Files

2817 lines
101 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 rdpgfx
import (
"encoding/binary"
"fmt"
"log/slog"
"runtime/debug"
"sync"
"sync/atomic"
"time"
"git.zeroonesoft.cn/golib/rdplib/plugin"
)
// regionPool reuses byte slices for progressive codec rectangle extraction,
// avoiding per-rectangle allocations that cause GC pressure.
var regionPool = sync.Pool{
New: func() any { return []byte(nil) },
}
const (
ChannelName = plugin.RDPGFX_DVC_CHANNEL_NAME
)
// suspendFrameAcknowledge is the queueDepth value defined in MS-RDPEGFX
// 2.2.2.8 that instructs the server to suspend sending new frames until
// the client sends a subsequent FRAME_ACKNOWLEDGE with a different value.
const suspendFrameAcknowledge uint32 = 0xFFFFFFFF
// RDPGFX Command IDs (MS-RDPEGFX 2.2.2)
const (
cmdidWireToSurface1 uint16 = 0x0001
cmdidWireToSurface2 uint16 = 0x0002
cmdidDeleteEncodingContext uint16 = 0x0003
cmdidSolidFill uint16 = 0x0004
cmdidSurfaceToSurface uint16 = 0x0005
cmdidSurfaceToCache uint16 = 0x0006
cmdidCacheToSurface uint16 = 0x0007
cmdidEvictCacheEntry uint16 = 0x0008
cmdidCreateSurface uint16 = 0x0009
cmdidDeleteSurface uint16 = 0x000A
cmdidStartFrame uint16 = 0x000B
cmdidEndFrame uint16 = 0x000C
cmdidFrameAcknowledge uint16 = 0x000D
cmdidResetGraphics uint16 = 0x000E
cmdidMapSurfaceToOutput uint16 = 0x000F
cmdidCacheImportOffer uint16 = 0x0010
cmdidCacheImportReply uint16 = 0x0011
cmdidCapsAdvertise uint16 = 0x0012
cmdidCapsConfirm uint16 = 0x0013
cmdidMapSurfaceToWindow uint16 = 0x0015
cmdidQoeFrameAcknowledge uint16 = 0x0016
cmdidMapSurfaceToScaledOutput uint16 = 0x0017
cmdidMapSurfaceToScaledWindow uint16 = 0x0018
)
// Pixel Formats
const (
pixelFormatXRGB8888 uint8 = 0x20
pixelFormatARGB8888 uint8 = 0x21
)
// Codec IDs (MS-RDPEGFX 2.2.2.1 / FreeRDP rdpgfx.h)
const (
codecUncompressed uint16 = 0x0000
codecCaVideo uint16 = 0x0003 // RDPGFX_CODECID_CAVIDEO (RemoteFX tiles)
codecPlanar uint16 = 0x0004
codecClear uint16 = 0x0008
codecProgressive uint16 = 0x0009
codecAVC420 uint16 = 0x000B
codecAVC444 uint16 = 0x000E
codecAVC444v2 uint16 = 0x000F
)
// Capability versions and flags
const (
capVersion8 uint32 = 0x00080004
capVersion81 uint32 = 0x00080105
capVersion10 uint32 = 0x000A0002
capVersion101 uint32 = 0x000A0100
capVersion102 uint32 = 0x000A0200
capVersion103 uint32 = 0x000A0301
capVersion104 uint32 = 0x000A0400
capVersion105 uint32 = 0x000A0502
capVersion106 uint32 = 0x000A0600
capVersion1061 uint32 = 0x000A0601
capVersion107 uint32 = 0x000A0701
capFlagThinClient uint32 = 0x00000001
capFlagSmallCache uint32 = 0x00000002
capFlagAVC420Enabled uint32 = 0x00000010 // v8.1: explicitly enable AVC420
capFlagAVCDisabled uint32 = 0x00000020 // v10+: disable AVC
)
const headerSize = 8
// BitmapUpdate represents a rendered bitmap region.
//
// Lifecycle: 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 BitmapUpdate struct {
DestLeft, DestTop, DestRight, DestBottom int
Width, Height int
Bpp int // bytes per pixel (always 4)
Data []byte // BGRA pixel data — see lifecycle note above
}
// bitmapBufPool reuses BGRA byte slices used to back BitmapUpdate.Data.
// Buffers are acquired with acquireBitmapBuf, handed to the onBitmap
// callback, and released with releaseBitmapBuf once the (synchronous)
// callback returns. This eliminates per-rectangle allocations on the
// hot CaVideo / AVC partial-blit paths.
var bitmapBufPool = sync.Pool{
New: func() any { return []byte(nil) },
}
// decodePkt is the message type for the async decode channel.
// pooled is true when data was acquired from bitmapBufPool; the receiver
// must call releaseBitmapBuf(data) after processing.
type decodePkt struct {
data []byte
pooled bool
}
func acquireBitmapBuf(size int) []byte {
if size <= 0 {
return nil
}
b := bitmapBufPool.Get().([]byte)
if cap(b) < size {
return make([]byte, size)
}
return b[:size]
}
func releaseBitmapBuf(b []byte) {
if b == nil {
return
}
//nolint:staticcheck // intentional pool of byte slices
bitmapBufPool.Put(b[:cap(b)])
}
// emitAndReleaseUpdates calls the onBitmap callback and then returns the
// pooled Data buffers of the supplied updates back to bitmapBufPool. All
// updates passed in must have Data acquired via acquireBitmapBuf.
func (g *GfxHandler) emitAndReleaseUpdates(updates []BitmapUpdate) {
if g.onBitmap != nil && len(updates) > 0 {
g.onBitmap(updates)
}
for i := range updates {
releaseBitmapBuf(updates[i].Data)
updates[i].Data = nil
}
}
type surface struct {
width, height uint16
format uint8
data []byte // BGRA, 4 bytes per pixel
outputX uint32
outputY uint32
mapped bool
// shadowStale is true when the CPU surface shadow (data) may not match the
// pixels currently shown on the GPU display. It is set when an AVC frame
// advances the GPU display (onNV12/onI420) without a corresponding BGRA
// shadow update (decoded == nil), and cleared by a full-surface blit. While
// stale, the region-only shadow blit fast path is bypassed in favour of a
// full blit so the shadow is fully repaired before the next partial update.
// Initialised true so the first frame on a fresh (zero-filled) surface always
// takes the full-blit path.
shadowStale bool
}
type vBarEntry struct {
pixels []byte // BGRA pixel data, 4 bytes per pixel
count int
}
type cacheEntry struct {
data []byte // BGRA pixel data
width, height int
key uint64 // 服务器提供的持久缓存键(SurfaceToCache),0 = 无
}
// GfxCacheEntry 是一条跨连接保留的位图缓存条目(MS-RDPEGFX 持久位图缓存)。
// Key 为服务器在 SurfaceToCache 中给出的 cacheKey(持久身份),Data 为当时
// 的 BGRA 像素副本。重连后经 CacheImportOffer 上报,服务器按前缀导入并
// 重新分配槽位,后续 CacheToSurface 即可直接回贴、无需重传像素。
type GfxCacheEntry struct {
Key uint64
Width, Height int
Bpp uint16
Data []byte
}
// GfxCacheStore 在浏览器侧持久化缓存条目(v1:页面内存,跨手动重连保留)。
// Persist 随每条 SurfaceToCache 调用;Export 在每次连接 caps 确认后调用一次,
// 返回待上报的条目(顺序即 CacheImportReply 的前缀导入顺序)。
// Get/Keys 供 bitmap 管线持久缓存(6.4b M2)按需查取与枚举键。
type GfxCacheStore interface {
Persist(key uint64, w, h int, bpp uint16, data []byte)
Export() []GfxCacheEntry
Get(key uint64) (GfxCacheEntry, bool)
Keys() []uint64
}
// GfxHandler implements the RDPGFX (MS-RDPEGFX) protocol.
type GfxHandler struct {
surfaces map[uint16]*surface
cacheEntries map[uint16]cacheEntry
cacheStore GfxCacheStore // 持久缓存桥(nil = 未启用)
offeredCache []GfxCacheEntry // 本连接 CacheImportOffer 的条目,等待 Reply 前缀映射
importOfferSent bool // 每连接只上报一次
clearCtx *clearCodecCtx
zgfx *zgfxContext
rfx *rfxDecoder
progressive *rfxProgressiveDecoder
// codecBytes[i] 累计 codecId=i 的表面位图字节数(WTS1+WTS2),用于带宽诊断
codecBytes [16]atomic.Int64
// cmdCounts[i] 累计 cmdId=i 的 PDU 条数(诊断服务端停止发帧用)
cmdCounts [64]atomic.Int64
h264dec H264Decoder
// h264dec2 is the auxiliary H.264 decoder used for AVC444v2 LC=2 chroma-upgrade
// frames. It decodes stream2, which carries chroma values for positions not
// covered by stream1's 4:2:0 quantiser. The decoded I420 planes are combined
// with the luma and chroma planes cached from the most recent LC=0/1 main-stream
// decode to reconstruct full 4:4:4 YUV before converting to BGRA.
h264dec2 H264Decoder
// avc444YPlane caches the luma (Y) and half-res chroma (U/V) planes from the
// last main-stream AVC444 decode, for use when an LC=2 chroma-upgrade frame arrives.
avc444YPlane avc444YPlane
// avc444IDRYPlane caches the luma and half-res chroma planes from the most
// recently decoded stream1 IDR frame. When a standalone LC=2 packet carries
// a stream2 IDR, it should be combined with the matching stream1 IDR luma
// (not the latest P-frame luma stored in avc444YPlane), so we keep this
// separate snapshot.
avc444IDRYPlane avc444YPlane
// lc2SampleLogged is set after the first LC=2 combine output has been
// sampled for green/pink colour diagnostics. Reset on each stream1 IDR so
// we can observe the combine quality at every GOP boundary.
lc2SampleLogged bool
lc2PFrameSampleLogged bool // logged first P-frame LC=2 combine after IDR
lc0SampleLogged bool // logged first LC=0 IDR frame pixel samples
// framesDecoded is accessed from both read and decode goroutines.
framesDecoded atomic.Uint32
// 帧间隔与单消息解码耗时的 EMA(微秒),DiagStats 输出用。
// 写方为 decode goroutine,读方为统计循环,CAS 循环更新。
frameIntvUs atomic.Int64
decUs atomic.Int64
lastFrameAt atomic.Int64 // unix nano,0 = 尚未收到帧
// 帧解码起点(unix ns):StartFrame 到 EndFrame 的墙钟时长进 QoE 上报
frameDecodeStart atomic.Int64
// SUSPEND 状态:队列深时置位(向服务器请求暂停发送),排空后复位
suspended atomic.Bool
// sessionW/H 为通道初始化时发送 RESET_GRAPHICS 用的会话尺寸
//(grdp.go 创建 handler 时注入)。
sessionW, sessionH uint32
sendFn func(data []byte)
onBitmap func([]BitmapUpdate)
// decodeCh receives decompressed PDU data for asynchronous decode.
decodeCh chan decodePkt
// ackCh is a buffered channel of serialized ACK PDUs. Every
// EndFrame ACK is enqueued here and the writeLoop goroutine sends
// each one to the server. The server tracks outstanding frames
// individually, so skipping ACKs causes it to stop sending.
ackCh chan []byte
// doneCh is closed by Close() to signal decodeLoop and writeLoop to exit.
doneCh chan struct{}
closeOnce sync.Once
// onDecoderBroken is called once when the H.264 decoder becomes permanently
// unrecoverable (all soft resets exhausted). The caller should reconnect
// the RDP session to create a fresh decoder.
onDecoderBroken func()
decoderBrokenNotified bool
// watchdogCh receives signals from background timers inside ffmpegDecoder
// when stall-probe or IDR-wait timeouts expire independently of server
// frame rate. decodeLoop selects on this channel so it calls
// maybeNotifyDecoderBroken even when no server frames are arriving.
watchdogCh chan struct{}
// lastDecodedFrame records when a visible AVC frame was last produced.
// Local-input watchdogs compare against this timestamp so recovery does
// not wait for a subsequent H.264 packet to arrive.
lastDecodedFrame atomic.Int64
inputWatchdogMu sync.Mutex
inputWatchdog *time.Timer
inputWatchdogNS int64
// lastLC2RecvTime records when the most recent AVC444 LC=2 frame arrived,
// regardless of whether it could be decoded. Used to detect the
// "server sending LC=2 only, aux decoder absent" deadlock.
lastLC2RecvTime atomic.Int64
// auxDecoderBrokenTimer fires after auxDecoderBrokenTimeout when h264dec2
// is nil. When it fires it signals watchdogCh so that decodeLoop can call
// maybeRenegotiateCapabilities and break the LC=2-only deadlock.
auxDecoderBrokenTimer *time.Timer
auxDecoderBrokenTimerMu sync.Mutex
// onKeyframeRequest is called to ask the server to send a fresh IDR
// keyframe. Optional: if nil, the decoder will wait for the next
// server-initiated keyframe.
onKeyframeRequest func()
// lastKeyframeRequest is the wall-clock time of the most recent keyframe
// request sent to the server. Used to rate-limit repeat requests.
lastKeyframeRequest time.Time
// softResetCount tracks how many in-place decoder resets have been
// attempted since the last server-triggered RESET_GRAPHICS.
softResetCount int
// noIDRSoftResetCount tracks soft resets triggered specifically by the
// no-IDR broken reason. This counter is kept separate from softResetCount
// so that a prior HW-stall reset (which increments softResetCount) does not
// consume the no-IDR recovery budget. Reset on RESET_GRAPHICS and on a
// successful frame decode.
noIDRSoftResetCount int
// usingSWFallback is set after a HW stall forces a switch to software
// decoding. Both h264dec and h264dec2 are created SW-only while this
// flag is true, avoiding repeated VideoToolbox stalls that would
// otherwise trigger a full RDP reconnect.
usingSWFallback bool
// swFallbackPrimed is set after a SW fallback decoder has been primed
// with the cached (stale) stream1 IDR. It stays armed for the whole SW
// fallback window — until a genuine fresh IDR resyncs the decoder
// (maybeCacheStream1IDR) or RESET_GRAPHICS — so that consecutive dropped
// frames caused by the stale prime can be detected and escalated to a
// reconnect instead of producing endless green/zero-UV frames. It is NOT
// cleared by a single successful decode: stale-prime corruption only
// appears once the following P-frames diverge from the missing reference.
swFallbackPrimed bool
// swFallbackDroppedCount counts dropped frames since the SW fallback
// decoder was primed. If it exceeds swFallbackDropLimit the decoder is
// declared broken and onDecoderBroken is fired.
swFallbackDroppedCount int
// swFallbackFirstDropTime records when the current run of consecutive
// stale-prime drops started. It bounds how long corruption may persist
// before escalating to a reconnect, giving the ForceRefresh resync IDR a
// chance to heal the picture first. Reset whenever a clean frame decodes.
swFallbackFirstDropTime time.Time
// lc2EverDecoded is set to true after the first successful AVC444 LC=2
// chroma-upgrade decode. maybeRenegotiateCapabilities uses this to
// distinguish "LC=2 was working and then broke" (reconnect needed) from
// "LC=2 never worked this session" (server may not support stream2 priming;
// gracefully degrade to LC=0 only without reconnecting).
lc2EverDecoded bool
// auxDecoderNoIDRRetries counts how many times maybeRenegotiateCapabilities
// has been called in the "stream2EverSeen but lc2EverDecoded=false" case
// this session. Each attempt sends a ForceRefresh; after
// auxDecoderMaxIDRRetries consecutive attempts the session degrades to LC=0
// only (no reconnect — the server consistently omits stream2 IDRs).
// Reset on RESET_GRAPHICS, when LC=2 successfully decodes, and whenever
// maybeRenegotiateCapabilities returns early due to no recent LC=2 activity
// (e.g. during a HW-decoder GOP-boundary stall) so a resumed burst gets
// fresh retries rather than inheriting the previous count.
auxDecoderNoIDRRetries int
// lc2PermanentlyDegraded is set when the server has not delivered a stream2
// IDR despite repeated keyframe requests and lc2EverDecoded is still false.
// Once set, LC=2 frames are silently skipped for the remainder of the session
// without arming the renegotiation timer — avoiding an endless reconnect loop.
// Cleared on RESET_GRAPHICS so a fresh AVC444 sequence gets a clean slate.
// Cleared in primeAuxDecoder on a stream2 IDR to allow late recovery.
lc2PermanentlyDegraded bool
// stream2EverSeen is set when a non-empty stream2 payload is observed inside
// an LC=0 packet. VirtualBox VRDE never includes stream2 in LC=0 packets,
// so stream2EverSeen stays false for the whole session. Windows does include
// stream2 in LC=0 IDRs, so stream2EverSeen becomes true as soon as the first
// LC=0 AVC444 packet arrives. maybeRenegotiateCapabilities uses this to
// distinguish "server never sends stream2" (VirtualBox → permanent degrade)
// from "stream2 seen but aux decoder not yet primed" (Windows, Chrome just
// launched → transient, wait for IDR rather than logging a WARN).
// Reset on RESET_GRAPHICS since the server starts a fresh AVC444 sequence.
stream2EverSeen bool
// lastStream1IDR caches the most recent stream1 H.264 IDR NAL data in
// Annex B format. When VideoToolbox stalls and the decoder falls back to
// software, this data is fed immediately to the new SW decoder so it can
// decode subsequent P-frames without waiting for the server to send a fresh
// IDR via ForceRefresh (which some servers, e.g. VirtualBox VRDE, ignore).
// Cleared on RESET_GRAPHICS to avoid feeding a stale IDR to a new pipeline.
lastStream1IDR []byte
// lastStream1IDRTime is the wall-clock time when the most recent stream1 IDR
// was cached, logged for diagnostics (e.g. how stale the priming IDR was).
// lastStream1IDRFrame is framesDecoded at the same instant, logged for diagnostics.
lastStream1IDRTime time.Time
lastStream1IDRFrame uint32
// avc444Disabled, when true, limits the CAPS_ADVERTISE to v8.0 and v8.1
// (AVC420 only). The server will never send AVC444/AVC444v2 frames, which
// avoids the LC=2 colour degradation seen with VirtualBox VRDE.
avc444Disabled bool
// avcDisabled, when true, advertises v10.x with AVC_DISABLED so the server
// keeps the RDPGFX channel on ClearCodec/RFX Progressive (no H.264 decode).
avcDisabled bool
// pduRecord, when non-nil, receives every wire-to-surface bitmap payload
// for offline replay (garbled-screen/bandwidth debugging harness).
pduRecord func(kind byte, codecId uint16, surfW, surfH, x, y, w, h uint32, payload []byte)
// queueDepthHint is a minimum queueDepth to report in FRAME_ACKNOWLEDGE
// PDUs. A higher value makes the server believe the client has a larger
// decode backlog, causing it to slow down or reduce encoding quality.
// 0 means "report the real queue length" (default, no throttling).
// See SetQueueDepthHint.
queueDepthHint atomic.Uint32
// singleUpdate is a pre-allocated one-element slice reused by emitBitmap
// and emitBitmapPooled to avoid a heap allocation on every BGRA frame.
// Safe: all emitBitmap* calls run on the single decode goroutine.
singleUpdate [1]BitmapUpdate
// updatesBuf is a pre-allocated slice reused by emitCaVideoRects and
// blitAndEmitAVCRegions to avoid make() on every multi-region frame.
// Safe: both functions run on the single decode goroutine and never call
// each other.
updatesBuf []BitmapUpdate
// avcStream1 and avcStream2 are pre-allocated AVC stream structs reused by
// fillAVC444Stream and the decode methods to avoid a heap allocation for
// the stream header (struct + regions slice) on every H.264 frame.
// Safe: all AVC decode calls run on the single decode goroutine.
avcStream1 avc420Stream
avcStream2 avc420Stream
// regionHintBuf is a pre-allocated slice reused by the SetRegionHint call
// sites to avoid a make([][4]uint16, ...) on every H.264 frame that carries
// dirty region metadata.
// Safe: used only on the single decode goroutine.
regionHintBuf [][4]uint16
// onH264Raw is called with raw H.264 NAL unit data when h264dec is nil
// (e.g. WASM builds without CGo). The caller can forward the data to a
// JavaScript WebCodecs VideoDecoder instead.
// destX, destY are the top-left canvas coordinates.
// regions 为扁平 [l,t,r,b,...](帧内坐标,右下开区间):服务器只保证
// 区域内像素有效,帧内其余像素未定义,绘制端必须只贴区域;空切片
// 表示整帧有效。
onH264Raw func(destX, destY, w, h int, isKey bool, data []byte, regions []int32)
// onI420 is called after a successful H.264 decode when I420 planar data
// is available. The caller can upload the planes to an SDL2 IYUV texture
// for GPU-accelerated YUV→RGB conversion, bypassing the CPU colour path.
// destX, destY are absolute canvas coordinates.
onI420 func(destX, destY, w, h int, y []byte, yStride int, u []byte, uStride int, v []byte, vStride int)
// onNV12 is like onI420 but receives native NV12 output. It is preferred
// by SDL2 clients when available because VideoToolbox commonly outputs
// NV12 and SDL_UpdateNVTexture can upload it directly.
onNV12 func(destX, destY, w, h int, y []byte, yStride int, uv []byte, uvStride int)
}
// NewGfxHandler creates a new RDPGFX handler.
func NewGfxHandler(onBitmap func([]BitmapUpdate)) *GfxHandler {
g := &GfxHandler{
surfaces: make(map[uint16]*surface),
cacheEntries: make(map[uint16]cacheEntry),
clearCtx: newClearCodecCtx(),
zgfx: newZgfxContext(),
rfx: newRfxDecoder(),
progressive: newRfxProgressiveDecoder(),
// h264dec2 starts nil; primeAuxDecoder creates it on the first stream2 IDR
// so it is always primed before decoding LC=2 P-frames.
onBitmap: onBitmap,
decodeCh: make(chan decodePkt, 64),
ackCh: make(chan []byte, 512),
doneCh: make(chan struct{}),
watchdogCh: make(chan struct{}, 4),
}
g.h264dec = newH264DecoderWithWatchdog(g.watchdogCh)
go g.decodeLoop()
go g.writeLoop()
return g
}
// inputStallSilentThreshold is the minimum time without a decoded frame before
// the local-input watchdog considers the HW decoder potentially stalled.
// Kept in the non-build-tagged file so both normal and !h264 builds compile.
// Must be kept in sync with avcHWReadyFreezeThreshold in h264_ffmpeg.go.
const inputStallSilentThreshold = 7 * time.Second
const localInputRecoveryGrace = 750 * time.Millisecond
// auxDecoderBrokenTimeout is the maximum time we wait for an LC=0 stream2 IDR
// to arrive and recreate the aux decoder (h264dec2) after it has been torn down.
// If LC=2 frames keep arriving beyond this window without an LC=0 IDR, the RDPGFX
// capabilities are re-advertised to force the server to issue RESET_GRAPHICS and
// restart the video pipeline with a fresh LC=0 IDR for both streams.
const auxDecoderBrokenTimeout = 10 * time.Second
// auxDecoderMaxIDRRetries is the maximum number of ForceRefresh keyframe
// requests sent while waiting for a stream2 IDR before giving up and
// reconnecting. Each attempt is spaced by auxDecoderBrokenTimeout (10 s).
// Windows servers typically begin a new GOP every 30–60 s, so 5 attempts
// (≤50 s) provides enough coverage for at least one natural IDR boundary.
const auxDecoderMaxIDRRetries = 5
// Close shuts down the GfxHandler's background goroutines.
// Safe to call multiple times; subsequent calls are no-ops.
//
// h264dec is intentionally NOT freed here: decodeLoop (goroutine 21) may be
// in the middle of avcodec_send_packet when Close is called from the transport
// goroutine, which would cause a use-after-free SIGSEGV. Instead, decodeLoop
// defers cleanup of h264dec so it always runs after the last Decode call.
func (g *GfxHandler) Close() {
g.closeOnce.Do(func() {
g.stopInputWatchdog()
g.stopAuxDecoderBrokenTimer()
close(g.doneCh)
})
}
// NotifyLocalInput tells the graphics pipeline that a real local input event
// was just sent to the server. If the decoder has already been silent longer
// than the HW stall threshold, arm a short watchdog so recovery no longer
// depends on the next H.264 packet arriving.
func (g *GfxHandler) NotifyLocalInput() {
if g.h264dec == nil {
return
}
// Don't arm the watchdog when the decoder is waiting for an IDR after a
// reset: it is already in a known recovery state, and arming the watchdog
// here would force-break the newly reset decoder 750 ms later when the IDR
// hasn't arrived yet, causing rapid cascading soft resets.
if g.h264dec.NeedsIDR() {
return
}
lastDecodedNS := g.lastDecodedFrame.Load()
if lastDecodedNS == 0 {
return
}
now := time.Now()
silentFor := now.Sub(time.Unix(0, lastDecodedNS))
if silentFor < inputStallSilentThreshold {
return
}
// Only arm the watchdog if the server has been actively sending video
// packets recently. A genuinely static screen (server sends nothing) is
// not a decoder stall and should not trigger a force-break.
if recvTime := g.h264dec.LastReceiveTime(); recvTime.IsZero() ||
time.Since(recvTime) >= inputStallSilentThreshold {
return
}
inputNS := now.UnixNano()
g.inputWatchdogMu.Lock()
g.inputWatchdogNS = inputNS
if g.inputWatchdog == nil {
g.inputWatchdog = time.AfterFunc(localInputRecoveryGrace, func() {
g.fireInputWatchdog(inputNS)
})
} else {
g.inputWatchdog.Reset(localInputRecoveryGrace)
}
g.inputWatchdogMu.Unlock()
slog.Debug("H.264: local input armed stall watchdog",
"silentFor", silentFor.Round(time.Millisecond))
}
func (g *GfxHandler) fireInputWatchdog(inputNS int64) {
g.inputWatchdogMu.Lock()
if g.inputWatchdogNS != inputNS {
g.inputWatchdogMu.Unlock()
return
}
g.inputWatchdog = nil
g.inputWatchdogMu.Unlock()
select {
case g.watchdogCh <- struct{}{}:
default:
}
}
func (g *GfxHandler) stopInputWatchdog() {
g.inputWatchdogMu.Lock()
if g.inputWatchdog != nil {
g.inputWatchdog.Stop()
g.inputWatchdog = nil
}
g.inputWatchdogNS = 0
g.inputWatchdogMu.Unlock()
}
// startAuxDecoderBrokenTimer arms a one-shot timer. If it fires before
// stopAuxDecoderBrokenTimer cancels it, it signals watchdogCh so that
// decodeLoop calls maybeRenegotiateCapabilities.
// Idempotent: only the first call while h264dec2 is nil takes effect.
func (g *GfxHandler) startAuxDecoderBrokenTimer() {
g.auxDecoderBrokenTimerMu.Lock()
defer g.auxDecoderBrokenTimerMu.Unlock()
if g.auxDecoderBrokenTimer == nil {
g.auxDecoderBrokenTimer = time.AfterFunc(auxDecoderBrokenTimeout, func() {
// Clear the pointer so startAuxDecoderBrokenTimer can re-arm
// the timer for subsequent retry attempts (e.g. after a
// ForceRefresh in case 2 of maybeRenegotiateCapabilities).
g.auxDecoderBrokenTimerMu.Lock()
g.auxDecoderBrokenTimer = nil
g.auxDecoderBrokenTimerMu.Unlock()
select {
case g.watchdogCh <- struct{}{}:
default:
}
})
}
}
func (g *GfxHandler) stopAuxDecoderBrokenTimer() {
g.auxDecoderBrokenTimerMu.Lock()
if g.auxDecoderBrokenTimer != nil {
g.auxDecoderBrokenTimer.Stop()
g.auxDecoderBrokenTimer = nil
}
g.auxDecoderBrokenTimerMu.Unlock()
}
// maybeRenegotiateCapabilities is called from decodeLoop when watchdogCh fires.
// If h264dec2 has been nil since the timer was armed AND the server has been
// actively sending LC=2 frames, we either reconnect (if LC=2 was previously
// working) or degrade gracefully to LC=0 only (if LC=2 never worked this session).
// Reconnecting mid-session when LC=2 never worked would just produce another
// identical cycle, since the server appears to not include stream2 in LC=0 IDRs.
func (g *GfxHandler) maybeRenegotiateCapabilities() {
if g.h264dec2 != nil {
return // aux decoder recovered while timer was in flight
}
if g.decoderBrokenNotified {
return // reconnect already in flight
}
if g.lc2PermanentlyDegraded {
return // already degraded to LC=0 only; buffered watchdog signals must not re-enter
}
// Only act when the server has recently been sending LC=2 frames —
// a genuinely idle server needs no intervention.
lastLC2NS := g.lastLC2RecvTime.Load()
if lastLC2NS == 0 || time.Since(time.Unix(0, lastLC2NS)) >= auxDecoderBrokenTimeout {
// No recent LC=2 activity (server idle, or HW-decoder GOP-boundary
// stall). Reset the retry counter so that when LC=2 resumes the
// next burst gets fresh ForceRefresh attempts rather than inheriting
// a stale count from a previous burst.
g.auxDecoderNoIDRRetries = 0
return
}
if !g.stream2EverSeen {
// stream2 has never appeared in any LC=0 packet this session. The server
// (e.g. VirtualBox VRDE) does not include stream2 in LC=0 IDRs and does
// not send standalone LC=2 IDR frames, so the aux decoder can never be
// primed. Reconnecting would reproduce the same failure. Degrade
// gracefully: LC=2 is silently skipped for the remainder of this session.
slog.Warn("H.264: server never sent stream2 in LC=0, LC=2 degraded to LC=0 only (no reconnect)")
g.lc2PermanentlyDegraded = true
return
}
if !g.lc2EverDecoded {
// stream2 has been seen in LC=0 packets (server supports LC=2), but
// the aux decoder has not yet produced a combined frame. The server
// may be sending only P-frame stream2 data in LC=0 packets — no IDR
// means primeAuxDecoder can never initialise h264dec2.
//
// Send a ForceRefresh keyframe request on each attempt so the server
// hopefully includes a stream2 IDR in the next LC=0 IDR packet.
// Windows servers typically begin a new GOP every 30–60 s; with
// auxDecoderMaxIDRRetries=5 (×10 s = 50 s) we cover at least one
// natural boundary before falling back to a reconnect.
//
// The retry counter is reset when the early-return above fires (no
// recent LC=2 activity, e.g. HW-decoder stall), so resumed bursts
// always start from attempt 1.
g.auxDecoderNoIDRRetries++
if g.auxDecoderNoIDRRetries <= auxDecoderMaxIDRRetries {
slog.Debug("H.264: aux decoder not primed — requesting keyframe to get stream2 IDR",
"attempt", g.auxDecoderNoIDRRetries, "maxRetries", auxDecoderMaxIDRRetries)
g.lastKeyframeRequest = time.Time{} // allow immediate send
g.maybeRequestKeyframe()
return
}
// The server has not delivered a stream2 IDR despite repeated ForceRefresh
// requests and LC=2 has never decoded successfully this session. Reconnecting
// reproduces the same failure because the server consistently omits stream2
// IDRs in LC=0 IDR packets. Degrade gracefully to LC=0-only instead.
slog.Warn("H.264: aux decoder never primed despite keyframe request — degrading to LC=0 only (no reconnect)",
"retries", g.auxDecoderNoIDRRetries)
g.lc2PermanentlyDegraded = true
return
}
// The timer firing is itself the auxDecoderBrokenTimeout signal. Do not
// gate on lastDecodedFrame: the main decoder (h264dec) may still be active
// and continuously updating lastDecodedFrame even while the aux decoder
// (h264dec2) is broken, which would prevent this function from ever
// triggering a reconnect and permanently lose LC=2 chroma enhancement.
slog.Debug("H.264: aux decoder absent, server sending LC=2 — triggering reconnect")
g.decoderBrokenNotified = true
if g.onDecoderBroken != nil {
go g.onDecoderBroken()
}
}
func (g *GfxHandler) noteSuccessfulDecode() {
g.lastDecodedFrame.Store(time.Now().UnixNano())
g.noIDRSoftResetCount = 0
// A single clean decode resets the *consecutive* drop counter, but does NOT
// disarm swFallbackPrimed. After a HW stall the SW decoder is primed with a
// stale cached IDR; the corruption it causes only manifests as the following
// P-frames diverge from the (missing) reference, so the very first frame can
// decode cleanly and then the picture rots into green/zero-chroma output.
// Clearing swFallbackPrimed here would permanently disable
// trackSWFallbackDroppedFrame's escalation and let that corruption persist.
// swFallbackPrimed is instead cleared only on a genuine fresh IDR resync
// (maybeCacheStream1IDR) or on RESET_GRAPHICS.
g.swFallbackDroppedCount = 0
g.swFallbackFirstDropTime = time.Time{}
g.stopInputWatchdog()
}
func (g *GfxHandler) maybeTriggerInputStall() {
g.inputWatchdogMu.Lock()
inputNS := g.inputWatchdogNS
g.inputWatchdogNS = 0
g.inputWatchdog = nil
g.inputWatchdogMu.Unlock()
if inputNS == 0 || g.h264dec == nil || g.h264dec.IsBroken() {
return
}
// Don't force-break a decoder that is already in IDR-wait state (e.g.
// after a soft reset). It is in a known recovery path; breaking it again
// would restart the soft-reset cycle unnecessarily.
if g.h264dec.NeedsIDR() {
return
}
// Don't force-break when the server has been idle (no video packets
// arriving). A static screen produces no frames but the decoder is
// healthy; only break when packets are flowing but output is absent.
if recvTime := g.h264dec.LastReceiveTime(); recvTime.IsZero() ||
time.Since(recvTime) >= inputStallSilentThreshold {
return
}
lastDecodedNS := g.lastDecodedFrame.Load()
if lastDecodedNS == 0 || lastDecodedNS >= inputNS {
return
}
silentFor := time.Since(time.Unix(0, lastDecodedNS))
if silentFor < inputStallSilentThreshold {
return
}
g.h264dec.ForceBroken(H264BrokenReasonHWStall)
slog.Debug("H.264: local input produced no new frame, marking decoder broken",
"silentFor", silentFor.Round(time.Millisecond),
"inputAgo", time.Since(time.Unix(0, inputNS)).Round(time.Millisecond))
}
// SetSendFunc sets the function used to send RDPGFX responses via DVC.
func (g *GfxHandler) SetSendFunc(fn func([]byte)) {
g.sendFn = fn
}
// SetDecoderBrokenCallback registers a function that is called once when the
// H.264 decoder becomes permanently unrecoverable (all soft resets exhausted).
// The callback should reconnect the RDP session so a fresh decoder can be
// created from scratch.
func (g *GfxHandler) SetDecoderBrokenCallback(fn func()) {
g.onDecoderBroken = fn
}
// SetKeyframeRequestFunc registers a function that is called after each
// soft decoder reset to ask the server for a fresh IDR keyframe. This
// speeds up recovery: without it the decoder waits for the server to
// spontaneously send a keyframe. A typical implementation calls
// pdu.SendRefreshRect with the current screen dimensions.
func (g *GfxHandler) SetKeyframeRequestFunc(fn func()) {
g.onKeyframeRequest = fn
}
// SetH264RawCallback registers a function that receives raw H.264 NAL unit
// data when the built-in decoder is unavailable (h264dec == nil). This
// allows the caller to hand off decoding to an external engine such as the
// browser WebCodecs VideoDecoder API.
//
// destX and destY are the top-left canvas coordinates of the decoded frame.
// isKey is true when the NAL data starts a new GOP (IDR frame).
func (g *GfxHandler) SetH264RawCallback(fn func(destX, destY, w, h int, isKey bool, data []byte, regions []int32)) {
g.onH264Raw = fn
}
// SetI420Callback registers a callback that receives I420 planar data when an
// H.264 frame is decoded and the underlying decoder supports I420 extraction.
// When set, H264 frames are NOT emitted via the normal OnBitmap path; the
// caller is responsible for rendering the I420 data directly (e.g. via an
// SDL2 IYUV texture). When the I420 fast path is used, the BGRA surface
// backing store is not updated for that frame.
// Set fn to nil to disable and revert to normal OnBitmap delivery.
func (g *GfxHandler) SetI420Callback(fn func(destX, destY, w, h int, y []byte, yStride int, u []byte, uStride int, v []byte, vStride int)) {
g.onI420 = fn
}
// SetNV12Callback registers a callback that receives native NV12 planar data
// when H.264 decoding produces NV12. When set for AVC420 frames, the normal
// OnBitmap path is bypassed for frames that can be delivered as NV12; callers
// should upload the Y and UV planes directly (for example with SDL2
// SDL_UpdateNVTexture). Set fn to nil to disable.
func (g *GfxHandler) SetNV12Callback(fn func(destX, destY, w, h int, y []byte, yStride int, uv []byte, uvStride int)) {
g.onNV12 = fn
}
// SetAVC444Disabled controls whether AVC444/AVC444v2 is advertised to the
// server. When disabled, CAPS_ADVERTISE only includes v8.0 and v8.1, so the
// server will encode frames using AVC420 (4:2:0) only and never send LC=2
// chroma-upgrade data. This avoids the colour degradation caused by servers
// (e.g. VirtualBox VRDE) that send LC=2 frames without including stream2 in
// LC=0 IDR packets. Must be called before the channel is opened.
func (g *GfxHandler) SetAVC444Disabled(v bool) {
g.avc444Disabled = v
}
// SetAVCDisabled advertises v10.x caps with RDPGFX_CAPS_FLAG_AVC_DISABLED while
// keeping the RDPGFX channel alive: the server then encodes with ClearCodec and
// RFX Progressive only ("RemoteFX mode"), never H.264. Must be called before
// the channel is opened.
func (g *GfxHandler) SetAVCDisabled(v bool) {
g.avcDisabled = v
}
// SetPduRecorder installs a callback that receives every wire-to-surface
// bitmap payload for offline replay analysis. Pass nil to disable.
func (g *GfxHandler) SetPduRecorder(fn func(kind byte, codecId uint16, surfW, surfH, x, y, w, h uint32, payload []byte)) {
g.pduRecord = fn
}
// Replay harness record kinds (see SetPduRecorder).
const (
ReplayKindWireToSurface1 byte = 1
ReplayKindWireToSurface2 byte = 2
ReplayKindCacheToSurface byte = 3
ReplayKindSurfaceToCache byte = 4
ReplayKindSolidFill byte = 5
ReplayKindSurfaceToSurf byte = 6
ReplayKindEvictCache byte = 7
ReplayKindResetGraphics byte = 8
)
// ReplaySurface ensures a replay surface with the given id/dimensions exists.
func (g *GfxHandler) ReplaySurface(id uint16, w, h uint16) {
g.onCreateSurface([]byte{
byte(id), byte(id >> 8),
byte(w), byte(w >> 8),
byte(h), byte(h >> 8),
0x20,
})
}
// ReplaySurfacePixels returns the live pixel buffer of a replay surface.
func (g *GfxHandler) ReplaySurfacePixels(id uint16) []byte {
if s, ok := g.surfaces[id]; ok {
return s.data
}
return nil
}
// ReplayPDU feeds one recorded payload through the normal decode path
// (offline replay harness for garbled-screen debugging).
func (g *GfxHandler) ReplayPDU(kind byte, codecId uint16, surfW, surfH uint16, x, y, w, h uint32, payload []byte) {
switch kind {
case ReplayKindWireToSurface1, ReplayKindWireToSurface2:
if kind == ReplayKindWireToSurface1 {
pdu := make([]byte, 17+len(payload))
binary.LittleEndian.PutUint16(pdu[0:], 0)
binary.LittleEndian.PutUint16(pdu[2:], codecId)
pdu[4] = 0x20
binary.LittleEndian.PutUint16(pdu[5:], uint16(x))
binary.LittleEndian.PutUint16(pdu[7:], uint16(y))
binary.LittleEndian.PutUint16(pdu[9:], uint16(x+w))
binary.LittleEndian.PutUint16(pdu[11:], uint16(y+h))
binary.LittleEndian.PutUint32(pdu[13:], uint32(len(payload)))
copy(pdu[17:], payload)
g.onWireToSurface1Decode(pdu)
return
}
pdu := make([]byte, 13+len(payload))
binary.LittleEndian.PutUint16(pdu[0:], 0)
binary.LittleEndian.PutUint16(pdu[2:], codecId)
binary.LittleEndian.PutUint32(pdu[4:], 0)
pdu[8] = 0x20
binary.LittleEndian.PutUint32(pdu[9:], uint32(len(payload)))
copy(pdu[13:], payload)
g.onWireToSurface2Decode(pdu)
case ReplayKindCacheToSurface:
g.onCacheToSurface(payload)
case ReplayKindSurfaceToCache:
g.onSurfaceToCache(payload)
case ReplayKindSolidFill:
g.onSolidFill(payload)
case ReplayKindSurfaceToSurf:
g.onSurfaceToSurface(payload)
case ReplayKindEvictCache:
g.onEvictCacheEntry(payload)
case ReplayKindResetGraphics:
g.onResetGraphics(payload)
}
}
// OnChannelCreated is called after the DVC CREATE_RSP has been sent.
// It sends RESET_GRAPHICS (MS-RDPEGFX 3.2.1.3:通道初始化后、caps 协商前,
// 客户端先发 RESET_GRAPHICS 声明会话尺寸) 再发 CAPS_ADVERTISE 启动管线。
// 服务器收到后重建表面并重发完整桌面帧——首次连接无害,断线重附会话时
// 则强制全量重绘,消除重连后的画面残留。
// 实测 Win10 19041 依赖该 reset 启动图形管线(省略后服务器不发 caps
// confirm 直接断连);Server 2025 对 reset 后的管线重建会崩(0x112F),
// 由上层对该错误做一次性 bitmap 回退兜底。
func (g *GfxHandler) OnChannelCreated() {
g.sendResetGraphics()
g.sendCapsAdvertise()
}
// SetSessionSize 注入会话尺寸,供通道初始化的 RESET_GRAPHICS 使用。
func (g *GfxHandler) SetSessionSize(width, height uint16) {
g.sessionW = uint32(width)
g.sessionH = uint32(height)
}
// sendResetGraphics sends RDPGFX_RESET_GRAPHICS_PDU: width(4) + height(4) +
// monitorCount(4)=1,负载共 12 字节。
// 注意:规范 2.2.2.10 要求整个 PDU 补齐到 340 字节(含 20 字节/显示器的
// MONITOR_DEF 数组),但实测 Win10 19041 只接受这个 12 字节负载的短格式
// (340 字节规范格式会令其静默拒绝 GFX 通道回退传统位图;全零 MONITOR_DEF
// 同样被拒)。Server 2025 对任何 reset 格式都会崩(0x112F),由前端
// ERRINFO_GRAPHICS_* 一次性 bitmap 回退兜底,见 app.js scheduleReconnect。
func (g *GfxHandler) sendResetGraphics() {
w, h := g.sessionW, g.sessionH
if w == 0 || h == 0 {
return
}
payload := make([]byte, 12)
binary.LittleEndian.PutUint32(payload[0:], w)
binary.LittleEndian.PutUint32(payload[4:], h)
binary.LittleEndian.PutUint32(payload[8:], 1)
g.sendPdu(cmdidResetGraphics, payload)
}
// sendCapsAdvertise sends RDPGFX_CAPS_ADVERTISE_PDU to the server.
// The client must advertise its capabilities before the server will
// send any graphics data (MS-RDPEGFX 2.2.3.1).
func (g *GfxHandler) sendCapsAdvertise() {
p := pduBufPool.Get().([]byte)[:0]
defer pduBufPool.Put(p[:0])
// AVC capsets are advertised when we can deliver decoded frames either
// in-process (h264dec) or by handing the raw NALs off to the embedder
// (onH264Raw, used by the WASM build to forward to WebCodecs). Without
// either, the v8.0+AVCDisabled fallback below forces the server to
// reject RDPGFX and use legacy bitmap PDUs.
if g.h264dec != nil || g.onH264Raw != nil {
if g.avc444Disabled {
// AVC444 关闭(AVC420-only):v8.0/v8.1 之外再广告 v10.0——
// Windows 服务器在 v10+ caps 下才启用 H.264 编码;只到 v10.0
//(不含 v10.1+)即不暗示 AVC444v2(LC=2 色度升级流),服务器
// 最高使用 AVC420。
p = binary.LittleEndian.AppendUint16(p, 3) // capsSetCount
// v8.0 — baseline fallback (no AVC)
p = binary.LittleEndian.AppendUint32(p, capVersion8)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagThinClient)
// v8.1 — AVC420 via explicit flag
p = binary.LittleEndian.AppendUint32(p, capVersion81)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagSmallCache|capFlagAVC420Enabled)
// v10.0 — H.264 启用门槛,仅 AVC420
p = binary.LittleEndian.AppendUint32(p, capVersion10)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagSmallCache|capFlagAVC420Enabled)
g.sendPdu(cmdidCapsAdvertise, p)
slog.Debug("RDPGFX: sent CAPS_ADVERTISE (v10.0..v8.0, AVC420 only)")
} else {
// Advertise capsets in ascending order (v8.0 → v10.7), matching
// rdpyqt / FreeRDP layout so servers pick the highest common version.
p = binary.LittleEndian.AppendUint16(p, 11) // capsSetCount
// v8.0 — baseline fallback (no AVC)
p = binary.LittleEndian.AppendUint32(p, capVersion8)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagThinClient)
// v8.1 — AVC420 via explicit flag
p = binary.LittleEndian.AppendUint32(p, capVersion81)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagSmallCache|capFlagAVC420Enabled)
// v10.0
p = binary.LittleEndian.AppendUint32(p, capVersion10)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagSmallCache)
// v10.1 — 16-byte capsData (12 zero bytes after flags)
p = binary.LittleEndian.AppendUint32(p, capVersion101)
p = binary.LittleEndian.AppendUint32(p, 16)
p = binary.LittleEndian.AppendUint32(p, 0)
p = binary.LittleEndian.AppendUint32(p, 0)
p = binary.LittleEndian.AppendUint32(p, 0)
p = binary.LittleEndian.AppendUint32(p, 0)
// v10.2
p = binary.LittleEndian.AppendUint32(p, capVersion102)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagSmallCache)
// v10.3
p = binary.LittleEndian.AppendUint32(p, capVersion103)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, 0)
// v10.4
p = binary.LittleEndian.AppendUint32(p, capVersion104)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagSmallCache)
// v10.5
p = binary.LittleEndian.AppendUint32(p, capVersion105)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagSmallCache)
// v10.6
p = binary.LittleEndian.AppendUint32(p, capVersion106)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagSmallCache)
// v10.6.1
p = binary.LittleEndian.AppendUint32(p, capVersion1061)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagSmallCache)
// v10.7
p = binary.LittleEndian.AppendUint32(p, capVersion107)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagSmallCache)
g.sendPdu(cmdidCapsAdvertise, p)
slog.Debug("RDPGFX: sent CAPS_ADVERTISE (v10.7..v8.0, AVC enabled)")
}
} else if g.avcDisabled {
// RemoteFX 模式:保留 RDPGFX 通道、显式禁用 AVC——服务端继续用
// ClearCodec + RFX Progressive 编码(capset 组合与 FreeRDP 无 H264
// 构建一致)。此前仅广告 v8.0+AVC_DISABLED(该标志自 v10 才在规范
// 中定义),服务器会拒绝整个 GFX 通道退化成传统 RLE 位图更新,
// 拖动窗口时带宽可达 80Mbps+。
// v10.3 特意不带 SMALL_CACHE(该版本固定 16MB 缓存槽)。
noAVC := capFlagSmallCache | capFlagAVCDisabled
p = binary.LittleEndian.AppendUint16(p, 11) // capsSetCount
// v8.0 / v8.1 — AVC_DISABLED 在 v8 未定义,仅 SMALL_CACHE
p = binary.LittleEndian.AppendUint32(p, capVersion8)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagSmallCache)
p = binary.LittleEndian.AppendUint32(p, capVersion81)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagSmallCache)
// v10.x — SMALL_CACHE | AVC_DISABLED
p = binary.LittleEndian.AppendUint32(p, capVersion10)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, noAVC)
// v10.1 — 16 字节布局(flags + 12 字节保留域)
p = binary.LittleEndian.AppendUint32(p, capVersion101)
p = binary.LittleEndian.AppendUint32(p, 16)
p = binary.LittleEndian.AppendUint32(p, 0)
p = binary.LittleEndian.AppendUint32(p, 0)
p = binary.LittleEndian.AppendUint32(p, 0)
p = binary.LittleEndian.AppendUint32(p, 0)
p = binary.LittleEndian.AppendUint32(p, capVersion102)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, noAVC)
p = binary.LittleEndian.AppendUint32(p, capVersion103)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, capFlagAVCDisabled)
p = binary.LittleEndian.AppendUint32(p, capVersion104)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, noAVC)
p = binary.LittleEndian.AppendUint32(p, capVersion105)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, noAVC)
p = binary.LittleEndian.AppendUint32(p, capVersion106)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, noAVC)
p = binary.LittleEndian.AppendUint32(p, capVersion1061)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, noAVC)
p = binary.LittleEndian.AppendUint32(p, capVersion107)
p = binary.LittleEndian.AppendUint32(p, 4)
p = binary.LittleEndian.AppendUint32(p, noAVC)
g.sendPdu(cmdidCapsAdvertise, p)
slog.Info("RDPGFX: sent CAPS_ADVERTISE (v10.7..v8.0, AVC disabled, RFX/Clear only)")
} else {
// 传统位图回退模式:v8.0 + AVC_DISABLED 使服务器"拒绝 RDPGFX 通道",
// 退回传统 RLE 位图更新。适用于图形管线损坏(任何 RDPGFX 会话都会
// 触发 ERRINFO 0x112F 崩溃)或不支持 H264 的服务器。
// 注意:AVC_DISABLED 只在 v8 caps 上安全(服务器会忽略非法组合)。
p = binary.LittleEndian.AppendUint16(p, 1) // capsSetCount
p = binary.LittleEndian.AppendUint32(p, capVersion8)
p = binary.LittleEndian.AppendUint32(p, 4) // capsDataLength
p = binary.LittleEndian.AppendUint32(p, capFlagThinClient|capFlagSmallCache|capFlagAVCDisabled)
g.sendPdu(cmdidCapsAdvertise, p)
slog.Debug("RDPGFX: sent CAPS_ADVERTISE (v8.0, AVC disabled → legacy bitmap fallback)")
}
}
// ZGFX segment descriptors (MS-RDPEGFX 2.2.4)
const (
zgfxSingle = 0xE0
zgfxMultipart = 0xE1
zgfxCompressedRDP8 = 0x04
)
// Process handles a complete RDPGFX payload (may contain multiple PDUs).
// Data arrives wrapped in ZGFX (RDP8 Bulk Compression) segments (MS-RDPEGFX 2.2.4).
//
// Called on the network read goroutine. Decompression happens here;
// the decompressed payload is then queued for asynchronous processing
// (including frame ACKs and decode) on the decode goroutine.
// This keeps the read goroutine free from any socket.Write calls that
// could cause TCP deadlock when both sides try to write simultaneously.
func (g *GfxHandler) Process(data []byte) {
defer func() {
if r := recover(); r != nil {
slog.Error("RDPGFX: panic in Process", "err", r)
}
}()
if len(data) < 1 {
return
}
var decompressed []byte
var decompPooled bool
descriptor := data[0]
switch descriptor {
case zgfxSingle:
if len(data) < 2 {
return
}
decompressed, decompPooled = g.decompressSegment(data[1:])
case zgfxMultipart:
decompressed, decompPooled = g.decompressMultipart(data[1:])
default:
slog.Warn("RDPGFX: unknown ZGFX descriptor", "descriptor", descriptor)
decompressed = data
}
if len(decompressed) == 0 {
return
}
// decompressSegment / decompressMultipart already return owned
// buffers (freshly allocated or copied from input) so we can hand
// the slice directly to the async decode goroutine.
pkt := decodePkt{data: decompressed, pooled: decompPooled}
// 阻塞式入队,绝不丢弃(协议正确性):被丢帧的区域服务器不会重发
//(FrameAck 已视为送达),表现为永久黑块/花屏。满时靠背压阻塞读
// 协程 → TCP 流控自然减速服务器;配合 onEndFrame 的 SUSPEND ACK
//(队列深时显式请求服务器暂停发送,排空后真实 queueDepth 恢复),
// 把内存占用约束在 decodeCh 容量级别。
g.decodeCh <- pkt
}
// ackDroppedFrames scans decompressed PDU data for EndFrame commands
// and sends ACKs for them. Called on the read goroutine when decodeCh
// is full and the message is being dropped. Without this, dropped
// EndFrames would leave the server's outstanding-frame count stuck,
// eventually causing it to stop sending entirely.
//
// queueDepth is set to suspendFrameAcknowledge (0xFFFFFFFF) so the
// server suspends sending new frames. As the decodeLoop drains the
// existing queue and sends ACKs with lower queueDepth values, the server
// will automatically resume (MS-RDPEGFX 2.2.2.8).
func (g *GfxHandler) ackDroppedFrames(pkt decodePkt) {
data := pkt.data
defer func() {
if pkt.pooled {
releaseBitmapBuf(data)
}
}()
for offset := 0; offset+headerSize <= len(data); {
cmdId := binary.LittleEndian.Uint16(data[offset:])
pduLength := binary.LittleEndian.Uint32(data[offset+4:])
if pduLength < uint32(headerSize) || int(pduLength) > len(data)-offset {
break
}
if cmdId == cmdidEndFrame {
pduData := data[offset+headerSize : offset+int(pduLength)]
if len(pduData) >= 4 {
g.sendFrameAck(binary.LittleEndian.Uint32(pduData), suspendFrameAcknowledge)
}
}
offset += int(pduLength)
}
}
// decompressSegment handles a single ZGFX segment (after the descriptor byte).
// First byte is RDP8_BULK_ENCODED_DATA header:
//
// bits 0-3: compression type (0x04 = RDP8)
// bit 5: PACKET_COMPRESSED (0x20)
//
// Always returns pooled=true: the returned slice was acquired from
// bitmapBufPool (directly or via Decompress) and must be released with
// releaseBitmapBuf once the caller is done with it.
func (g *GfxHandler) decompressSegment(seg []byte) ([]byte, bool) {
if len(seg) < 1 {
return nil, false
}
header := seg[0]
payload := seg[1:]
if header&0x20 != 0 {
// Acquire a pool buffer as the initial backing for Decompress output.
// Decompress may grow beyond it; the returned slice (not buf) must be
// released by the caller. Any over-small buf that gets replaced is
// abandoned to GC — the pool converges to the right size over time.
buf := acquireBitmapBuf(len(payload) * 3)
return g.zgfx.Decompress(payload, buf), true
}
g.zgfx.historyWrite(payload)
// Return a pooled copy: payload aliases the caller's network buffer, which
// will be reused on the next read. Callers hand the slice off to the async
// decode goroutine and must own the memory.
buf := acquireBitmapBuf(len(payload))
copy(buf, payload)
return buf, true
}
// decompressMultipart handles ZGFX multipart segments and returns the
// concatenated decompressed data (without processing PDUs).
// Returns a slice acquired from bitmapBufPool; caller must release it.
func (g *GfxHandler) decompressMultipart(data []byte) ([]byte, bool) {
if len(data) < 6 {
return nil, false
}
// Direct slice indexing — avoids bytes.NewReader and per-field io.ReadFull.
segCount := binary.LittleEndian.Uint16(data[0:])
uncompSize := binary.LittleEndian.Uint32(data[2:])
offset := 6
// Pre-allocate to the advertised uncompressed size to avoid repeated
// buffer growths as each segment is appended.
buf := acquireBitmapBuf(int(uncompSize))
result := buf[:0]
for range segCount {
if offset+4 > len(data) {
break
}
segSize := int(binary.LittleEndian.Uint32(data[offset:]))
offset += 4
if offset+segSize > len(data) {
break
}
segData := data[offset : offset+segSize]
offset += segSize
raw, rawPooled := g.decompressSegment(segData)
if raw != nil {
result = append(result, raw...)
if rawPooled {
releaseBitmapBuf(raw)
}
}
}
if len(result) == 0 {
releaseBitmapBuf(buf)
return nil, false
}
// If result grew beyond buf, buf was abandoned; result is the new owner.
return result, true
}
// decodeLoop runs in a dedicated goroutine, reading decompressed PDU data
// from decodeCh and dispatching all processing — including frame ACKs and
// heavy decode work. Keeping socket.Write calls off the read goroutine
// avoids TCP deadlock (where both sides try to write while neither reads).
// It automatically restarts on panic, unless Close() has been called.
//
// decodeLoop owns the h264dec lifecycle: it is the sole caller of Decode()
// and it frees h264dec on final exit (when doneCh is closed) to avoid a
// use-after-free race with Close() freeing the AVCodecContext from a
// different goroutine while Decode() holds it.
func (g *GfxHandler) decodeLoop() {
defer func() {
if r := recover(); r != nil {
slog.Error("RDPGFX: panic in decodeLoop, restarting", "err", r, "stack", string(debug.Stack()))
select {
case <-g.doneCh:
// Shutting down — free h264dec resources on this final exit.
if g.h264dec != nil {
g.h264dec.Close()
g.h264dec = nil
}
if g.h264dec2 != nil {
g.h264dec2.Close()
g.h264dec2 = nil
}
default:
go g.decodeLoop()
}
return
}
// Normal exit triggered by doneCh being closed.
if g.h264dec != nil {
g.h264dec.Close()
g.h264dec = nil
}
if g.h264dec2 != nil {
g.h264dec2.Close()
g.h264dec2 = nil
}
}()
slog.Debug("RDPGFX: decodeLoop started")
for {
select {
case <-g.doneCh:
return
case pkt := <-g.decodeCh:
g.decodePDUs(pkt.data)
if pkt.pooled {
releaseBitmapBuf(pkt.data)
}
case <-g.watchdogCh:
// A background timer in ffmpegDecoder fired because the stall-probe
// or IDR-wait timeout expired while the server was sending no frames
// (static screen → near-0 fps), or because local input produced no
// new frame while the decoder had already been silent too long.
slog.Debug("H.264: watchdog triggered, checking decoder state")
g.maybeTriggerInputStall()
g.maybeRenegotiateCapabilities()
g.maybeNotifyDecoderBroken()
}
}
}
// decodePDUs processes all PDUs in decompressed data.
// Every PDU is decoded: silently dropping CaVideo/progressive updates under
// backpressure permanently diverges the canvas (the server considers acked
// frames delivered and never resends them — the corruption manifests as
// persistent noise after a fast window drag). Pacing is the server's job
// via the queueDepth we report in FRAME_ACKNOWLEDGE.
func (g *GfxHandler) decodePDUs(data []byte) {
start := time.Now()
defer func() {
// 单条消息解码耗时的 EMA(α=0.3):诊断解码吞吐是否跟不上到达速率
if cost := time.Since(start).Microseconds(); cost > 0 {
for {
cur := g.decUs.Load()
next := cur * 7 / 10
if next == 0 {
next = cost
} else {
next += cost * 3 / 10
}
if g.decUs.CompareAndSwap(cur, next) {
return
}
}
}
}()
for offset := 0; offset+headerSize <= len(data); {
cmdId := binary.LittleEndian.Uint16(data[offset:])
pduLength := binary.LittleEndian.Uint32(data[offset+4:])
if pduLength < uint32(headerSize) || int(pduLength) > len(data)-offset {
break
}
pduData := data[offset+headerSize : offset+int(pduLength)]
g.dispatchDecode(cmdId, pduData)
offset += int(pduLength)
}
}
// dispatchDecode routes a single PDU.
func (g *GfxHandler) dispatchDecode(cmdId uint16, data []byte) {
if int(cmdId) < len(g.cmdCounts) {
g.cmdCounts[cmdId].Add(1)
}
switch cmdId {
case cmdidCapsConfirm:
g.onCapsConfirm(data)
case cmdidResetGraphics:
g.onResetGraphics(data)
case cmdidCreateSurface:
g.onCreateSurface(data)
case cmdidDeleteSurface:
g.onDeleteSurface(data)
case cmdidMapSurfaceToOutput:
g.onMapSurfaceToOutput(data)
case cmdidStartFrame:
// 记录帧解码起点,供 QoE 周期上报(timeDiffSE)使用
g.frameDecodeStart.Store(time.Now().UnixNano())
case cmdidSurfaceToSurface:
g.onSurfaceToSurface(data)
case cmdidSurfaceToCache:
g.onSurfaceToCache(data)
case cmdidEndFrame:
g.onEndFrame(data) // always ACK, even under backpressure
case cmdidWireToSurface1:
g.onWireToSurface1Decode(data)
case cmdidWireToSurface2:
g.onWireToSurface2Decode(data)
case cmdidSolidFill:
g.onSolidFill(data)
case cmdidCacheToSurface:
g.onCacheToSurface(data)
case cmdidEvictCacheEntry:
g.onEvictCacheEntry(data)
case cmdidCacheImportOffer:
g.onCacheImportOffer()
case cmdidMapSurfaceToWindow, cmdidMapSurfaceToScaledWindow:
// ignored — we don't support per-window mapping
case cmdidCacheImportReply:
g.onCacheImportReply(data)
case cmdidDeleteEncodingContext, cmdidQoeFrameAcknowledge:
// no client state to maintain for these
case cmdidMapSurfaceToScaledOutput:
g.onMapSurfaceToScaledOutput(data)
default:
slog.Debug("RDPGFX: unhandled cmd", "cmdId", cmdId)
}
}
// writeLoop runs in a dedicated goroutine. It reads serialized ACK
// PDUs from ackCh and sends each one via sendFn. Every ACK must reach
// the server — the server tracks outstanding frames individually and
// stops sending if ACKs are missing. Automatically restarts on panic,
// unless Close() has been called.
func (g *GfxHandler) writeLoop() {
defer func() {
if r := recover(); r != nil {
slog.Error("RDPGFX: panic in writeLoop, restarting", "err", r)
select {
case <-g.doneCh:
// Shut down; do not restart.
default:
go g.writeLoop()
}
}
}()
for {
select {
case <-g.doneCh:
return
case pdu := <-g.ackCh:
if g.sendFn != nil {
g.sendFn(pdu)
}
ackPDUPool.Put(pdu)
}
}
}
// sendPdu sends a PDU synchronously. Used for rare control messages
// (CapsAdvertise, CacheImportReply) that must not be dropped.
// pduBufPool reuses scratch byte slices for assembling outbound PDU frames,
// avoiding per-call heap allocations on the sendPdu hot path.
var pduBufPool = sync.Pool{
New: func() any { return make([]byte, 0, headerSize+256) },
}
// ackPDUPool reuses the fixed-size 20-byte slices used for FRAME_ACKNOWLEDGE
// PDUs (~60/s during video), eliminating per-frame heap allocations.
var ackPDUPool = sync.Pool{
New: func() any { return make([]byte, 20) },
}
func (g *GfxHandler) sendPdu(cmdId uint16, payload []byte) {
if g.sendFn == nil {
return
}
buf := pduBufPool.Get().([]byte)
buf = buf[:0]
buf = binary.LittleEndian.AppendUint16(buf, cmdId)
buf = binary.LittleEndian.AppendUint16(buf, 0) // flags
buf = binary.LittleEndian.AppendUint32(buf, uint32(headerSize+len(payload)))
buf = append(buf, payload...)
g.sendFn(buf)
pduBufPool.Put(buf[:0])
}
// --- Command Handlers ---
// DebugSurfacePixel 诊断:返回第一个 mapped surface 上 (x,y) 的 BGRA 原始值
func (g *GfxHandler) DebugSurfacePixel(x, y int) (uint8, uint8, uint8, uint8, bool) {
for _, s := range g.surfaces {
if !s.mapped {
continue
}
if x < 0 || y < 0 || x >= int(s.width) || y >= int(s.height) {
return 0, 0, 0, 0, false
}
i := (y*int(s.width) + x) * 4
return s.data[i], s.data[i+1], s.data[i+2], s.data[i+3], true
}
return 0, 0, 0, 0, false
}
func (g *GfxHandler) onCapsConfirm(data []byte) {
if len(data) < 12 {
slog.Debug("RDPGFX: CAPS_CONFIRM received (short)")
return
}
version := binary.LittleEndian.Uint32(data[0:])
dataLen := binary.LittleEndian.Uint32(data[4:])
flags := uint32(0)
if dataLen >= 4 {
flags = binary.LittleEndian.Uint32(data[8:])
}
slog.Info("RDPGFX: CAPS_CONFIRM", "version", fmt.Sprintf("0x%08X", version), "flags", fmt.Sprintf("0x%08X", flags))
// caps 交换完成后上报持久缓存条目(MS-RDPEGFX 2.2.2.16)。此时连接的
// 其他 GFX 状态尚未开始流动,是最安全的上报时点。
g.sendCacheImportOffer()
}
func (g *GfxHandler) onResetGraphics(data []byte) {
if len(data) < 12 {
return
}
if g.pduRecord != nil {
g.pduRecord(ReplayKindResetGraphics, 0, 0, 0, 0, 0, 0, 0, data)
}
w := binary.LittleEndian.Uint32(data[0:])
h := binary.LittleEndian.Uint32(data[4:])
// RESET_GRAPHICS 意味着服务端图形子系统重启(常见于 0x112F 内部错误),
// 后续帧流是否恢复依赖此处的状态清理,值得始终留痕。
slog.Info("RDPGFX: RESET_GRAPHICS", "w", w, "h", h)
g.surfaces = make(map[uint16]*surface)
g.clearCtx = newClearCodecCtx()
g.framesDecoded.Store(0)
g.softResetCount = 0
g.noIDRSoftResetCount = 0
g.decoderBrokenNotified = false
g.lc2EverDecoded = false
g.stream2EverSeen = false
g.auxDecoderNoIDRRetries = 0
g.lc2PermanentlyDegraded = false
g.lastKeyframeRequest = time.Time{}
g.lastStream1IDR = g.lastStream1IDR[:0]
g.lastStream1IDRTime = time.Time{}
g.lastStream1IDRFrame = 0
g.swFallbackPrimed = false
g.swFallbackDroppedCount = 0
g.swFallbackFirstDropTime = time.Time{}
g.lastDecodedFrame.Store(0)
g.stopInputWatchdog()
if g.h264dec != nil {
g.h264dec.Close()
g.h264dec = newH264DecoderWithWatchdog(g.watchdogCh)
}
if g.h264dec2 != nil {
g.h264dec2.Close()
// Keep h264dec2 nil; primeAuxDecoder will recreate it on the next stream2 IDR
// so the fresh decoder is always primed before receiving LC=2 P-frames.
g.h264dec2 = nil
}
g.avc444YPlane = avc444YPlane{}
g.avc444IDRYPlane = avc444YPlane{}
g.progressive.Reset()
}
func (g *GfxHandler) onCreateSurface(data []byte) {
if len(data) < 7 {
return
}
id := binary.LittleEndian.Uint16(data[0:])
w := binary.LittleEndian.Uint16(data[2:])
h := binary.LittleEndian.Uint16(data[4:])
f := data[6]
slog.Debug("RDPGFX: CREATE_SURFACE", "id", id, "w", w, "h", h)
g.surfaces[id] = &surface{
width: w, height: h, format: f,
data: make([]byte, int(w)*int(h)*4),
shadowStale: true,
}
}
func (g *GfxHandler) onDeleteSurface(data []byte) {
if len(data) < 2 {
return
}
id := binary.LittleEndian.Uint16(data)
delete(g.surfaces, id)
}
func (g *GfxHandler) onMapSurfaceToOutput(data []byte) {
if len(data) < 12 {
return
}
id := binary.LittleEndian.Uint16(data[0:])
// data[2:4] = reserved
ox := binary.LittleEndian.Uint32(data[4:])
oy := binary.LittleEndian.Uint32(data[8:])
slog.Debug("RDPGFX: MAP_SURFACE", "id", id, "ox", ox, "oy", oy)
if s, ok := g.surfaces[id]; ok {
s.outputX = ox
s.outputY = oy
s.mapped = true
}
}
func (g *GfxHandler) onMapSurfaceToScaledOutput(data []byte) {
if len(data) < 20 {
return
}
id := binary.LittleEndian.Uint16(data[0:])
// data[2:4] = reserved
ox := binary.LittleEndian.Uint32(data[4:])
oy := binary.LittleEndian.Uint32(data[8:])
// data[12:16] = targetWidth, data[16:20] = targetHeight (unused)
slog.Debug("RDPGFX: MAP_SURFACE_SCALED", "id", id, "ox", ox, "oy", oy)
if s, ok := g.surfaces[id]; ok {
s.outputX = ox
s.outputY = oy
s.mapped = true
}
}
// sendFrameAck builds and queues a FRAME_ACKNOWLEDGE PDU.
// Safe to call from any goroutine (uses atomic framesDecoded).
// The PDU is serialized directly into a 20-byte slice to avoid
// the two bytes.Buffer allocations the previous implementation required.
//
// queueDepth is reported to the server so it can adjust encoding quality
// and frame rate based on the client's decode backlog. Pass
// suspendFrameAcknowledge (0xFFFFFFFF) to ask the server to suspend new
// frames until a subsequent ACK with a lower value is received.
func (g *GfxHandler) sendFrameAck(frameId uint32, queueDepth uint32) {
decoded := g.framesDecoded.Add(1)
// 8-byte RDPGFX header + 12-byte FRAME_ACKNOWLEDGE payload = 20 bytes.
pdu := ackPDUPool.Get().([]byte)
binary.LittleEndian.PutUint16(pdu[0:], cmdidFrameAcknowledge)
// pdu[2:4] = flags (0) — zero value
binary.LittleEndian.PutUint16(pdu[2:], 0)
binary.LittleEndian.PutUint32(pdu[4:], 20) // total PDU length
binary.LittleEndian.PutUint32(pdu[8:], queueDepth)
binary.LittleEndian.PutUint32(pdu[12:], frameId)
binary.LittleEndian.PutUint32(pdu[16:], decoded)
select {
case g.ackCh <- pdu:
default:
ackPDUPool.Put(pdu)
slog.Warn("RDPGFX: ackCh full, ACK dropped")
}
}
func (g *GfxHandler) onEndFrame(data []byte) {
if len(data) < 4 {
return
}
// 帧间隔 EMA(α=0.3):诊断服务器帧率与传输节奏
now := time.Now().UnixNano()
if prev := g.lastFrameAt.Swap(now); prev != 0 {
intv := (now - prev) / 1000 // µs
for {
cur := g.frameIntvUs.Load()
next := cur * 7 / 10
if next == 0 {
next = intv
} else {
next += intv * 3 / 10
}
if g.frameIntvUs.CompareAndSwap(cur, next) {
break
}
}
}
realDepth := uint32(len(g.decodeCh))
if hint := g.queueDepthHint.Load(); hint > realDepth {
realDepth = hint
}
frameId := binary.LittleEndian.Uint32(data)
// 背压:队列深时显式 SUSPEND(queueDepth=0xFFFFFFFF 请求服务器暂停
// 发送),排空到低水位后以真实 queueDepth 恢复(MS-RDPEGFX 2.2.2.8)。
// 这是协议规定的限速通道——配合读协程的阻塞式入队,保证任何已发送
// 的帧都会被完整解码上屏,不存在"ACK 了却没画"的黑块来源。
switch depth := uint32(len(g.decodeCh)); {
case depth > 48:
g.suspended.Store(true)
case depth <= 16:
g.suspended.Store(false)
}
if g.suspended.Load() {
g.sendFrameAck(frameId, suspendFrameAcknowledge)
return
}
g.sendFrameAck(frameId, realDepth)
// QoE 周期上报(每 30 帧):mstsc/FreeRDP 用它向服务器反馈解码时延,
// 服务器据此做码率/画质自适应。cmdId=0x16,布局对齐 FreeRDP
// rdpgfx_send_qoe_frame_acknowledge_pdu(header8 + frameId4 +
// timestamp4 + timeDiffSE2 + timeDiffEDR2 = 20 字节,与 ack 同尺寸,
// 复用 ackCh 写循环)。
if frameId%30 == 0 {
now := time.Now().UnixNano()
startNs := g.frameDecodeStart.Load()
if startNs == 0 {
startNs = now
}
diffSE := (now - startNs) / int64(time.Millisecond)
if diffSE < 0 {
diffSE = 0
}
if diffSE > 65000 {
diffSE = 65000
}
edr := g.decUs.Load() / 1000
if edr > 65000 {
edr = 65000
}
pdu := ackPDUPool.Get().([]byte)
binary.LittleEndian.PutUint16(pdu[0:], cmdidQoeFrameAcknowledge)
binary.LittleEndian.PutUint16(pdu[2:], 0)
binary.LittleEndian.PutUint32(pdu[4:], 20)
binary.LittleEndian.PutUint32(pdu[8:], frameId)
binary.LittleEndian.PutUint32(pdu[12:], uint32(startNs/int64(time.Millisecond)))
binary.LittleEndian.PutUint16(pdu[16:], uint16(diffSE))
binary.LittleEndian.PutUint16(pdu[18:], uint16(edr))
select {
case g.ackCh <- pdu:
default:
ackPDUPool.Put(pdu)
}
}
}
// DiagStats 返回实时诊断指标:累计解码帧数、当前解码队列深度、
// 帧间隔(毫秒)与单消息解码耗时 EMA(微秒)。解码耗时必须保留
// 微秒精度——WebCodecs 路径单消息只有几百微秒,取整到毫秒恒为 0。
// 统计循环定期输出,用于判断解码吞吐是否跟不上到达速率、服务器帧率是否异常。
func (g *GfxHandler) DiagStats() (frames, qdepth, fintvMs, decUs int64) {
return int64(g.framesDecoded.Load()),
int64(len(g.decodeCh)),
g.frameIntvUs.Load() / 1000,
g.decUs.Load()
}
// SetQueueDepthHint sets a minimum queueDepth to report in FRAME_ACKNOWLEDGE
// PDUs (MS-RDPEGFX 2.2.2.8). The server uses this value to pace its frame
// rate and encoding quality: a larger value signals that the client's decode
// queue is full, causing the server to slow down or reduce quality.
//
// A hint of 0 (the default) means "report the real queue length".
// Values in the range 10–100 are typical for moderate throttling.
// Use suspendFrameAcknowledge (0xFFFFFFFF) to pause the stream entirely
// (the stream resumes automatically when the hint is cleared).
// CodecStats returns cumulative surface-bitmap bytes per codec id
// (WTS1 + WTS2 combined). Useful for bandwidth diagnostics: e.g. a session
// dominated by codec 3 (RemoteFX) vs 0x0B/0x0E (AVC420/444) vs 8 (ClearCodec).
// Keys 100+cmdId carry the per-command PDU counters (100+0..100+63).
func (g *GfxHandler) CodecStats() map[uint16]int64 {
out := make(map[uint16]int64, len(g.codecBytes))
for i := range g.codecBytes {
if v := g.codecBytes[i].Load(); v != 0 {
out[uint16(i)] = v
}
}
for i := range g.cmdCounts {
if v := g.cmdCounts[i].Load(); v != 0 {
out[uint16(100+i)] = v
}
}
return out
}
func (g *GfxHandler) SetQueueDepthHint(depth uint32) {
g.queueDepthHint.Store(depth)
}
// onWireToSurface1Decode handles RDPGFX_WIRE_TO_SURFACE_PDU_1 (MS-RDPEGFX 2.2.2.1).
func (g *GfxHandler) onWireToSurface1Decode(data []byte) {
if len(data) < 17 {
return
}
// Parse fixed header fields via direct binary indexing (avoids bytes.NewReader
// and per-field io.ReadFull overhead on the hot H.264 path).
surfId := binary.LittleEndian.Uint16(data[0:])
codecId := binary.LittleEndian.Uint16(data[2:])
pixFmt := data[4]
left := binary.LittleEndian.Uint16(data[5:])
top := binary.LittleEndian.Uint16(data[7:])
right := binary.LittleEndian.Uint16(data[9:])
bottom := binary.LittleEndian.Uint16(data[11:])
bmpLen := binary.LittleEndian.Uint32(data[13:])
if int(bmpLen) > len(data)-17 {
return
}
// 服务端在拖动时会夹杂 pixFmt 非法(0x62 等)的垃圾微型更新,
// FreeRDP 对非法 pixFmt 直接拒绝 PDU;照单全收会在屏幕上画出
// 彩色杂线(拖动花屏的组成部分)。
if pixFmt != pixelFormatXRGB8888 && pixFmt != pixelFormatARGB8888 {
return
}
bmpData := data[17 : 17+int(bmpLen)]
if slog.Default().Enabled(nil, slog.LevelDebug) {
slog.Debug("RDPGFX: WTS1", "surfId", surfId, "codecId", codecId,
"w", right-left, "h", bottom-top, "bmpLen", bmpLen)
}
w := int(right - left)
h := int(bottom - top)
if w <= 0 || h <= 0 {
return
}
if int(codecId) < len(g.codecBytes) {
g.codecBytes[codecId].Add(int64(bmpLen))
}
s, ok := g.surfaces[surfId]
if !ok {
return
}
if g.pduRecord != nil {
g.pduRecord(1, codecId, uint32(s.width), uint32(s.height),
uint32(left), uint32(top), uint32(w), uint32(h), bmpData)
}
// CaVideo (0x0003) carries RFX tile-encoded data; decode onto the
// persistent surface buffer like the progressive codec in WTS2.
if codecId == codecCaVideo {
rects := g.rfx.Decode(bmpData, int(left), int(top), s.data, int(s.width), int(s.height))
g.emitCaVideoRects(s, rects)
return
}
var decoded []byte
var avcRegions []avcRect
owned := false // true ⇒ decoded buffer is from bitmapBufPool and must be released
switch codecId {
case codecUncompressed:
decoded = decodeUncompressed(bmpData, w, h, pixFmt)
owned = true
case codecPlanar:
decoded = decodePlanar(bmpData, w, h)
owned = true
case codecAVC420:
destX := int(s.outputX) + int(left)
destY := int(s.outputY) + int(top)
if g.onNV12 != nil {
var ownedAVC bool
var nv12 *H264FrameNV12
decoded, nv12, avcRegions, ownedAVC = g.decodeAVC420WithNV12(bmpData, destX, destY, w, h)
owned = ownedAVC
if nv12 != nil {
if decoded != nil {
blitToSurface(s, int(left), int(top), w, h, decoded)
if owned {
releaseBitmapBuf(decoded)
}
} else {
// Display advanced without a shadow update; mark stale so a
// later full-surface (WTS2) frame fully repairs the shadow.
s.shadowStale = true
}
g.onNV12(destX, destY, w, h, nv12.Y, nv12.YStride, nv12.UV, nv12.UVStride)
return
}
// NV12 unavailable; fall through to BGRA emit.
} else if g.onI420 != nil {
var ownedAVC bool
var i420 *H264FrameI420
decoded, i420, avcRegions, ownedAVC = g.decodeAVC420WithI420(bmpData, destX, destY, w, h)
owned = ownedAVC
if i420 != nil {
if decoded != nil {
blitToSurface(s, int(left), int(top), w, h, decoded)
if owned {
releaseBitmapBuf(decoded)
}
} else {
s.shadowStale = true
}
g.onI420(destX, destY, w, h, i420.Y, i420.YStride, i420.U, i420.UStride, i420.V, i420.VStride)
return
}
// I420 unavailable (nil frame or unsupported format); fall through to BGRA emit.
} else {
var ownedAVC bool
decoded, avcRegions, ownedAVC = g.decodeAVC420(bmpData, destX, destY, w, h)
owned = ownedAVC
}
case codecAVC444, codecAVC444v2:
destX := int(s.outputX) + int(left)
destY := int(s.outputY) + int(top)
if g.onNV12 != nil {
var ownedAVC bool
var nv12 *H264FrameNV12
decoded, nv12, avcRegions, ownedAVC = g.decodeAVC444WithNV12(bmpData, destX, destY, w, h)
owned = ownedAVC
if nv12 != nil {
if decoded != nil {
blitToSurface(s, int(left), int(top), w, h, decoded)
if owned {
releaseBitmapBuf(decoded)
}
} else {
// Display advanced without a shadow update; mark stale so a
// later full-surface frame fully repairs the shadow.
s.shadowStale = true
}
g.onNV12(destX, destY, w, h, nv12.Y, nv12.YStride, nv12.UV, nv12.UVStride)
return
}
// nv12 == nil: LC=2 chroma frame or decoder stall.
// decoded may contain combined BGRA; fall through to emit it.
} else if g.onI420 != nil {
var ownedAVC bool
var i420 *H264FrameI420
decoded, i420, avcRegions, ownedAVC = g.decodeAVC444WithI420(bmpData, destX, destY, w, h)
owned = ownedAVC
if i420 != nil {
if decoded != nil {
blitToSurface(s, int(left), int(top), w, h, decoded)
if owned {
releaseBitmapBuf(decoded)
}
} else {
s.shadowStale = true
}
g.onI420(destX, destY, w, h, i420.Y, i420.YStride, i420.U, i420.UStride, i420.V, i420.VStride)
return
}
// i420 == nil: LC=2 or decoder unavailable; decoded may contain BGRA.
} else {
var ownedAVC bool
decoded, avcRegions, ownedAVC = g.decodeAVC444(bmpData, destX, destY, w, h)
owned = ownedAVC
}
case codecClear:
// ClearCodec 承载 UI/文字/图标等内容(现代服务器主力编码器),
// 解码器维护跨帧 vBar 缓存,输出 BGRA 位图
decoded = g.clearCtx.decode(bmpData, w, h)
owned = true
case codecProgressive:
// Progressive 瓦片按 surface 网格绝对定位,直接解到持久 surface 缓冲
//(与 WTS2 分支相同)。服务器可能经任一 WTS 消息投递 progressive。
rects := g.progressive.Decode(bmpData, s.data, int(s.width), int(s.height))
stride := int(s.width) * 4
for _, rc := range rects {
needed := rc.w * rc.h * 4
region := regionPool.Get().([]byte)
if cap(region) < needed {
region = make([]byte, needed)
} else {
region = region[:needed]
}
rowBytes := rc.w * 4
for row := 0; row < rc.h; row++ {
srcOff := (rc.y+row)*stride + rc.x*4
dstOff := row * rowBytes
if srcOff+rowBytes <= len(s.data) {
copy(region[dstOff:dstOff+rowBytes], s.data[srcOff:srcOff+rowBytes])
}
}
g.emitBitmap(s, rc.x, rc.y, rc.w, rc.h, region)
regionPool.Put(region)
}
return
default:
slog.Warn("RDPGFX: unsupported codec in WTS1", "codecId", codecId, "surfId", surfId, "w", w, "h", h, "bmpLen", bmpLen)
return
}
if decoded == nil {
return
}
if len(avcRegions) > 0 && shouldUseAVCRegions(avcRegions, w, h) {
g.blitAndEmitAVCRegions(s, int(left), int(top), w, h, decoded, avcRegions)
if owned {
releaseBitmapBuf(decoded)
}
return
}
blitToSurface(s, int(left), int(top), w, h, decoded)
if owned {
g.emitBitmapPooled(s, int(left), int(top), w, h, decoded)
} else {
g.emitBitmap(s, int(left), int(top), w, h, decoded)
}
}
// onWireToSurface2Decode handles RDPGFX_WIRE_TO_SURFACE_PDU_2 (MS-RDPEGFX 2.2.2.2).
func (g *GfxHandler) onWireToSurface2Decode(data []byte) {
if len(data) < 13 {
return
}
// Parse fixed header fields via direct binary indexing (avoids bytes.NewReader
// and per-field io.ReadFull overhead on the hot H.264 path).
surfId := binary.LittleEndian.Uint16(data[0:])
codecId := binary.LittleEndian.Uint16(data[2:])
codecCtxId := binary.LittleEndian.Uint32(data[4:])
pixFmt := data[8]
bmpLen := binary.LittleEndian.Uint32(data[9:])
if int(bmpLen) > len(data)-13 {
return
}
if pixFmt != pixelFormatXRGB8888 && pixFmt != pixelFormatARGB8888 {
return
}
bmpData := data[13 : 13+int(bmpLen)]
s, ok := g.surfaces[surfId]
if !ok {
return
}
w := int(s.width)
h := int(s.height)
if g.pduRecord != nil {
g.pduRecord(2, codecId, uint32(s.width), uint32(s.height), 0, 0,
uint32(w), uint32(h), bmpData)
}
if slog.Default().Enabled(nil, slog.LevelDebug) {
slog.Debug("RDPGFX: WTS2", "surfId", surfId, "codecId", codecId,
"w", w, "h", h, "bmpLen", bmpLen)
}
if int(codecId) < len(g.codecBytes) {
g.codecBytes[codecId].Add(int64(bmpLen))
}
var decoded []byte
switch codecId {
case codecUncompressed:
decoded = decodeUncompressed(bmpData, w, h, pixFmt)
blitToSurface(s, 0, 0, w, h, decoded)
g.emitBitmapPooled(s, 0, 0, w, h, decoded)
case codecPlanar:
decoded = decodePlanar(bmpData, w, h)
blitToSurface(s, 0, 0, w, h, decoded)
g.emitBitmapPooled(s, 0, 0, w, h, decoded)
case codecClear:
clearDecoded := g.clearCtx.decode(bmpData, w, h)
blitToSurface(s, 0, 0, w, h, clearDecoded)
g.emitBitmapPooled(s, 0, 0, w, h, clearDecoded)
case codecCaVideo:
rects := g.rfx.Decode(bmpData, 0, 0, s.data, w, h)
g.emitCaVideoRects(s, rects)
case codecAVC420:
destX := int(s.outputX)
destY := int(s.outputY)
if g.onNV12 != nil {
decoded, nv12, avcRegions, ownedAVC := g.decodeAVC420WithNV12(bmpData, destX, destY, w, h)
if nv12 != nil {
if decoded != nil {
// GPU drives the display via onNV12; the CPU shadow only needs
// the dirty regions when it is already in sync. A full blit is
// forced when the shadow is stale (a prior GPU-only frame
// advanced the display without updating it).
if !s.shadowStale && len(avcRegions) > 0 && shouldUseAVCRegions(avcRegions, w, h) {
g.blitAVCRegionsToSurface(s, 0, 0, w, h, decoded, avcRegions)
} else {
blitToSurface(s, 0, 0, w, h, decoded)
s.shadowStale = false
}
if ownedAVC {
releaseBitmapBuf(decoded)
}
} else {
// Display advanced (onNV12) without a BGRA shadow update; mark
// the shadow stale so the next decoded frame fully repairs it.
s.shadowStale = true
}
g.onNV12(destX, destY, w, h, nv12.Y, nv12.YStride, nv12.UV, nv12.UVStride)
} else if decoded != nil {
// NV12 unavailable; fall back to BGRA emit.
if len(avcRegions) > 0 && shouldUseAVCRegions(avcRegions, w, h) {
g.blitAndEmitAVCRegions(s, 0, 0, w, h, decoded, avcRegions)
if ownedAVC {
releaseBitmapBuf(decoded)
}
} else {
blitToSurface(s, 0, 0, w, h, decoded)
if ownedAVC {
g.emitBitmapPooled(s, 0, 0, w, h, decoded)
} else {
g.emitBitmap(s, 0, 0, w, h, decoded)
}
}
}
} else if g.onI420 != nil {
decoded, i420, avcRegions, ownedAVC := g.decodeAVC420WithI420(bmpData, destX, destY, w, h)
if i420 != nil {
if decoded != nil {
if !s.shadowStale && len(avcRegions) > 0 && shouldUseAVCRegions(avcRegions, w, h) {
g.blitAVCRegionsToSurface(s, 0, 0, w, h, decoded, avcRegions)
} else {
blitToSurface(s, 0, 0, w, h, decoded)
s.shadowStale = false
}
if ownedAVC {
releaseBitmapBuf(decoded)
}
} else {
s.shadowStale = true
}
g.onI420(destX, destY, w, h, i420.Y, i420.YStride, i420.U, i420.UStride, i420.V, i420.VStride)
} else if decoded != nil {
// I420 unavailable; fall back to BGRA emit.
if len(avcRegions) > 0 && shouldUseAVCRegions(avcRegions, w, h) {
g.blitAndEmitAVCRegions(s, 0, 0, w, h, decoded, avcRegions)
if ownedAVC {
releaseBitmapBuf(decoded)
}
} else {
blitToSurface(s, 0, 0, w, h, decoded)
if ownedAVC {
g.emitBitmapPooled(s, 0, 0, w, h, decoded)
} else {
g.emitBitmap(s, 0, 0, w, h, decoded)
}
}
}
} else {
decoded, avcRegions, ownedAVC := g.decodeAVC420(bmpData, destX, destY, w, h)
if decoded != nil {
if len(avcRegions) > 0 && shouldUseAVCRegions(avcRegions, w, h) {
g.blitAndEmitAVCRegions(s, 0, 0, w, h, decoded, avcRegions)
if ownedAVC {
releaseBitmapBuf(decoded)
}
} else {
blitToSurface(s, 0, 0, w, h, decoded)
if ownedAVC {
g.emitBitmapPooled(s, 0, 0, w, h, decoded)
} else {
g.emitBitmap(s, 0, 0, w, h, decoded)
}
}
}
}
case codecAVC444, codecAVC444v2:
destX := int(s.outputX)
destY := int(s.outputY)
if g.onI420 != nil {
decoded, i420, avcRegions, ownedAVC := g.decodeAVC444WithI420(bmpData, destX, destY, w, h)
if i420 != nil {
if decoded != nil {
if !s.shadowStale && len(avcRegions) > 0 && shouldUseAVCRegions(avcRegions, w, h) {
g.blitAVCRegionsToSurface(s, 0, 0, w, h, decoded, avcRegions)
} else {
blitToSurface(s, 0, 0, w, h, decoded)
s.shadowStale = false
}
if ownedAVC {
releaseBitmapBuf(decoded)
}
} else {
s.shadowStale = true
}
g.onI420(destX, destY, w, h, i420.Y, i420.YStride, i420.U, i420.UStride, i420.V, i420.VStride)
} else if decoded != nil {
if len(avcRegions) > 0 && shouldUseAVCRegions(avcRegions, w, h) {
g.blitAndEmitAVCRegions(s, 0, 0, w, h, decoded, avcRegions)
if ownedAVC {
releaseBitmapBuf(decoded)
}
} else {
blitToSurface(s, 0, 0, w, h, decoded)
if ownedAVC {
g.emitBitmapPooled(s, 0, 0, w, h, decoded)
} else {
g.emitBitmap(s, 0, 0, w, h, decoded)
}
}
}
} else {
decoded, avcRegions, ownedAVC := g.decodeAVC444(bmpData, destX, destY, w, h)
if decoded != nil {
if len(avcRegions) > 0 && shouldUseAVCRegions(avcRegions, w, h) {
g.blitAndEmitAVCRegions(s, 0, 0, w, h, decoded, avcRegions)
if ownedAVC {
releaseBitmapBuf(decoded)
}
} else {
blitToSurface(s, 0, 0, w, h, decoded)
if ownedAVC {
g.emitBitmapPooled(s, 0, 0, w, h, decoded)
} else {
g.emitBitmap(s, 0, 0, w, h, decoded)
}
}
}
}
case codecProgressive:
// Decode tiles directly onto the persistent surface buffer.
rects := g.progressive.Decode(bmpData, s.data, w, h)
for _, rc := range rects {
needed := rc.w * rc.h * 4
region := regionPool.Get().([]byte)
if cap(region) < needed {
region = make([]byte, needed)
} else {
region = region[:needed]
}
stride := w * 4
rowBytes := rc.w * 4
for row := 0; row < rc.h; row++ {
srcOff := (rc.y+row)*stride + rc.x*4
dstOff := row * rowBytes
if srcOff+rowBytes <= len(s.data) {
copy(region[dstOff:dstOff+rowBytes], s.data[srcOff:srcOff+rowBytes])
}
}
g.emitBitmap(s, rc.x, rc.y, rc.w, rc.h, region)
regionPool.Put(region)
}
default:
slog.Debug("RDPGFX: WTS2 unsupported codec", "codecId", codecId, "ctxId", codecCtxId)
return
}
}
func (g *GfxHandler) onSolidFill(data []byte) {
if len(data) < 8 {
return
}
if g.pduRecord != nil {
g.pduRecord(ReplayKindSolidFill, 0, 0, 0, 0, 0, 0, 0, data)
}
surfId := binary.LittleEndian.Uint16(data[0:])
cb := data[2]
cg := data[3]
cr := data[4]
// data[5] = XA (ignored)
fillCount := binary.LittleEndian.Uint16(data[6:])
s, ok := g.surfaces[surfId]
if !ok {
return
}
stride := int(s.width) * 4
// Pre-compose a single BGRA pixel as a uint32 for one-shot writes.
pixelU32 := uint32(cb) | uint32(cg)<<8 | uint32(cr)<<16 | uint32(0xFF)<<24
offset := 8
for range fillCount {
if offset+8 > len(data) {
break
}
left := binary.LittleEndian.Uint16(data[offset:])
top := binary.LittleEndian.Uint16(data[offset+2:])
right := binary.LittleEndian.Uint16(data[offset+4:])
bottom := binary.LittleEndian.Uint16(data[offset+6:])
offset += 8
w := int(right - left)
h := int(bottom - top)
if w <= 0 || h <= 0 {
continue
}
// Clamp to surface bounds
yEnd := min(int(bottom), int(s.height))
xEnd := min(int(right), int(s.width))
// Fill the first row with PutUint32 (single 32-bit store per pixel),
// then replicate it to subsequent rows with copy().
rowStart := int(top)*stride + int(left)*4
rowBytes := (xEnd - int(left)) * 4
if rowStart+rowBytes <= len(s.data) {
row := s.data[rowStart : rowStart+rowBytes]
for x := 0; x+4 <= rowBytes; x += 4 {
binary.LittleEndian.PutUint32(row[x:], pixelU32)
}
for y := int(top) + 1; y < yEnd; y++ {
dst := y*stride + int(left)*4
if dst+rowBytes <= len(s.data) {
copy(s.data[dst:dst+rowBytes], row)
}
}
}
if s.mapped && g.onBitmap != nil {
// Build fill data: fill first row, then replicate (doubling).
fillData := acquireBitmapBuf(w * h * 4)
rowW := w * 4
for x := 0; x+4 <= rowW; x += 4 {
binary.LittleEndian.PutUint32(fillData[x:], pixelU32)
}
// Doubling copy: O(log h) memmoves instead of h linear copies.
filled := rowW
total := rowW * h
for filled*2 <= total {
copy(fillData[filled:filled*2], fillData[:filled])
filled *= 2
}
if filled < total {
copy(fillData[filled:total], fillData[:total-filled])
}
destL := int(s.outputX) + int(left)
destT := int(s.outputY) + int(top)
g.singleUpdate[0] = BitmapUpdate{
DestLeft: destL, DestTop: destT,
DestRight: destL + w - 1, DestBottom: destT + h - 1,
Width: w, Height: h, Bpp: 4, Data: fillData,
}
g.emitAndReleaseUpdates(g.singleUpdate[:])
}
}
}
// onSurfaceToSurface handles RDPGFX_SURFACE_TO_SURFACE_PDU (MS-RDPEGFX
// 2.2.2.5). Layout: srcId(2) dstId(2) rectSrc(8) destPtsCount(2) destPts[]
// (4 bytes each). rectSrc on the source surface is copied to each
// destination point — this is how the server scrolls/blits content during
// window drags, and same-surface copies may overlap, so the row order must
// adapt to the copy direction.
func (g *GfxHandler) onSurfaceToSurface(data []byte) {
// Minimum payload: srcId(2)+dstId(2)+rectSrc(8)+destPtsCount(2) = 14 bytes
if len(data) < 14 {
return
}
if g.pduRecord != nil {
g.pduRecord(ReplayKindSurfaceToSurf, 0, 0, 0, 0, 0, 0, 0, data)
}
srcId := binary.LittleEndian.Uint16(data[0:])
dstId := binary.LittleEndian.Uint16(data[2:])
left := int(binary.LittleEndian.Uint16(data[4:]))
top := int(binary.LittleEndian.Uint16(data[6:]))
right := int(binary.LittleEndian.Uint16(data[8:]))
bottom := int(binary.LittleEndian.Uint16(data[10:]))
destCount := int(binary.LittleEndian.Uint16(data[12:]))
src, srcOk := g.surfaces[srcId]
dst, dstOk := g.surfaces[dstId]
if !srcOk || !dstOk {
return
}
w := right - left
h := bottom - top
if w <= 0 || h <= 0 {
return
}
srcStride := int(src.width) * 4
dstStride := int(dst.width) * 4
rowBytes := w * 4
offset := 14
for i := 0; i < destCount; i++ {
if offset+4 > len(data) {
break
}
dstX := int(binary.LittleEndian.Uint16(data[offset:]))
dstY := int(binary.LittleEndian.Uint16(data[offset+2:]))
offset += 4
if left+rowBytes/4 > int(src.width) || top+h > int(src.height) ||
dstX+w > int(dst.width) || dstY+h > int(dst.height) {
continue
}
// 同面滚动拷贝:目标在源下方时自下而上逐行复制,避免未读源行
// 被提前覆盖(行内水平重叠由 copy 的 memmove 语义保证)。
first, last, step := 0, h, 1
if src == dst && dstY > top {
first, last, step = h-1, -1, -1
}
for row := first; row != last; row += step {
s := (top+row)*srcStride + left*4
d := (dstY+row)*dstStride + dstX*4
copy(dst.data[d:d+rowBytes], src.data[s:s+rowBytes])
}
// emitBitmap 的 Data 必须是 w×h 的矩形位图:把拷贝结果抽出来再发。
// 直接传整面 dst.data 会让 JS 端按 w×h 读表面左上角内容贴到
// 目标位置——拖动时窗口内部被桌面角落内容覆盖,即花屏条纹。
rectBuf := acquireBitmapBuf(w * h * 4)
for row := 0; row < h; row++ {
srcOff := (dstY+row)*dstStride + dstX*4
dstOff := row * rowBytes
copy(rectBuf[dstOff:dstOff+rowBytes], dst.data[srcOff:srcOff+rowBytes])
}
g.emitBitmapPooled(dst, dstX, dstY, w, h, rectBuf)
}
}
// onSurfaceToCache copies a surface region into a cache slot
// (MS-RDPEGFX 2.2.2.6). Without this the cache stays empty and every
// subsequent CACHE_TO_SURFACE blits nothing — a major source of missing
// UI elements on Win11 servers, which lean on the cache heavily.
func (g *GfxHandler) onSurfaceToCache(data []byte) {
if len(data) < 20 {
return
}
if g.pduRecord != nil {
g.pduRecord(ReplayKindSurfaceToCache, 0, 0, 0, 0, 0, 0, 0, data)
}
surfId := binary.LittleEndian.Uint16(data[0:])
// 持久缓存键:跨重连的身份标识,CacheImportOffer 上报它而非槽位
cacheKey := binary.LittleEndian.Uint64(data[2:])
cacheSlot := binary.LittleEndian.Uint16(data[10:])
left := int(binary.LittleEndian.Uint16(data[12:]))
top := int(binary.LittleEndian.Uint16(data[14:]))
right := int(binary.LittleEndian.Uint16(data[16:]))
bottom := int(binary.LittleEndian.Uint16(data[18:]))
s, ok := g.surfaces[surfId]
if !ok {
return
}
w := right - left
h := bottom - top
if w <= 0 || h <= 0 || left < 0 || top < 0 ||
right > int(s.width) || bottom > int(s.height) {
return
}
rowBytes := w * 4
entry := cacheEntry{key: cacheKey, width: w, height: h, data: make([]byte, w*h*4)}
for row := 0; row < h; row++ {
srcOff := (top+row)*int(s.width)*4 + left*4
dstOff := row * rowBytes
copy(entry.data[dstOff:dstOff+rowBytes], s.data[srcOff:srcOff+rowBytes])
}
g.cacheEntries[cacheSlot] = entry
// 交浏览器侧持久化:重连时以 CacheImportOffer 上报该键(MS-RDPEGFX
// 持久位图缓存)。key=0 视为服务器未提供有效键,跳过。
if g.cacheStore != nil && cacheKey != 0 {
g.cacheStore.Persist(cacheKey, w, h, 32, entry.data)
}
}
func (g *GfxHandler) onCacheToSurface(data []byte) {
if len(data) < 6 {
return
}
if g.pduRecord != nil {
g.pduRecord(ReplayKindCacheToSurface, 0, 0, 0, 0, 0, 0, 0, data)
}
cacheSlot := binary.LittleEndian.Uint16(data[0:])
surfId := binary.LittleEndian.Uint16(data[2:])
destCount := binary.LittleEndian.Uint16(data[4:])
ce, hasCE := g.cacheEntries[cacheSlot]
s, hasSurf := g.surfaces[surfId]
offset := 6
for range destCount {
if offset+4 > len(data) {
break
}
dx := binary.LittleEndian.Uint16(data[offset:])
dy := binary.LittleEndian.Uint16(data[offset+2:])
offset += 4
if hasCE && hasSurf {
blitToSurface(s, int(dx), int(dy), ce.width, ce.height, ce.data)
g.emitBitmap(s, int(dx), int(dy), ce.width, ce.height, ce.data)
}
}
}
func (g *GfxHandler) onEvictCacheEntry(data []byte) {
if len(data) < 2 {
return
}
if g.pduRecord != nil {
g.pduRecord(ReplayKindEvictCache, 0, 0, 0, 0, 0, 0, 0, data)
}
slot := binary.LittleEndian.Uint16(data)
delete(g.cacheEntries, slot)
}
func (g *GfxHandler) onCacheImportOffer() {
var p [2]byte // importedEntriesCount = 0 (little-endian zero)
g.sendPdu(cmdidCacheImportReply, p[:])
}
// maxCacheImportEntries 限制单条 CacheImportOffer 的条目数。规范上限 5461
// (cacheEntriesCount < 0x1556),静态桌面的瓦片数百量级即够,同时把
// 浏览器侧存储与导出开销控制在几 MB 内。
const maxCacheImportEntries = 512
// sendCacheImportOffer 在 caps 确认后向服务器上报客户端持久缓存中仍持有的
// 条目(cacheKey + bitmapLength)。服务器按前缀导入并在 CacheImportReply
// 中回分配的槽位;导入的条目随后可被 CacheToSurface 直接回贴,静态内容
// 在重连后无需重传(MS-RDPEGFX 2.2.2.16/2.2.2.17)。
func (g *GfxHandler) sendCacheImportOffer() {
if g.importOfferSent || g.cacheStore == nil {
return
}
g.importOfferSent = true
entries := g.cacheStore.Export()
if len(entries) == 0 {
return
}
if len(entries) > maxCacheImportEntries {
entries = entries[:maxCacheImportEntries]
}
payload := make([]byte, 2, 2+12*len(entries))
binary.LittleEndian.PutUint16(payload, uint16(len(entries)))
for _, e := range entries {
payload = binary.LittleEndian.AppendUint64(payload, e.Key)
payload = binary.LittleEndian.AppendUint32(payload, uint32(len(e.Data)))
}
g.offeredCache = entries
g.sendPdu(cmdidCacheImportOffer, payload)
slog.Info("RDPGFX: cache import offer", "entries", len(entries))
}
// onCacheImportReply 处理 CacheImportOffer 的应答:importedEntriesCount 为
// 前缀语义——上报条目的前 N 条被导入,cacheSlots[i] 是第 i 条的新槽位。
// 像素数据客户端本就持有,直接重建槽位映射即可。
func (g *GfxHandler) onCacheImportReply(data []byte) {
if len(data) < 2 {
return
}
n := int(binary.LittleEndian.Uint16(data))
if n > len(g.offeredCache) {
n = len(g.offeredCache)
}
if len(data) < 2+2*n { // 槽位数组被截断:只取完整部分
n = (len(data) - 2) / 2
}
imported := 0
for i := 0; i < n; i++ {
e := g.offeredCache[i]
// 任何尺寸/长度不一致的条目绝不入缓存——错误回贴就是花屏
if e.Width <= 0 || e.Height <= 0 || len(e.Data) != e.Width*e.Height*4 {
continue
}
slot := binary.LittleEndian.Uint16(data[2+i*2:])
g.cacheEntries[slot] = cacheEntry{key: e.Key, width: e.Width, height: e.Height, data: e.Data}
imported++
}
g.offeredCache = nil
slog.Info("RDPGFX: cache import reply", "imported", imported, "of", n)
}
// SetPersistentCacheStore 安装持久缓存桥;必须在连接建立前调用。
func (g *GfxHandler) SetPersistentCacheStore(s GfxCacheStore) {
g.cacheStore = s
}
// --- Helpers ---
// emitCaVideoRects copies decoded RemoteFX tile regions from the surface
// pixel buffer into individual BitmapUpdate slices and emits them.
// Used by both onWireToSurface1Decode and onWireToSurface2Decode.
func (g *GfxHandler) emitCaVideoRects(s *surface, rects []rfxRect) {
if !s.mapped || g.onBitmap == nil || len(rects) == 0 {
return
}
g.updatesBuf = g.updatesBuf[:0]
stride := int(s.width) * 4
for _, rc := range rects {
needed := rc.w * rc.h * 4
region := acquireBitmapBuf(needed)
rowBytes := rc.w * 4
for row := 0; row < rc.h; row++ {
srcOff := (rc.y+row)*stride + rc.x*4
dstOff := row * rowBytes
if srcOff+rowBytes <= len(s.data) {
copy(region[dstOff:dstOff+rowBytes], s.data[srcOff:srcOff+rowBytes])
}
}
destL := int(s.outputX) + rc.x
destT := int(s.outputY) + rc.y
g.updatesBuf = append(g.updatesBuf, BitmapUpdate{
DestLeft: destL, DestTop: destT,
DestRight: destL + rc.w - 1, DestBottom: destT + rc.h - 1,
Width: rc.w, Height: rc.h, Bpp: 4, Data: region,
})
}
g.emitAndReleaseUpdates(g.updatesBuf)
}
func blitToSurface(s *surface, x, y, w, h int, src []byte) {
stride := int(s.width) * 4
// Full-width fast path: when x==0 and w==surface.width the entire region
// is contiguous in both src and s.data — replace h row-copies with one.
if x == 0 && w == int(s.width) && y >= 0 && y+h <= int(s.height) {
dstOff := y * stride
n := h * stride
if n <= len(src) && dstOff+n <= len(s.data) {
copy(s.data[dstOff:dstOff+n], src[:n])
return
}
}
for row := range h {
dy := y + row
if dy < 0 || dy >= int(s.height) {
continue
}
srcOff := row * w * 4
dstOff := dy*stride + x*4
n := w * 4
if dstOff >= 0 && dstOff+n <= len(s.data) && srcOff+n <= len(src) {
copy(s.data[dstOff:dstOff+n], src[srcOff:srcOff+n])
}
}
}
// emitBitmapPooled is like emitBitmap but releases `decoded` back to
// bitmapBufPool after the synchronous onBitmap callback returns. Use this
// for codec output buffers that the GfxHandler owns end-to-end (currently
// uncompressed and planar).
func (g *GfxHandler) emitBitmapPooled(s *surface, x, y, w, h int, decoded []byte) {
if !s.mapped || g.onBitmap == nil {
releaseBitmapBuf(decoded)
return
}
destL := int(s.outputX) + x
destT := int(s.outputY) + y
g.singleUpdate[0] = BitmapUpdate{
DestLeft: destL, DestTop: destT,
DestRight: destL + w - 1, DestBottom: destT + h - 1,
Width: w, Height: h, Bpp: 4, Data: decoded,
}
g.emitAndReleaseUpdates(g.singleUpdate[:])
}
func (g *GfxHandler) emitBitmap(s *surface, x, y, w, h int, decoded []byte) {
if !s.mapped || g.onBitmap == nil {
return
}
destL := int(s.outputX) + x
destT := int(s.outputY) + y
g.singleUpdate[0] = BitmapUpdate{
DestLeft: destL, DestTop: destT,
DestRight: destL + w - 1, DestBottom: destT + h - 1,
Width: w, Height: h, Bpp: 4, Data: decoded,
}
g.onBitmap(g.singleUpdate[:])
g.singleUpdate[0].Data = nil // release reference; decoded is not pooled
}
// --- Codec: Uncompressed ---
func decodeUncompressed(data []byte, w, h int, pixFmt uint8) []byte {
out := acquireBitmapBuf(w * h * 4)
n := w * h * 4
if len(data) >= n {
copy(out, data[:n])
} else {
copy(out[:len(data)], data)
// Zero the unfilled tail in case the slice was reused from the pool.
clear(out[len(data):n])
}
return out
}
// --- Codec: Planar (RDP 6.0 Bitmap Codec, MS-RDPEGDI 2.2.2.5) ---
func decodePlanar(data []byte, w, h int) []byte {
if len(data) < 1 {
return acquireBitmapBuf(w * h * 4)
}
header := data[0]
// 头部标志位(FreeRDP include/freerdp/codec/planar.h):
// bit3 CS 色度子采样、bit4 RLE、bit5 NA(无 alpha 平面)、bits0-2 CLL。
rle := (header >> 4) & 1
noAlpha := (header >> 5) & 1
planeSize := w * h
offset := 1
var alphaPlane, redPlane, greenPlane, bluePlane []byte
defer func() {
releasePlaneBuf(alphaPlane)
releasePlaneBuf(redPlane)
releasePlaneBuf(greenPlane)
releasePlaneBuf(bluePlane)
}()
if rle == 0 {
if noAlpha == 0 {
alphaPlane, offset = readRawPlane(data, offset, planeSize)
}
redPlane, offset = readRawPlane(data, offset, planeSize)
greenPlane, offset = readRawPlane(data, offset, planeSize)
bluePlane, offset = readRawPlane(data, offset, planeSize)
} else {
if noAlpha == 0 {
alphaPlane, offset = decodePlanarRLEPlane(data, offset, w, h)
}
redPlane, offset = decodePlanarRLEPlane(data, offset, w, h)
greenPlane, offset = decodePlanarRLEPlane(data, offset, w, h)
bluePlane, offset = decodePlanarRLEPlane(data, offset, w, h)
}
_ = offset
out := acquireBitmapBuf(planeSize * 4)
// Hoist the per-pixel nil/length checks: clamp each plane to
// `planeSize` (zero-fill missing planes) so the inner loop has no
// branches and the bounds checks are eliminated.
rp := planeOrZero(redPlane, planeSize)
gp := planeOrZero(greenPlane, planeSize)
bp := planeOrZero(bluePlane, planeSize)
ap := alphaPlane
hasAlpha := ap != nil && len(ap) >= planeSize
if hasAlpha {
ap = ap[:planeSize]
for i := range planeSize {
j := i * 4
out[j] = bp[i]
out[j+1] = gp[i]
out[j+2] = rp[i]
out[j+3] = ap[i]
}
} else {
for i := range planeSize {
j := i * 4
out[j] = bp[i]
out[j+1] = gp[i]
out[j+2] = rp[i]
out[j+3] = 0xFF
}
}
return out
}
// planeOrZero returns a slice of exactly `size` bytes, either the input
// plane (truncated if longer) or a zero-filled buffer when the plane is
// nil or short. Used to drop per-pixel nil/bounds checks in decodePlanar.
func planeOrZero(plane []byte, size int) []byte {
if len(plane) >= size {
return plane[:size]
}
out := make([]byte, size)
copy(out, plane)
return out
}
func readRawPlane(data []byte, offset, size int) ([]byte, int) {
plane := acquirePlaneBuf(size)
end := min(offset+size, len(data))
if offset < end {
copy(plane, data[offset:end])
if end-offset < size {
clear(plane[end-offset:])
}
} else {
clear(plane)
}
return plane, offset + size
}
// decodePlanarRLEPlane 解码一个 RLE 压缩的像素平面
// (FreeRDP planar_decompress_plane_rle)。控制字节:低 4 位 runLen、
// 高 4 位 rawLen;runLen==1 → run=raw+16、runLen==2 → run=raw+32
// (此时 rawLen 清零)。首行携带绝对像素值,后续行携带相对上一行
// 同列像素的有符号 delta(最低位为符号位),行程重复最后一个像素
// 的 delta。
func decodePlanarRLEPlane(data []byte, offset, w, h int) ([]byte, int) {
out := acquirePlaneBuf(w * h)
clamp := func(v int) byte {
if v < 0 {
return 0
}
if v > 255 {
return 255
}
return byte(v)
}
for y := 0; y < h; y++ {
rowStart := y * w
prevRow := rowStart - w
for x := 0; x < w; {
if offset >= len(data) {
return out, offset
}
ctrl := data[offset]
offset++
runLen := int(ctrl & 0x0F)
rawLen := int((ctrl >> 4) & 0x0F)
switch runLen {
case 1:
runLen = rawLen + 16
rawLen = 0
case 2:
runLen = rawLen + 32
rawLen = 0
}
if y == 0 {
// 首行:绝对像素值
var last byte
for i := 0; i < rawLen; i++ {
if offset >= len(data) || x+i >= w {
return out, offset
}
last = data[offset]
out[rowStart+x+i] = last
offset++
}
x += rawLen
for i := 0; i < runLen && x < w; i++ {
out[rowStart+x] = last
x++
}
continue
}
// 后续行:相对上一行的有符号 delta
lastDelta := 0
for i := 0; i < rawLen; i++ {
if offset >= len(data) || x+i >= w {
return out, offset
}
dv := data[offset]
offset++
if dv&1 != 0 {
lastDelta = -int(dv>>1) - 1
} else {
lastDelta = int(dv >> 1)
}
out[rowStart+x+i] = clamp(int(out[prevRow+x+i]) + lastDelta)
}
x += rawLen
// 行程延续最后一个 delta(相对各自位置上一行的像素)
for i := 0; i < runLen && x < w; i++ {
out[rowStart+x] = clamp(int(out[prevRow+x]) + lastDelta)
x++
}
}
}
return out, offset
}