2817 lines
101 KiB
Go
2817 lines
101 KiB
Go
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
|
||
}
|