package grdp import ( "errors" "fmt" "image" "log/slog" "net" "os" "runtime/debug" "strings" "sync" "sync/atomic" "time" "unsafe" "git.zeroonesoft.cn/golib/rdplib/plugin" "git.zeroonesoft.cn/golib/rdplib/plugin/cliprdr" "git.zeroonesoft.cn/golib/rdplib/plugin/rdpdr" "git.zeroonesoft.cn/golib/rdplib/plugin/drdynvc" "git.zeroonesoft.cn/golib/rdplib/plugin/rdpedisp" "git.zeroonesoft.cn/golib/rdplib/plugin/rdpgfx" "git.zeroonesoft.cn/golib/rdplib/plugin/rdpsnd" "git.zeroonesoft.cn/golib/rdplib/core" "git.zeroonesoft.cn/golib/rdplib/protocol/nla" "git.zeroonesoft.cn/golib/rdplib/protocol/pdu" "git.zeroonesoft.cn/golib/rdplib/protocol/sec" "git.zeroonesoft.cn/golib/rdplib/protocol/t125" "git.zeroonesoft.cn/golib/rdplib/protocol/t125/gcc" "git.zeroonesoft.cn/golib/rdplib/protocol/tpkt" "git.zeroonesoft.cn/golib/rdplib/protocol/x224" ) // mouseCoalescer holds all state for mouse-move coalescing. // High-frequency move events are collapsed into at most one network PDU per // mouseCoalesceInterval, with the latest position always winning. type mouseCoalescer struct { mu sync.Mutex pending bool x, y int timer *time.Timer lastTx time.Time pdu pdu.PointerEvent pduBuf [1]pdu.InputEventsInterface } // wheelCoalescer holds all state for wheel-scroll coalescing (vertical and // horizontal). Rapid scroll events are accumulated over mouseCoalesceInterval // and sent as a single PDU whose rotation value is the sum of all deltas in // that window. accum/haccum are stored in RDP WHEEL_DELTA units (120 per // physical notch); haccum is the horizontal axis (positive = scroll right). type wheelCoalescer struct { mu sync.Mutex accum float64 haccum float64 timer *time.Timer lastTx time.Time pdu pdu.PointerEvent pduBuf [1]pdu.InputEventsInterface } // stubChannel is a no-op virtual channel handler for channels the server // expects to be present (e.g. rdpdr, cliprdr) but that we don't process. type stubChannel struct { name string option uint32 sender core.ChannelSender } func (s *stubChannel) GetType() (string, uint32) { return s.name, s.option } func (s *stubChannel) Sender(f core.ChannelSender) { s.sender = f } func (s *stubChannel) Process(data []byte) {} type RdpClient struct { hostPort string // ip:port width int height int kbdLayout uint32 keyboardType uint32 keyboardSubType uint32 timezoneName string timezoneBias int tpkt *tpkt.TPKT x224 *x224.X224 mcs *t125.MCSClient sec *sec.Client pdu *pdu.Client channels *plugin.Channels eventReady atomic.Bool decompressPool sync.Pool // pools []uint8 buffers for bitmap decompression flipLinePool sync.Pool // pools line-sized []uint8 buffers for bitmap vertical flip bmpEmitLogged int // 临时诊断:前 10 次位图事件计数 closed atomic.Bool // credentials stored for reconnection domain string user string password string // stored callbacks for re-registration on reconnect onErrorFn func(e error) onCloseFn func() onSuccessFn func() onReadyFn func() onBitmapPaintFn func([]Bitmap) onPointerHideFn func() onPointerCachedFn func(uint16) onPointerDefaultFn func() onPointerUpdateFn func(uint16, uint16, uint16, uint16, uint16, uint16, []byte, []byte) onAudioFn func(rdpsnd.AudioFormat, []byte) onAudioResetFn func() onH264RawFn func(destX, destY, w, h int, isKey bool, data []byte, regions []int32) onH264I420Fn func(destX, destY, w, h int, y []byte, yStride int, u []byte, uStride int, v []byte, vStride int) onH264NV12Fn func(destX, destY, w, h int, y []byte, yStride int, uv []byte, uvStride int) onDecoderBrokenFn func() // clipboard callbacks and handler onClipboardFn func(text string) // remote → local getClipboardFn func() string // local → remote onClipboardImageFn func(png []byte) // remote → local(PNG) onClipboardHTMLFn func(string) // remote → local(HTML Format) getClipboardHTMLFn func() string // local → remote(HTML Format,空串=无) getClipboardImgFn func() []byte // local → remote(PNG,无图返回 nil) // 文件剪贴板(CF_HDROP + FileContentsRequest/Response,MS-RDPECLIP §3.1.5.4) onClipboardFilesFn func(names []string) // remote → local:远端剪贴板含文件 onClipboardFileDataFn func(index int, name string, data []byte) // 单文件下载完成/失败 onFileProgressFn func(index int, received, total int64) // 下载进度(每段一次) // certVerifierFn:TLS 服务器证书 SHA-256 指纹校验(TOFU),nil = 不校验 certVerifierFn func(sha256Fp []byte) error // 位图缓存(stage6 6.4b M1 会话内缓存):CacheBitmapV2 存入, // MemBlt 按键引用回贴。键 = cacheId<<16 | cacheIndex。 bitmapCache map[uint32]*Bitmap bitmapCacheFIFO []uint32 bmpCacheStores atomic.Uint64 bmpCacheHits atomic.Uint64 // gfxCacheStore:GFX/bitmap 两条管线共用的持久缓存桥(6.4/6.4b)。 gfxCacheStore rdpgfx.GfxCacheStore // 连接参数(mstsc 基本选项): // audioMode: "local"(默认,本机播放)| "none"(不播放)| "remote"(远端播放) audioMode string // colorDepth: 0/32 = 默认(WANT_32BPP_SESSION),16/24 生效于传统位图管线 colorDepth int // perfFlags/perfFlagsSet: Client Info performanceFlags;未设置保持库默认 perfFlags uint32 perfFlagsSet bool // rdpsndHandler 在 doLogin 内创建(audioMode != "remote" 时), // 供同连接周期的 AUDIO_PLAYBACK DVC 适配器复用。 rdpsndHandler *rdpsnd.Handler cliprdrHandler *cliprdr.CliprdrHandler // 驱动器重定向(MS-RDPEFS):SetDriveRedirect(true) 后 doLogin 注册 // 完整 rdpdr 处理器(替换 stub),设备内容由异步 Filesystem 桥提供。 driveRedirect bool driveLabel string driveFS rdpdr.Filesystem // Login 前暂存,处理器创建时补挂 rdpdrHandler *rdpdr.Handler // reconnectMu serialises concurrent Reconnect() calls. // reconnecting is also set during async server redirects to suppress // user-facing callbacks while the transport is being re-established. reconnectMu sync.Mutex reconnecting atomic.Bool // mouse and wheel hold all coalescing state for pointer input. mouse mouseCoalescer wheel wheelCoalescer // gfxHandler is the active RDPGFX handler; nil when not connected. // Stored here so closeTransport() can stop its goroutines. gfxHandler *rdpgfx.GfxHandler // avc444Disabled, when true, limits CAPS_ADVERTISE to v8.1 so the server // uses AVC420 only. Set via DisableAVC444() before Connect(); preserved // across reconnects. avc444Disabled bool // gfxNoAVC, when true, advertises v10.x caps with AVC_DISABLED so the // server keeps using the RDPGFX channel with ClearCodec/RFX Progressive // instead of H.264. Set via NoAVC() before Connect(). gfxNoAVC bool // pduRecorderFn backs SetPduRecorder: the wasm layer may arm recording // before the client (and its GfxHandler) exists; the callback is attached // at GfxHandler creation. pduRecorderFn func(kind byte, codecId uint16, sw, sh, x, y, w, h uint32, payload []byte) // rejectGFX, when true, rejects the RDPGFX dynvc channel entirely so the // server falls back to legacy bitmap updates (compat mode). rejectGFX bool // dispHandler is the active MS-RDPEDISP handler; nil when not connected. // Used by SetResolution to send MONITOR_LAYOUT PDUs. dispHandler *rdpedisp.Handler dialer func(hostPort string) (net.Conn, error) } const mouseCoalesceInterval = 16 * time.Millisecond // Bitmap is a single rendered region delivered to the OnBitmap callback. // // Lifecycle: Bitmap.Data is borrowed from an internal buffer pool and is // only valid for the duration of the synchronous OnBitmap callback. After // the callback returns the slice may be returned to the pool and overwritten // by subsequent updates. Callers that need to retain the pixels (e.g. to // hand them to an asynchronous paint goroutine) MUST copy the bytes before // the callback returns. type Bitmap struct { DestLeft int DestTop int DestRight int DestBottom int Width int Height int BitsPerPixel int Data []byte } // FillRGBA converts the bitmap's pixel data to RGBA format, writing into dst. // If dst is nil or has the wrong dimensions a new *image.RGBA is allocated. // Callers that process tiles of stable dimensions can reuse the same *image.RGBA // across frames to avoid repeated heap allocations: // // var tile *image.RGBA // tile = bm.FillRGBA(tile) func (bm *Bitmap) FillRGBA(dst *image.RGBA) *image.RGBA { if dst == nil || dst.Bounds().Dx() != bm.Width || dst.Bounds().Dy() != bm.Height { dst = image.NewRGBA(image.Rect(0, 0, bm.Width, bm.Height)) } pix := dst.Pix data := bm.Data // Per-format specialised loops avoid a per-pixel switch and let the // compiler hoist bounds checks and emit tight, branch-free inner code. switch bm.BitsPerPixel { case 1: // 16-bit RGB555 stored big-endian in two bytes. n := len(pix) >> 2 if len(data) < n*2 { n = len(data) / 2 } rgb555BatchToRGBA(pix, data, n) case 2: // 16-bit RGB565 stored big-endian in two bytes. n := len(pix) >> 2 if len(data) < n*2 { n = len(data) / 2 } rgb565BatchToRGBA(pix, data, n) default: // 24/32-bit BGR(A) → RGBA with stride = bm.BitsPerPixel. stride := bm.BitsPerPixel n := len(pix) >> 2 if len(data) < n*stride { n = len(data) / stride } if stride == 4 { // BGRA32 is the common case; use the SIMD-accelerated path. bgr32BatchToRGBA(pix, data, n) } else { // BGR24 (stride==3) and any other depth: scalar fallback. // Write each pixel as a single 32-bit store to let the compiler // vectorise the loop (avoids 4 separate byte stores per pixel). for i := range n { s := i * stride *(*uint32)(unsafe.Pointer(&pix[i*4])) = uint32(data[s+2]) | uint32(data[s+1])<<8 | uint32(data[s])<<16 | 0xFF000000 } } } return dst } // RGBA converts the bitmap pixel data to an *image.RGBA. // A new *image.RGBA is allocated on each call. If the caller processes tiles // of the same dimensions across frames, prefer FillRGBA to avoid allocations. func (bm *Bitmap) RGBA() *image.RGBA { return bm.FillRGBA(nil) } // SwapRB swaps the red and blue byte of every 32-bit pixel in p in place and // forces the alpha byte to 0xFF, converting between RGBA and BGRA byte order // (the operation is its own inverse). It reuses the SIMD-accelerated batch // converter (SSE2 on amd64, NEON on arm64), so callers that need BGRA pixels // for an SDL_PIXELFORMAT_BGRA32 texture can convert an *image.RGBA buffer with // a single vectorised pass instead of a scalar per-byte swap loop. len(p) must // be a multiple of 4; any trailing bytes are ignored. func SwapRB(p []byte) { bgr32BatchToRGBA(p, p, len(p)/4) } // FillBGRA converts the bitmap pixel data into packed BGRA32 (4 bytes/pixel, // B at offset 0) and writes it into dst, growing dst if necessary. The caller // may pass a previously returned slice to reuse the allocation across calls // (common for tiled rendering where many same-sized bitmaps are converted in a // loop). The returned slice has length == bm.Width * bm.Height * 4. // // Unlike RGBA() + SwapRB, FillBGRA avoids an intermediate *image.RGBA // allocation and the second-pass R/B swap for BGR24/BGRA32 inputs. For // BGR24 the output is a single-pass direct pack; for BGRA32 it is a memcopy. func (bm *Bitmap) FillBGRA(dst []byte) []byte { n := bm.Width * bm.Height need := n * 4 if cap(dst) < need { dst = make([]byte, need) } dst = dst[:need] data := bm.Data switch bm.BitsPerPixel { case 2: // 16-bit RGB565: convert to RGBA then swap R↔B in-place → BGRA. if len(data) < n*2 { n = len(data) / 2 } rgb565BatchToRGBA(dst, data, n) bgr32BatchToRGBA(dst, dst, n) default: stride := bm.BitsPerPixel if stride < 3 { stride = 4 // treat unknown as BGRA32 } if stride == 4 { // BGRA32 source — already in the target format; bulk copy. if len(data) < n*4 { n = len(data) / 4 } copy(dst[:n*4], data[:n*4]) } else { // BGR24 source — pack B,G,R directly with A=0xFF; no R/B swap. if len(data) < n*3 { n = len(data) / 3 } for i := range n { s := i * 3 *(*uint32)(unsafe.Pointer(&dst[i*4])) = uint32(data[s]) | uint32(data[s+1])<<8 | uint32(data[s+2])<<16 | 0xFF000000 } } } return dst } func NewRdpClient(host string, width, height int, dialer func(string) (net.Conn, error)) *RdpClient { g := &RdpClient{ hostPort: host, width: width, height: height, kbdLayout: uint32(gcc.US), keyboardType: uint32(gcc.KT_IBM_101_102_KEYS), keyboardSubType: 0, timezoneName: "UTC", dialer: dialer, decompressPool: sync.Pool{ New: func() any { return []uint8(nil) }, }, flipLinePool: sync.Pool{ New: func() any { return []uint8(nil) }, }, } // Point the cached single-element slices at the cached PDU fields so // sendMouseMoveLocked / sendWheelLocked need no per-call allocations. g.mouse.pduBuf[0] = &g.mouse.pdu g.wheel.pduBuf[0] = &g.wheel.pdu return g } var keyboardLayoutMap = map[string]uint32{ "ARABIC": uint32(gcc.ARABIC), "BULGARIAN": uint32(gcc.BULGARIAN), "CHINESE_US_KEYBOARD": uint32(gcc.CHINESE_US_KEYBOARD), "CZECH": uint32(gcc.CZECH), "DANISH": uint32(gcc.DANISH), "GERMAN": uint32(gcc.GERMAN), "GREEK": uint32(gcc.GREEK), "US": uint32(gcc.US), "SPANISH": uint32(gcc.SPANISH), "FINNISH": uint32(gcc.FINNISH), "FRENCH": uint32(gcc.FRENCH), "HEBREW": uint32(gcc.HEBREW), "HUNGARIAN": uint32(gcc.HUNGARIAN), "ICELANDIC": uint32(gcc.ICELANDIC), "ITALIAN": uint32(gcc.ITALIAN), "JAPANESE": uint32(gcc.JAPANESE), "KOREAN": uint32(gcc.KOREAN), "DUTCH": uint32(gcc.DUTCH), "NORWEGIAN": uint32(gcc.NORWEGIAN), } var keyboardTypeMap = map[string]uint32{ "IBM_PC_XT_83_KEY": uint32(gcc.KT_IBM_PC_XT_83_KEY), "OLIVETTI": uint32(gcc.KT_OLIVETTI), "IBM_PC_AT_84_KEY": uint32(gcc.KT_IBM_PC_AT_84_KEY), "IBM_101_102_KEYS": uint32(gcc.KT_IBM_101_102_KEYS), "NOKIA_1050": uint32(gcc.KT_NOKIA_1050), "NOKIA_9140": uint32(gcc.KT_NOKIA_9140), "JAPANESE": uint32(gcc.KT_JAPANESE), } // SetKeyboardLayout sets the keyboard layout by name (e.g. "US", "FRENCH"). // Must be called before Login. func (g *RdpClient) SetKeyboardLayout(layout string) { if v, ok := keyboardLayoutMap[strings.ToUpper(layout)]; ok { g.kbdLayout = v } else { slog.Warn("Unknown keyboard layout, falling back to US", "layout", layout) g.kbdLayout = uint32(gcc.US) } } // SetKeyboardType sets the keyboard type by name (e.g. "IBM_101_102_KEYS"). // Must be called before Login. func (g *RdpClient) SetKeyboardType(keyboardType string) { if v, ok := keyboardTypeMap[strings.ToUpper(keyboardType)]; ok { g.keyboardType = v } else { slog.Warn("Unknown keyboard type, falling back to IBM_101_102_KEYS", "keyboardType", keyboardType) g.keyboardType = uint32(gcc.KT_IBM_101_102_KEYS) } } // SetTimezone sets the client timezone reported in the Client Info PDU. // name is a Windows timezone registry key name (e.g. "UTC", "China Standard // Time"); biasMinutes = UTC minus local time in minutes (e.g. -480 for UTC+8). // Must be called before Login. func (g *RdpClient) SetTimezone(name string, biasMinutes int) { g.timezoneName = name g.timezoneBias = biasMinutes } // DisableAVC444 prevents the client from advertising AVC444/AVC444v2 support. // When called before Login, the RDPGFX CAPS_ADVERTISE is limited to v8.1 // (AVC420 only), so the server will never send LC=2 chroma-upgrade frames. // This avoids the colour distortion seen with VirtualBox VRDE, which sends // LC=2 data but does not include stream2 in LC=0 IDR packets. // The setting is preserved across automatic reconnects. func (g *RdpClient) DisableAVC444() *RdpClient { g.avc444Disabled = true return g } // RejectGFXChannel rejects the RDPGFX dynamic channel entirely, forcing the // server to fall back to legacy bitmap updates. Useful against servers whose // graphics pipeline misbehaves. Must be called before Login. func (g *RdpClient) RejectGFXChannel() { g.rejectGFX = true } // NoAVC keeps the RDPGFX channel alive while forbidding H.264: CAPS_ADVERTISE // includes v10.x sets with RDPGFX_CAPS_FLAG_AVC_DISABLED, so the server encodes // with ClearCodec + RFX Progressive ("RemoteFX mode"). Must be called before // Login; preserved across automatic reconnects. func (g *RdpClient) NoAVC() *RdpClient { g.gfxNoAVC = true return g } // SetGfxCacheStore installs the persistent bitmap cache store (MS-RDPEGFX // CacheImport): SurfaceToCache entries go to the store, and each connection // offers persisted entries to the server after the caps exchange. Must be // called after client creation, before the GFX channel finishes negotiation. func (g *RdpClient) SetGfxCacheStore(s rdpgfx.GfxCacheStore) { g.gfxCacheStore = s if h := g.gfxHandler; h != nil { h.SetPersistentCacheStore(s) } } // SetPduRecorder installs a callback receiving every wire-to-surface bitmap // payload for offline replay analysis (garbled-screen debugging harness). func (g *RdpClient) SetPduRecorder(fn func(kind byte, codecId uint16, surfW, surfH, x, y, w, h uint32, payload []byte)) { g.pduRecorderFn = fn if h := g.gfxHandler; h != nil { h.SetPduRecorder(fn) } } func bpp(BitsPerPixel uint16) int { switch BitsPerPixel { case 15, 16: return 2 case 24: return 3 case 32: return 4 default: slog.Error("invalid bitmap data format", "BitsPerPixel", BitsPerPixel) return 0 } } // mouseButtonFlag returns the PTRFLAGS constant for button index 0/1/2. func mouseButtonFlag(button int) uint16 { switch button { case 0: return pdu.PTRFLAGS_BUTTON1 case 2: return pdu.PTRFLAGS_BUTTON2 case 1: return pdu.PTRFLAGS_BUTTON3 default: return pdu.PTRFLAGS_MOVE } } func (g *RdpClient) Login(domain string, user string, password string) error { slog.Debug("Login", "Host", g.hostPort, "domain", domain, "user", user) g.domain = domain g.user = user g.password = password return g.doLogin(nil) } // doLogin establishes an RDP connection. // When routingToken is non-nil it replaces the username cookie in the // x224 Connection Request (required for Server Redirection). func (g *RdpClient) doLogin(routingToken []byte) error { conn, err := g.dialer(g.hostPort) if err != nil { return fmt.Errorf("[dial err] %v", err) } host, _, _ := net.SplitHostPort(g.hostPort) socketLayer := core.NewSocketLayer(conn, host) socketLayer.SetCertVerifier(g.certVerifierFn) g.tpkt = tpkt.New(socketLayer, nla.NewNTLMv2(g.domain, g.user, g.password)) g.x224 = x224.New(g.tpkt) g.mcs = t125.NewMCSClient(g.x224, g.kbdLayout, g.keyboardType, g.keyboardSubType) g.sec = sec.NewClient(g.mcs) if g.perfFlagsSet { g.sec.SetPerformanceFlags(g.perfFlags) } // 声音模式的协议层声明(与通道注册行为互补,对齐 mstsc) switch g.audioMode { case "none": g.sec.SetNoAudioPlayback() case "remote": g.sec.SetRemoteConsoleAudio() } g.pdu = pdu.NewClient(g.sec) g.channels = plugin.NewChannels(g.sec) // Wire user-registered callbacks now that g.pdu is initialised. // This allows callers to invoke On* methods before Login. g.reregisterCallbacks() // Wire RemoteFX surface decoder so the pdu layer can decode // codecID=3 in surface bitmap commands without importing rdpgfx. pdu.DecodeRemoteFX = rdpgfx.DecodeSurfaceRFX g.mcs.SetClientDesktop(uint16(g.width), uint16(g.height)) // 客户端名随机化(RDPDR-2 假设实验):M1 验收成功时 ClientName 为 wasm // os.Hostname 的固定值 "js";引入随机名后所有会话的 \\tsclient\<共享名> // 打开均报"试图访问无效的地址"且 rdpdr 零 IRP。ClientName 是成功/失败 // 之间唯一的客户端侧系统差异——临时禁用随机化以隔离变量。 if os.Getenv("WEBRDP_RANDOM_CLIENT_NAME") == "1" { cn := fmt.Sprintf("web%06x", uint64(time.Now().UnixNano())&0xffffff) g.mcs.SetClientName(cn) slog.Info("client name", "name", cn) } if g.colorDepth != 0 { g.mcs.SetSessionColorDepth(g.colorDepth) } // Register channels in order: rdpdr, rdpsnd, cliprdr, drdynvc // (matching the channel order that Windows servers expect) // rdpdr (Device Redirection) — stub, required for server to enable audio。 // 启用驱动器重定向时注册完整处理器(宣告文件系统设备并处理 IO), // 否则保持 stub。 if g.driveRedirect { label := g.driveLabel if label == "" { label = "local" } // DosName 与 DeviceData 须同名(服务端建 \\tsclient\<共享名> UNC // 映射的键,FreeRDP 同款);label 同时作卷标。 // (2026-09-13 凌晨实验结论:DosName 固定 "C" 与 label 同名行为 // 一致——rdpdr 通道同样在验证序列后被服务端关闭,DosName 已彻底 // 排除,见 doc/RDPDR-2.md 与 doc/history/stage6-plan.md。) g.rdpdrHandler = rdpdr.NewHandler(label, label) if g.driveFS != nil { g.rdpdrHandler.SetFilesystem(g.driveFS) } g.channels.Register(g.rdpdrHandler) } else { g.channels.Register(&stubChannel{name: "rdpdr", option: plugin.CHANNEL_OPTION_INITIALIZED | plugin.CHANNEL_OPTION_ENCRYPT_RDP | plugin.CHANNEL_OPTION_COMPRESS_RDP}) } g.mcs.SetClientDeviceRedirection() // RDPSND (Audio Output) handler — static virtual channel + DVC paths // 声音重定向三模式(对应 mstsc 远程音频播放): // remote(远端播放):不注册 rdpsnd / AUDIO_PLAYBACK 通道, // 服务器音频走本机扬声器; // none(不播放):照常注册协商,但 wave 数据丢弃——服务器认为 // 音频已被重定向而保持静音; // local(本机播放,默认):注册并播放。 switch g.audioMode { case "remote": // 不注册任何音频通道 default: rdpsndHandler := rdpsnd.NewHandler(func(format rdpsnd.AudioFormat, data []byte) { if g.onAudioFn != nil { g.onAudioFn(format, data) } }) if g.audioMode == "none" { rdpsndHandler.SetMuted(true) } rdpsndHandler.SetAudioResetCallback(func() { if g.onAudioResetFn != nil { g.onAudioResetFn() } }) g.channels.Register(rdpsndHandler) g.mcs.SetClientSoundProtocol() // 音频 DVC 适配器在 doLogin 后半段注册(需要 rdpsndHandler 存活), // 用局部暂存传递;remote 模式下两通道同样不注册。 g.rdpsndHandler = rdpsndHandler } // cliprdr (Clipboard) — cross-platform text clipboard handler cliprdrHandler := cliprdr.NewHandler( func(text string) { if g.onClipboardFn != nil { g.onClipboardFn(text) } }, func() string { if g.getClipboardFn != nil { return g.getClipboardFn() } return "" }, ) cliprdrHandler.SetImageCallbacks( func(png []byte) { if g.onClipboardImageFn != nil { g.onClipboardImageFn(png) } }, func() []byte { if g.getClipboardImgFn != nil { return g.getClipboardImgFn() } return nil }, ) cliprdrHandler.SetHTMLCallbacks( func(html string) { if g.onClipboardHTMLFn != nil { g.onClipboardHTMLFn(html) } }, func() string { if g.getClipboardHTMLFn != nil { return g.getClipboardHTMLFn() } return "" }, ) cliprdrHandler.SetFileCallbacks( func(names []string) { if g.onClipboardFilesFn != nil { g.onClipboardFilesFn(names) } }, func(index int, name string, data []byte) { if g.onClipboardFileDataFn != nil { g.onClipboardFileDataFn(index, name, data) } }, func(index int, received, total int64) { if g.onFileProgressFn != nil { g.onFileProgressFn(index, received, total) } }, ) g.cliprdrHandler = cliprdrHandler g.channels.Register(cliprdrHandler) g.mcs.SetClientClipboard() // drdynvc (Dynamic Virtual Channels) dvcClient := drdynvc.NewDvcClient() g.channels.Register(dvcClient) g.mcs.SetClientDynvcProtocol() // RDPGFX (Graphics Pipeline) handler gfxHandler := rdpgfx.NewGfxHandler(func(updates []rdpgfx.BitmapUpdate) { if g.onBitmapPaintFn == nil { return } bs := make([]Bitmap, len(updates)) for i, u := range updates { bs[i] = Bitmap{ DestLeft: u.DestLeft, DestTop: u.DestTop, DestRight: u.DestRight, DestBottom: u.DestBottom, Width: u.Width, Height: u.Height, BitsPerPixel: u.Bpp, Data: u.Data, } } g.onBitmapPaintFn(bs) }) gfxHandler.SetDecoderBrokenCallback(func() { slog.Debug("H.264 decoder broken") if g.onDecoderBrokenFn != nil { g.onDecoderBrokenFn() } }) gfxHandler.SetKeyframeRequestFunc(func() { slog.Debug("H.264: requesting keyframe via force refresh") if g.pdu != nil { // SendRefreshRect is silently ignored by Windows servers while // an H.264 video stream is active. Use the suppress→allow // toggle (SendForceRefresh) which mstsc/FreeRDP rely on to // reliably trigger a fresh IDR. See protocol/pdu/pdu.go. g.pdu.SendForceRefresh(uint16(g.width), uint16(g.height)) } }) if g.onH264RawFn != nil { gfxHandler.SetH264RawCallback(g.onH264RawFn) } if g.onH264I420Fn != nil { gfxHandler.SetI420Callback(g.onH264I420Fn) } if g.onH264NV12Fn != nil { gfxHandler.SetNV12Callback(g.onH264NV12Fn) } if g.avc444Disabled { gfxHandler.SetAVC444Disabled(true) } if g.gfxNoAVC { gfxHandler.SetAVCDisabled(true) } if g.pduRecorderFn != nil { gfxHandler.SetPduRecorder(g.pduRecorderFn) } g.gfxHandler = gfxHandler // 持久缓存桥:SetGfxCacheStore 在 handler 创建前调用,必须在这里补挂 // (此前 store 在这里丢失,GFX 持久缓存实际从未生效——服务端 0 条 // SurfaceToCache 掩盖了这一点)。 if g.gfxCacheStore != nil && !g.rejectGFX { gfxHandler.SetPersistentCacheStore(g.gfxCacheStore) } // bitmap 管线持久缓存(6.4b M2):把跨会话持有的键注册给 finalize, // 服务器广告 HOST SUPPORT 时经 PERSISTENT_KEY_LIST 上报。 if g.rejectGFX && g.gfxCacheStore != nil { g.pdu.SetPersistentKeyList(g.gfxCacheStore.Keys()) } gfxHandler.SetSessionSize(uint16(g.width), uint16(g.height)) if g.rejectGFX { dvcClient.RegisterRejectedChannel(rdpgfx.ChannelName) } else { dvcClient.RegisterHandler(rdpgfx.ChannelName, gfxHandler) } // RDPEDISP (Display Update Virtual Channel) handler — allows requesting // a resolution change while connected (MS-RDPEDISP). Pass 0,0 so no // initial MONITOR_LAYOUT PDU is sent: some servers' graphics pipeline // fails (ERRINFO_GRAPHICS_SUBSYSTEM_FAILED 0x112F) when the desktop is // resized during GFX surface setup. Resolution changes go through // SetResolution() instead. dispHandler := rdpedisp.NewHandler(0, 0) g.dispHandler = dispHandler dvcClient.RegisterHandler(rdpedisp.ChannelName, dispHandler) // Reject Video Optimized Remoting (VOR) channels so the server keeps // sending video through the RDPGFX pipeline which we do handle. // Without this, the server detects video playback (e.g. YouTube) and // switches to VOR channels that we don't implement, causing the video // to freeze while audio continues. dvcClient.RegisterRejectedChannel("Microsoft::Windows::RDS::Video::Control::v08.01") dvcClient.RegisterRejectedChannel("Microsoft::Windows::RDS::Video::Data::v08.01") dvcClient.RegisterRejectedChannel("Microsoft::Windows::RDS::Geometry::v08.01") // Register DVC audio handlers for both the lossless and lossy variants. // gnome-remote-desktop requests AUDIO_PLAYBACK_LOSSY_DVC first; if it is // rejected, gnome-remote-desktop triggers its SVC fallback path which also // sets prevent_dvc_initialization=true, silently blocking AUDIO_PLAYBACK_DVC // as well — leaving the client with no audio at all. // By accepting both channels with the same rdpsnd handler, format negotiation // (which only advertises PCM) ensures PCM is used regardless of which channel // gnome-remote-desktop chooses. // remote(远端播放)模式不注册,让服务器走本机音频。 if g.rdpsndHandler != nil { dvcClient.RegisterHandler("AUDIO_PLAYBACK_DVC", rdpsnd.NewDvcAdapter(g.rdpsndHandler)) dvcClient.RegisterHandler("AUDIO_PLAYBACK_LOSSY_DVC", rdpsnd.NewDvcAdapter(g.rdpsndHandler)) } g.sec.SetUser(g.user) g.sec.SetPwd(g.password) g.sec.SetDomain(g.domain) // 时区:dynamic DST 键名为空会导致服务器 0x112F 断连;默认发 UTC g.sec.SetClientTimezone(g.timezoneName, g.timezoneBias) g.tpkt.SetFastPathListener(g.sec) g.sec.SetFastPathListener(g.pdu) g.sec.SetChannelSender(g.mcs) g.channels.SetChannelSender(g.sec) // Wire fast-path output: pdu → sec → tpkt. This enables the much // shorter Fast-Path Client Input PDU framing for mouse/keyboard events // (MS-RDPBCGR §2.2.8.1.2). Use is gated at runtime both by capability // negotiation in the PDU layer and by sec.SendFastPath itself, which // refuses when legacy RDP encryption is in effect. g.sec.SetFastPathSender(g.tpkt) g.pdu.SetFastPathSender(g.sec) g.x224.SetRequestedProtocol(x224.PROTOCOL_SSL | x224.PROTOCOL_HYBRID) if routingToken != nil { g.x224.SetRoutingToken(routingToken) } else { g.x224.SetUsername(g.user) } err = g.x224.Connect() if err != nil { return fmt.Errorf("[x224 connect err] %v", err) } // Wait for the RDP handshake to complete or fail. // Events arrive asynchronously from the TPKT read goroutine. type connResult struct { err error redirect *pdu.ServerRedirectionPDU } ch := make(chan connResult, 4) send := func(r connResult) { select { case ch <- r: default: } } // readyFired is set by the "ready" callback. All emitter callbacks // run synchronously on the TPKT read goroutine, so no mutex needed. readyFired := false g.pdu.On("ready", func() { g.eventReady.Store(true) readyFired = true send(connResult{}) }) g.pdu.On("error", func(err error) { if !readyFired { send(connResult{err: err}) } else { // Mid-session error: stop accepting input so we don't // try to write to the now-dead transport. g.eventReady.Store(false) } }) // Redirect may arrive before or after "ready". // Before ready: send to channel for synchronous handling. // After ready: launch async goroutine (GNOME Remote Desktop // sends redirect ~5s after the GFX retry's "ready"). g.pdu.Once("redirect", func(redir *pdu.ServerRedirectionPDU) { if !readyFired { send(connResult{redirect: redir}) } else { go g.handleRedirect(redir) } }) // DeactivateAllPDU during an active session means the server is // reactivating (e.g. desktop resize). Pause input until "ready" // fires again after the reactivation handshake completes. g.pdu.On("deactivateAll", func() { g.eventReady.Store(false) }) select { case r := <-ch: if r.err != nil { g.tpkt.Close() return fmt.Errorf("[connection err] %v", r.err) } if r.redirect != nil { slog.Debug("Server redirect", "loadBalanceInfo", string(r.redirect.LoadBalanceInfo)) g.tpkt.Close() g.eventReady.Store(false) return g.doLogin(r.redirect.LoadBalanceInfo) } // "ready" received — session established. return nil case <-time.After(30 * time.Second): g.tpkt.Close() return fmt.Errorf("[connection timeout]") } } // handleRedirect handles a Server Redirection PDU that arrives after // "ready" (e.g. GNOME Remote Desktop). Runs asynchronously. func (g *RdpClient) handleRedirect(redir *pdu.ServerRedirectionPDU) { slog.Debug("Async server redirect", "loadBalanceInfo", string(redir.LoadBalanceInfo)) g.reconnecting.Store(true) g.tpkt.Close() g.eventReady.Store(false) err := g.doLogin(redir.LoadBalanceInfo) g.reconnecting.Store(false) if err != nil { slog.Error("handleRedirect: login failed", "err", err) if g.onErrorFn != nil { g.onErrorFn(err) } return } g.reregisterCallbacks() } func (g *RdpClient) Width() int { return g.width } func (g *RdpClient) Height() int { return g.height } func (g *RdpClient) OnError(f func(e error)) *RdpClient { g.onErrorFn = f if g.pdu != nil { g.pdu.On("error", func(e error) { if !g.reconnecting.Load() { f(e) } }) } return g } func (g *RdpClient) OnClose(f func()) *RdpClient { g.onCloseFn = f if g.pdu != nil { g.pdu.On("close", func() { if !g.reconnecting.Load() { f() } }) } return g } func (g *RdpClient) OnSuccess(f func()) *RdpClient { g.onSuccessFn = f if g.sec != nil { g.sec.On("success", f) } return g } func (g *RdpClient) OnReady(f func()) *RdpClient { g.onReadyFn = f if g.pdu != nil { g.pdu.On("ready", f) } return g } // OnBitmap registers a callback for bitmap update events. // For compressed bitmaps, Bitmap.Data is borrowed from an internal pool and // is valid only for the duration of the paint call. If you need to retain // the raw pixel data beyond paint, copy it or call bm.RGBA() inside paint. func (g *RdpClient) OnBitmap(paint func([]Bitmap)) *RdpClient { g.onBitmapPaintFn = paint if g.pdu == nil { return g } g.pdu.On("bitmap", func(rectangles []pdu.BitmapData) { // 16/24bpp 位图路径曾发生 panic 直接杀死整个 wasm 程序; // recover 保证连接存活并留下完整堆栈用于定位。 defer func() { if r := recover(); r != nil { slog.Error("bitmap update panic", "err", r, "stack", string(debug.Stack())) } }() bs := make([]Bitmap, 0, len(rectangles)) var pooled [][]uint8 // track buffers borrowed from pool for idx, v := range rectangles { data := v.BitmapDataStream wireBpp := v.BitsPerPixel if wireBpp == 0 { // Win10 在低色深会话(经 postBeta2ColorDepth 协商)里把 // bitsPerPixel 置 0:按会话色深处理 wireBpp = uint16(g.colorDepth) if wireBpp == 0 { wireBpp = 32 } } Bpp := bpp(wireBpp) if Bpp == 0 { slog.Error("bitmap rect with invalid bpp", "idx", idx, "count", len(rectangles), "rect", fmt.Sprintf("%+v", v), "allRects", fmt.Sprintf("%+v", rectangles)) continue } if v.Flags&pdu.BITMAP_NO_PROCESSING != 0 { // Surface command: data is already decoded top-down BGRA } else if v.IsCompress() { buf := g.decompressPool.Get().([]uint8) var ok bool buf, ok = core.DecompressInto(v.BitmapDataStream, buf, int(v.Width), int(v.Height), Bpp) if !ok { // 解码失败的矩形绝不能上屏:缓冲里是半解码+池化残留, // 画出来就是噪块。跳过并等服务器后续更新修复该区域。 g.decompressPool.Put(buf) slog.Warn("skip undecodable bitmap rect", "dx", v.DestLeft, "dy", v.DestTop, "dr", v.DestRight, "db", v.DestBottom, "w", v.Width, "h", v.Height, "bpp", Bpp, "flags", fmt.Sprintf("0x%04X", v.Flags), "len", len(v.BitmapDataStream)) continue } data = buf pooled = append(pooled, buf) } else { // Uncompressed bitmaps are bottom-up; flip to top-down. stride := int(v.Width) * Bpp h := int(v.Height) tmp := g.flipLinePool.Get().([]byte) if cap(tmp) < stride { tmp = make([]byte, stride) } else { tmp = tmp[:stride] } for y := 0; y < h/2; y++ { top := y * stride bot := (h - 1 - y) * stride copy(tmp, data[top:top+stride]) copy(data[top:top+stride], data[bot:bot+stride]) copy(data[bot:bot+stride], tmp) } g.flipLinePool.Put(tmp[:cap(tmp)]) } b := Bitmap{int(v.DestLeft), int(v.DestTop), int(v.DestRight), int(v.DestBottom), int(v.Width), int(v.Height), Bpp, data} bs = append(bs, b) } if g.bmpEmitLogged < 10 { g.bmpEmitLogged++ slog.Warn("BITMAP_EMIT", "rects", len(bs)) } paint(bs) for _, buf := range pooled { g.decompressPool.Put(buf[:cap(buf)]) } }) // 位图缓存(stage6 6.4b M1):CacheBitmapV2 存入,MemBlt 引用回贴。 // 订单与位图更新同为 fast-path PDU,按到达顺序同步处理,服务器保证 // MemBlt 先于对应 CacheBitmapV2 的乱序不存在。 g.bitmapCache = make(map[uint32]*Bitmap) ordersSeen := 0 g.pdu.On("orders", func(orderPdus []pdu.OrderPdu) { ordersSeen++ if ordersSeen <= 5 { kinds := make(map[string]int) for i := range orderPdus { o := &orderPdus[i] switch { case o.CacheBitmapV2 != nil: kinds["CacheBitmapV2"]++ case o.Primary != nil && o.Primary.Data != nil: kinds[fmt.Sprintf("%T", o.Primary.Data)]++ default: kinds["empty"]++ } } slog.Warn("ORDERS event", "n", ordersSeen, "pdus", len(orderPdus), "kinds", fmt.Sprintf("%v", kinds)) } defer func() { if r := recover(); r != nil { slog.Error("orders update panic", "err", r, "stack", string(debug.Stack())) } }() for i := range orderPdus { o := &orderPdus[i] if cb := o.CacheBitmapV2; cb != nil { g.storeBitmapCacheV2(cb) continue } if o.Primary != nil { if mb, ok := o.Primary.Data.(*pdu.Memblt); ok { g.drawMemblt(mb) } } } }) return g } func (g *RdpClient) OnPointerHide(f func()) *RdpClient { g.onPointerHideFn = f if g.pdu != nil { g.pdu.On("pointer_hide", f) } return g } func (g *RdpClient) OnPointerCached(f func(uint16)) *RdpClient { g.onPointerCachedFn = f if g.pdu != nil { g.pdu.On("pointer_cached", f) } return g } func (g *RdpClient) OnPointerDefault(f func()) *RdpClient { g.onPointerDefaultFn = f if g.pdu != nil { g.pdu.On("pointer_default", f) } return g } // bitmapCacheMaxEntries 单元总量上限(超出按 FIFO 逐出,防御服务器 // 引用已逐出条目时缓存无限增长)。 const bitmapCacheMaxEntries = 4096 // storeBitmapCacheV2 解压 CacheBitmapV2 次级订单的位图并按 // cacheId<<16|cacheIndex 存入会话内缓存(DO_NOT_CACHE 条目不存)。 func (g *RdpClient) storeBitmapCacheV2(cb *pdu.CacheBitmapV2Order) { if cb.CacheIndex == 0x7FFF || cb.BitmapWidth == 0 || cb.BitmapHeight == 0 { return } // 持久键(6.4b M2):服务器在 CacheBitmapV2 上携带跨会话键。 // 带键且零数据 = 服务器指示客户端从持久库回填该单元。 persistentKey := uint64(0) if cb.Flags&pdu.CBR2_PERSISTENT_KEY_PRESENT != 0 { persistentKey = uint64(cb.Key2)<<32 | uint64(cb.Key1) } Bpp := bpp(uint16(cb.BitmapBpp)) if Bpp == 0 { return } w, h := int(cb.BitmapWidth), int(cb.BitmapHeight) if cb.BitmapLength == 0 && persistentKey != 0 { if g.gfxCacheStore == nil { return } e, ok := g.gfxCacheStore.Get(persistentKey) if !ok || e.Width != w || e.Height != h || e.Bpp != uint16(Bpp) || len(e.Data) != w*h*int(Bpp) { slog.Debug("bmpcache: persistent backfill miss", "key", persistentKey, "cacheId", cb.CacheId, "idx", cb.CacheIndex, "w", w, "h", h) return } entry := &Bitmap{Width: w, Height: h, BitsPerPixel: Bpp, Data: e.Data} g.bitmapCachePut(uint32(cb.CacheId)<<16|uint32(cb.CacheIndex), entry) return } stride := w * Bpp var data []byte var ok bool if cb.Compressed { buf := g.decompressPool.Get().([]byte) var out []byte out, ok = core.DecompressInto(cb.BitmapDataStream, buf, w, h, Bpp) if !ok { g.decompressPool.Put(buf) slog.Debug("bmpcache: skip undecodable cache bitmap", "cacheId", cb.CacheId, "idx", cb.CacheIndex, "w", w, "h", h, "bpp", Bpp, "len", len(cb.BitmapDataStream)) return } data = out } else { data = append([]byte(nil), cb.BitmapDataStream...) } // RLE 位图自底向上,翻转为自顶向下(与 bitmap 更新路径一致) mirrorRows(data, stride, h) entry := &Bitmap{Width: w, Height: h, BitsPerPixel: Bpp, Data: data} g.bitmapCachePut(uint32(cb.CacheId)<<16|uint32(cb.CacheIndex), entry) // 带持久键的条目跨会话持久化(像素副本,键 = key2<<32|key1) if persistentKey != 0 && g.gfxCacheStore != nil { g.gfxCacheStore.Persist(persistentKey, w, h, uint16(Bpp), data) } stores := g.bmpCacheStores.Add(1) if stores == 1 || stores%500 == 0 { slog.Info("bmpcache: stored", "stores", stores, "total", len(g.bitmapCache), "cacheId", cb.CacheId, "idx", cb.CacheIndex, "w", w, "h", h, "bpp", Bpp) } } // bitmapCachePut 写入会话内缓存单元(cacheId<<16|cacheIndex → 位图), // 超过 bitmapCacheMaxEntries 按插入序 FIFO 逐出。 func (g *RdpClient) bitmapCachePut(key uint32, entry *Bitmap) { if _, exists := g.bitmapCache[key]; !exists { g.bitmapCacheFIFO = append(g.bitmapCacheFIFO, key) if len(g.bitmapCacheFIFO) > bitmapCacheMaxEntries { evict := g.bitmapCacheFIFO[0] g.bitmapCacheFIFO = g.bitmapCacheFIFO[1:] delete(g.bitmapCache, evict) } } g.bitmapCache[key] = entry } // drawMemblt 处理 MemBlt 主订单:从缓存取位图按目标坐标回贴。 // 缓存缺失时跳过(保持既有画面,等服务器后续更新修复该区域)。 func (g *RdpClient) drawMemblt(mb *pdu.Memblt) { if g.onBitmapPaintFn == nil { return } key := uint32(mb.CacheId)<<16 | uint32(mb.CacheIdx) b, ok := g.bitmapCache[key] if !ok || b == nil { slog.Debug("bmpcache: memblt cache miss", "cacheId", mb.CacheId, "idx", mb.CacheIdx) return } out := *b out.DestLeft, out.DestTop = int(mb.X), int(mb.Y) out.DestRight, out.DestBottom = int(mb.X)+int(mb.Cx), int(mb.Y)+int(mb.Cy) hits := g.bmpCacheHits.Add(1) if hits == 1 || hits%500 == 0 { slog.Info("bmpcache: hit", "hits", hits, "cacheId", mb.CacheId, "idx", mb.CacheIdx, "x", mb.X, "y", mb.Y, "w", mb.Cx, "h", mb.Cy) } g.onBitmapPaintFn([]Bitmap{out}) } // mirrorRows 将 stride 对齐的行序图像原地垂直镜像(bottom-up ↔ top-down)。 func mirrorRows(buf []byte, stride, height int) { tmp := make([]byte, stride) for y := 0; y < height/2; y++ { top := y * stride bot := (height - 1 - y) * stride if top+stride <= len(buf) && bot+stride <= len(buf) { copy(tmp, buf[top:top+stride]) copy(buf[top:top+stride], buf[bot:bot+stride]) copy(buf[bot:bot+stride], tmp) } } } func (g *RdpClient) OnPointerUpdate(f func(uint16, uint16, uint16, uint16, uint16, uint16, []byte, []byte)) *RdpClient { g.onPointerUpdateFn = f if g.pdu != nil { g.pdu.On("pointer_update", func(p *pdu.FastPathUpdatePointerPDU) { w := int(p.Width) h := int(p.Height) // xorBpp 由线路直接携带(TS_POINTER_NEW 首字段,FreeRDP // update_read_pointer_new 校验 1≤xorBpp≤32)。 xorBpp := int(p.XorBpp) if xorBpp == 0 { xorBpp = 1 } slog.Debug("OnPointerUpdate", "cacheIdx", p.CacheIdx, "xorBpp", p.XorBpp, "hotX", p.HotX, "hotY", p.HotY, "w", p.Width, "h", p.Height, "andLen", p.MaskLen, "xorLen", p.XorLen) // 掩码行序:按 MS-RDPBCGR/FreeRDP(vFlip = xorBpp != 1),彩色 // 指针(24/32bpp)自底向上存储,翻转为自顶向下;1bpp 单色指针 // 本身自顶向下。stride 均按 2 字节对齐。 flip := xorBpp != 1 xorStride := ((w*xorBpp + 15) / 16) * 2 andStride := ((w + 15) / 16) * 2 var xorData []byte if len(p.Data) > 0 && h > 0 && w > 0 { xorData = make([]byte, len(p.Data)) copy(xorData, p.Data) if flip { mirrorRows(xorData, xorStride, h) } } else { xorData = p.Data } var andMask []byte if len(p.Mask) > 0 && h > 0 && w > 0 { andMask = make([]byte, len(p.Mask)) copy(andMask, p.Mask) if flip { mirrorRows(andMask, andStride, h) } } else { andMask = p.Mask } f(p.CacheIdx, uint16(xorBpp), p.HotX, p.HotY, p.Width, p.Height, andMask, xorData) }) } return g } // OnAudio registers a callback for server audio data. // The callback receives the AudioFormat describing the PCM data and the raw audio bytes. // Must be called before Login. func (g *RdpClient) OnAudio(f func(rdpsnd.AudioFormat, []byte)) *RdpClient { g.onAudioFn = f return g } // OnAudioReset registers a callback that is called when the server closes the // audio channel (e.g. media seek or stream restart). The application should // flush its audio playback buffer so that stale audio does not keep playing. // Must be called before Login. func (g *RdpClient) OnAudioReset(f func()) *RdpClient { g.onAudioResetFn = f return g } // OnH264Raw registers a callback that receives raw H.264 NAL unit data when // the built-in decoder is unavailable (e.g. WASM builds without CGo). // destX, destY are the top-left canvas coordinates; isKey flags an IDR frame. // regions 为扁平 [l,t,r,b,...](帧内坐标,右下开区间)的脏矩形:服务器只 // 保证区域内像素有效,绘制端必须只上屏区域;空切片表示整帧有效。 // The caller owns data and may retain it beyond the callback. func (g *RdpClient) OnH264Raw(fn func(destX, destY, w, h int, isKey bool, data []byte, regions []int32)) *RdpClient { g.onH264RawFn = fn return g } // OnH264I420 registers a callback that receives decoded H.264 frames in I420 // planar format (Y, U, V planes with associated strides). When set, the // decoded frame is NOT delivered via OnBitmap; the caller is responsible for // rendering it directly (e.g. via an SDL2 IYUV texture for GPU-accelerated // YUV→RGB conversion). When I420 extraction is unavailable for a frame // (e.g. non-YUV420P/NV12 formats), grdp falls back to OnBitmap delivery. // destX, destY are top-left canvas coordinates; w, h are frame dimensions. // The plane slices are only valid for the duration of the callback; copy them // if they need to be retained beyond the callback's return. func (g *RdpClient) OnH264I420(fn func(destX, destY, w, h int, y []byte, yStride int, u []byte, uStride int, v []byte, vStride int)) *RdpClient { g.onH264I420Fn = fn if g.gfxHandler != nil { g.gfxHandler.SetI420Callback(fn) } return g } // OnH264NV12 registers a callback that receives decoded H.264 frames in NV12 // format (Y plane plus interleaved UV plane). This is the fastest SDL2 path // on platforms whose hardware decoder already outputs NV12 (notably macOS // VideoToolbox), because callers can upload the planes directly with an NV12 // texture and avoid NV12->I420 deinterleaving in grdp. When NV12 extraction // is unavailable for a frame, grdp falls back to OnBitmap delivery. // destX, destY are top-left canvas coordinates; w, h are frame dimensions. // The plane slices are only valid for the duration of the callback; copy them // if they need to be retained beyond the callback's return. func (g *RdpClient) OnH264NV12(fn func(destX, destY, w, h int, y []byte, yStride int, uv []byte, uvStride int)) *RdpClient { g.onH264NV12Fn = fn if g.gfxHandler != nil { g.gfxHandler.SetNV12Callback(fn) } return g } // OnDecoderBroken registers a callback that is invoked when the H.264 decoder // enters an unrecoverable state (all hard-reset attempts exhausted). When // this callback is set, grdp does NOT automatically call Reconnect; the // application is responsible for deciding when to reconnect (e.g. via its // own stall watchdog). If no callback is registered, grdp falls back to // the previous behaviour of reconnecting immediately. func (g *RdpClient) OnDecoderBroken(f func()) *RdpClient { g.onDecoderBrokenFn = f return g } // OnClipboard registers callbacks for bidirectional clipboard sharing. // // - onRemote is called with the text when the RDP server's clipboard // content is received (server → client). // - getLocal is called to retrieve the current local clipboard text // when the server requests it (client → server). // // Must be called before Login. func (g *RdpClient) OnClipboard(onRemote func(text string), getLocal func() string) *RdpClient { g.onClipboardFn = onRemote g.getClipboardFn = getLocal return g } // OnClipboardImage registers the callback invoked with PNG-encoded bytes when // the RDP server's clipboard image is received (server → client). Must be // called before Login. func (g *RdpClient) OnClipboardImage(onRemoteImage func(png []byte)) *RdpClient { g.onClipboardImageFn = onRemoteImage return g } // SetClipboardImageProvider registers the provider used to answer server // requests for the local clipboard image (client → server). The provider // returns PNG-encoded bytes, or nil when no image is on the local clipboard. // Must be called before Login. func (g *RdpClient) SetClipboardImageProvider(getImage func() []byte) *RdpClient { g.getClipboardImgFn = getImage return g } // OnClipboardHTML registers the callback invoked when the remote clipboard // HTML content (HTML Format, fragment already extracted) arrives. func (g *RdpClient) OnClipboardHTML(onRemoteHTML func(html string)) *RdpClient { g.onClipboardHTMLFn = onRemoteHTML return g } // SetClipboardHTMLProvider registers the provider used to answer server // requests for the local clipboard HTML (HTML Format). The provider returns // raw HTML (no CF_HTML envelope), or "" when no HTML is on the local // clipboard. Must be called before Login. func (g *RdpClient) SetClipboardHTMLProvider(getHTML func() string) *RdpClient { g.getClipboardHTMLFn = getHTML return g } // OnClipboardFiles registers the callback invoked when the remote clipboard // holds files (CF_HDROP path list received, server → client). Files are NOT // fetched automatically; call RequestRemoteFile for each file to download. func (g *RdpClient) OnClipboardFiles(onRemoteFiles func(names []string)) *RdpClient { g.onClipboardFilesFn = onRemoteFiles return g } // OnClipboardFileData registers the callback invoked when a file requested // via RequestRemoteFile has been fully received. data is nil on failure. func (g *RdpClient) OnClipboardFileData(fn func(index int, name string, data []byte)) *RdpClient { g.onClipboardFileDataFn = fn return g } // OnFileTransferProgress registers the callback invoked after every received // chunk of an in-flight RequestRemoteFile transfer. func (g *RdpClient) OnFileTransferProgress(fn func(index int, received, total int64)) *RdpClient { g.onFileProgressFn = fn return g } // SetLocalFiles stages files as the local clipboard file content (client → // server). The server will see a CF_HDROP format and fetch bytes on paste. func (g *RdpClient) SetLocalFiles(files []cliprdr.LocalFile) { if g.cliprdrHandler != nil { g.cliprdrHandler.SetLocalFiles(files) } } // ClearLocalFiles removes locally staged clipboard files. func (g *RdpClient) ClearLocalFiles() { if g.cliprdrHandler != nil { g.cliprdrHandler.ClearLocalFiles() } } // RequestRemoteFile starts downloading file `index` from the remote clipboard // file list; progress and completion arrive via the registered callbacks. func (g *RdpClient) RequestRemoteFile(index int) error { if g.cliprdrHandler == nil { return errors.New("client not connected") } return g.cliprdrHandler.RequestRemoteFile(index) } // SetCertVerifier registers a TOFU (trust-on-first-use) verifier for the // server's TLS certificate. The callback receives the SHA-256 of the leaf // certificate DER; return an error to abort the connection. Must be called // before Login. func (g *RdpClient) SetCertVerifier(fn func(sha256Fp []byte) error) *RdpClient { g.certVerifierFn = fn return g } // NotifyClipboardChanged tells the server that the local clipboard has // changed. The UI should call this when it detects a system clipboard // change (e.g. via polling or a platform clipboard-change signal). func (g *RdpClient) NotifyClipboardChanged() { if g.cliprdrHandler != nil { g.cliprdrHandler.OnLocalClipboardChanged() } } func (g *RdpClient) notifyGfxLocalInput() { if gfx := g.gfxHandler; gfx != nil { gfx.NotifyLocalInput() } } // newScancodeEvent builds a TS_SCANCODE_EVENT from the 0xE0xx convention // used throughout this codebase: the E0 prefix must travel as // KBDFLAGS_EXTENDED with the 8-bit make code in KeyCode — the slow-path // serializer sends KeyCode verbatim, and a raw 0xE0xx value there is an // invalid scancode the server silently drops (observed: Delete did // nothing on the remote). func newScancodeEvent(sc int, release bool) *pdu.ScancodeKeyEvent { p := &pdu.ScancodeKeyEvent{} if sc&0xFF00 == 0xE000 { p.KeyboardFlags |= pdu.KBDFLAGS_EXTENDED sc &= 0xFF } p.KeyCode = uint16(sc) if release { p.KeyboardFlags |= pdu.KBDFLAGS_RELEASE } return p } func (g *RdpClient) KeyUp(sc int) { if !g.eventReady.Load() { return } slog.Debug("KeyUp", "sc", sc) g.flushMouseMove() g.flushWheel() p := newScancodeEvent(sc, true) g.pdu.SendInputEvents(pdu.INPUT_EVENT_SCANCODE, []pdu.InputEventsInterface{p}) g.notifyGfxLocalInput() } // SendUnicodeText sends text as RDP Unicode input events (MS-RDPBCGR // RDP_INPUT_UNICODE, TS_UNICODE_EVENT). Used for IME-committed text and // clipboard paste, which have no meaningful scancode representation. // Each character is sent as a press+release pair; characters outside the // BMP (no UTF-16 code unit mapping) are skipped. func (g *RdpClient) SendUnicodeText(text string) { if !g.eventReady.Load() { return } g.flushMouseMove() g.flushWheel() const maxEventsPerPDU = 14 // fast-path input PDU hard limit is 15 events events := make([]pdu.InputEventsInterface, 0, maxEventsPerPDU) flush := func() { if len(events) == 0 { return } g.pdu.SendInputEvents(pdu.INPUT_EVENT_UNICODE, events) events = events[:0] } for _, r := range text { if r == '\r' || r == '\n' { // CR/LF 走 Unicode 事件会被服务端当不可打印字符丢弃(实测 cmd // 收不到回车,粘贴多行文本/IME 提交带回车的文本无法执行),转成 // Enter 扫描码按下+释放。扫描码必须走 INPUT_EVENT_SCANCODE PDU: // 与 Unicode 事件混在同一 PDU 里服务端会按 TS_UNICODE_EVENT // 解析出控制字符并丢弃(实测 Enter 静默丢失)。 flush() g.pdu.SendInputEvents(pdu.INPUT_EVENT_SCANCODE, []pdu.InputEventsInterface{ newScancodeEvent(0x001C, false), newScancodeEvent(0x001C, true), }) continue } if r > 0xFFFF || (r >= 0xD800 && r <= 0xDFFF) { continue } u := uint16(r) events = append(events, &pdu.UnicodeKeyEvent{Unicode: u}, &pdu.UnicodeKeyEvent{Unicode: u, KeyboardFlags: pdu.KBDFLAGS_RELEASE}, ) if len(events) >= maxEventsPerPDU { flush() } } flush() g.notifyGfxLocalInput() } func (g *RdpClient) KeyDown(sc int) { if !g.eventReady.Load() { return } slog.Debug("KeyDown", "sc", sc) g.flushMouseMove() g.flushWheel() p := newScancodeEvent(sc, false) g.pdu.SendInputEvents(pdu.INPUT_EVENT_SCANCODE, []pdu.InputEventsInterface{p}) g.notifyGfxLocalInput() } // MouseMove queues a mouse-move event. Successive moves within // mouseCoalesceInterval are collapsed: only the latest (x,y) is sent. The // first move in a burst is sent immediately so the server sees no extra // latency for a single isolated motion. func (g *RdpClient) MouseMove(x, y int) { if !g.eventReady.Load() { return } g.mouse.mu.Lock() g.mouse.x = x g.mouse.y = y g.mouse.pending = true now := time.Now() since := now.Sub(g.mouse.lastTx) if since >= mouseCoalesceInterval { // Throttle window has elapsed — send right away. g.sendMouseMoveLocked(now) g.mouse.mu.Unlock() return } // Within throttle window: schedule a flush for the remainder of it // (unless one is already scheduled). if g.mouse.timer == nil { delay := mouseCoalesceInterval - since g.mouse.timer = time.AfterFunc(delay, g.flushMouseMoveTimer) } g.mouse.mu.Unlock() } // flushMouseMove sends any pending mouse-move event synchronously. Called // before any non-move input event to preserve server-side ordering. func (g *RdpClient) flushMouseMove() { g.mouse.mu.Lock() if g.mouse.timer != nil { g.mouse.timer.Stop() g.mouse.timer = nil } if g.mouse.pending { g.sendMouseMoveLocked(time.Now()) } g.mouse.mu.Unlock() } // flushMouseMoveTimer is the time.AfterFunc callback. Acquires the lock // itself and sends whatever's pending. func (g *RdpClient) flushMouseMoveTimer() { g.mouse.mu.Lock() g.mouse.timer = nil if g.mouse.pending && g.eventReady.Load() { g.sendMouseMoveLocked(time.Now()) } g.mouse.mu.Unlock() } // sendMouseMoveLocked must be called with mouse.mu held. func (g *RdpClient) sendMouseMoveLocked(now time.Time) { g.mouse.pdu.PointerFlags = pdu.PTRFLAGS_MOVE g.mouse.pdu.XPos = uint16(g.mouse.x) g.mouse.pdu.YPos = uint16(g.mouse.y) g.mouse.pending = false g.mouse.lastTx = now g.pdu.SendInputEvents(pdu.INPUT_EVENT_MOUSE, g.mouse.pduBuf[:]) } // MouseWheel sends a vertical scroll event to the remote desktop. // delta is the rotation amount in physical notches (1.0 = one click of a // scroll wheel = Windows WHEEL_DELTA). Fractional values are accepted for // smooth / high-resolution input devices such as trackpads. // Positive values scroll up (away from the user); negative values scroll down. func (g *RdpClient) MouseWheel(delta float64) { if !g.eventReady.Load() { return } slog.Debug("MouseWheel", "delta", delta) g.flushMouseMove() // Convert notch count to RDP WHEEL_DELTA units (120 per notch). const wheelDelta = 120 g.wheel.mu.Lock() g.wheel.accum += delta * wheelDelta if g.wheel.accum == 0 { // Opposite deltas cancelled out; nothing to send. g.wheel.mu.Unlock() return } now := time.Now() since := now.Sub(g.wheel.lastTx) if since >= mouseCoalesceInterval { g.sendWheelLocked(now) g.wheel.mu.Unlock() return } if g.wheel.timer == nil { delay := mouseCoalesceInterval - since g.wheel.timer = time.AfterFunc(delay, g.flushWheelTimer) } g.wheel.mu.Unlock() } // MouseHWheel sends a horizontal scroll event (trackpad two-finger horizontal // pan or tilt-wheel). delta is the rotation amount in physical notches; // positive values scroll right, negative scroll left. Shares the vertical // axis's coalescing window and timer. func (g *RdpClient) MouseHWheel(delta float64) { if !g.eventReady.Load() { return } g.flushMouseMove() const wheelDelta = 120 g.wheel.mu.Lock() g.wheel.haccum += delta * wheelDelta if g.wheel.haccum == 0 { g.wheel.mu.Unlock() return } now := time.Now() since := now.Sub(g.wheel.lastTx) if since >= mouseCoalesceInterval { g.sendWheelLocked(now) g.wheel.mu.Unlock() return } if g.wheel.timer == nil { delay := mouseCoalesceInterval - since g.wheel.timer = time.AfterFunc(delay, g.flushWheelTimer) } g.wheel.mu.Unlock() } // flushWheel sends any pending wheel event synchronously. Called before any // non-wheel input event to preserve server-side ordering. func (g *RdpClient) flushWheel() { g.wheel.mu.Lock() if g.wheel.timer != nil { g.wheel.timer.Stop() g.wheel.timer = nil } if g.wheel.accum != 0 { g.sendWheelLocked(time.Now()) } g.wheel.mu.Unlock() } // flushWheelTimer is the time.AfterFunc callback for wheel coalescing. func (g *RdpClient) flushWheelTimer() { g.wheel.mu.Lock() g.wheel.timer = nil if g.wheel.accum != 0 && g.eventReady.Load() { g.sendWheelLocked(time.Now()) } g.wheel.mu.Unlock() } // sendWheelLocked must be called with wheel.mu held. // Modelled on FreeRDP's send_mouse_wheel in client/SDL/SDL2/sdl_touch.cpp. func (g *RdpClient) sendWheelLocked(now time.Time) { // Truncate the accumulated float to a whole WHEEL_DELTA integer; keep the // fractional remainder so sub-notch trackpad movements aren't discarded. iaccum := int(g.wheel.accum) g.wheel.accum -= float64(iaccum) haccum := int(g.wheel.haccum) g.wheel.haccum -= float64(haccum) g.wheel.lastTx = now // sendAxis emits 0xFF-capped wheel events for one axis. The WheelRotation // field is 9 bits. Bits 0–7 hold the unsigned magnitude (max 0xFF per // event); bit 8 is the sign (PTRFLAGS_WHEEL_NEGATIVE). For negative values // the receiver computes -(0x100 - bits[0:7]), so we must store the 9-bit // two's-complement form, not the raw magnitude (same loop as FreeRDP). sendAxis := func(baseFlags uint16, v int) { negative := v < 0 if negative { v = -v } if negative { baseFlags |= uint16(pdu.PTRFLAGS_WHEEL_NEGATIVE) } for v > 0 { cval := min(v, 0xFF) v -= cval if negative { g.wheel.pdu.PointerFlags = (baseFlags & 0xFF00) | uint16(0x100-cval) } else { g.wheel.pdu.PointerFlags = baseFlags | uint16(cval) } g.pdu.SendInputEvents(pdu.INPUT_EVENT_MOUSE, g.wheel.pduBuf[:]) } } if iaccum != 0 { sendAxis(uint16(pdu.PTRFLAGS_WHEEL), iaccum) } if haccum != 0 { // 水平轴:正值 = 向右滚动(不设 NEGATIVE 位 = 远端向右)。 sendAxis(uint16(pdu.PTRFLAGS_HWHEEL), haccum) } if iaccum != 0 || haccum != 0 { g.notifyGfxLocalInput() } } func (g *RdpClient) MouseUp(button int, x, y int) { if !g.eventReady.Load() { return } slog.Debug("MouseUp", "x", x, "y", y, "button", button) g.flushMouseMove() g.flushWheel() p := &pdu.PointerEvent{} p.PointerFlags = mouseButtonFlag(button) p.XPos = uint16(x) p.YPos = uint16(y) g.pdu.SendInputEvents(pdu.INPUT_EVENT_MOUSE, []pdu.InputEventsInterface{p}) g.notifyGfxLocalInput() } func (g *RdpClient) MouseDown(button int, x, y int) { if !g.eventReady.Load() { return } slog.Debug("MouseDown", "x", x, "y", y, "button", button) g.flushMouseMove() g.flushWheel() p := &pdu.PointerEvent{} p.PointerFlags = pdu.PTRFLAGS_DOWN | mouseButtonFlag(button) p.XPos = uint16(x) p.YPos = uint16(y) g.pdu.SendInputEvents(pdu.INPUT_EVENT_MOUSE, []pdu.InputEventsInterface{p}) g.notifyGfxLocalInput() } // SetResolution requests a desktop resolution change via the MS-RDPEDISP // Display Update Virtual Channel. The server will reshape the desktop to the // given dimensions and send a fresh RDPGFX ResetGraphics command. // // width must be even and both width and height must be >= 200. // This method is a no-op when the RDPEDISP channel has not been established // (e.g. when the server does not support it). func (g *RdpClient) SetResolution(width, height int) { if g.dispHandler == nil { slog.Warn("SetResolution: RDPEDISP channel not available") return } w := uint32(width) if w%2 != 0 { w++ } w = max(w, 200) h := uint32(height) h = max(h, 200) g.dispHandler.SendMonitorLayout([]rdpedisp.Monitor{ { Flags: rdpedisp.MonitorFlagPrimary, Left: 0, Top: 0, Width: w, Height: h, PhysicalWidth: 0, PhysicalHeight: 0, Orientation: 0, DesktopScaleFactor: 100, DeviceScaleFactor: 100, }, }) slog.Debug("SetResolution", "width", w, "height", h) } // SetQueueDepthHint controls the frame-rate and encoding quality reported to // the server via the RDPGFX FRAME_ACKNOWLEDGE queueDepth field // (MS-RDPEGFX 2.2.2.8). // // A higher value signals a larger client decode backlog, causing the server to // slow down or reduce H.264/RFX encoding quality. 0 (default) means "report // the real decode-queue length" — no artificial throttling. // DebugSurfacePixel 诊断:读第一个 mapped surface 上 (x,y) 的 BGRA func (g *RdpClient) DebugSurfacePixel(x, y int) (uint8, uint8, uint8, uint8, bool) { return g.gfxHandler.DebugSurfacePixel(x, y) } // GfxDiagStats 返回图形管线实时诊断指标:累计解码帧数、解码队列深度、 // 帧间隔(毫秒)与单消息解码耗时 EMA(微秒)。见 GfxHandler.DiagStats。 func (g *RdpClient) GfxDiagStats() (frames, qdepth, fintvMs, decUs int64) { return g.gfxHandler.DiagStats() } // CodecStats exposes cumulative surface-bitmap bytes per codec id for // bandwidth diagnostics. Returns nil before a graphics session exists. func (g *RdpClient) CodecStats() map[uint16]int64 { if g.gfxHandler != nil { return g.gfxHandler.CodecStats() } return nil } // Typical values: 0 = off, 10–50 = moderate throttle, 100+ = heavy throttle. // Use 0xFFFFFFFF to pause new frames entirely (the stream resumes when hint is // reduced or cleared). func (g *RdpClient) SetQueueDepthHint(depth uint32) { if g.gfxHandler != nil { g.gfxHandler.SetQueueDepthHint(depth) } } // SetAudioMode 设置声音重定向模式(mstsc 远程音频播放): // "local"(本机播放,默认)| "none"(不播放:协商后丢弃,服务器静音)| // "remote"(远端播放:不注册音频通道,服务器本机扬声器出声)。 // 必须在 Login 前调用。 func (g *RdpClient) SetAudioMode(mode string) { g.audioMode = mode } // SetDriveRedirect 启用驱动器重定向(MS-RDPEFS):远端将看到一个只读的 // 重定向设备,内容由 SetFilesystem 桥接的异步文件系统提供。必须在 Login // 前调用。 func (g *RdpClient) SetDriveRedirect(enabled bool) { g.driveRedirect = enabled } // SetDriveLabel 设置重定向卷的卷标(远端 Explorer 显示名)。必须在 Login // 前调用;空值取 "local"。 func (g *RdpClient) SetDriveLabel(label string) { g.driveLabel = label } // SetDriveFilesystem 挂接异步文件系统桥。处理器在 Login(doLogin)内才 // 创建——这里先暂存字段,创建时补挂(同 SetGfxCacheStore 的时序教训)。 func (g *RdpClient) SetDriveFilesystem(fs rdpdr.Filesystem) { g.driveFS = fs } // Rdpdr 返回驱动器重定向处理器(未启用时为 nil)。桥接层用它接收完成回调。 func (g *RdpClient) Rdpdr() *rdpdr.Handler { return g.rdpdrHandler } // SetSessionColorDepth 请求会话颜色位数(16/24/32,其它值按 32 处理)。 // 仅影响传统位图管线;RDPGFX 会话表面恒为 32bpp。必须在 Login 前调用。 func (g *RdpClient) SetSessionColorDepth(bpp int) { g.colorDepth = bpp } // SetPerformanceFlags 覆盖 Client Info PDU 的 performanceFlags(体验选项)。 // 禁用类位置位 = 关闭(PERF_DISABLE_WALLPAPER 等),0 = 视觉全开默认。 // 必须在 Login 前调用。 func (g *RdpClient) SetPerformanceFlags(flags uint32) { g.perfFlags = flags g.perfFlagsSet = true } // RequestKeyframe asks the server to send a fresh full-screen IDR keyframe via // the SuppressOutput off→on toggle (SendForceRefresh). It is the same request // the RDPGFX decoder issues internally when it stalls, exposed publicly so the // frontend's black-screen watchdog can recover an initial all-black session — // where the server sent its first IDR before the desktop finished painting // (decoded as a black warm-up frame and dropped) and then went idle — without // the cost of a full reconnect. Safe to call from the render loop goroutine. func (g *RdpClient) RequestKeyframe() { if g.closed.Load() { return } if g.pdu != nil { g.pdu.SendForceRefresh(uint16(g.width), uint16(g.height)) } } func (g *RdpClient) Reconnect(width, height int) error { if g.closed.Load() { return fmt.Errorf("client is closed") } g.reconnectMu.Lock() defer g.reconnectMu.Unlock() g.reconnecting.Store(true) defer func() { g.reconnecting.Store(false) }() slog.Debug("Reconnect", "width", width, "height", height) g.closeTransport() g.width = width g.height = height g.eventReady.Store(false) const maxRetries = 3 for attempt := 1; attempt <= maxRetries; attempt++ { // No delay on the first attempt: the transport was already closed above // so the server has already started session teardown. Use exponential // backoff (1s, 2s) only for retries after a failed login. delay := time.Duration(0) if attempt > 1 { delay = time.Duration(1<