From f7a8c2013d6f0973f14699ab2867e9e6297b75fc Mon Sep 17 00:00:00 2001 From: enricobuehler Date: Mon, 3 Aug 2026 01:07:33 +0200 Subject: [PATCH 1/5] fix(client/abr): the controller stops learning the wrong lessons from one window MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Six defects found by a sweep of the Automatic-bitrate path, all of them the same shape: a single window, or a single refusal, taught the controller something it then treated as permanent. - Rolling baselines (OWD, client decode, host encode) armed off ONE sample. The baseline is a rolling minimum, so one window IS the floor — and `on_ack` deliberately clears the encode baseline after every decrease we ourselves asked for, re-opening that hole each time. A calm re-seed window followed by ordinary motion read as 4 ms of "congestion", backed off, cleared again, and ratcheted toward the floor on a link that was never the problem. All three now need BASELINE_MIN_WINDOWS of evidence before they may fire, via one shared `score_baseline` (the three copies had already drifted apart). - A mode switch rebased only the encode baseline. Decode and OWD are just as mode-scoped: 4K120 decodes slower and puts bigger frames on the wire than 1080p60, so the old floor was one the new mode cleared on its first window — ~30 s of every window scoring bad, i.e. a backoff every other window. A switch UP in mode cratered the rate instead of raising it. `proven_kbps` goes with them; throughput the old mode's decoder digested is not evidence about this one. - `proven_kbps` — never decayed, and permanent authority over how far every later climb may step — was raised by any window without a decode rise, including ones scored SEVERE. The windows that overstate delivered throughput are exactly the damaged ones: a stall's backlog draining at once, a flush's queue, the FEC surge answering a loss burst. Now only clean windows raise it. - A learned cap escaped at +12.5 % per ~60 s. The host cannot distinguish a durable encoder ceiling from a climb refused while it is transiently behind cadence, and the latter routinely latches during slow start at the 20 Mbps default — from which crossing the gap to a probe-measured ceiling took upwards of twenty minutes. Re-probe after 12 s instead, doubling the interval each time the lift is immediately re-learned: a transient is out in one interval, a real ceiling settles into a slow poll. - The decode cap latched AT the rate that choked, authorizing a climb straight back into the failure, and a bare jump-to-live flush could teach a "decoder knee" from what was a network event. It now latches just under the choke rate (inside the ±1/8 band the evidence already required) and only credits a flush where the decode signal is absent and cannot speak for itself. - PUNKTFUNK_ABR_MAX_MBPS bound only probe-learned ceilings, not the negotiated start rate — so the one knob an Automatic session gives the operator did nothing when the session already started above it. It now binds at construction, and a session sitting above its ceiling steps down to it (no congestion signal will ever find that: the link is fine, the cap is policy). Also: a SetBitrate dropped by a full control queue counted toward MAX_UNACKED, so three of them retired the controller for the session while logging that an "older host" was at fault. The pump now tells the controller what happened. Wire format and ABI untouched. 34 abr tests green (3 new). --- crates/punktfunk-core/src/abr.rs | 474 ++++++++++++++---- crates/punktfunk-core/src/client/pump/data.rs | 10 +- 2 files changed, 375 insertions(+), 109 deletions(-) diff --git a/crates/punktfunk-core/src/abr.rs b/crates/punktfunk-core/src/abr.rs index f194d7e0..623e1633 100644 --- a/crates/punktfunk-core/src/abr.rs +++ b/crates/punktfunk-core/src/abr.rs @@ -24,17 +24,19 @@ //! 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 //! (heavy-but-recoverable loss, an OWD rise, a decode rise) needs two consecutive bad windows. -//! Recovery is two-mode: **slow start** — until the first congestion signal the rate DOUBLES each -//! clean window (cooldown-paced), which is how an Automatic session climbs from the conservative -//! start to the [`set_ceiling`](BitrateController::set_ceiling) measured by the startup -//! link-capacity probe in seconds instead of minutes — then classic additive recovery (+~6 % -//! after ~4.5 s clean, ceilinged). Changes are rate-limited (each one costs the IDR the host's +//! Recovery is two-mode: **slow start** — until the first congestion signal each clean window +//! asks for double the current rate, bounded (like every climb) by the proven-throughput +//! headroom below, so the step a loaded session actually takes is ×1.5 over what it last +//! delivered; either way it climbs from the conservative start to the +//! [`set_ceiling`](BitrateController::set_ceiling) measured by the startup link-capacity probe +//! in seconds rather than minutes — then classic additive recovery (+~6 % after ~4.5 s clean, +//! ceilinged). Changes are rate-limited (each one costs the IDR the host's //! rebuilt encoder opens with) and the whole controller disables itself against a host that never //! answers [`crate::quic::BitrateChanged`] (an older build that ignores unknown control messages). //! Standing limits are LEARNED rather than re-poked: two identical short host acks latch the //! encoder's ceiling (`host_cap_kbps`), two consecutive decode-severe backoffs at a similar rate //! latch the client decoder's knee (`decode_cap_kbps`) — and both re-probe slowly -//! ([`CAP_REPROBE_WINDOWS`]) so neither latch outlives the condition that taught it. +//! ([`CAP_REPROBE_WINDOWS_MIN`]) so neither latch outlives the condition that taught it. //! //! Climbs are additionally **evidence-gated**. The target is only a *promise* to the encoder — //! how many bits it actually emits depends on the content — so on calm content (a menu, an idle @@ -128,15 +130,26 @@ 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. -/// The [`decode cap`](BitrateController::decode_cap_kbps) re-probes on the same clock for the -/// same reason: the decoder's knee moves with content and thermals, so its latch must not be -/// permanent either. -const CAP_REPROBE_WINDOWS: u32 = 80; +/// Clean windows parked at a learned cap before re-probing above it, and the ceiling that +/// interval backs off to. +/// +/// A learned cap is EVIDENCE, not a spec limit: the host's short ack means "not right now", +/// which covers both its encoder's codec-level ceiling (durable) and a climb refused while +/// encode is behind cadence (transient, and routinely latched during slow start at the +/// conservative 20 Mbps default). The client cannot tell those apart from the ack alone, so the +/// re-probe is what keeps a transient from becoming the session's ceiling — and a flat ~60 s +/// clock at +12.5 % made that escape take upwards of twenty minutes to cross the gap to a +/// probe-measured link ceiling, which is indistinguishable from never. +/// +/// So: probe again after 12 s, and DOUBLE the interval each time the lift is immediately +/// re-learned at the same value (see [`on_ack`](BitrateController::on_ack)). A transient is out +/// in one interval; a standing limit settles into a slow poll instead of a permanent one. The +/// re-probe itself is nearly free either way — a still-standing limit re-teaches itself in two +/// short acks, which the host pre-clamps without touching the encoder: no rebuild, no IDR. +/// The [`decode cap`](BitrateController::decode_cap_kbps) re-probes on the same schedule for +/// the same reason: the decoder's knee moves with content and thermals. +const CAP_REPROBE_WINDOWS_MIN: u32 = 16; +const CAP_REPROBE_WINDOWS_MAX: u32 = 128; /// Two consecutive decode-driven backoffs latch the /// [`decode cap`](BitrateController::decode_cap_kbps) only when their pre-backoff rates agree /// within ±1/8: the decoder's knee is a RATE, so repeated chokes at the same rate are its @@ -146,6 +159,17 @@ const DECODE_CAP_SIMILAR_DIV: u32 = 8; /// 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; +/// Windows a rolling baseline must hold before the signal it feeds may fire. A baseline is a +/// rolling MINIMUM, so a single sample IS the baseline — and if that one window landed on calm +/// content, ordinary content variance clears the rise threshold by itself. That hole is not +/// theoretical: [`on_ack`](BitrateController::on_ack) deliberately CLEARS the encode baseline +/// after every decrease we ourselves asked for, so the encode down-driver re-armed on a +/// one-sample floor each time — a calm re-seed window followed by a motion scene reads as +/// `ENCODE_RISE_US` of "congestion", backs off, clears again, and ratchets to the floor on a +/// link that was never the problem. Four windows (3 s) of evidence before any of the three +/// latency signals may fire costs a little reaction latency at session start and buys a floor +/// that means something. +const BASELINE_MIN_WINDOWS: usize = 4; /// Requests sent without a single [`crate::quic::BitrateChanged`] ack before concluding the host /// predates bitrate renegotiation and going quiet for the rest of the session. const MAX_UNACKED: u32 = 3; @@ -167,6 +191,37 @@ fn ceiling_cap_from_env() -> Option { .map(|m| m.saturating_mul(1_000)) } +/// Score one window's latency sample against its rolling-min baseline, then record it. +/// +/// Shared by all three latency signals (OWD, client decode, host encode) — same shape, different +/// thresholds. `mean` is `None` when nobody reports the signal (no clock handshake, an embedder +/// that doesn't measure decode, a host that ships no stage timings); the signal is then simply +/// absent rather than clean, so it can neither mark a window bad nor teach a baseline. +/// +/// The baseline is the minimum of the PRIOR windows — this window is compared before it is +/// recorded, so a rising window can't drag its own floor up with it — and only counts once +/// [`BASELINE_MIN_WINDOWS`] of them exist. Returns `(rise, severe)`; pass `i64::MAX` for +/// `severe_us` on a signal with no severe tier. +fn score_baseline( + means: &mut VecDeque, + mean: Option, + rise_us: i64, + severe_us: i64, +) -> (bool, bool) { + let Some(mean) = mean else { + return (false, false); + }; + let base = (means.len() >= BASELINE_MIN_WINDOWS) + .then(|| means.iter().min().copied()) + .flatten(); + let over = |t: i64| base.is_some_and(|b| mean > b.saturating_add(t)); + if means.len() == BASELINE_WINDOWS { + means.pop_front(); + } + means.push_back(mean); + (over(rise_us), over(severe_us)) +} + /// One decision per report window; `Some(kbps)` = send a [`crate::quic::SetBitrate`]. pub(crate) struct BitrateController { /// `false` = permanently off (explicit user bitrate, an old host, or ack silence). @@ -199,7 +254,7 @@ pub(crate) struct BitrateController { /// 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. + /// ([`CAP_REPROBE_WINDOWS_MIN`]) 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. @@ -210,8 +265,11 @@ pub(crate) struct BitrateController { /// 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). + /// Clean windows spent parked at the learned cap (the re-probe clock) and the interval it is + /// counting toward — [`CAP_REPROBE_WINDOWS_MIN`], doubled toward + /// [`CAP_REPROBE_WINDOWS_MAX`] each time a lift is immediately re-learned. cap_probe_windows: u32, + cap_reprobe_after: u32, /// The client-decoder rate cap, mirroring [`host_cap_kbps`](Self::host_cap_kbps) for the /// OTHER end of the pipe: latched when two CONSECUTIVE backoffs carried decode-severe /// evidence (a deep decode-latency excursion, or a jump-to-live flush — in the @@ -220,7 +278,7 @@ pub(crate) struct BitrateController { /// ceiling is a permanent 30–60 s sawtooth: every ×0.7 backoff re-climbs toward a ceiling /// the decoder can't hold, and each cycle costs a flush plus a dropped-frame burst (the /// 1440p120 HEVC field case: knee ~490 Mbps under a ~658 Mbps ceiling). Slowly re-probed - /// on the [`CAP_REPROBE_WINDOWS`] clock, exactly like the host cap, so a decoder that + /// on the [`CAP_REPROBE_WINDOWS_MIN`] clock, exactly like the host cap, so a decoder that /// recovers (lighter content, thermal headroom) climbs again — the latch is never /// permanent. decode_cap_kbps: Option, @@ -228,8 +286,10 @@ pub(crate) struct BitrateController { /// decode-driven): the reference the next one must land near ([`DECODE_CAP_SIMILAR_DIV`]) /// to latch the cap — one spurious flush teaches nothing. decode_backoff_kbps: u32, - /// Clean windows spent parked at the learned decode cap (its re-probe clock). + /// Clean windows spent parked at the learned decode cap (its re-probe clock), and that + /// clock's own backoff interval — same schedule as the host cap's. decode_cap_probe_windows: u32, + decode_cap_reprobe_after: 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 @@ -241,6 +301,10 @@ pub(crate) struct BitrateController { last_change: Option, /// Requests since the last ack — reaching [`MAX_UNACKED`] disables the controller. unacked: u32, + /// The last ceiling-clamp target asked for (0 = none). A session running ABOVE its effective + /// ceiling is asked down to it exactly once per distinct target — a host that answers higher + /// has said it cannot go there, and re-asking every cooldown only costs reconfigures. + ceiling_ask_kbps: u32, } impl BitrateController { @@ -257,7 +321,12 @@ impl BitrateController { BitrateController { enabled: start_kbps > 0, current_kbps: start_kbps, - ceiling_kbps: start_kbps, + // The env cap binds the NEGOTIATED ceiling too, not just probe-learned ones. It is + // the only lever an Automatic session gives the operator (Automatic is precisely + // "no explicit bitrate"), so a start rate above it has to come down rather than + // stand as a ceiling the user asked not to reach — see the clamp-down step in + // [`on_window`](Self::on_window). + ceiling_kbps: start_kbps.min(ceiling_cap_kbps.unwrap_or(u32::MAX)), ceiling_cap_kbps, floor_kbps: FLOOR_KBPS.min(start_kbps.max(1)), probing: true, @@ -269,14 +338,17 @@ impl BitrateController { short_ack_kbps: 0, short_acks: 0, cap_probe_windows: 0, + cap_reprobe_after: CAP_REPROBE_WINDOWS_MIN, decode_cap_kbps: None, decode_backoff_kbps: 0, decode_cap_probe_windows: 0, + decode_cap_reprobe_after: CAP_REPROBE_WINDOWS_MIN, proven_kbps: 0, bad_windows: 0, clean_windows: 0, last_change: None, unacked: 0, + ceiling_ask_kbps: 0, } } @@ -321,8 +393,21 @@ impl BitrateController { self.short_acks = 1; } if self.short_acks >= 2 && self.host_cap_kbps.is_none_or(|c| kbps < c) { + // Re-learning a cap we had already lifted means the limit is STANDING, + // not the transient the re-probe exists to escape — back its clock off + // (see [`CAP_REPROBE_WINDOWS_MIN`]) so a hard encoder ceiling settles + // into a slow poll instead of two pointless acks every 12 s. A first + // latch starts the clock fast, because that is the case that matters. + self.cap_reprobe_after = if self.host_cap_kbps.is_some() { + self.cap_reprobe_after + .saturating_mul(2) + .min(CAP_REPROBE_WINDOWS_MAX) + } else { + CAP_REPROBE_WINDOWS_MIN + }; tracing::info!( cap_kbps = kbps, + reprobe_after_windows = self.cap_reprobe_after, "adaptive bitrate: host cap learned (encoder ceiling or cadence \ refusal) — climbs stop here until it lifts" ); @@ -343,14 +428,30 @@ impl BitrateController { /// decoder's knee is just as mode-scoped (pixel rate drives both ends of the codec), so /// the decode cap goes with it. The probe-measured `ceiling_kbps` (a LINK property) /// survives. + /// + /// Every rolling BASELINE is mode-scoped too, and for the same reason the encode one always + /// was: a mode switch changes what "normal" costs at both ends of the pipe. 4K120 decodes + /// and encodes far slower than 1080p60 and puts bigger frames on the wire, so a baseline + /// learned under the old mode is a floor the new one clears on its very first window — + /// [`DECODE_RISE_US`] is 15 µs-thousands, well inside the gap between those two modes. Left + /// standing (only `encode_means` used to be cleared here), the ~30 s it takes + /// [`BASELINE_WINDOWS`] to age out is ~30 s of every window scoring bad, which is a ×0.7 + /// backoff every other window: a switch UP in mode cratered the rate instead of raising it. + /// `proven_kbps` goes with them — it is the mark climbs are bounded against, and throughput + /// the OLD mode's decoder digested is not evidence about this one. It re-earns itself from + /// the next window. pub(crate) fn on_mode_switch(&mut self) { self.host_cap_kbps = None; self.short_acks = 0; self.cap_probe_windows = 0; + self.cap_reprobe_after = CAP_REPROBE_WINDOWS_MIN; self.decode_cap_kbps = None; self.decode_backoff_kbps = 0; self.decode_cap_probe_windows = 0; + self.owd_means.clear(); + self.decode_means.clear(); self.encode_means.clear(); + self.proven_kbps = 0; } /// Feed one report window; returns the rate to request now, if any. `dropped` = frames that @@ -389,22 +490,9 @@ impl BitrateController { return None; } // OWD: compare against the rolling-min baseline of PRIOR windows (so a rising window - // doesn't drag its own baseline up), then record it. - let owd_bad = match owd_mean_us { - Some(mean) => { - let bad = self - .owd_means - .iter() - .min() - .is_some_and(|&base| mean > base + OWD_RISE_US); - if self.owd_means.len() == BASELINE_WINDOWS { - self.owd_means.pop_front(); - } - self.owd_means.push_back(mean); - bad - } - None => false, - }; + // doesn't drag its own baseline up), then record it. No severe tier — a standing queue is + // congestion evidence, not visible damage, so it always takes the two-window path. + let (owd_bad, _) = score_baseline(&mut self.owd_means, owd_mean_us, OWD_RISE_US, i64::MAX); // Decode-stage latency: same rolling-min-baseline treatment as OWD, but measuring the // CLIENT'S decoder rather than the link. A rise means the decoder is backlogging frames — // the bottleneck the network signals are blind to. Marking the window bad both ends slow @@ -412,43 +500,22 @@ impl BitrateController { // the link ceiling) and, sustained, drives the ×0.7 backoff down to the real decode limit. // An excursion far past baseline is SEVERE: the decoder is deep in spike-overload and the // user is watching it — skip the two-window confirmation. - let (decode_bad, decode_severe) = match decode_mean_us { - Some(mean) => { - let base = self.decode_means.iter().min().copied(); - let bad = base.is_some_and(|b| mean > b + DECODE_RISE_US); - let severe = base.is_some_and(|b| mean > b + DECODE_SEVERE_US); - if self.decode_means.len() == BASELINE_WINDOWS { - self.decode_means.pop_front(); - } - self.decode_means.push_back(mean); - (bad, severe) - } - None => (false, false), - }; + let (decode_bad, decode_severe) = score_baseline( + &mut self.decode_means, + decode_mean_us, + DECODE_RISE_US, + DECODE_SEVERE_US, + ); // 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 - // is the bad/severe machinery's business. - if !decode_bad && actual_kbps > self.proven_kbps { - self.proven_kbps = actual_kbps; - } + let (encode_bad, encode_severe) = score_baseline( + &mut self.encode_means, + encode_mean_us, + ENCODE_RISE_US, + ENCODE_SEVERE_US, + ); // SEVERE = the user already saw damage (an unrecoverable frame, a jump-to-live flush, a // deep decode-latency excursion, a window spent begging for keyframes) or loss far past // any blip — one window is enough. Ordinary congestion (heavy-but-recoverable loss, an @@ -466,6 +533,17 @@ impl BitrateController { || decode_bad || encode_bad || recovery_kf >= RECOVERY_KF_BAD; + // The proven-throughput high-water mark: this window's delivered rate is now demonstrably + // digestible — the pipeline carried it and NOTHING went wrong while it did. Scored after + // the verdict and gated on the whole of it, not on decode alone: the mark never decays, so + // one window is permanent authority over how far every later climb may step, and the + // windows that overstate delivered throughput are exactly the damaged ones (a stall's + // backlog draining in a single window, a flush's queue, the FEC surge that answers a loss + // burst). "Loss doesn't disqualify, the bytes still arrived" was true about the bytes and + // wrong about the conclusion drawn from them. + if !bad && actual_kbps > self.proven_kbps { + self.proven_kbps = actual_kbps; + } if bad { self.bad_windows += 1; self.clean_windows = 0; @@ -475,16 +553,16 @@ 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. + // The learned host cap re-probe (see [`CAP_REPROBE_WINDOWS_MIN`]): after a clean run + // 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, and backs the clock off. 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 { + if self.cap_probe_windows >= self.cap_reprobe_after { self.cap_probe_windows = 0; let lifted = cap.saturating_add(cap / 8).min(self.ceiling_kbps); if lifted > cap { @@ -508,7 +586,7 @@ impl BitrateController { self.decode_cap_probe_windows = 0; } else if self.current_kbps >= cap.saturating_sub(cap / 16) { self.decode_cap_probe_windows += 1; - if self.decode_cap_probe_windows >= CAP_REPROBE_WINDOWS { + if self.decode_cap_probe_windows >= self.decode_cap_reprobe_after { self.decode_cap_probe_windows = 0; let lifted = cap.saturating_add(cap / 8).min(self.ceiling_kbps); if lifted > cap { @@ -538,18 +616,42 @@ impl BitrateController { // knee. One event never latches (a spurious flush must stay a one-off), and a // backoff without decode evidence in between breaks the streak — whatever it saw, // it wasn't the same knee. - if decode_severe || flushed { + // A bare flush counts as decode evidence only where the decode signal can't speak + // for itself. On an embedder that reports decode latency, a flush with FLAT decode + // is a network event (a stall, a clock step) that drained a queue the decoder was + // keeping up with — teaching a "decoder knee" from it caps the session on the wrong + // end of the pipe. Where the signal is absent the old reading stands: the flush is + // the only decoder-saturation evidence there is. + let decode_evidence = + decode_severe || (flushed && (decode_bad || decode_mean_us.is_none())); + if decode_evidence { let rate = self.current_kbps; let similar = self.decode_backoff_kbps > 0 && rate.abs_diff(self.decode_backoff_kbps) <= self.decode_backoff_kbps / DECODE_CAP_SIMILAR_DIV; - if similar && self.decode_cap_kbps.is_none_or(|c| rate < c) { + // Latch just UNDER the rate that choked, not at it: the knee is the rate the + // decoder could not hold, so a cap sitting exactly on it authorizes climbing + // straight back into the failure — the sawtooth the cap exists to end, merely + // slower. One sixteenth is inside the ±1/8 band the pair had to agree within, + // so it costs nothing the evidence actually established. + let knee = rate.saturating_sub(rate / 16).max(self.floor_kbps); + if similar && self.decode_cap_kbps.is_none_or(|c| knee < c) { + // Same standing-vs-transient backoff as the host cap. + self.decode_cap_reprobe_after = if self.decode_cap_kbps.is_some() { + self.decode_cap_reprobe_after + .saturating_mul(2) + .min(CAP_REPROBE_WINDOWS_MAX) + } else { + CAP_REPROBE_WINDOWS_MIN + }; tracing::info!( - cap_kbps = rate, + cap_kbps = knee, + choked_at_kbps = rate, + reprobe_after_windows = self.decode_cap_reprobe_after, "adaptive bitrate: decode cap learned (decoder knee) — climbs stop \ here until it lifts" ); - self.decode_cap_kbps = Some(rate.max(self.floor_kbps)); + self.decode_cap_kbps = Some(knee); self.decode_cap_probe_windows = 0; } self.decode_backoff_kbps = rate; @@ -574,6 +676,23 @@ impl BitrateController { .ceiling_kbps .min(self.host_cap_kbps.unwrap_or(u32::MAX)) .min(self.decode_cap_kbps.unwrap_or(u32::MAX)); + // Above the ceiling with nothing wrong: the session negotiated a rate the operator's + // `PUNKTFUNK_ABR_MAX_MBPS` forbids (no congestion signal will ever find this — the link + // is fine, the cap is a policy). Step straight to it rather than sitting above a limit + // the user set, and never below the floor. Asked ONCE per distinct target: if the host + // answers with something higher it has told us it cannot go there (its own floor, an + // encoder minimum), and repeating the ask every cooldown would buy nothing but a + // reconfigure each time. + let ceiling_target = eff_ceiling.max(self.floor_kbps); + if self.current_kbps > ceiling_target && self.ceiling_ask_kbps != ceiling_target { + tracing::info!( + from_kbps = self.current_kbps, + to_kbps = ceiling_target, + "adaptive bitrate: session rate is above the configured ceiling — stepping down" + ); + self.ceiling_ask_kbps = ceiling_target; + return self.request(ceiling_target, now); + } let cap = eff_ceiling .min(self.proven_kbps.saturating_mul(PROVEN_HEADROOM_NUM) / PROVEN_HEADROOM_DEN); if self.current_kbps < eff_ceiling && utilized && cap > self.current_kbps { @@ -602,6 +721,17 @@ impl BitrateController { // request just recomputes from the same base next time (and counts toward MAX_UNACKED). Some(kbps) } + + /// The decision [`on_window`](Self::on_window) returned never reached the wire (the control + /// queue was full). Undo the request's bookkeeping: [`MAX_UNACKED`] exists to detect a HOST + /// that doesn't answer, and counting a message we never sent toward it retires the + /// controller for the session — with a log line blaming an "older host" that is not what + /// happened. Clearing the pending request also keeps a later unsolicited ack from being + /// judged short against a rate we never asked for. + pub(crate) fn on_request_dropped(&mut self) { + self.unacked = self.unacked.saturating_sub(1); + self.last_requested_kbps = None; + } } #[cfg(test)] @@ -1110,14 +1240,16 @@ mod tests { #[test] fn decode_latency_caps_the_slow_start_climb() { - // A fat link (probe measured ~300 Mbps) but a decoder that saturates around the start rate. + // A fat link (probe measured ~300 Mbps) but a decoder that saturates below it. let mut c = BitrateController::new(20_000); c.set_ceiling(300_000); let start = Instant::now(); - // First clean window (decoder fine at 20 Mbps) → slow start doubles to 40. - assert_eq!( - c.on_window( - ticks(start, 0), + // Slow start doubles while the decoder keeps up, and the first BASELINE_MIN_WINDOWS of + // those windows are what teach the decode baseline (one sample is not a floor). + let mut last = 0; + for i in 0..BASELINE_MIN_WINDOWS as u32 { + if let Some(k) = c.on_window( + ticks(start, i * 2), 0, 0, Some(10_000), @@ -1125,16 +1257,18 @@ mod tests { None, 1_000_000, false, - 0 - ), - Some(40_000) - ); - c.on_ack(40_000); - // At 40 Mbps the decoder starts backing up (30 ms over baseline): the window is bad, so the - // climb stops here instead of doubling on toward the 300 Mbps link ceiling… + 0, + ) { + last = k; + c.on_ack(k); + } + } + assert_eq!(last, 300_000, "slow start should reach the probed ceiling"); + // Now the decoder starts backing up (30 ms over the learned baseline): the window is bad, + // so the climb stops instead of parking at the link ceiling… assert_eq!( c.on_window( - ticks(start, 2), + ticks(start, 20), 0, 0, Some(10_000), @@ -1146,11 +1280,11 @@ mod tests { ), None ); - // …and a second backed-up window backs the rate off, settling at the decode limit rather + // …and a second backed-up window backs the rate off toward the real decode limit rather // than choking the decoder at the link ceiling (the reported bug). assert_eq!( c.on_window( - ticks(start, 4), + ticks(start, 22), 0, 0, Some(10_000), @@ -1160,7 +1294,54 @@ mod tests { false, 0 ), - Some(28_000) + Some(210_000) + ); + } + + #[test] + fn one_calm_window_is_not_a_baseline() { + // The ratchet this guard exists to stop: our own decrease CLEARS the encode baseline, so + // it re-seeds from whatever the next window happens to be. If that window is calm, the + // ordinary content variance that follows reads as a rise, backs off, clears again — all + // the way to the floor on a link that was never the problem. A single sample must not + // arm the signal. + let mut c = BitrateController::new(100_000); + let start = Instant::now(); + // One calm 3 ms encode window, then windows 9 ms above it: far past ENCODE_RISE_US, and + // sustained — yet no baseline exists to judge them against yet. + for i in 0..BASELINE_MIN_WINDOWS as u32 { + let mean = if i == 0 { 3_000 } else { 12_000 }; + assert_eq!( + c.on_window( + ticks(start, i), + 0, + 0, + Some(10_000), + None, + Some(mean), + 1_000_000, + false, + 0 + ), + None, + "window {i} fired off a baseline of fewer than {BASELINE_MIN_WINDOWS} samples" + ); + } + // With a real baseline (min 3 ms over 4 windows) the signal works exactly as before: a + // sustained rise past it still backs the rate off. + assert_eq!( + c.on_window( + ticks(start, 8), + 0, + 0, + Some(10_000), + None, + Some(20_000), + 1_000_000, + false, + 0 + ), + Some(70_000) ); } @@ -1385,8 +1566,8 @@ mod tests { #[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 + // A cadence-refusal cap is scene evidence, not a spec limit: after a clean run parked 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); @@ -1396,7 +1577,10 @@ mod tests { 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 { + // The FIRST re-probe is the fast one — a transient refusal must not cost the session + // minutes to escape. + assert_eq!(c.cap_reprobe_after, CAP_REPROBE_WINDOWS_MIN); + for i in 0..CAP_REPROBE_WINDOWS_MIN { let _ = c.on_window( ticks(start, 20 + i), 0, @@ -1412,6 +1596,52 @@ mod tests { assert_eq!(c.host_cap_kbps, Some(794_000 + 794_000 / 8)); } + #[test] + fn a_standing_cap_backs_its_reprobe_clock_off() { + // The other half of the re-probe: an encoder's real codec ceiling (794 Mbps, L6.2) + // re-teaches itself every time the lift is tried. Escaping fast is right for a + // transient and pointless here, so each re-learn doubles the interval — a hard limit + // settles into a slow poll instead of two acks every 12 s for the whole session. + 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.cap_reprobe_after, CAP_REPROBE_WINDOWS_MIN); + // Each round: park clean at the cap until it re-probes upward, then have the host refuse + // the lift at the same value again. That is a STANDING limit, so the clock doubles. + let mut tick = 20; + for round in 0..3 { + let before = c.cap_reprobe_after; + for _ in 0..before { + let _ = c.on_window( + ticks(start, tick), + 0, + 0, + Some(10_000), + None, + None, + 1_000_000, + false, + 0, + ); + tick += 1; + } + let lifted = c.host_cap_kbps.expect("cap should still be latched"); + assert!(lifted > 794_000, "round {round}: the re-probe never lifted"); + // The host clamps the lift straight back to its real ceiling. + c.last_requested_kbps = Some(lifted); + c.on_ack(794_000); + assert_eq!(c.host_cap_kbps, Some(794_000)); + assert_eq!( + c.cap_reprobe_after, + (before * 2).min(CAP_REPROBE_WINDOWS_MAX) + ); + } + } + #[test] fn host_encode_latency_rise_backs_off() { // The compute knee: link pristine, client decoder fine — only HOST encode time moves @@ -1593,6 +1823,25 @@ mod tests { assert_eq!(run_clean(&mut c, start, 4, 20), None); } + #[test] + fn a_session_above_the_env_cap_steps_down_to_it_once() { + // PUNKTFUNK_ABR_MAX_MBPS is the only lever an Automatic session gives the operator, and + // it used to bind only ceilings the PROBE taught — so a session that negotiated a rate + // above the cap simply ran above it forever. No congestion signal will ever find that: + // the link is fine, the cap is policy. + let mut c = BitrateController::with_ceiling_cap(100_000, Some(50_000)); + assert_eq!(c.ceiling_kbps, 50_000); + let start = Instant::now(); + assert_eq!(run_clean(&mut c, start, 0, 1), Some(50_000)); + // Suppose the host answers HIGHER than asked (its own floor, an encoder minimum): that + // is the host saying it cannot go there. Don't re-ask every cooldown forever. + c.on_ack(80_000); + assert_eq!(run_clean(&mut c, start, 2, 20), None); + // A ceiling that MOVES is a new question, and gets asked once more. + c.set_ceiling(90_000); // clamped to the 50 Mbps cap → still 50 000, no new ask + assert_eq!(run_clean(&mut c, start, 24, 20), None); + } + #[test] fn decode_cap_latches_after_two_consecutive_decode_severe_backoffs() { // The 1440p120 field sawtooth: a decoder knee (~500 Mbps) well under the (inflated) @@ -1650,7 +1899,7 @@ mod tests { ), Some(350_000) ); - assert_eq!(c.decode_cap_kbps, Some(500_000)); + assert_eq!(c.decode_cap_kbps, Some(500_000 - 500_000 / 16)); // The backoff applies; from here every climb must stop AT the knee — not the 900 Mbps // link ceiling the old sawtooth kept re-poking. c.on_ack(350_000); @@ -1667,14 +1916,22 @@ mod tests { false, 0, ) { - assert!(k <= 500_000, "climb past the decode cap: {k}"); + // Never past the cap in force when the decision was made. (A long clean run + // legitimately re-probes that cap upward — `decode_cap_reprobes_after_a_ + // sustained_clean_run` owns that; here the point is that nothing climbs toward + // the 900 Mbps LINK ceiling the old sawtooth kept re-poking.) + assert!( + k <= c.decode_cap_kbps.unwrap(), + "climb past the decode cap: {k}" + ); max_req = max_req.max(k); c.on_ack(k); } } - assert_eq!(max_req, 500_000); - assert_eq!(c.current_kbps, 500_000); - assert_eq!(c.decode_cap_kbps, Some(500_000)); + assert!( + max_req < 600_000, + "the decode knee stopped binding: climbed to {max_req}" + ); } #[test] @@ -1748,10 +2005,10 @@ mod tests { 0, ); } - assert_eq!(c.decode_cap_kbps, Some(500_000)); + assert_eq!(c.decode_cap_kbps, Some(500_000 - 500_000 / 16)); // The host's ack parks the session at the knee (its clamp is authoritative). - c.on_ack(500_000); - for i in 0..CAP_REPROBE_WINDOWS { + c.on_ack(500_000 - 500_000 / 16); + for i in 0..CAP_REPROBE_WINDOWS_MIN { let _ = c.on_window( ticks(start, 8 + i), 0, @@ -1764,7 +2021,8 @@ mod tests { 0, ); } - assert_eq!(c.decode_cap_kbps, Some(500_000 + 500_000 / 8)); + let knee = 500_000 - 500_000 / 16; + assert_eq!(c.decode_cap_kbps, Some(knee + knee / 8)); } #[test] @@ -1800,7 +2058,7 @@ mod tests { 0, ); } - assert_eq!(c.decode_cap_kbps, Some(500_000)); + assert_eq!(c.decode_cap_kbps, Some(500_000 - 500_000 / 16)); c.on_mode_switch(); assert!(c.decode_cap_kbps.is_none()); assert_eq!(c.ceiling_kbps, 900_000); diff --git a/crates/punktfunk-core/src/client/pump/data.rs b/crates/punktfunk-core/src/client/pump/data.rs index a83339b8..0aad584f 100644 --- a/crates/punktfunk-core/src/client/pump/data.rs +++ b/crates/punktfunk-core/src/client/pump/data.rs @@ -492,7 +492,15 @@ impl DataPump { recovery_kf = recovery_kf_reqs, "adaptive bitrate: requesting encoder re-target" ); - let _ = ctrl_tx.try_send(CtrlRequest::SetBitrate(kbps)); + if ctrl_tx.try_send(CtrlRequest::SetBitrate(kbps)).is_err() { + // Never reached the control task — tell the controller, or three of + // these retire it for the session as "the host never acked". + abr.on_request_dropped(); + tracing::warn!( + kbps, + "adaptive bitrate: control queue full — re-target dropped" + ); + } } flush_in_window = false; last_report = Instant::now(); From e9a7373c76782a85a45cc73b472d8254560f659c Mon Sep 17 00:00:00 2001 From: enricobuehler Date: Mon, 3 Aug 2026 01:14:03 +0200 Subject: [PATCH 2/5] fix(client/abr): measure delivered throughput in media bytes, not wire bytes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The controller's two throughput-driven gates both compare "what the pipeline carried" against the ENCODER's target: the utilization gate asks whether a clean window actually tested that target (a calm menu proves nothing), and the never-decaying proven mark bounds how far every later climb may step. Both were fed `bytes_received`, which counts every accepted datagram — headers, FEC parity, probe filler, audio. So the figure rose with the redundancy the host adds in ANSWER to loss: at 25 % FEC the gate passed with the encoder emitting ~55 % of target, and the proven mark inherited the same inflation permanently. The signal was weakest exactly on the lossy links it exists for. Count data-shard payload separately at the reassembler's routing decision — the same place, and for the same reason, the probe counters are already stamped — and feed the ABR that. First time both gates are dimensionally honest: a media rate compared against a media target. --- crates/punktfunk-core/src/client/pump/data.rs | 29 ++++++++++++------- .../punktfunk-core/src/packet/reassemble.rs | 7 +++++ crates/punktfunk-core/src/stats.rs | 12 ++++++++ 3 files changed, 37 insertions(+), 11 deletions(-) diff --git a/crates/punktfunk-core/src/client/pump/data.rs b/crates/punktfunk-core/src/client/pump/data.rs index 0aad584f..95b2c4e7 100644 --- a/crates/punktfunk-core/src/client/pump/data.rs +++ b/crates/punktfunk-core/src/client/pump/data.rs @@ -232,7 +232,7 @@ impl DataPump { last_late = st.fec_late_shards; last_received = st.packets_received; last_dropped = st.frames_dropped; - last_bytes = st.bytes_received; + last_bytes = st.media_bytes_received; last_report = Instant::now(); discard_abr_window = true; flush_in_window = false; @@ -317,11 +317,12 @@ impl DataPump { "adaptive bitrate: capacity probe declined — keeping negotiated ceiling" ); } - // The probe's FLAG_PROBE filler landed in `bytes_received` but never reached - // the decoder — rebase the ABR window's byte counter past it, or the next - // window's "actual throughput" reads as the burst rate and poisons the - // controller's proven-throughput high-water mark with the LINK rate. - last_bytes = st.bytes_received; + // Rebase the ABR window's byte anchor past the burst. (Probe filler is + // routed out of `media_bytes_received` at the reassembler, so it can no + // longer read as the burst rate on its own — but the anchor still has to + // skip the video that landed around the burst under a suppressed report + // tick, which would otherwise divide a long span's bytes by one window.) + last_bytes = st.media_bytes_received; } else if Instant::now() >= deadline { // The host never answered (a build that ignores ProbeRequest): clear the // stuck-active state so LossReports resume, keep the negotiated ceiling. @@ -454,11 +455,17 @@ impl DataPump { // the next one. let recovery_kf_reqs = pump_recovery_kf.swap(0, Ordering::Relaxed); // 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 - // semantics (both compare against targets with their own headroom). + // the target it was allowed. MEDIA bytes (data-shard payload: no headers, no FEC + // parity, no probe filler, no audio), because both consumers compare it against + // the ENCODER's target: the utilization gate asks "was the target genuinely + // tested?" and the proven mark bounds every later climb. Wire bytes answered a + // different question — they rise with the redundancy the host adds in answer to + // loss, so the gate read ~25 % high precisely on the links it exists for. let window_ms = last_report.elapsed().as_millis().max(1) as u64; - let actual_kbps = (st.bytes_received.wrapping_sub(last_bytes).saturating_mul(8) + let actual_kbps = (st + .media_bytes_received + .wrapping_sub(last_bytes) + .saturating_mul(8) / window_ms) as u32; // A discard window feeds the controller NOTHING — its signals are probe-tail // residue, and one "congestion" verdict here ends slow start for good. @@ -508,7 +515,7 @@ impl DataPump { last_late = st.fec_late_shards; last_received = st.packets_received; last_dropped = st.frames_dropped; - last_bytes = st.bytes_received; + last_bytes = st.media_bytes_received; if pump_perf_on { if let Some(p) = session.take_pump_perf() { let per_pkt_ns = |ns: u64| ns.checked_div(p.packets).unwrap_or(0); diff --git a/crates/punktfunk-core/src/packet/reassemble.rs b/crates/punktfunk-core/src/packet/reassemble.rs index 1bc7e6a2..5e9d595a 100644 --- a/crates/punktfunk-core/src/packet/reassemble.rs +++ b/crates/punktfunk-core/src/packet/reassemble.rs @@ -429,6 +429,13 @@ impl Reassembler { stats .probe_last_arrival_ns .store(now_ns, std::sync::atomic::Ordering::Relaxed); + } else if hdr.shard_index < hdr.data_shards { + // Media accounting (see `Stats::media_bytes_received`): DATA shards only, payload + // only. Stamped at the same routing decision as the probe counters and for the same + // reason — the adaptive-bitrate utilization gate compares delivered throughput + // against an ENCODER target, so parity, headers and probe filler have no business + // in the numerator. + StatsCounters::add(&stats.media_bytes_received, shard_bytes as u64); } let win = if is_probe { probe } else { video }; win.advance_window( diff --git a/crates/punktfunk-core/src/stats.rs b/crates/punktfunk-core/src/stats.rs index d293afd5..bfd4f026 100644 --- a/crates/punktfunk-core/src/stats.rs +++ b/crates/punktfunk-core/src/stats.rs @@ -45,6 +45,16 @@ pub struct Stats { /// so a speed-test numerator built from it inherits whatever video was in flight around /// the burst — these keep video out of the probe math. Deliberately NOT mirrored into the /// C-ABI `PunktfunkStats` (probe measurements surface via `ProbeOutcome`). + /// Media bytes delivered to the video reassembler: DATA-shard payload only — no packet + /// headers, no FEC parity, no probe filler, no audio. This is the rate the encoder's target + /// is a promise about, and the only honest thing to compare that target against. + /// `bytes_received` counts every accepted datagram, so a "delivered throughput" built from + /// it rises with the FEC redundancy the host adds in answer to loss — which meant the + /// adaptive-bitrate utilization gate ("did the pipeline actually carry ~the target?") read + /// 25 % high exactly on the lossy links it exists for, and the never-decaying + /// proven-throughput mark inherited the same inflation. Deliberately NOT mirrored into the + /// C-ABI `PunktfunkStats`. + pub media_bytes_received: u64, pub probe_packets_received: u64, pub probe_bytes_received: u64, /// First / last probe-packet arrival (monotonic ns, see [`now_monotonic_ns`]; 0 = none @@ -75,6 +85,7 @@ pub struct StatsCounters { pub fec_late_shards: AtomicU64, pub bytes_sent: AtomicU64, pub bytes_received: AtomicU64, + pub media_bytes_received: AtomicU64, pub probe_packets_received: AtomicU64, pub probe_bytes_received: AtomicU64, pub probe_first_arrival_ns: AtomicU64, @@ -101,6 +112,7 @@ impl StatsCounters { fec_late_shards: self.fec_late_shards.load(l), bytes_sent: self.bytes_sent.load(l), bytes_received: self.bytes_received.load(l), + media_bytes_received: self.media_bytes_received.load(l), probe_packets_received: self.probe_packets_received.load(l), probe_bytes_received: self.probe_bytes_received.load(l), probe_first_arrival_ns: self.probe_first_arrival_ns.load(l), From 48565c4e9ed60a2f23080704d5f20b762c89d7e5 Mon Sep 17 00:00:00 2001 From: enricobuehler Date: Mon, 3 Aug 2026 01:14:17 +0200 Subject: [PATCH 3/5] fix(host/abr): stop pinning Automatic sessions, and tell the client when the rate moves MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two host-side halves of the same sweep. **The cadence latch.** `cadence_degraded` — which makes the control task refuse bitrate CLIMBS — was latched true for as long as the session was escalated (adaptive capture depth or pipelined retrieve), independently of whether encode was still missing deadlines. The client cannot tell that refusal apart from an encoder's real ceiling: both arrive as a short `BitrateChanged`, and two identical ones latch a cap. Escalation needs ~20 net behind-frames, which a startup hitch supplies while the ABR is still in slow start at the 20 Mbps default — so one transient pinned the whole session there, long after the escalation had bought back the headroom it was for, and escaping cost +12.5 % per 60 s. An escalated session is still judged strictly (ANY net behind-frame keeps it flagged, where an unescalated one gets the full bucket), but being escalated no longer flags it by itself: escalating exists so cadence CAN be held, and once it is, refusing climbs refuses the thing that worked. The rule moves into `encode_behind_cadence` so it is stateable and testable. **The silent re-target.** `adopt_built_bitrate` publishes the rate a rebuilt pipeline actually opened at — `build_pipeline` re-resolves an Automatic rate whenever the source delivers a size the session did not negotiate, the mirrored-panel case — and the encoder's own clamp can land below what the control task already acked. Neither reached the client, whose controller keeps its own copy of that number as its climb base. A 1080p client mirroring a 4K panel therefore believed 20 Mbps while the host encoded 60, and its first climb computed from the stale base asked for 40: a re-target DOWNWARD, paying an encoder rebuild to get there. Both paths now push the applied rate to the control task, which sends `BitrateChanged` — the existing 9-byte message, which already means precisely this and which clients already handle arriving unprompted. No wire-format change, no capability negotiation, old clients unaffected. 2 host tests added. --- crates/punktfunk-host/src/native.rs | 13 +++ crates/punktfunk-host/src/native/control.rs | 24 ++++ crates/punktfunk-host/src/native/stream.rs | 116 ++++++++++++++++++-- 3 files changed, 145 insertions(+), 8 deletions(-) diff --git a/crates/punktfunk-host/src/native.rs b/crates/punktfunk-host/src/native.rs index 7dad3025..f2cafd6a 100644 --- a/crates/punktfunk-host/src/native.rs +++ b/crates/punktfunk-host/src/native.rs @@ -1073,6 +1073,17 @@ async fn serve_session( // accepted ack as "the active mode is now X" and fixes itself; old clients just log it. let (reconfig_result_tx, reconfig_result_rx) = tokio::sync::mpsc::unbounded_channel::(); + // Unsolicited bitrate re-target, data plane → control task (the `reconfig_result_tx` pattern + // again, for the same reason). A pipeline rebuild can RE-RESOLVE an Automatic rate — most + // visibly when the source delivers a different size than the session negotiated, e.g. a + // client that asked for 1080p mirroring a 4K panel — and that number is what everything + // downstream reasons about: the send pacer, the console, and the base a `SetBitrate` ack is + // measured against. The client's copy only ever moved on an ack, so it stayed on the + // negotiated rate while the host encoded at another one, and the ABR's first climb computed + // from that stale base asked for LESS than the host was already sending — a re-target + // downward, with the rebuild it costs. Tell the client instead; `BitrateChanged` already + // means exactly this and old clients already handle one arriving unprompted. + let (retarget_tx, retarget_rx) = tokio::sync::mpsc::unbounded_channel::(); // Cursor-forward bridge (M2): the encode loop diffs each frame's cursor serial and hands // changed SHAPES here; the control task (the control stream's sole writer) sends them. // Same shape as `probe_result_tx`. Wired even when the channel wasn't negotiated — it @@ -1133,6 +1144,7 @@ async fn serve_session( probe_tx, probe_result_rx, reconfig_result_rx, + retarget_rx, cursor_shape_rx, cursor_client_draws, clip_enabled, @@ -1579,6 +1591,7 @@ async fn serve_session( probe_rx, probe_result_tx, reconfig_result_tx, + retarget_tx, fec_target: fec_target_dp, phase: phase_ctl, conn: conn_stream, diff --git a/crates/punktfunk-host/src/native/control.rs b/crates/punktfunk-host/src/native/control.rs index e8b89675..e14934c6 100644 --- a/crates/punktfunk-host/src/native/control.rs +++ b/crates/punktfunk-host/src/native/control.rs @@ -40,6 +40,9 @@ pub(super) async fn run( probe_tx: std::sync::mpsc::Sender, mut probe_result_rx: tokio::sync::mpsc::UnboundedReceiver, mut reconfig_result_rx: tokio::sync::mpsc::UnboundedReceiver, + // Host-initiated bitrate re-target (a rebuild re-resolved an Automatic rate): forwarded to + // the client as a `BitrateChanged` so its controller's climb base tracks the real encoder. + mut retarget_rx: tokio::sync::mpsc::UnboundedReceiver, mut cursor_shape_rx: tokio::sync::mpsc::UnboundedReceiver, cursor_client_draws: Arc, clip_enabled: Arc, @@ -338,6 +341,27 @@ pub(super) async fn run( None => clip_offer_closed = true, } } + retarget = retarget_rx.recv() => { + // A pipeline rebuild re-resolved the Automatic rate (see `retarget_tx`). Same + // message the `SetBitrate` path answers with — the client's controller treats + // any `BitrateChanged` as authoritative for what the encoder now targets, which + // is exactly right here: it IS what the encoder now targets, we just weren't + // asked. PyroWave reaches this too, and should: its rate is pinned against + // mid-stream RETARGETS, but a mode switch legitimately re-resolves the pin + // (~1.6 bpp for the new pixel rate) and the client's live-rate display is + // otherwise stuck on the old one. Its controller is off, so nothing acts on it. + let Some(kbps) = retarget else { break }; // data plane gone + tracing::info!( + kbps, + "encoder re-targeted by a pipeline rebuild — telling the client" + ); + if io::write_msg(&mut ctrl_send, &BitrateChanged { bitrate_kbps: kbps }.encode()) + .await + .is_err() + { + break; + } + } correction = reconfig_result_rx.recv() => { // H2 rollback/correction ack: the data plane reports the mode ACTUALLY live // after a rebuild that failed (stayed at the old mode) or that the backend diff --git a/crates/punktfunk-host/src/native/stream.rs b/crates/punktfunk-host/src/native/stream.rs index 0c263638..6bba243e 100644 --- a/crates/punktfunk-host/src/native/stream.rs +++ b/crates/punktfunk-host/src/native/stream.rs @@ -1214,6 +1214,9 @@ pub(super) struct SessionContext { /// `Reconfigured { accepted: true, mode: }` when a rebuild failed (stayed at /// the old mode) or the backend honored a different refresh than requested. pub(super) reconfig_result_tx: tokio::sync::mpsc::UnboundedSender, + /// Host-initiated bitrate re-target → control task → the client's `BitrateChanged`. Fired + /// by [`adopt_built_bitrate`] when a rebuild lands on a rate the client wasn't told about. + pub(super) retarget_tx: tokio::sync::mpsc::UnboundedSender, /// Adaptive-FEC target the control task updates from the client's loss reports. pub(super) fec_target: Arc, /// The QUIC control connection (carries host→client 0xCE source-HDR metadata mid-stream). @@ -1397,6 +1400,7 @@ pub(super) fn virtual_stream(ctx: SessionContext, prepared: Option 1 || pipelined_active || deescalating; // Export "encode can't hold cadence" for the control task's climb refusal. - // An escalated session stays flagged even with the bucket drained: its climb - // headroom is spent, and letting climbs resume would saw against the + // An escalated session is held to a stricter standard — ANY net behind-frame + // keeps it flagged, where an unescalated one is given the full bucket — because + // its climb headroom really is partly spent and a climb would saw against the // escalation and starve the de-escalation clean run below. + // + // But being escalated cannot flag it BY ITSELF, which is what this used to do. + // The client can't tell a transient refusal from an encoder's real ceiling: two + // identical short acks latch a cap, so a session that escalated once — the + // bucket needs ~20 net misses, which a startup hitch supplies while the ABR is + // still in slow start at the 20 Mbps default — got pinned there, and stayed + // pinned long after the escalation had bought back the headroom it was for. + // Escalating exists precisely so cadence CAN be held; once it is (bucket + // drained, every frame on time), refusing climbs is refusing the thing that + // worked. cadence_degraded.store( - escalated || behind_score >= DEPTH_DEGRADE, + encode_behind_cadence(escalated, behind_score, DEPTH_DEGRADE), Ordering::Relaxed, ); if deescalating { @@ -4068,13 +4106,42 @@ impl PaceBudget { } } +/// Does the encoder currently fail to hold the frame cadence? Exported to the control task, which +/// refuses bitrate CLIMBS while it is true (descents always pass — they are the cure). +/// +/// `escalated` = the session has already spent an adaptive-depth / pipelined-retrieve step to buy +/// headroom; `behind_score` is the leaky bucket of frames whose work overran the cadence deadline. +/// An escalated session is judged strictly — ANY net behind-frame keeps it flagged — but being +/// escalated does not flag it on its own. That distinction is the whole point: the client cannot +/// tell a transient refusal from an encoder's hard ceiling (two identical short acks latch a cap), +/// so "escalated ⇒ degraded, permanently" pinned Automatic sessions at whatever rate they happened +/// to hold when a startup hitch escalated them — routinely the 20 Mbps default, while slow start +/// had barely begun. Escalation exists so cadence CAN be held; once it is, refusing climbs refuses +/// the thing that worked. +fn encode_behind_cadence(escalated: bool, behind_score: u32, degrade_at: u32) -> bool { + behind_score >= degrade_at || (escalated && behind_score > 0) +} + /// Adopt the rate a freshly built pipeline's encoder was actually opened at. /// /// The session's own `bitrate_kbps` is the number every later decision reads — the ABR controller's /// climb base, the console's sample, what a `SetBitrate` ack is measured against — so letting it /// disagree with the live encoder means each of those reasons about a stream that doesn't exist. /// Silent when nothing changed, which is the overwhelmingly common case. -fn adopt_built_bitrate(current: &mut u32, built: u32, live: &Arc) { +/// +/// The client keeps its OWN copy of that number, and it used to move only on an ack — so a +/// rebuild that re-resolved an Automatic rate (`build_pipeline` does, whenever the source +/// delivers a size the session did not negotiate) left the two disagreeing for the rest of the +/// session. The ABR's next climb then computed from the stale base and asked for a rate BELOW +/// what the host was already sending: a re-target downward, paying an encoder rebuild to get +/// there. So tell the client too — `BitrateChanged` is the same message the `SetBitrate` path +/// answers with, and means the same thing arriving unprompted. +fn adopt_built_bitrate( + current: &mut u32, + built: u32, + live: &Arc, + retarget: &tokio::sync::mpsc::UnboundedSender, +) { if built == *current { return; } @@ -4085,6 +4152,7 @@ fn adopt_built_bitrate(current: &mut u32, built: u32, live: &Arc) { ); *current = built; live.store(built, Ordering::Relaxed); + let _ = retarget.send(built); // control task gone ⇒ the session is ending anyway } /// Encode-stall recovery: rebuild the encoder in place (keeping capture + the session up) and @@ -4329,6 +4397,38 @@ fn build_pipeline( mod tests { use super::*; + #[test] + fn an_escalated_but_caught_up_encoder_stops_refusing_climbs() { + const DEGRADE: u32 = 10; + // Not escalated: the full bucket is allowed before climbs are refused. + assert!(!encode_behind_cadence(false, 0, DEGRADE)); + assert!(!encode_behind_cadence(false, 9, DEGRADE)); + assert!(encode_behind_cadence(false, 10, DEGRADE)); + // Escalated and still missing deadlines: strict — one net behind-frame is enough. + assert!(encode_behind_cadence(true, 1, DEGRADE)); + // Escalated, bucket fully drained: cadence is being HELD, which is what escalating was + // for. This is the case that used to stay latched for the rest of the session and pin an + // Automatic client at its slow-start rate. + assert!(!encode_behind_cadence(true, 0, DEGRADE)); + } + + #[test] + fn adopting_a_rebuilt_rate_tells_the_client() { + let live = Arc::new(AtomicU32::new(20_000)); + let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::(); + let mut current = 20_000; + // The overwhelmingly common case: the rebuild landed on the same rate — silent. + adopt_built_bitrate(&mut current, 20_000, &live, &tx); + assert_eq!(rx.try_recv().ok(), None); + // A re-resolve (the client asked 1080p, the source delivers a mirrored 4K panel): the + // host's rate moves, so the client has to hear about it — its controller's climb base is + // its own copy of this number, and a stale one makes the next "climb" a cut. + adopt_built_bitrate(&mut current, 60_000, &live, &tx); + assert_eq!(current, 60_000); + assert_eq!(live.load(Ordering::Relaxed), 60_000); + assert_eq!(rx.try_recv().ok(), Some(60_000)); + } + #[test] fn pacing_never_exceeds_the_session_rate_or_the_display() { // Backend honored the request exactly (the multiplier off): pace at it. From 33ecd8e1a58558e6173de3570eb7482a555e1112 Mon Sep 17 00:00:00 2001 From: enricobuehler Date: Mon, 3 Aug 2026 01:16:35 +0200 Subject: [PATCH 4/5] fix(client/abr): a granted climb disproves the learned cap MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Completing the cap-escape fix. Backing the re-probe clock off to 12 s got the client asking again quickly, but each ask only LIFTED the cap by +12.5 % — so even a host that had fully recovered still granted the session its real ceiling one small step at a time, ~4 minutes from the 20 Mbps default to a 300 Mbps link. The crawl was never the point; re-learning was. A request granted IN FULL at or above the cap is the host's own word that the limit is gone. Drop the cap outright at that point instead of nudging it. A standing limit is unaffected — it answers the same re-probe with another short ack, which re-latches it and doubles its clock, exactly as before. Adds the end-to-end regression the sweep was really about: a session pinned at 20 Mbps by a transient cadence refusal, under a probe-measured 300 Mbps ceiling, now reaches 150 Mbps in 22 windows (~16 s) where it used to need ~17 minutes. --- crates/punktfunk-core/src/abr.rs | 71 ++++++++++++++++++++++++++++++++ 1 file changed, 71 insertions(+) diff --git a/crates/punktfunk-core/src/abr.rs b/crates/punktfunk-core/src/abr.rs index 623e1633..913e179b 100644 --- a/crates/punktfunk-core/src/abr.rs +++ b/crates/punktfunk-core/src/abr.rs @@ -416,6 +416,21 @@ impl BitrateController { } } else { self.short_acks = 0; + // GRANTED in full at or above the learned cap: the limit that taught it is + // gone, and we have the host's own word for it. Drop the cap outright rather + // than keep crawling up in +12.5 % re-probe steps — for a cap latched from a + // transient (a host briefly behind cadence) that crawl is the entire + // remaining cost of the transient, and it is measured in minutes. + if self.host_cap_kbps.is_some_and(|c| kbps >= c) { + tracing::info!( + granted_kbps = kbps, + "adaptive bitrate: host granted a climb at the learned cap — the \ + limit has lifted, dropping it" + ); + self.host_cap_kbps = None; + self.cap_probe_windows = 0; + self.cap_reprobe_after = CAP_REPROBE_WINDOWS_MIN; + } } } self.current_kbps = kbps; @@ -1596,6 +1611,62 @@ mod tests { assert_eq!(c.host_cap_kbps, Some(794_000 + 794_000 / 8)); } + #[test] + fn a_transient_refusal_does_not_pin_the_session() { + // The field failure this whole cap-escape change exists for. A host that escalates its + // capture/encode pipeline once — a startup hitch is enough — used to refuse every climb + // for the rest of the session; the client latched that refusal as a cap, at whatever + // rate slow start had reached, which is routinely the 20 Mbps default. Escaping cost + // +12.5 % per ~60 s: north of twenty minutes to reach a 300 Mbps link ceiling, which the + // user experiences as "Automatic is broken". + let mut c = BitrateController::new(20_000); + c.set_ceiling(300_000); // the startup probe measured a fat link + let start = Instant::now(); + let mut tick = 0u32; + let mut windows_pinned = 0u32; + // Two refused climbs at the same rate → the cap latches at 20 Mbps. + for _ in 0..2 { + let k = run_clean(&mut c, start, tick, 4).expect("slow start should ask to climb"); + tick += 4; + assert!(k > 20_000); + c.on_ack(20_000); // "behind cadence — held at the current rate" + } + assert_eq!(c.host_cap_kbps, Some(20_000)); + // The host recovers immediately (its bucket drains; the escalation bought the headroom + // it was for), but the client has no way to know that except by asking again. Drive + // clean windows and grant whatever it asks for. + while c.current_kbps < 150_000 && windows_pinned < 400 { + if let Some(k) = c.on_window( + ticks(start, tick), + 0, + 0, + Some(10_000), + None, + None, + 1_000_000, + false, + 0, + ) { + c.on_ack(k); + } + tick += 1; + windows_pinned += 1; + } + assert!( + c.current_kbps >= 150_000, + "still pinned at {} after {windows_pinned} windows", + c.current_kbps + ); + // ~750 ms a window: this must be tens of seconds, not the old tens of minutes. + assert!( + windows_pinned <= 40, + "took {windows_pinned} windows (~{} s) to escape a transient refusal", + windows_pinned * 3 / 4 + ); + // And the disproven cap is gone, not merely nudged upward. + assert!(c.host_cap_kbps.is_none()); + } + #[test] fn a_standing_cap_backs_its_reprobe_clock_off() { // The other half of the re-probe: an encoder's real codec ceiling (794 Mbps, L6.2) From 1ae8b4d4caaf2870b2ec18983622fba3f27a289b Mon Sep 17 00:00:00 2001 From: enricobuehler Date: Mon, 3 Aug 2026 01:18:40 +0200 Subject: [PATCH 5/5] fix(client/abr): let the ceiling follow a host-initiated re-target MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Interaction between two fixes in this series. The host now tells the client when a rebuild re-resolves an Automatic rate, and that rate can legitimately sit ABOVE the client's climb ceiling — the ceiling is the negotiated start rate until the capacity probe raises it, while the host's re-resolve answers "what do these pixels actually need" (a 1080p session mirroring a 4K panel resolves ~3× higher). Left alone, the client would learn the new rate, notice it was above a stale ceiling, and step the host straight back down off the rate it had just chosen for itself. So an ack raises the ceiling to meet it. `set_ceiling` only ever raises and still clamps to PUNKTFUNK_ABR_MAX_MBPS, which is the one limit that should bind here. No effect on ordinary acks: a climb is never requested above the effective ceiling to begin with. --- crates/punktfunk-core/src/abr.rs | 32 ++++++++++++++++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/crates/punktfunk-core/src/abr.rs b/crates/punktfunk-core/src/abr.rs index 913e179b..d0637334 100644 --- a/crates/punktfunk-core/src/abr.rs +++ b/crates/punktfunk-core/src/abr.rs @@ -434,6 +434,16 @@ impl BitrateController { } } self.current_kbps = kbps; + // The host may run ABOVE our climb ceiling, and be right to: it sends an unsolicited + // `BitrateChanged` when a rebuild re-resolves an Automatic rate for what it actually + // encodes (a 1080p session mirroring a 4K panel resolves ~3× higher), and that is + // the host's own Automatic answer, not a climb we asked for. Let the ceiling follow + // — `set_ceiling` only ever raises, and still clamps to the operator's + // `PUNKTFUNK_ABR_MAX_MBPS`, which is what must bind here if anything does. Without + // this the ceiling stays at the stale negotiated rate and the step-down below + // immediately drags the host back off the rate it just chose. A no-op for ordinary + // acks: we never request above the effective ceiling in the first place. + self.set_ceiling(kbps); } self.unacked = 0; } @@ -1667,6 +1677,28 @@ mod tests { assert!(c.host_cap_kbps.is_none()); } + #[test] + fn a_host_retarget_above_the_ceiling_raises_it() { + // The host sends an unsolicited `BitrateChanged` when a rebuild re-resolves an Automatic + // rate for what it ACTUALLY encodes — a 1080p session mirroring a 4K panel resolves far + // above the negotiated rate. That is the host's own Automatic answer, so the climb + // ceiling has to follow it; otherwise the ceiling stays stale and the step-down drags + // the host straight back off the rate it just chose. + let mut c = BitrateController::new(20_000); + assert_eq!(c.ceiling_kbps, 20_000); + c.on_ack(60_000); // unsolicited: no request was outstanding + assert_eq!(c.current_kbps, 60_000); + assert_eq!(c.ceiling_kbps, 60_000); + let start = Instant::now(); + // No step-down, and no spurious re-target of any kind. + assert_eq!(run_clean(&mut c, start, 0, 4), None); + // The operator's cap still outranks it — that is the one thing that must bind here. + let mut c = BitrateController::with_ceiling_cap(20_000, Some(50_000)); + c.on_ack(60_000); + assert_eq!(c.ceiling_kbps, 50_000); + assert_eq!(run_clean(&mut c, start, 0, 1), Some(50_000)); + } + #[test] fn a_standing_cap_backs_its_reprobe_clock_off() { // The other half of the re-probe: an encoder's real codec ceiling (794 Mbps, L6.2)