package pdu import ( "bytes" "encoding/binary" "fmt" "log/slog" "sync" "sync/atomic" "git.zeroonesoft.cn/golib/rdplib/core" "git.zeroonesoft.cn/golib/rdplib/emission" "git.zeroonesoft.cn/golib/rdplib/protocol/t125/gcc" ) var readerPool = sync.Pool{ New: func() any { return new(bytes.Reader) }, } // LastServerErrorInfo 记录服务器最近一次通过 SET_ERROR_INFO PDU 报告的错误码。 // 服务器断开前通常先发送该 PDU(如 0x0000112F = ERRINFO_GRAPHICS_SUBSYSTEM_FAILED), // 上层据此判断断连原因并做降级重试。 var LastServerErrorInfo atomic.Uint32 // fastPathBufPool reuses byte slices for serializing fast-path input PDUs. // Capacity 128 covers the maximum frame: 1 + 7*15 = 106 bytes. var fastPathBufPool = sync.Pool{ New: func() any { return make([]byte, 0, 128) }, } type PDULayer struct { bitmapCachePersistKeys []uint64 emission.Emitter transport core.Transport sharedId uint32 userId uint16 channelId uint16 serverCapabilities map[CapsType]Capability clientCapabilities map[CapsType]Capability fastPathSender core.FastPathSender // serverFastPathInput is set after capability exchange when both sides // advertise INPUT_FLAG_FASTPATH_INPUT, allowing client input to be sent // using the much shorter fast-path framing (MS-RDPBCGR §2.2.8.1.2). serverFastPathInput bool demandActivePDU *DemandActivePDU mppc *core.MppcDecompressor } func NewPDULayer(t core.Transport) *PDULayer { p := &PDULayer{ Emitter: *emission.NewEmitter(), transport: t, sharedId: 0x103EA, serverCapabilities: map[CapsType]Capability{ CAPSTYPE_GENERAL: &GeneralCapability{ ProtocolVersion: 0x0200, }, CAPSTYPE_BITMAP: &BitmapCapability{ Receive1BitPerPixel: 0x0001, Receive4BitsPerPixel: 0x0001, Receive8BitsPerPixel: 0x0001, BitmapCompressionFlag: 0x0001, MultipleRectangleSupport: 0x0001, }, CAPSTYPE_ORDER: &OrderCapability{ DesktopSaveXGranularity: 1, DesktopSaveYGranularity: 20, MaximumOrderLevel: 1, OrderFlags: NEGOTIATEORDERSUPPORT, DesktopSaveSize: 480 * 480, }, CAPSTYPE_POINTER: &PointerCapability{ColorPointerCacheSize: 20}, CAPSTYPE_INPUT: &InputCapability{}, CAPSTYPE_VIRTUALCHANNEL: &VirtualChannelCapability{}, CAPSTYPE_FONT: &FontCapability{SupportFlags: 0x0001}, CAPSTYPE_COLORCACHE: &ColorCacheCapability{CacheSize: 0x0006}, CAPSTYPE_SHARE: &ShareCapability{}, }, clientCapabilities: map[CapsType]Capability{ CAPSTYPE_GENERAL: &GeneralCapability{ ProtocolVersion: 0x0200, }, CAPSTYPE_BITMAP: &BitmapCapability{ Receive1BitPerPixel: 0x0001, Receive4BitsPerPixel: 0x0001, Receive8BitsPerPixel: 0x0001, BitmapCompressionFlag: 0x0001, MultipleRectangleSupport: 0x0001, }, CAPSTYPE_ORDER: &OrderCapability{ DesktopSaveXGranularity: 1, DesktopSaveYGranularity: 20, MaximumOrderLevel: 1, OrderFlags: NEGOTIATEORDERSUPPORT, DesktopSaveSize: 480 * 480, TextANSICodePage: 0x4e4, }, CAPSTYPE_CONTROL: &ControlCapability{0, 0, 2, 2}, CAPSTYPE_ACTIVATION: &WindowActivationCapability{}, CAPSTYPE_POINTER: &PointerCapability{1, 20, 20}, CAPSTYPE_SHARE: &ShareCapability{}, CAPSTYPE_COLORCACHE: &ColorCacheCapability{6, 0}, CAPSTYPE_SOUND: &SoundCapability{0x0001, 0}, CAPSTYPE_INPUT: &InputCapability{}, CAPSTYPE_FONT: &FontCapability{0x0001, 0}, CAPSTYPE_BRUSH: &BrushCapability{BRUSH_COLOR_8x8}, CAPSTYPE_GLYPHCACHE: &GlyphCapability{}, CAPSETTYPE_BITMAP_CODECS: newClientBitmapCodecsCapability(), CAPSTYPE_BITMAPCACHE_REV2: &BitmapCache2Capability{ BitmapCachePersist: 2, CachesNum: 5, BmpC0Cells: 0x258, BmpC1Cells: 0x258, BmpC2Cells: 0x800, BmpC3Cells: 0x1000, BmpC4Cells: 0x800, }, CAPSTYPE_VIRTUALCHANNEL: &VirtualChannelCapability{0, 1600}, CAPSETTYPE_MULTIFRAGMENTUPDATE: &MultiFragmentUpdate{0x3F0000}, CAPSTYPE_RAIL: &RemoteProgramsCapability{ RailSupportLevel: RAIL_LEVEL_SUPPORTED | RAIL_LEVEL_SHELL_INTEGRATION_SUPPORTED | RAIL_LEVEL_LANGUAGE_IME_SYNC_SUPPORTED | RAIL_LEVEL_SERVER_TO_CLIENT_IME_SYNC_SUPPORTED | RAIL_LEVEL_HIDE_MINIMIZED_APPS_SUPPORTED | RAIL_LEVEL_WINDOW_CLOAKING_SUPPORTED | RAIL_LEVEL_HANDSHAKE_EX_SUPPORTED | RAIL_LEVEL_DOCKED_LANGBAR_SUPPORTED, }, CAPSETTYPE_LARGE_POINTER: &LargePointerCapability{1}, CAPSETTYPE_COMPDESK: &DesktopCompositionCapability{ CompDeskSupportLevel: 1, // COMPDESK_SUPPORTED }, CAPSETTYPE_SURFACE_COMMANDS: &SurfaceCommandsCapability{ CmdFlags: SURFCMDS_SET_SURFACE_BITS | SURFCMDS_STREAM_SURFACE_BITS | SURFCMDS_FRAME_MARKER, }, CAPSSETTYPE_FRAME_ACKNOWLEDGE: &FrameAcknowledgeCapability{2}, }, mppc: core.NewMppcDecompressor(), } t.On("close", func() { p.Emit("close") }).On("error", func(err error) { p.Emit("error", err) }) return p } func (p *PDULayer) sendPDU(message PDUMessage) { pdu := NewPDU(p.userId, message) p.transport.Write(pdu.serialize()) } func (p *PDULayer) sendDataPDU(message DataPDUData) { dataPdu := NewDataPDU(message, p.sharedId) p.sendPDU(dataPdu) } func (p *PDULayer) SetFastPathSender(f core.FastPathSender) { p.fastPathSender = f } type Client struct { *PDULayer clientCoreData *gcc.ClientCoreData buff *bytes.Buffer // fragCode 记住分片更新首片的 updateCode:续片的 updateCode 无意义 //(Win10 发送的是载荷字节),重组完成后必须用它来解析。 fragCode uint8 // fragActive/fragPoisoned:分片重组进行中标记;解压失败的首片会把 // 整个更新标记为中毒,其剩余分片被吞掉(防止孤儿分片拼出垃圾)。 fragActive bool fragPoisoned bool // decErrLogs:MPPC 解压失败日志限次计数器。 decErrLogs int // surfResetDumped:SURFCMDS 重置全量转储限次(离线分析用)。 surfResetDumped int // bmpRectLogs:16bpp 矩形取证日志限次计数器。 bmpRectLogs int // bitmapCachePersistKeys:客户端持久位图缓存键(6.4b M2), // finalize 时经 maybeSendPersistentKeyList 上报服务器。 bitmapCachePersistKeys []uint64 // bmpFrame:fast-path 位图跨段累积缓冲(Win10 会把一个位图更新拆到 // 多个 fast-path PDU,按矩形结构完整解析后整体消费)。 bmpFrame []byte // surfBuf:fast-path 表面命令跨段累积缓冲(按命令流自定界解析)。 surfBuf []byte // surfHdrLen:SET_SURFACE_BITS 头部长度(21=经典 / 22=Win10),由本 // 连接第一条命令判定后锁定。服务器布局会话内恒定;逐帧启发式在 // inclusive/exclusive 双兼容下会偶发选错 ±1 字节,错位累积到批次 // 末尾即"未知命令类型→整段重置"——表现为周期性花屏+断流。 surfHdrLen int } func NewClient(t core.Transport) *Client { c := &Client{ PDULayer: NewPDULayer(t), buff: &bytes.Buffer{}, } c.transport.Once("connect", c.connect) return c } func (c *Client) connect(data *gcc.ClientCoreData, userId uint16, channelId uint16) { slog.Debug("pdu connect", "userId", userId, "channelId", channelId) c.clientCoreData = data c.userId = userId c.channelId = channelId c.transport.Once("data", c.recvDemandActivePDU) } func (c *Client) recvDemandActivePDU(s []byte) { r := readerPool.Get().(*bytes.Reader) r.Reset(s) defer readerPool.Put(r) pdu, err := readPDU(r, c.mppc) if err != nil { slog.Error("recvDemandActivePDU", "err", err) return } if pdu.ShareCtrlHeader.PDUType != PDUTYPE_DEMANDACTIVEPDU { if pdu.ShareCtrlHeader.PDUType == PDUTYPE_DEACTIVATEALLPDU { // Per [MS-RDPBCGR] the server may send DeactivateAllPDU before // DemandActivePDU (e.g. GNOME RDP after RDPGFX capability // exchange). Stay on the same connection and keep waiting, // exactly as FreeRDP does. slog.Debug("received DeactivateAllPDU while waiting for DemandActivePDU; continuing to wait") c.transport.Once("data", c.recvDemandActivePDU) return } if pdu.ShareCtrlHeader.PDUType == PDUTYPE_SERVER_REDIR_PKT { if redir, ok := pdu.Message.(*ServerRedirectionPDU); ok { c.Emit("redirect", redir) } return } slog.Debug("ignore message during connection sequence", "type", pdu.ShareCtrlHeader.PDUType) c.transport.Once("data", c.recvDemandActivePDU) return } c.sharedId = pdu.Message.(*DemandActivePDU).SharedId c.demandActivePDU = pdu.Message.(*DemandActivePDU) for _, caps := range c.demandActivePDU.CapabilitySets { slog.Debug("serverCaps", "type", caps.Type(), "value", caps) c.serverCapabilities[caps.Type()] = caps } if ic, ok := c.serverCapabilities[CAPSTYPE_INPUT].(*InputCapability); ok { c.serverFastPathInput = ic.Flags&INPUT_FLAG_FASTPATH_INPUT != 0 } c.sendConfirmActivePDU() c.sendClientFinalizeSynchronizePDU() c.transport.Once("data", c.recvServerSynchronizePDU) } func (c *Client) sendConfirmActivePDU() { pdu := NewConfirmActivePDU() generalCapa := c.clientCapabilities[CAPSTYPE_GENERAL].(*GeneralCapability) generalCapa.OSMajorType = OSMAJORTYPE_WINDOWS generalCapa.OSMinorType = OSMINORTYPE_WINDOWS_NT generalCapa.GeneralCompressionTypes = 0x0002 // PACKET_COMPR_TYPE_64K: advertise MPPC-64K support generalCapa.ExtraFlags = LONG_CREDENTIALS_SUPPORTED | NO_BITMAP_COMPRESSION_HDR | FASTPATH_OUTPUT_SUPPORTED | AUTORECONNECT_SUPPORTED generalCapa.RefreshRectSupport = 1 generalCapa.SuppressOutputSupport = 1 bitmapCapa := c.clientCapabilities[CAPSTYPE_BITMAP].(*BitmapCapability) bitmapCapa.PreferredBitsPerPixel = 32 bitmapCapa.DesktopWidth = c.clientCoreData.DesktopWidth bitmapCapa.DesktopHeight = c.clientCoreData.DesktopHeight bitmapCapa.DesktopResizeFlag = 0x0001 orderCapa := c.clientCapabilities[CAPSTYPE_ORDER].(*OrderCapability) orderCapa.OrderFlags = NEGOTIATEORDERSUPPORT | ZEROBOUNDSDELTASSUPPORT | COLORINDEXSUPPORT | ORDERFLAGS_EXTRA_FLAGS orderCapa.OrderSupportExFlags |= ORDERFLAGS_EX_ALTSEC_FRAME_MARKER_SUPPORT orderCapa.OrderSupport[TS_NEG_DSTBLT_INDEX] = 1 orderCapa.OrderSupport[TS_NEG_PATBLT_INDEX] = 1 orderCapa.OrderSupport[TS_NEG_SCRBLT_INDEX] = 1 // MEMBLT + 位图缓存 v2(stage6 6.4b M1 会话内缓存):服务器把重复 // 位图存入客户端缓存单元,后续以 MemBlt 引用回贴,替代全量重发。 orderCapa.OrderSupport[TS_NEG_MEMBLT_INDEX] = 1 //orderCapa.OrderSupport[TS_NEG_LINETO_INDEX] = 1 //orderCapa.OrderSupport[TS_NEG_MEM3BLT_INDEX] = 1 //orderCapa.OrderSupport[TS_NEG_POLYLINE_INDEX] = 1 /*orderCapa.OrderSupport[TS_NEG_MULTIOPAQUERECT_INDEX] = 1 orderCapa.OrderSupport[TS_NEG_GLYPH_INDEX_INDEX] = 1 //orderCapa.OrderSupport[TS_NEG_DRAWNINEGRID_INDEX] = 1 orderCapa.OrderSupport[TS_NEG_SAVEBITMAP_INDEX] = 1 orderCapa.OrderSupport[TS_NEG_POLYGON_SC_INDEX] = 1 orderCapa.OrderSupport[TS_NEG_POLYGON_CB_INDEX] = 1 orderCapa.OrderSupport[TS_NEG_ELLIPSE_SC_INDEX] = 1 orderCapa.OrderSupport[TS_NEG_ELLIPSE_CB_INDEX] = 1*/ //orderCapa.OrderSupport[TS_NEG_FAST_GLYPH_INDEX] = 1 inputCapa := c.clientCapabilities[CAPSTYPE_INPUT].(*InputCapability) inputCapa.Flags = INPUT_FLAG_SCANCODES | INPUT_FLAG_MOUSEX | INPUT_FLAG_UNICODE | INPUT_FLAG_FASTPATH_INPUT | INPUT_FLAG_FASTPATH_INPUT2 inputCapa.KeyboardLayout = c.clientCoreData.KbdLayout inputCapa.KeyboardType = c.clientCoreData.KeyboardType inputCapa.KeyboardSubType = c.clientCoreData.KeyboardSubType inputCapa.KeyboardFunctionKey = c.clientCoreData.KeyboardFnKeys inputCapa.ImeFileName = c.clientCoreData.ImeFileName glyphCapa := c.clientCapabilities[CAPSTYPE_GLYPHCACHE].(*GlyphCapability) /*glyphCapa.GlyphCache[0] = cacheEntry{254, 4} glyphCapa.GlyphCache[1] = cacheEntry{254, 4} glyphCapa.GlyphCache[2] = cacheEntry{254, 8} glyphCapa.GlyphCache[3] = cacheEntry{254, 8} glyphCapa.GlyphCache[4] = cacheEntry{254, 16} glyphCapa.GlyphCache[5] = cacheEntry{254, 32} glyphCapa.GlyphCache[6] = cacheEntry{254, 64} glyphCapa.GlyphCache[7] = cacheEntry{254, 128} glyphCapa.GlyphCache[8] = cacheEntry{254, 256} glyphCapa.GlyphCache[9] = cacheEntry{64, 2048} glyphCapa.FragCache = 0x01000100*/ glyphCapa.SupportLevel = GLYPH_SUPPORT_NONE // 位图缓存 v2 单元(MS-RDPBCGR 2.2.7.1.8):单元尺寸取 mstsc 量级。 // CacheFlags 持久键支持位留待 6.4b M2(配合 PERSISTENT_KEY_LIST)。 rev2, ok := c.clientCapabilities[CAPSTYPE_BITMAPCACHE_REV2].(*BitmapCacheRev2Capability) if !ok || rev2 == nil { rev2 = &BitmapCacheRev2Capability{} c.clientCapabilities[CAPSTYPE_BITMAPCACHE_REV2] = rev2 } rev2.CacheCells = [5]BitmapCacheV2CellInfo{ {NumEntries: 600, MaxCellSize: 256}, {NumEntries: 1024, MaxCellSize: 512}, {NumEntries: 4096, MaxCellSize: 1024}, {}, {}, } // v1 位图缓存 caps(0x0004)与 Rev2 一并广告——与 mstsc 行为对齐, // 部分服务器以 v1 caps 存在与否决定是否启用缓存订单。 v1, ok := c.clientCapabilities[CAPSTYPE_BITMAPCACHE].(*BitmapCacheCapability) if !ok || v1 == nil { v1 = &BitmapCacheCapability{} c.clientCapabilities[CAPSTYPE_BITMAPCACHE] = v1 } pdu.SharedId = c.sharedId for _, v := range c.clientCapabilities { slog.Debug("clientCaps", "type", v.Type(), "value", v) pdu.CapabilitySets = append(pdu.CapabilitySets, v) } pdu.NumberCapabilities = uint16(len(pdu.CapabilitySets)) pdu.LengthSourceDescriptor = c.demandActivePDU.LengthSourceDescriptor pdu.SourceDescriptor = c.demandActivePDU.SourceDescriptor pdu.LengthCombinedCapabilities = c.demandActivePDU.LengthCombinedCapabilities c.sendPDU(pdu) } func (c *Client) sendClientFinalizeSynchronizePDU() { c.sendDataPDU(NewSynchronizeDataPDU(c.channelId)) c.sendDataPDU(&ControlDataPDU{Action: CTRLACTION_COOPERATE}) c.sendDataPDU(&ControlDataPDU{Action: CTRLACTION_REQUEST_CONTROL}) c.maybeSendPersistentKeyList() c.sendDataPDU(&FontListDataPDU{ListFlags: 0x0003, EntrySize: 0x0032}) } func (c *Client) recvServerSynchronizePDU(s []byte) { r := readerPool.Get().(*bytes.Reader) r.Reset(s) defer readerPool.Put(r) pdu, err := readPDU(r, c.mppc) if err != nil { slog.Error("recvServerSynchronizePDU", "err", err) return } dataPdu, ok := pdu.Message.(*DataPDU) if !ok || dataPdu.Header.PDUType2 != PDUTYPE2_SYNCHRONIZE { if ok { slog.Error("recvServerSynchronizePDU ignore datapdu", "type2", dataPdu.Header.PDUType2) } else { slog.Error("recvServerSynchronizePDU ignore message", "type", pdu.ShareCtrlHeader.PDUType) } slog.Debug("recvServerSynchronizePDU dataPdu", "pdu", &dataPdu) c.transport.Once("data", c.recvServerSynchronizePDU) return } c.transport.Once("data", c.recvServerControlCooperatePDU) } func (c *Client) recvServerControlCooperatePDU(s []byte) { r := readerPool.Get().(*bytes.Reader) r.Reset(s) defer readerPool.Put(r) pdu, err := readPDU(r, c.mppc) if err != nil { slog.Error("recvServerControlCooperatePDU", "err", err) return } dataPdu, ok := pdu.Message.(*DataPDU) if !ok || dataPdu.Header.PDUType2 != PDUTYPE2_CONTROL { if ok { slog.Error("recvServerControlCooperatePDU ignore datapdu", "type2", dataPdu.Header.PDUType2) } else { slog.Error("recvServerControlCooperatePDU ignore message", "type", pdu.ShareCtrlHeader.PDUType) } c.transport.Once("data", c.recvServerControlCooperatePDU) return } if dataPdu.Data.(*ControlDataPDU).Action != CTRLACTION_COOPERATE { slog.Error("recvServerControlCooperatePDU ignore", "action", dataPdu.Data.(*ControlDataPDU).Action) c.transport.Once("data", c.recvServerControlCooperatePDU) return } c.transport.Once("data", c.recvServerControlGrantedPDU) } func (c *Client) recvServerControlGrantedPDU(s []byte) { r := readerPool.Get().(*bytes.Reader) r.Reset(s) defer readerPool.Put(r) pdu, err := readPDU(r, c.mppc) if err != nil { slog.Error("recvServerControlGrantedPDU", "err", err) return } dataPdu, ok := pdu.Message.(*DataPDU) if !ok || dataPdu.Header.PDUType2 != PDUTYPE2_CONTROL { if ok { slog.Error("recvServerControlGrantedPDU ignore datapdu", "type2", dataPdu.Header.PDUType2) } else { slog.Error("recvServerControlGrantedPDU ignore message", "type", pdu.ShareCtrlHeader.PDUType) } c.transport.Once("data", c.recvServerControlGrantedPDU) return } if dataPdu.Data.(*ControlDataPDU).Action != CTRLACTION_GRANTED_CONTROL { slog.Error("recvServerControlGrantedPDU ignore", "action", dataPdu.Data.(*ControlDataPDU).Action) c.transport.Once("data", c.recvServerControlGrantedPDU) return } c.transport.Once("data", c.recvServerFontMapPDU) } func (c *Client) recvServerFontMapPDU(s []byte) { r := readerPool.Get().(*bytes.Reader) r.Reset(s) defer readerPool.Put(r) pdu, err := readPDU(r, c.mppc) if err != nil { slog.Error("recvServerFontMapPDU", "err", err) return } dataPdu, ok := pdu.Message.(*DataPDU) if !ok || dataPdu.Header.PDUType2 != PDUTYPE2_FONTMAP { if ok { slog.Error("recvServerFontMapPDU ignore datapdu", "type2", dataPdu.Header.PDUType2) } else { slog.Error("recvServerFontMapPDU ignore message", "type", pdu.ShareCtrlHeader.PDUType) } return } c.transport.On("data", c.recvPDU) // Tell the server we're ready to receive display updates (MS-RDPBCGR 2.2.11.3.1) slog.Debug("Sending SuppressOutput (ALLOW_DISPLAY_UPDATES)") c.sendDataPDU(&SuppressOutputPDU{ AllowDisplayUpdates: 1, Right: c.clientCoreData.DesktopWidth - 1, Bottom: c.clientCoreData.DesktopHeight - 1, }) c.Emit("ready") } func (c *Client) recvPDU(s []byte) { r := readerPool.Get().(*bytes.Reader) r.Reset(s) defer readerPool.Put(r) if r.Len() > 0 { p, err := readPDU(r, c.mppc) if err != nil { slog.Error("recvPDU", "err", err) return } if p.ShareCtrlHeader.PDUType == PDUTYPE_DEACTIVATEALLPDU { // Server is reactivating the session (e.g. desktop resize). // Signal callers to pause input until "ready" fires again. slog.Debug("received DeactivateAllPDU during active session, waiting for reactivation") c.Emit("deactivateAll") c.transport.Once("data", c.recvDemandActivePDU) } else if p.ShareCtrlHeader.PDUType == PDUTYPE_SERVER_REDIR_PKT { if redir, ok := p.Message.(*ServerRedirectionPDU); ok { c.Emit("redirect", redir) } } else if p.ShareCtrlHeader.PDUType == PDUTYPE_DATAPDU { d := p.Message.(*DataPDU) if d.Header.PDUType2 == PDUTYPE2_UPDATE { up := d.Data.(*UpdateDataPDU) p := up.Udata if up.UpdateType == FASTPATH_UPDATETYPE_BITMAP { c.Emit("bitmap", p.(*BitmapUpdateDataPDU).Rectangles) } else if up.UpdateType == FASTPATH_UPDATETYPE_ORDERS { c.Emit("orders", p.(*FastPathOrdersPDU).OrderPdus) } } else if d.Header.PDUType2 == PDUTYPE2_POINTER { pp := d.Data.(*PointerDataPDU) if pp.Pdata != nil { switch pp.MessageType { case TS_PTRUPDATE_TYPE_CACHED: c.Emit("pointer_cached", pp.Pdata.(*FastPathUpdateCachedPDU).CacheIdx) case TS_PTRUPDATE_TYPE_POINTER: c.Emit("pointer_update", pp.Pdata.(*FastPathUpdatePointerPDU)) } } if pp.MessageType == TS_PTRUPDATE_TYPE_SYSTEM { c.Emit("pointer_hide") } } else if d.Header.PDUType2 == PDUTYPE2_SET_ERROR_INFO_PDU { // 服务器在断开连接前会发送错误码,不解析就无法定位断连原因 ei := d.Data.(*ErrorInfoDataPDU) LastServerErrorInfo.Store(ei.ErrorInfo) slog.Warn("server SET_ERROR_INFO", "code", fmt.Sprintf("0x%08X", ei.ErrorInfo)) } } } } func (c *Client) RecvFastPath(secFlag byte, s []byte) { r := readerPool.Get().(*bytes.Reader) r.Reset(s) defer readerPool.Put(r) for r.Len() > 0 { updateHeader, err := core.ReadUInt8(r) if err != nil { return } updateCode := updateHeader & 0x0f fragmentation := updateHeader & 0x30 // compression 占 header 的高 2 位:字段值 0b10(左移后 0x80)表示 // 后随 1 字节 compressionFlags。此前误拿位掩码值与字段值 0x2 直接 // 比较,恒为假 → 压缩标志字节从未被消费,整条 fast-path 流错位。 compression := updateHeader & 0xC0 var compressionFlags uint8 = 0 if compression == FASTPATH_OUTPUT_COMPRESSION_USED<<6 { compressionFlags, err = core.ReadUInt8(r) if err != nil { return } } size, err := core.ReadUint16LE(r) if err != nil { return } // Read exactly `size` bytes for this update's payload. payload, err := core.ReadBytes(int(size), r) if err != nil { return } slog.Debug("RecvFastPath", "Code", FastPathUpdateType(updateCode), "compressionFlags", compressionFlags, "fragmentation", fragmentation, "size", size) // 先解压再重组:压缩以每个 fast-path 分片为独立单位(各分片头部 // 携带各自的 compressionFlags),共享同一条 MPPC 历史流;拼接压缩 // 字节后一次性解压必然产生垃圾。 decompressed := payload decFailed := false if compressionFlags != 0 && c.mppc != nil { out, err := c.mppc.Decompress(compressionFlags, payload) if err != nil { decFailed = true if c.decErrLogs < 4 { c.decErrLogs++ slog.Warn("RecvFastPath: MPPC decompression failed", "err", err, "cf", fmt.Sprintf("%02X", compressionFlags), "frag", fragmentation, "size", size) } } else { decompressed = out } } // 分片重组:续片的 updateCode 不可信(MS-RDPBCGR 2.2.9.1.1.3.1 // 规定忽略),解析时使用首片记住的 code。解压失败的分片会把所在 // 更新标记为中毒:吞掉其剩余分片,避免孤儿分片拼出垃圾更新。 if fragmentation != FASTPATH_FRAGMENT_SINGLE { if fragmentation == FASTPATH_FRAGMENT_FIRST { c.buff.Reset() c.fragCode = updateCode c.fragActive = true c.fragPoisoned = decFailed } else if !c.fragActive || c.fragPoisoned { continue // 无首片的孤儿分片 / 中毒更新的剩余分片 } c.buff.Write(decompressed) if fragmentation != FASTPATH_FRAGMENT_LAST { continue } c.fragActive = false if c.fragPoisoned { continue } payload = c.buff.Bytes() updateCode = c.fragCode } else if decFailed { continue } else { // 单片更新:解析对象必须是解压后的结果(此前漏赋值导致 // 嗅探/解析拿到原始压缩数据,单片位图全部报废)。 payload = decompressed } // Surface Commands:Win10 会把一条表面命令拆到多个 fast-path PDU //(实测 frag 位 FIRST/NEXT/LAST 使用规范),按命令流自定界累积解析,命令完整即发射。 if updateCode == FASTPATH_UPDATETYPE_SURFCMDS { if decFailed { continue } c.surfBuf = append(c.surfBuf, payload...) for { result, consumed, needMore, valid, dropped, hdrLen := parseSurfaceCommandsIncremental(c.surfBuf, c.surfHdrLen) if c.surfHdrLen == 0 && hdrLen != 0 { c.surfHdrLen = hdrLen // 首条命令锁定头部布局(会话内恒定) } if !valid { if c.surfResetDumped < 2 { c.surfResetDumped++ slog.Warn("SURFCMDS stream reset full dump", "n", c.surfResetDumped, "bufLen", len(c.surfBuf), "hdrLen", c.surfHdrLen, "hex", fmt.Sprintf("% X", c.surfBuf)) } slog.Warn("SURFCMDS stream reset", "bufLen", len(c.surfBuf), "hex", fmt.Sprintf("% X", c.surfBuf[:min(24, len(c.surfBuf))])) c.surfBuf = c.surfBuf[:0] c.surfHdrLen = 0 // 布局判定作废,下批重新锁定 break } if len(result.Rects) > 0 { c.Emit("bitmap", result.Rects) } for _, fid := range result.FrameIDs { c.sendDataPDU(&FrameAcknowledgeDataPDU{FrameID: fid}) } c.surfBuf = c.surfBuf[consumed:] if dropped { // 尾部未知结构:已解析前缀照常上屏,仅丢弃尾巴。 // 实测 Win10 整屏重绘批次末尾带 11 字节未文档化尾巴, // 整批报废会造成周期性花屏+断流。 slog.Debug("SURFCMDS tail dropped", "consumed", consumed, "bufLen", len(c.surfBuf)) c.surfBuf = c.surfBuf[:0] break } if needMore || consumed == 0 || len(c.surfBuf) < 2 { break } } if len(c.surfBuf) > 8<<20 { c.surfBuf = c.surfBuf[:0] // 防御:异常累积封顶 } continue } // fast-path 位图:Win10 会把一个位图更新拆到多个 fast-path PDU //(首段携带 updateType + numberRectangles + 前几个矩形,后续段 // 继续补充矩形数据)。将解压结果持续累积,一旦能完整解析出全部 // 矩形就发射。非位图更新仍走通用分片状态机。 // 同样必须追加 payload(重组缓冲)而非 decompressed,理由同上。 if updateCode == FASTPATH_UPDATETYPE_BITMAP { if decFailed { continue } c.bmpFrame = append(c.bmpFrame, payload...) for { rects, consumed, complete, valid := parseBitmapFrame(c.bmpFrame) if !valid { slog.Warn("BITMAP frame reset", "bufLen", len(c.bmpFrame), "hex", fmt.Sprintf("% X", c.bmpFrame[:min(24, len(c.bmpFrame))])) c.bmpFrame = c.bmpFrame[:0] break } if !complete { break } if len(rects) > 0 { if c.bmpRectLogs < 24 { for i, rc := range rects { if c.bmpRectLogs >= 24 { break } c.bmpRectLogs++ slog.Debug("BITMAP rect", "i", i, "n", len(rects), "dx", rc.DestLeft, "dy", rc.DestTop, "dr", rc.DestRight, "db", rc.DestBottom, "w", rc.Width, "h", rc.Height, "bpp", rc.BitsPerPixel, "flags", fmt.Sprintf("0x%04X", rc.Flags), "len", len(rc.BitmapDataStream)) } } c.Emit("bitmap", rects) } c.bmpFrame = c.bmpFrame[consumed:] if len(c.bmpFrame) < 4 { break } } continue } if updateCode == FASTPATH_UPDATETYPE_POINTER { // 取证:指针更新原始字节,核对 xorBpp 后各字段对齐。 slog.Debug("PTR raw", "hex", fmt.Sprintf("% X", payload[:min(32, len(payload))])) } pr := bytes.NewReader(payload) p, err := readFastPathUpdatePDU(pr, updateCode) if err != nil { slog.Warn("readFastPathUpdatePDU:", "Code", FastPathUpdateType(updateCode), "err", err) continue } if updateCode == FASTPATH_UPDATETYPE_BITMAP { c.Emit("bitmap", p.Data.(*FastPathBitmapUpdateDataPDU).Rectangles) } else if updateCode == FASTPATH_UPDATETYPE_COLOR { c.Emit("color", p.Data.(*FastPathColorPdu)) } else if updateCode == FASTPATH_UPDATETYPE_ORDERS { c.Emit("orders", p.Data.(*FastPathOrdersPDU).OrderPdus) } else if updateCode == FASTPATH_UPDATETYPE_PTR_NULL { c.Emit("pointer_hide") } else if updateCode == FASTPATH_UPDATETYPE_PTR_DEFAULT { c.Emit("pointer_default") } else if updateCode == FASTPATH_UPDATETYPE_PTR_POSITION { pp := p.Data.(*FastPathPointerPositionPDU) c.Emit("pointer_position", pp.X, pp.Y) } else if updateCode == FASTPATH_UPDATETYPE_CACHED { c.Emit("pointer_cached", p.Data.(*FastPathUpdateCachedPDU).CacheIdx) } else if updateCode == FASTPATH_UPDATETYPE_POINTER { c.Emit("pointer_update", p.Data.(*FastPathUpdatePointerPDU)) } } } type InputEventsInterface interface { Serialize() []byte } // fastPathEncoder is implemented by input event types that know how to // produce their Fast-Path Input wire encoding (MS-RDPBCGR §2.2.8.1.2.2). type fastPathEncoder interface { FastPathEncode(buf []byte) []byte } func (c *Client) SendInputEvents(msgType uint16, events []InputEventsInterface) { if c.serverFastPathInput && c.fastPathSender != nil && c.canSendFastPathInput(events) { if c.sendFastPathInputEvents(events) { return } // Fall back to slow-path on send failure (e.g. legacy encryption). } p := &ClientInputEventPDU{} p.NumEvents = uint16(len(events)) p.SlowPathInputEvents = make([]SlowPathInputEvent, 0, p.NumEvents) for _, in := range events { seria := in.Serialize() s := SlowPathInputEvent{0, msgType, len(seria), seria} p.SlowPathInputEvents = append(p.SlowPathInputEvents, s) } c.sendDataPDU(p) } // canSendFastPathInput reports whether every event in the batch implements // the fast-path encoder. Falls back to slow-path if any event type doesn't // (currently just SynchronizeEvent, which the client never sends). func (c *Client) canSendFastPathInput(events []InputEventsInterface) bool { if len(events) == 0 || len(events) > 15 { return false } for _, e := range events { if _, ok := e.(fastPathEncoder); !ok { return false } } return true } func (c *Client) sendFastPathInputEvents(events []InputEventsInterface) bool { buf := fastPathBufPool.Get().([]byte) buf = buf[:0] buf = append(buf, byte(len(events))) for _, e := range events { buf = e.(fastPathEncoder).FastPathEncode(buf) } _, err := c.fastPathSender.SendFastPath(0, buf) fastPathBufPool.Put(buf[:cap(buf)]) if err != nil { // Disable for the rest of the session so we don't keep paying the // failed-attempt cost on every input event. c.serverFastPathInput = false slog.Warn("fast-path input disabled, falling back to slow-path", "err", err) return false } return true } // SendRefreshRect requests the server to redraw the given screen rectangle. // This causes the server to send a full refresh (including a new H.264 IDR) // for the specified region, which is useful after a decoder reset. func (c *Client) SendRefreshRect(width, height uint16) { slog.Debug("PDU: SendRefreshRect", "w", width, "h", height) c.sendDataPDU(&RefreshRectPDU{ NumberOfAreas: 1, Right: width - 1, Bottom: height - 1, }) } // SendForceRefresh asks the server for a complete display repaint by toggling // SuppressOutput off→on. Per MS-RDPBCGR 2.2.11.3.1, sending ALLOW_DISPLAY_UPDATES // after SUPPRESS_DISPLAY_UPDATES forces the server to send a fresh full-screen // update — for the RDPGFX H.264 pipeline this means a new IDR frame, which is // what we need to recover after a hardware-decoder hard reset. Plain // SendRefreshRect is sometimes silently ignored by Windows servers while a // video stream is active; this is the reliable fallback used by mstsc/FreeRDP. func (c *Client) SendForceRefresh(width, height uint16) { slog.Debug("PDU: SendForceRefresh (suppress→allow)", "w", width, "h", height) c.sendDataPDU(&SuppressOutputPDU{ AllowDisplayUpdates: 0x00, // SUPPRESS_DISPLAY_UPDATES }) c.sendDataPDU(&SuppressOutputPDU{ AllowDisplayUpdates: 0x01, // ALLOW_DISPLAY_UPDATES Right: width - 1, Bottom: height - 1, }) } // bitmapDropLogs:垃圾帧丢弃诊断日志限次计数器。 var bitmapDropLogs = 0 // parseBitmapFrame 尝试把累积缓冲解析为完整的 fast-path 位图更新 // (updateType(2) + numberRectangles(2) + 矩形数组)。 // 返回 valid=false 表示缓冲开头不是合法位图结构(应丢弃); // complete=false 表示数据不足(应继续累积); // complete=true 时 consumed 为本次更新占用的字节数,尾部可能残留 // 下一更新的开头。 func parseBitmapFrame(buf []byte) (rects []BitmapData, consumed int, complete bool, valid bool) { if len(buf) < 4 { return nil, 0, false, true } if u16 := binary.LittleEndian.Uint16(buf[0:2]); u16 != 0x0001 { // UPDATE_TYPE_BITMAP return nil, 0, false, false } nr := int(binary.LittleEndian.Uint16(buf[2:4])) if nr > 4096 { if bitmapDropLogs < 3 { bitmapDropLogs++ slog.Warn("bitmap frame dropped: implausible rect count", "nr", nr, "len", len(buf), "prefix", fmt.Sprintf("% X", buf[:min(24, len(buf))])) } return nil, 0, false, false } rects = make([]BitmapData, 0, nr) pos := 4 for i := 0; i < nr; i++ { if pos+18 > len(buf) { return nil, 0, false, true } r := BitmapData{} r.DestLeft = binary.LittleEndian.Uint16(buf[pos:]) r.DestTop = binary.LittleEndian.Uint16(buf[pos+2:]) r.DestRight = binary.LittleEndian.Uint16(buf[pos+4:]) r.DestBottom = binary.LittleEndian.Uint16(buf[pos+6:]) r.Width = binary.LittleEndian.Uint16(buf[pos+8:]) r.Height = binary.LittleEndian.Uint16(buf[pos+10:]) r.BitsPerPixel = binary.LittleEndian.Uint16(buf[pos+12:]) r.Flags = binary.LittleEndian.Uint16(buf[pos+14:]) bl := int(binary.LittleEndian.Uint16(buf[pos+16:])) hdr := 0 if r.Flags&BITMAP_COMPRESSION != 0 && r.Flags&NO_BITMAP_COMPRESSION_HDR == 0 { if pos+18+8 > len(buf) { return nil, 0, false, true } bl = int(binary.LittleEndian.Uint16(buf[pos+18+2 : pos+18+4])) hdr = 8 } if bl < 0 || pos+18+hdr+bl > len(buf) { return nil, 0, false, true } r.BitmapDataStream = buf[pos+18+hdr : pos+18+hdr+bl] pos += 18 + hdr + bl rects = append(rects, r) } return rects, pos, true, true } // SetPersistentKeyList 注册客户端持久位图缓存持有的键(bitmap 管线)。 // 须在 DemandActive 之前调用;finalize 时仅当服务器广告了 // CAPSTYPE_BITMAPCACHE_HOSTSUPPORT 且键非空才实际发送。 func (p *PDULayer) SetPersistentKeyList(keys []uint64) { p.bitmapCachePersistKeys = keys } // maybeSendPersistentKeyList 在 Connection Finalization 序列 // (REQUEST_CONTROL 之后、FONT_LIST 之前)发送持久缓存键列表, // 位置与 FreeRDP/MS-RDPBCGR §2.2.1.17 一致。 func (p *PDULayer) maybeSendPersistentKeyList() { if len(p.bitmapCachePersistKeys) == 0 { return } hs, ok := p.serverCapabilities[CAPSTYPE_BITMAPCACHE_HOSTSUPPORT].(*BitmapCacheHostSupportCapability) if !ok || hs.CacheVersion < 1 { slog.Info("bmpcache: server lacks bitmap cache host support, skip persistent key list") return } var cells [5]uint16 if rev2, ok := p.clientCapabilities[CAPSTYPE_BITMAPCACHE_REV2].(*BitmapCacheRev2Capability); ok { for i := range rev2.CacheCells { if i < len(cells) { cells[i] = rev2.CacheCells[i].NumEntries } } } keys := p.bitmapCachePersistKeys if len(keys) > 2042 { // FreeRDP 同款上限:>2042 条会触发服务器错误 keys = keys[:2042] slog.Info("bmpcache: truncating persistent key list", "kept", len(keys)) } pdu := NewPDU(p.userId, NewPersistentKeyListPDU(p.sharedId, keys, cells)) p.transport.Write(pdu.serialize()) slog.Info("bmpcache: sent persistent key list", "keys", len(keys)) }