From f3c3a9427bd402300ca39d640e2a1fbee36676b1 Mon Sep 17 00:00:00 2001 From: enricobuehler Date: Sat, 25 Jul 2026 01:26:23 +0200 Subject: [PATCH] feat(client/abr): learn the host's rate cap from short acks, and read host encode time as a signal MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two client-side halves of the §ABR-overdrive fix (the 4K120 sessions that climbed to 1 Gbps against a 794 Mbps encoder and a 8.33 ms budget): - Short-ack cap learning: the host now resolves a climb it can't serve to what the encoder actually runs at (its codec-level ceiling, or the current rate while encode is behind cadence). Two consecutive identical short acks latch that value as a host cap the climb logic folds into its ceiling — one short ack stays a transient (a failed rebuild also acks short once). The cap is mode-scoped (cleared on an accepted mode switch, tracked via a mode generation counter the control task bumps) and re-probes one step after ~60 s parked clean, so a heavy-scene refusal can't quietly cap the whole session. - Host-encode-latency down-driver: the per-AU 0xCF `encode_us` the host already ships (and the overlay already draws) now feeds the controller through its own window accumulator — the overlay channel is lossy and embedder-drained, so the ABR gets a dedicated mirror of the decode-latency path. Baseline-relative like the decode signal (an escalated host reports encode_us inflated by ~a frame of queue depth; an absolute budget threshold would read permanently red), with the baseline rebased after our own decreases so one backoff doesn't train-fire into the floor. This is the only signal that can push an already-too-high rate back under the encoder's compute knee — host climb refusal stops the climb, but nothing else descends on a clean LAN. Co-Authored-By: Claude Fable 5 --- crates/punktfunk-core/src/abr.rs | 587 ++++++++++++++++-- .../src/client/frame_channel.rs | 14 + crates/punktfunk-core/src/client/pump.rs | 10 + .../src/client/pump/control_task.rs | 6 + crates/punktfunk-core/src/client/pump/data.rs | 33 +- .../src/client/pump/datagram_task.rs | 8 + 6 files changed, 616 insertions(+), 42 deletions(-) diff --git a/crates/punktfunk-core/src/abr.rs b/crates/punktfunk-core/src/abr.rs index 7cc53cf7..708b829e 100644 --- a/crates/punktfunk-core/src/abr.rs +++ b/crates/punktfunk-core/src/abr.rs @@ -13,7 +13,13 @@ //! its rolling baseline: standing queue growth, the *pre-loss* signature of a saturated link //! (bufferbloat) — this is the early-warning signal loss-based control lacks; //! - **a jump-to-live flush** — the pump discarded its backlog, the strongest "we were behind" -//! evidence there is. +//! evidence there is; +//! - **host-encode-latency rise** — the host's per-AU 0xCF `encode_us` climbing above its rolling +//! baseline: the ENCODER falling behind its frame budget (the compute knee), the one failure a +//! fat LAN never surfaces as loss/OWD/decode. Paired with the host's own climb refusal (a +//! behind-cadence host acks climbs at the current rate) and short-ack cap learning +//! ([`BitrateController::on_ack`]), this is what stops an Automatic session from driving the +//! encoder off a cliff the network could carry. //! //! AIMD shape: a SEVERE window (an unrecoverable frame, a flush, ≥6 % loss, or a decode-latency //! excursion far past baseline) backs off ×0.7 immediately; ordinary congestion @@ -94,6 +100,25 @@ const UTILIZATION_DEN: u64 = 4; /// cap always leaves ≥ ~12 % of climbing room — the two gates can't deadlock. const PROVEN_HEADROOM_NUM: u32 = 3; const PROVEN_HEADROOM_DEN: u32 = 2; +/// How far the window's mean HOST-ENCODE latency (the 0xCF `HostStages::encode_us` the host +/// already ships per AU) may rise above its rolling baseline before the window is bad. This is +/// the down-driver for the ENCODER's compute knee — the failure loss/OWD/decode are all blind +/// to: on a fat LAN the controller can climb to a rate the link carries fine but the ASIC +/// can't encode inside the frame budget (4K120 HEVC at ~800 Mbps ≈ 9.3 ms against 8.33), and +/// the only symptom is encode time. Baseline-RELATIVE on purpose: an escalated host reports +/// encode_us inflated by its retrieve-queue depth (~a frame), so an absolute budget threshold +/// would read permanently-red and drive the rate to the floor; a rise above the session's own +/// baseline survives that offset. ~half a 120 Hz frame budget of standing rise is real. +const ENCODE_RISE_US: i64 = 4_000; +/// Host-encode latency this far above baseline (≈1.5 × a 120 Hz budget) is SEVERE — the encode +/// queue is growing past the knee; skip the two-window confirmation. +const ENCODE_SEVERE_US: i64 = 12_000; +/// Clean windows parked at the learned [`host cap`](BitrateController::host_cap_kbps) before +/// re-probing above it (~60 s at the 750 ms tick). A cadence-refusal cap is scene-dependent +/// evidence, not a spec limit — without a re-probe, one heavy scene would cap the whole +/// session. A still-standing limit just re-teaches itself in two short acks, which the host +/// pre-clamps without touching the encoder — the re-probe costs no rebuild, no IDR. +const CAP_REPROBE_WINDOWS: u32 = 80; /// Rolling window (in 750 ms report windows, ~30 s) whose minimum mean is the OWD baseline. /// Long enough to remember the uncongested floor, short enough to follow genuine path changes. const BASELINE_WINDOWS: usize = 40; @@ -121,6 +146,27 @@ pub(crate) struct BitrateController { /// keeping-up baseline. Empty on embedders that don't report decode latency (the decode /// signal is then simply absent — identical to the pre-decode-signal behavior). decode_means: VecDeque, + /// Recent window mean host-encode latencies (µs, from the 0xCF datagrams); rolling-min + /// baseline like the decode signal. Cleared whenever OUR OWN rate decrease changes the + /// encode regime (see [`on_ack`](Self::on_ack)) and on a mode switch. + encode_means: VecDeque, + /// The host-taught rate cap (§ABR overdrive): latched when the host acks BELOW what we + /// asked twice consecutively at the same value — its encoder's codec-level ceiling, or a + /// climb refusal while host encode can't hold cadence. Kept apart from `ceiling_kbps` so + /// the probe-measured link authority survives a mode switch's reset. Slowly re-probed + /// ([`CAP_REPROBE_WINDOWS`]) so scene-dependent evidence can't cap the session forever. + host_cap_kbps: Option, + /// The rate the last [`request`](Self::request) asked for — the reference an ack is judged + /// short against. Taken (not kept) by the ack, so one request is judged at most once. + last_requested_kbps: Option, + /// Consecutive short-ack streak: the value and how many times in a row it was acked. Two + /// identical short acks latch [`host_cap_kbps`](Self::host_cap_kbps) — one can be a + /// transient (a failed host rebuild keeping the old rate); the host's resolves are + /// deterministic min()s, so a persistent limit reproduces exactly. + short_ack_kbps: u32, + short_acks: u32, + /// Clean windows spent parked at the learned cap (the re-probe clock). + cap_probe_windows: u32, /// Proven throughput: the session's highest windowed ACTUAL delivered rate seen with flat /// decode latency — the known-good high-water mark climbs are bounded against. Never decays; /// shrinking capacity (thermals, a heavier scene) is the reactive decode signal's job. On @@ -147,6 +193,12 @@ impl BitrateController { probing: true, owd_means: VecDeque::with_capacity(BASELINE_WINDOWS), decode_means: VecDeque::with_capacity(BASELINE_WINDOWS), + encode_means: VecDeque::with_capacity(BASELINE_WINDOWS), + host_cap_kbps: None, + last_requested_kbps: None, + short_ack_kbps: 0, + short_acks: 0, + cap_probe_windows: 0, proven_kbps: 0, bad_windows: 0, clean_windows: 0, @@ -168,22 +220,68 @@ impl BitrateController { /// The host's [`crate::quic::BitrateChanged`] ack: its clamp is authoritative for what the /// encoder now targets, and any ack proves the host renegotiates (resets the silence counter). + /// + /// A SHORT ack (below what we asked) is the host telling us about a limit the network + /// signals can't see — its encoder's codec-level ceiling, or a climb refusal while encode + /// can't hold cadence. Two consecutive short acks at the SAME value latch it as + /// [`host_cap_kbps`](Self::host_cap_kbps), stopping the AIMD sawtooth from re-poking a + /// limit the host already refused; ONE is not enough — a failed host rebuild also acks + /// short once, and latching a transient would cap the session on a hiccup. pub(crate) fn on_ack(&mut self, kbps: u32) { if kbps > 0 { + if kbps < self.current_kbps { + // Our own decrease changes the encode-time regime (less work per frame; on an + // escalated host the queue offset shifts too) — judging the new regime against + // the old baseline would train-fire the encode down-driver. Re-seed it. + self.encode_means.clear(); + } + if let Some(req) = self.last_requested_kbps.take() { + if kbps < req { + if self.short_ack_kbps == kbps { + self.short_acks += 1; + } else { + self.short_ack_kbps = kbps; + self.short_acks = 1; + } + if self.short_acks >= 2 && self.host_cap_kbps.is_none_or(|c| kbps < c) { + tracing::info!( + cap_kbps = kbps, + "adaptive bitrate: host cap learned (encoder ceiling or cadence \ + refusal) — climbs stop here until it lifts" + ); + self.host_cap_kbps = Some(kbps.max(self.floor_kbps)); + self.cap_probe_windows = 0; + } + } else { + self.short_acks = 0; + } + } self.current_kbps = kbps; } self.unacked = 0; } + /// An accepted mode switch: the encoder's ceiling and compute knee are properties of the + /// MODE (4K120 caps where 1080p60 never would) — drop the mode-scoped learned state. The + /// probe-measured `ceiling_kbps` (a LINK property) survives. + pub(crate) fn on_mode_switch(&mut self) { + self.host_cap_kbps = None; + self.short_acks = 0; + self.cap_probe_windows = 0; + self.encode_means.clear(); + } + /// Feed one report window; returns the rate to request now, if any. `dropped` = frames that /// went FEC-unrecoverable in the window, `loss_ppm` the window's [`crate::quic::LossReport`] /// figure, `owd_mean_us` the window's mean skew-corrected capture→received latency (`None` /// without a clock handshake), `decode_mean_us` the window's mean client decode-stage latency - /// (`None` on an embedder that doesn't report it — the signal is then absent), `actual_kbps` - /// the window's ACTUAL delivered throughput (wire bytes received ÷ window — what the pipeline - /// really carried, as opposed to the target it was allowed; feeds the utilization climb gate - /// and the proven-throughput high-water mark), `flushed` = the pump's jump-to-live fired in - /// the window. + /// (`None` on an embedder that doesn't report it — the signal is then absent), + /// `encode_mean_us` the window's mean HOST encode-stage latency (from the per-AU 0xCF + /// datagrams; `None` on an old host that doesn't send them), `actual_kbps` the window's + /// ACTUAL delivered throughput (wire bytes received ÷ window — what the pipeline really + /// carried, as opposed to the target it was allowed; feeds the utilization climb gate and + /// the proven-throughput high-water mark), `flushed` = the pump's jump-to-live fired in the + /// window. #[allow(clippy::too_many_arguments)] pub(crate) fn on_window( &mut self, @@ -192,6 +290,7 @@ impl BitrateController { loss_ppm: u32, owd_mean_us: Option, decode_mean_us: Option, + encode_mean_us: Option, actual_kbps: u32, flushed: bool, ) -> Option { @@ -242,6 +341,23 @@ impl BitrateController { } None => (false, false), }; + // Host-encode latency: the same rolling-min-baseline treatment, measuring the HOST'S + // encoder — the compute-knee down-driver (see [`ENCODE_RISE_US`]). This is the only + // signal that can push an already-too-high rate back under the knee: the host refuses + // further climbs while behind cadence, but nothing else ever DESCENDS on a clean LAN. + let (encode_bad, encode_severe) = match encode_mean_us { + Some(mean) => { + let base = self.encode_means.iter().min().copied(); + let bad = base.is_some_and(|b| mean > b + ENCODE_RISE_US); + let severe = base.is_some_and(|b| mean > b + ENCODE_SEVERE_US); + if self.encode_means.len() == BASELINE_WINDOWS { + self.encode_means.pop_front(); + } + self.encode_means.push_back(mean); + (bad, severe) + } + None => (false, false), + }; // The proven-throughput high-water mark: this window's delivered rate is now demonstrably // digestible (decode latency stayed flat while it was carried). Loss doesn't disqualify — // the bytes that DID arrive still went through the decoder; what loss means for the rate @@ -253,8 +369,9 @@ impl BitrateController { // deep decode-latency excursion) or loss far past any blip — one window is enough. // Ordinary congestion (heavy-but-recoverable loss, an OWD rise, a decode-latency rise) // still needs two consecutive windows. - let severe = dropped > 0 || flushed || loss_ppm >= SEVERE_LOSS_PPM || decode_severe; - let bad = severe || loss_ppm >= HEAVY_LOSS_PPM || owd_bad || decode_bad; + let severe = + dropped > 0 || flushed || loss_ppm >= SEVERE_LOSS_PPM || decode_severe || encode_severe; + let bad = severe || loss_ppm >= HEAVY_LOSS_PPM || owd_bad || decode_bad || encode_bad; if bad { self.bad_windows += 1; self.clean_windows = 0; @@ -264,6 +381,29 @@ impl BitrateController { self.clean_windows += 1; self.bad_windows = 0; } + // The learned host cap re-probe (see [`CAP_REPROBE_WINDOWS`]): after ~60 s of clean + // windows parked at the cap, lift it one step (+12.5 %, ceiling-bounded) so a + // scene-dependent refusal can't quietly cap the whole session — a still-standing limit + // just re-latches from the next pair of short acks, at zero encoder cost. + if let Some(cap) = self.host_cap_kbps { + if bad { + self.cap_probe_windows = 0; + } else if self.current_kbps >= cap.saturating_sub(cap / 16) { + self.cap_probe_windows += 1; + if self.cap_probe_windows >= CAP_REPROBE_WINDOWS { + self.cap_probe_windows = 0; + let lifted = cap.saturating_add(cap / 8).min(self.ceiling_kbps); + if lifted > cap { + tracing::debug!( + from_kbps = cap, + to_kbps = lifted, + "adaptive bitrate: re-probing above the learned host cap" + ); + self.host_cap_kbps = Some(lifted); + } + } + } + } let cooled = self .last_change .is_none_or(|t| now.duration_since(t) >= CHANGE_COOLDOWN); @@ -284,10 +424,14 @@ impl BitrateController { // utilized window after a long-enough clean run climbs immediately. let utilized = actual_kbps as u64 * UTILIZATION_DEN >= self.current_kbps as u64 * UTILIZATION_NUM; - let cap = self + // The effective ceiling folds in the host-taught cap: the probe measured the LINK, but + // the host's short acks measured the ENCODER — whichever binds first is the limit. + let eff_ceiling = self .ceiling_kbps + .min(self.host_cap_kbps.unwrap_or(u32::MAX)); + let cap = eff_ceiling .min(self.proven_kbps.saturating_mul(PROVEN_HEADROOM_NUM) / PROVEN_HEADROOM_DEN); - if self.current_kbps < self.ceiling_kbps && utilized && cap > self.current_kbps { + if self.current_kbps < eff_ceiling && utilized && cap > self.current_kbps { // Slow start: double on every cooled clean window until the first congestion signal // (this is how an Automatic session reaches a probe-measured ceiling in seconds). // Congestion avoidance: +~6 % after a sustained clean run. @@ -308,6 +452,7 @@ impl BitrateController { fn request(&mut self, kbps: u32, now: Instant) -> Option { self.last_change = Some(now); self.unacked += 1; + self.last_requested_kbps = Some(kbps); // `current_kbps` is NOT updated here — the host's ack is authoritative. A lost/ignored // request just recomputes from the same base next time (and counts toward MAX_UNACKED). Some(kbps) @@ -331,7 +476,16 @@ mod tests { fn run_clean(c: &mut BitrateController, start: Instant, from: u32, n: u32) -> Option { let mut out = None; for i in from..from + n { - out = c.on_window(ticks(start, i), 0, 0, Some(10_000), None, 1_000_000, false); + out = c.on_window( + ticks(start, i), + 0, + 0, + Some(10_000), + None, + None, + 1_000_000, + false, + ); if out.is_some() { return out; } @@ -345,7 +499,7 @@ mod tests { let mut c = BitrateController::new(0); let now = Instant::now(); assert_eq!( - c.on_window(now, 5, 900_000, Some(500_000), None, 1_000_000, true), + c.on_window(now, 5, 900_000, Some(500_000), None, None, 1_000_000, true), None ); } @@ -356,22 +510,58 @@ mod tests { let start = Instant::now(); // Heavy-but-recoverable loss (2–6 %) is ORDINARY: one window is a blip — no reaction. assert_eq!( - c.on_window(ticks(start, 0), 0, 25_000, None, None, 1_000_000, false), + c.on_window( + ticks(start, 0), + 0, + 25_000, + None, + None, + None, + 1_000_000, + false + ), None ); // The second consecutive bad window backs off ×0.7. assert_eq!( - c.on_window(ticks(start, 1), 0, 25_000, None, None, 1_000_000, false), + c.on_window( + ticks(start, 1), + 0, + 25_000, + None, + None, + None, + 1_000_000, + false + ), Some(14_000) ); c.on_ack(14_000); // Still bad after the cooldown → another ×0.7 step from the ACKED rate. assert_eq!( - c.on_window(ticks(start, 6), 0, 25_000, None, None, 1_000_000, false), + c.on_window( + ticks(start, 6), + 0, + 25_000, + None, + None, + None, + 1_000_000, + false + ), None ); // bad #1 again assert_eq!( - c.on_window(ticks(start, 7), 0, 25_000, None, None, 1_000_000, false), + c.on_window( + ticks(start, 7), + 0, + 25_000, + None, + None, + None, + 1_000_000, + false + ), Some(9_800) ); } @@ -382,19 +572,28 @@ mod tests { let mut c = BitrateController::new(20_000); let start = Instant::now(); assert_eq!( - c.on_window(ticks(start, 0), 1, 0, None, None, 1_000_000, false), + c.on_window(ticks(start, 0), 1, 0, None, None, None, 1_000_000, false), Some(14_000) ); // …and so does a jump-to-live flush. let mut c = BitrateController::new(20_000); assert_eq!( - c.on_window(ticks(start, 0), 0, 0, None, None, 1_000_000, true), + c.on_window(ticks(start, 0), 0, 0, None, None, None, 1_000_000, true), Some(14_000) ); // …and ≥6 % window loss. let mut c = BitrateController::new(20_000); assert_eq!( - c.on_window(ticks(start, 0), 0, 80_000, None, None, 1_000_000, false), + c.on_window( + ticks(start, 0), + 0, + 80_000, + None, + None, + None, + 1_000_000, + false + ), Some(14_000) ); } @@ -404,18 +603,18 @@ mod tests { let mut c = BitrateController::new(20_000); let start = Instant::now(); assert_eq!( - c.on_window(ticks(start, 0), 1, 0, None, None, 1_000_000, false), + c.on_window(ticks(start, 0), 1, 0, None, None, None, 1_000_000, false), Some(14_000) ); c.on_ack(14_000); // A severe window INSIDE the 1.5 s cooldown (tick 1 = 750 ms) → held; at the cooldown // boundary (tick 2 = 1.5 s) it fires. assert_eq!( - c.on_window(ticks(start, 1), 1, 0, None, None, 1_000_000, false), + c.on_window(ticks(start, 1), 1, 0, None, None, None, 1_000_000, false), None ); assert_eq!( - c.on_window(ticks(start, 2), 1, 0, None, None, 1_000_000, false), + c.on_window(ticks(start, 2), 1, 0, None, None, None, 1_000_000, false), Some(9_800) ); } @@ -426,17 +625,17 @@ mod tests { let start = Instant::now(); // ×0.7 of 6000 = 4200 < floor → clamped to 5000. assert_eq!( - c.on_window(ticks(start, 0), 1, 0, None, None, 1_000_000, false), + c.on_window(ticks(start, 0), 1, 0, None, None, None, 1_000_000, false), Some(5_000) ); c.on_ack(5_000); // At the floor, further bad windows request nothing. assert_eq!( - c.on_window(ticks(start, 6), 1, 0, None, None, 1_000_000, false), + c.on_window(ticks(start, 6), 1, 0, None, None, None, 1_000_000, false), None ); assert_eq!( - c.on_window(ticks(start, 7), 1, 0, None, None, 1_000_000, false), + c.on_window(ticks(start, 7), 1, 0, None, None, None, 1_000_000, false), None ); } @@ -446,7 +645,7 @@ mod tests { let mut c = BitrateController::new(20_000); let start = Instant::now(); assert_eq!( - c.on_window(ticks(start, 0), 1, 0, None, None, 1_000_000, false), + c.on_window(ticks(start, 0), 1, 0, None, None, None, 1_000_000, false), Some(14_000) ); c.on_ack(14_000); @@ -469,9 +668,16 @@ mod tests { // Every cooled clean window doubles until the ceiling caps the climb, then quiet. let mut got = Vec::new(); for i in 0..14 { - if let Some(k) = - c.on_window(ticks(start, i), 0, 0, Some(10_000), None, 1_000_000, false) - { + if let Some(k) = c.on_window( + ticks(start, i), + 0, + 0, + Some(10_000), + None, + None, + 1_000_000, + false, + ) { c.on_ack(k); got.push(k); } @@ -485,20 +691,47 @@ mod tests { c.set_ceiling(300_000); let start = Instant::now(); assert_eq!( - c.on_window(ticks(start, 0), 0, 0, Some(10_000), None, 1_000_000, false), + c.on_window( + ticks(start, 0), + 0, + 0, + Some(10_000), + None, + None, + 1_000_000, + false + ), Some(40_000) ); c.on_ack(40_000); // Severe window → immediate ×0.7, and slow start is over. assert_eq!( - c.on_window(ticks(start, 2), 1, 0, Some(10_000), None, 1_000_000, false), + c.on_window( + ticks(start, 2), + 1, + 0, + Some(10_000), + None, + None, + 1_000_000, + false + ), Some(28_000) ); c.on_ack(28_000); // Clean again — but the next climb is additive, after the 6-window clean run. let mut next = None; for i in 3..12 { - next = c.on_window(ticks(start, i), 0, 0, Some(10_000), None, 1_000_000, false); + next = c.on_window( + ticks(start, i), + 0, + 0, + Some(10_000), + None, + None, + 1_000_000, + false, + ); if next.is_some() { assert!(i >= 8, "additive climb must wait for the clean run"); break; @@ -512,7 +745,7 @@ mod tests { let mut c = BitrateController::new(0); c.set_ceiling(1_000_000); assert_eq!( - c.on_window(Instant::now(), 0, 0, None, None, 1_000_000, false), + c.on_window(Instant::now(), 0, 0, None, None, None, 1_000_000, false), None ); let mut c = BitrateController::new(20_000); @@ -527,17 +760,44 @@ mod tests { // Establish a ~10 ms baseline over a few clean windows. for i in 0..4 { assert_eq!( - c.on_window(ticks(start, i), 0, 0, Some(10_000), None, 1_000_000, false), + c.on_window( + ticks(start, i), + 0, + 0, + Some(10_000), + None, + None, + 1_000_000, + false + ), None ); } // Delay climbs 40 ms above baseline with ZERO loss — bufferbloat. Two windows → back off. assert_eq!( - c.on_window(ticks(start, 4), 0, 0, Some(50_000), None, 1_000_000, false), + c.on_window( + ticks(start, 4), + 0, + 0, + Some(50_000), + None, + None, + 1_000_000, + false + ), None ); assert_eq!( - c.on_window(ticks(start, 5), 0, 0, Some(52_000), None, 1_000_000, false), + c.on_window( + ticks(start, 5), + 0, + 0, + Some(52_000), + None, + None, + 1_000_000, + false + ), Some(14_000) ); } @@ -557,6 +817,7 @@ mod tests { 0, Some(10_000), Some(8_000), + None, 1_000_000, false ), @@ -572,6 +833,7 @@ mod tests { 0, Some(10_000), Some(38_000), + None, 1_000_000, false ), @@ -584,6 +846,7 @@ mod tests { 0, Some(10_000), Some(40_000), + None, 1_000_000, false ), @@ -605,6 +868,7 @@ mod tests { 0, Some(10_000), Some(8_000), + None, 1_000_000, false ), @@ -620,6 +884,7 @@ mod tests { 0, Some(10_000), Some(38_000), + None, 1_000_000, false ), @@ -634,6 +899,7 @@ mod tests { 0, Some(10_000), Some(40_000), + None, 1_000_000, false ), @@ -657,6 +923,7 @@ mod tests { 0, Some(10_000), Some(8_000), + None, 2_000, false ), @@ -673,6 +940,7 @@ mod tests { 0, Some(10_000), Some(8_000), + None, 18_000, false ), @@ -696,6 +964,7 @@ mod tests { 0, Some(10_000), Some(8_000), + None, 20_000, false ), @@ -710,6 +979,7 @@ mod tests { 0, Some(10_000), Some(8_000), + None, 30_000, false ), @@ -732,6 +1002,7 @@ mod tests { 0, Some(10_000), Some(8_000), + None, 20_000, false ), @@ -741,7 +1012,16 @@ mod tests { // A long calm stretch (2 % utilization, decoder idle): the controller stays silent. for i in 2..30 { assert_eq!( - c.on_window(ticks(start, i), 0, 0, Some(10_000), Some(4_000), 600, false), + c.on_window( + ticks(start, i), + 0, + 0, + Some(10_000), + Some(4_000), + None, + 600, + false + ), None ); } @@ -761,6 +1041,7 @@ mod tests { 0, Some(10_000), Some(8_000), + None, 1_000_000, false ), @@ -776,6 +1057,7 @@ mod tests { 0, Some(10_000), Some(60_000), + None, 1_000_000, false ), @@ -783,6 +1065,235 @@ mod tests { ); } + #[test] + fn two_identical_short_acks_latch_the_host_cap() { + // The 4K120 field failure: the encoder ceilings at ~794 Mbps while the link carries + // more — the host acks short. TWO identical short acks teach the cap; climbs then stop + // poking a limit the host already refused (the rebuild-storm driver). + let mut c = BitrateController::new(400_000); + c.set_ceiling(1_400_000); + let start = Instant::now(); + assert_eq!(run_clean(&mut c, start, 0, 1), Some(800_000)); + // First short ack: current follows (authoritative), but one short ack is not a cap. + c.on_ack(794_000); + assert!(c.host_cap_kbps.is_none()); + // The next climb overshoots again and is short-acked at the SAME value: latch. + assert_eq!(run_clean(&mut c, start, 10, 1), Some(1_400_000)); + c.on_ack(794_000); + assert_eq!(c.host_cap_kbps, Some(794_000)); + // Parked AT the learned cap, nothing left to climb to — no more requests. + assert_eq!(run_clean(&mut c, start, 20, 12), None); + } + + #[test] + fn one_short_ack_is_a_transient_not_a_cap() { + // A failed host rebuild acks short once (it kept the old rate) — latching THAT would + // cap the session on a driver hiccup. The streak must survive only identical repeats. + let mut c = BitrateController::new(400_000); + c.set_ceiling(1_400_000); + let start = Instant::now(); + assert_eq!(run_clean(&mut c, start, 0, 1), Some(800_000)); + c.on_ack(400_000); // rebuild failed, host kept the old rate + assert!(c.host_cap_kbps.is_none()); + // The retry applies fully: streak broken, still no cap, full authority kept. + assert_eq!(run_clean(&mut c, start, 10, 1), Some(800_000)); + c.on_ack(800_000); + assert!(c.host_cap_kbps.is_none()); + } + + #[test] + fn mode_switch_clears_the_learned_cap() { + let mut c = BitrateController::new(400_000); + c.set_ceiling(1_400_000); + let start = Instant::now(); + assert_eq!(run_clean(&mut c, start, 0, 1), Some(800_000)); + c.on_ack(794_000); + assert_eq!(run_clean(&mut c, start, 10, 1), Some(1_400_000)); + c.on_ack(794_000); + assert_eq!(c.host_cap_kbps, Some(794_000)); + // 4K120's ceiling means nothing at the new mode — the cap must not survive the switch + // (the probe-measured link ceiling does). + c.on_mode_switch(); + assert!(c.host_cap_kbps.is_none()); + assert_eq!(c.ceiling_kbps, 1_400_000); + } + + #[test] + fn learned_cap_reprobes_after_a_sustained_clean_run() { + // A cadence-refusal cap is scene evidence, not a spec limit: after ~60 s parked clean + // at the cap, lift one step so a one-time heavy scene can't cap the session forever. A + // still-standing limit just re-latches from the next short-ack pair, at zero cost. + let mut c = BitrateController::new(400_000); + c.set_ceiling(1_400_000); + let start = Instant::now(); + assert_eq!(run_clean(&mut c, start, 0, 1), Some(800_000)); + c.on_ack(794_000); + assert_eq!(run_clean(&mut c, start, 10, 1), Some(1_400_000)); + c.on_ack(794_000); + assert_eq!(c.host_cap_kbps, Some(794_000)); + for i in 0..CAP_REPROBE_WINDOWS { + let _ = c.on_window( + ticks(start, 20 + i), + 0, + 0, + Some(10_000), + None, + None, + 1_000_000, + false, + ); + } + assert_eq!(c.host_cap_kbps, Some(794_000 + 794_000 / 8)); + } + + #[test] + fn host_encode_latency_rise_backs_off() { + // The compute knee: link pristine, client decoder fine — only HOST encode time moves + // (the 4K120 case: ~9.3 ms against an 8.33 ms budget shows up nowhere else). Two risen + // windows → ×0.7, exactly like an OWD/decode rise. + let mut c = BitrateController::new(20_000); + let start = Instant::now(); + for i in 0..4 { + assert_eq!( + c.on_window( + ticks(start, i), + 0, + 0, + Some(10_000), + None, + Some(7_000), + 1_000_000, + false + ), + None + ); + } + assert_eq!( + c.on_window( + ticks(start, 4), + 0, + 0, + Some(10_000), + None, + Some(11_500), + 1_000_000, + false + ), + None + ); + assert_eq!( + c.on_window( + ticks(start, 6), + 0, + 0, + Some(10_000), + None, + Some(12_000), + 1_000_000, + false + ), + Some(14_000) + ); + } + + #[test] + fn deep_encode_excursion_is_severe() { + // Encode time shooting ≈1.5 frame budgets over baseline = the queue is growing past + // the knee right now — no two-window confirmation. + let mut c = BitrateController::new(20_000); + let start = Instant::now(); + for i in 0..4 { + assert_eq!( + c.on_window( + ticks(start, i), + 0, + 0, + Some(10_000), + None, + Some(7_000), + 1_000_000, + false + ), + None + ); + } + assert_eq!( + c.on_window( + ticks(start, 4), + 0, + 0, + Some(10_000), + None, + Some(20_000), + 1_000_000, + false + ), + Some(14_000) + ); + } + + #[test] + fn rate_decrease_rebases_the_encode_baseline() { + // After OUR OWN decrease the encode regime legitimately changes (less work per frame; + // an escalated host's reported encode_us also carries a queue offset) — the old + // baseline must not train-fire repeated backoffs down to the floor. + let mut c = BitrateController::new(20_000); + let start = Instant::now(); + for i in 0..4 { + let _ = c.on_window( + ticks(start, i), + 0, + 0, + Some(10_000), + None, + Some(7_000), + 1_000_000, + false, + ); + } + let _ = c.on_window( + ticks(start, 4), + 0, + 0, + Some(10_000), + None, + Some(12_000), + 1_000_000, + false, + ); + assert_eq!( + c.on_window( + ticks(start, 6), + 0, + 0, + Some(10_000), + None, + Some(12_500), + 1_000_000, + false + ), + Some(14_000) + ); + // The decrease applies → rebase. The new regime's ~15 ms means (an escalated host's + // queue offset) would be far over the OLD 7 ms baseline, but must now read clean. + c.on_ack(14_000); + for i in 8..11 { + assert_eq!( + c.on_window( + ticks(start, i), + 0, + 0, + Some(10_000), + None, + Some(15_000), + 1_000_000, + false + ), + None + ); + } + } + #[test] fn ack_silence_disables_the_controller() { let mut c = BitrateController::new(20_000); @@ -791,7 +1302,7 @@ mod tests { let mut i = 0; // Keep every window bad and never ack: exactly MAX_UNACKED requests, then silence. while i < 60 { - if c.on_window(ticks(start, i), 1, 0, None, None, 1_000_000, false) + if c.on_window(ticks(start, i), 1, 0, None, None, None, 1_000_000, false) .is_some() { sent += 1; diff --git a/crates/punktfunk-core/src/client/frame_channel.rs b/crates/punktfunk-core/src/client/frame_channel.rs index 39d7e88d..32c49e8c 100644 --- a/crates/punktfunk-core/src/client/frame_channel.rs +++ b/crates/punktfunk-core/src/client/frame_channel.rs @@ -224,6 +224,20 @@ pub(crate) struct DecodeLatAcc { pub(crate) count: u32, } +/// Host encode-stage latency accumulator — [`DecodeLatAcc`]'s mirror for the HOST side of the +/// pipeline. The datagram task adds one sample per 0xCF `HostStages::encode_us` (host encoder +/// submit → bitstream ready) and the pump drains a window mean into +/// [`crate::abr::BitrateController::on_window`]'s encode signal. Host encode time was measured, +/// shipped and drawn on the overlay, but never an ABR input — which is how a fat-LAN Automatic +/// session drove the encoder past its compute knee with nothing to stop it (§ABR overdrive). +/// Its own accumulator rather than the overlay's `host_timing` channel: that channel is a lossy +/// `try_send` the embedder may never drain, and the controller must not depend on it. +#[derive(Default)] +pub(crate) struct EncodeLatAcc { + pub(crate) sum_us: u64, + pub(crate) count: u32, +} + /// The pre-decode video hand-off from the data-plane pump to the embedder. Unlike the side planes /// (self-contained samples that drop the newest on overflow), video AUs are reference-chained under the /// host's infinite GOP: dropping ANY frame mid-stream corrupts every dependent frame until the next diff --git a/crates/punktfunk-core/src/client/pump.rs b/crates/punktfunk-core/src/client/pump.rs index 74fa2055..aee9b3a3 100644 --- a/crates/punktfunk-core/src/client/pump.rs +++ b/crates/punktfunk-core/src/client/pump.rs @@ -112,6 +112,12 @@ pub(super) async fn run_pump(args: WorkerArgs) { // Adaptive bitrate ack slot: the control task parks the latest BitrateChanged here; the // pump's controller drains it on its report tick (`take()` — an ack is consumed once). let bitrate_ack: Arc>> = Arc::new(Mutex::new(None)); + // Host-encode-latency accumulator (the ABR encode signal, see [`EncodeLatAcc`]): the + // datagram task adds one sample per 0xCF; the pump drains a window mean per report tick. + let encode_lat = Arc::new(Mutex::new(super::frame_channel::EncodeLatAcc::default())); + // Bumped by the control task on every accepted mode switch (the `clock_gen` pattern): the + // pump resets the controller's mode-scoped learned state (host cap, encode baseline). + let mode_gen = Arc::new(AtomicU32::new(0)); // Control task (see [`control_task`]): the handshake stream stays open for mid-stream // renegotiation, speed tests, clock re-sync, and clipboard metadata. @@ -128,6 +134,7 @@ pub(super) async fn run_pump(args: WorkerArgs) { clock_gen: clock_gen.clone(), clip_event_tx: clip_event_tx.clone(), cursor_shape_tx, + mode_gen: mode_gen.clone(), } .run(), ); @@ -141,6 +148,7 @@ pub(super) async fn run_pump(args: WorkerArgs) { hidout_tx, hdr_meta_tx, host_timing_tx, + encode_lat.clone(), cursor_state_tx, )); @@ -176,6 +184,8 @@ pub(super) async fn run_pump(args: WorkerArgs) { clock_offset, clock_gen, decode_lat, + encode_lat, + mode_gen, frames_dropped, fec_recovered, bitrate_ack, diff --git a/crates/punktfunk-core/src/client/pump/control_task.rs b/crates/punktfunk-core/src/client/pump/control_task.rs index f236baba..b34d1078 100644 --- a/crates/punktfunk-core/src/client/pump/control_task.rs +++ b/crates/punktfunk-core/src/client/pump/control_task.rs @@ -24,6 +24,10 @@ pub(super) struct ControlTask { /// Host cursor shapes ([`CursorShape`], sent on pointer-bitmap change) → the embedder's /// shape plane ([`NativeClient::next_cursor_shape`]). pub(super) cursor_shape_tx: std::sync::mpsc::SyncSender, + /// Bumped on every ACCEPTED mode switch (the `clock_gen` pattern): the pump watches it and + /// resets the bitrate controller's mode-scoped learned state — the encoder ceiling / compute + /// knee it was taught belong to the OLD mode. + pub(super) mode_gen: Arc, } impl ControlTask { @@ -40,6 +44,7 @@ impl ControlTask { clock_gen, clip_event_tx, cursor_shape_tx, + mode_gen, } = self; // Mid-stream clock re-sync (see [`ClockResync`]): a batch runs every // CLOCK_RESYNC_INTERVAL and whenever the pump asks (CtrlRequest::ClockResync after @@ -88,6 +93,7 @@ impl ControlTask { if let Ok(ack) = Reconfigured::decode(&msg) { if ack.accepted { *mode_slot.lock().unwrap() = ack.mode; + mode_gen.fetch_add(1, Ordering::Relaxed); tracing::info!(mode = ?ack.mode, "host accepted mode switch"); } else { tracing::warn!(active = ?ack.mode, "host rejected mode switch"); diff --git a/crates/punktfunk-core/src/client/pump/data.rs b/crates/punktfunk-core/src/client/pump/data.rs index c28722a5..8b0cae8b 100644 --- a/crates/punktfunk-core/src/client/pump/data.rs +++ b/crates/punktfunk-core/src/client/pump/data.rs @@ -19,6 +19,12 @@ pub(super) struct DataPump { pub(super) clock_offset: Arc, pub(super) clock_gen: Arc, pub(super) decode_lat: Arc>, + /// Host encode-stage latency window accumulator (the ABR encode signal — see + /// [`super::super::frame_channel::EncodeLatAcc`]); fed by the datagram task. + pub(super) encode_lat: Arc>, + /// Accepted-mode-switch generation (control task bumps): a change resets the controller's + /// mode-scoped learned state ([`BitrateController::on_mode_switch`]). + pub(super) mode_gen: Arc, pub(super) frames_dropped: Arc, pub(super) fec_recovered: Arc, pub(super) bitrate_ack: Arc>>, @@ -41,6 +47,8 @@ impl DataPump { clock_offset: pump_clock_offset, clock_gen: pump_clock_gen, decode_lat: pump_decode_lat, + encode_lat: pump_encode_lat, + mode_gen: pump_mode_gen, frames_dropped, fec_recovered, bitrate_ack, @@ -127,6 +135,7 @@ impl DataPump { let mut clock_detector_armed = true; let mut resync_wanted = false; let mut seen_clock_gen = pump_clock_gen.load(Ordering::Relaxed); + let mut seen_mode_gen = pump_mode_gen.load(Ordering::Relaxed); // Standing-latency bleed (see StandingLatency): the third detector, for the small, // constant, loss-free OWD elevation the two jump-to-live detectors deliberately // tolerate (< QUEUE_HIGH frames, < FLUSH_LATENCY behind) — a sub-frame standing @@ -338,9 +347,15 @@ impl DataPump { ); } } - // Adaptive bitrate: drain any host ack first (its clamp is authoritative), then - // feed the controller this window's congestion signals; a decision becomes a - // SetBitrate on the control stream. + // Adaptive bitrate: an accepted mode switch first (it invalidates the + // mode-scoped learned state), then drain any host ack (its clamp is + // authoritative), then feed the controller this window's congestion signals; a + // decision becomes a SetBitrate on the control stream. + let mg = pump_mode_gen.load(Ordering::Relaxed); + if mg != seen_mode_gen { + seen_mode_gen = mg; + abr.on_mode_switch(); + } if let Some(acked) = bitrate_ack.lock().unwrap().take() { abr.on_ack(acked); } @@ -356,6 +371,14 @@ impl DataPump { *acc = DecodeLatAcc::default(); (count > 0).then(|| (sum / count as u64) as i64) }; + // Same drain for the host-encode window (0xCF `encode_us` via the datagram + // task) — `None` on an old host that doesn't send stage timings. + let encode_mean_us = { + let mut acc = pump_encode_lat.lock().unwrap(); + let (sum, count) = (acc.sum_us, acc.count); + *acc = Default::default(); + (count > 0).then(|| (sum / count as u64) as i64) + }; // The window's ACTUAL delivered throughput — what the pipeline really carried, vs // the target it was allowed. Wire bytes (headers + FEC) slightly overstate the // media rate the decoder ingests; acceptable for the climb gate / proven-mark @@ -369,17 +392,19 @@ impl DataPump { loss_ppm, owd_mean_us, decode_mean_us, + encode_mean_us, actual_kbps, flush_in_window, ) { // Log the window's signals alongside the decision so an on-glass session can - // tell a decode-driven re-target (the new signal — decode_mean_us elevated with + // tell a decode-/encode-driven re-target (the new signals — elevated with // loss/OWD flat) from a network-driven one. tracing::info!( kbps, loss_ppm, owd_mean_us = owd_mean_us.unwrap_or(-1), decode_mean_us = decode_mean_us.unwrap_or(-1), + encode_mean_us = encode_mean_us.unwrap_or(-1), actual_kbps, flushed = flush_in_window, "adaptive bitrate: requesting encoder re-target" diff --git a/crates/punktfunk-core/src/client/pump/datagram_task.rs b/crates/punktfunk-core/src/client/pump/datagram_task.rs index 904f2144..b62c8dd7 100644 --- a/crates/punktfunk-core/src/client/pump/datagram_task.rs +++ b/crates/punktfunk-core/src/client/pump/datagram_task.rs @@ -14,6 +14,9 @@ pub(super) async fn run( hidout_tx: std::sync::mpsc::SyncSender, hdr_meta_tx: std::sync::mpsc::SyncSender, host_timing_tx: std::sync::mpsc::SyncSender, + // The ABR encode signal's accumulator (see [`EncodeLatAcc`]) — fed HERE, not off + // `host_timing_tx`: that channel is the overlay's, lossy and embedder-drained. + encode_lat: Arc>, cursor_state_tx: std::sync::mpsc::SyncSender, ) { // Per-pad reorder gate for v2 rumble envelopes (the seq analog of the host's gamepad-state @@ -74,6 +77,11 @@ pub(super) async fn run( } Some(&crate::quic::HOST_TIMING_MAGIC) => { if let Some(t) = crate::quic::decode_host_timing_datagram(&d) { + if let Some(s) = &t.stages { + let mut acc = encode_lat.lock().unwrap(); + acc.sum_us += s.encode_us as u64; + acc.count += 1; + } let _ = host_timing_tx.try_send(t); } }