//go:build !js && !ios && !android package server import ( "net" "sync" "sync/atomic" "time" ) // SessionTick is one sampling slice of a VNC session's wire activity. // BytesOut / Writes / FBUs are deltas observed during this tick; // Max* fields are the high-water marks observed during this tick (reset // at the start of the next). Period is the wall-clock duration covered // (typically sessionTickInterval, shorter for the final flush). type SessionTick struct { Period time.Duration BytesOut uint64 Writes uint64 FBUs uint64 MaxFBUBytes uint64 MaxFBURects uint64 MaxWriteBytes uint64 WriteNanos uint64 // FBUsTracked is false when the connection carried framebuffer updates it // could not see the boundaries of, which is the case for a service-mode // proxy relaying bytes to a per-session agent. FBUs, MaxFBUBytes and // MaxFBURects then mean "unknown", not "none". FBUsTracked bool } // sessionTickInterval is how often metricsConn emits a SessionTick. One // second covers roughly one FBU round-trip at typical client request // cadences during steady-state activity. const sessionTickInterval = time.Second // metricsConn wraps a net.Conn and tracks per-session byte / write / FBU // counters. Updates are atomic so the cost is a few atomic ops per Write // (well under 100 ns), negligible against the syscall itself, so the wrap // is always installed. A goroutine emits a SessionTick to the recorder // every sessionTickInterval (only when the tick has activity to report); // a final partial-tick flush runs on Close. type metricsConn struct { net.Conn recorder func(SessionTick) // framed is true when the session writing through this connection marks // its FramebufferUpdate boundaries with beginFBU/endFBU. framed bool bytesOut atomic.Uint64 writes atomic.Uint64 writeNanos atomic.Uint64 largestPkt atomic.Uint64 fbus atomic.Uint64 // inFBU is true between beginFBU and endFBU, so only the writes that make // up a FramebufferUpdate are counted against it. inFBU atomic.Bool fbuBytes atomic.Uint64 fbuRects atomic.Uint64 maxFBUBytes atomic.Uint64 maxFBURects atomic.Uint64 tickMu sync.Mutex tickStart time.Time tickPrevB uint64 tickPrevW uint64 tickPrevF uint64 tickPrevNS uint64 // busyMu guards the sliding window used by BusyFraction. busyMu sync.Mutex busyLastTime time.Time busyLastNanos uint64 busyFraction float64 // writeMu is held shared by each Write for its whole duration and taken // exclusively by Close before the final snapshot, so that snapshot counts // every write that was in flight when the connection was closed. writeMu sync.RWMutex closeOnce sync.Once done chan struct{} } func newMetricsConn(c net.Conn, recorder func(SessionTick)) net.Conn { return newMetricsConnFramed(c, recorder, true) } // newProxyMetricsConn wraps a connection whose bytes are relayed to a session // running elsewhere, the service-mode agent. Nothing on this side sees where // one FramebufferUpdate ends and the next begins, so its ticks carry byte and // write counts only and report the FBU fields as untracked. func newProxyMetricsConn(c net.Conn, recorder func(SessionTick)) net.Conn { return newMetricsConnFramed(c, recorder, false) } func newMetricsConnFramed(c net.Conn, recorder func(SessionTick), framed bool) net.Conn { m := &metricsConn{ Conn: c, recorder: recorder, framed: framed, tickStart: time.Now(), done: make(chan struct{}), } if recorder != nil { go m.tickLoop() } return m } // tickLoop emits a SessionTick every sessionTickInterval until done. // Empty ticks (no writes since the last tick) are skipped. func (m *metricsConn) tickLoop() { t := time.NewTicker(sessionTickInterval) defer t.Stop() for { select { case <-m.done: return case <-t.C: m.flushTick(false) } } } // flushTick computes deltas since the last tick, resets the per-tick max // trackers, and emits a SessionTick to the recorder. final=true forces // emission even if no writes happened (used at session close to record // the trailing partial period). func (m *metricsConn) flushTick(final bool) { m.tickMu.Lock() defer m.tickMu.Unlock() b := m.bytesOut.Load() w := m.writes.Load() f := m.fbus.Load() ns := m.writeNanos.Load() db := b - m.tickPrevB dw := w - m.tickPrevW df := f - m.tickPrevF dns := ns - m.tickPrevNS m.tickPrevB, m.tickPrevW, m.tickPrevF, m.tickPrevNS = b, w, f, ns maxFBU := m.maxFBUBytes.Swap(0) maxRects := m.maxFBURects.Swap(0) maxPkt := m.largestPkt.Swap(0) period := time.Since(m.tickStart) m.tickStart = time.Now() if dw == 0 && !final { return } m.recorder(SessionTick{ Period: period, BytesOut: db, Writes: dw, FBUs: df, MaxFBUBytes: maxFBU, MaxFBURects: maxRects, MaxWriteBytes: maxPkt, WriteNanos: dns, FBUsTracked: m.framed, }) } // BusyFraction reports the fraction of recent wall time that Write spent // blocked in the underlying socket, as an exponentially smoothed value in // [0, 1]. Approximates downstream backpressure: persistent values near 1 // mean the socket cannot keep up with the encoder's output. Callers can // throttle JPEG quality or skip frames in response. func (m *metricsConn) BusyFraction() float64 { now := time.Now() ns := m.writeNanos.Load() m.busyMu.Lock() defer m.busyMu.Unlock() if m.busyLastTime.IsZero() { m.busyLastTime = now m.busyLastNanos = ns return 0 } period := now.Sub(m.busyLastTime) if period < 50*time.Millisecond { return m.busyFraction } delta := ns - m.busyLastNanos sample := float64(delta) / float64(period.Nanoseconds()) if sample > 1 { sample = 1 } const alpha = 0.4 m.busyFraction = alpha*sample + (1-alpha)*m.busyFraction m.busyLastTime = now m.busyLastNanos = ns return m.busyFraction } // beginFBU records that the writes which follow belong to a new // FramebufferUpdate carrying rects rectangles, closing off the accounting for // the previous one. // // The encoder says this rather than the wrapper inferring it from the payload. // A FramebufferUpdate's message type is 0, and so is the leading byte of plenty // of other traffic: the four zero bytes of a successful security result, and // any rect body whose x coordinate is below 256. Sniffing counted each of those // as a new update, inflating the FBU count and cutting the byte and rect totals // of the update actually in flight into pieces. func (m *metricsConn) beginFBU(rects int) { m.flushFBUMax() m.inFBU.Store(true) m.fbus.Add(1) if rects > 0 { m.fbuRects.Add(uint64(rects)) } } // endFBU closes the update beginFBU opened and folds its totals into the // per-tick maxima straight away, so a frame is accounted in the tick it was // sent in rather than whenever the next one happens to start. Writes outside a // begin/end pair — the handshake, clipboard traffic, bell — are not part of any // update and are left out of the FBU byte total. func (m *metricsConn) endFBU() { m.flushFBUMax() m.inFBU.Store(false) } func (m *metricsConn) Write(p []byte) (int, error) { m.writeMu.RLock() defer m.writeMu.RUnlock() t0 := time.Now() n, err := m.Conn.Write(p) m.writeNanos.Add(uint64(time.Since(t0).Nanoseconds())) m.bytesOut.Add(uint64(n)) m.writes.Add(1) if m.inFBU.Load() { m.fbuBytes.Add(uint64(n)) } if uint64(n) > m.largestPkt.Load() { m.largestPkt.Store(uint64(n)) } return n, err } // flushFBUMax folds the bytes and rects accumulated for the FBU that just // ended into the per-tick high-water marks, then resets the accumulators // for the next FBU. func (m *metricsConn) flushFBUMax() { if b := m.fbuBytes.Swap(0); b > m.maxFBUBytes.Load() { m.maxFBUBytes.Store(b) } if r := m.fbuRects.Swap(0); r > m.maxFBURects.Load() { m.maxFBURects.Store(r) } } // Close closes the connection and records the final partial tick. The socket is // closed first, which is what unblocks a Write stuck on a peer that stopped // reading; the snapshot is taken only once those writes have returned, so the // last frame and the last write of a session are counted in it rather than // lost. func (m *metricsConn) Close() error { closed := false var err error m.closeOnce.Do(func() { closed = true close(m.done) err = m.Conn.Close() m.writeMu.Lock() defer m.writeMu.Unlock() if m.recorder == nil { return } m.flushFBUMax() m.flushTick(true) }) if !closed { return m.Conn.Close() } return err }