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();