Files
rdplib/protocol/pdu/pdu.go
T

936 lines
34 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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))
}