Merge branch 'main' into feat/android-pad-audio
ci / web (pull_request) Successful in 1m24s
android / android (pull_request) Successful in 3m55s
apple / swift (pull_request) Canceled after 0s
apple / screenshots (pull_request) Canceled after 0s
ci / rust (pull_request) Canceled after 5m1s
ci / rust-arm64 (pull_request) Canceled after 3m22s
ci / docs-site (pull_request) Canceled after 1m35s
windows / build (aarch64-pc-windows-msvc) (pull_request) Canceled after 0s
windows / build (x86_64-pc-windows-msvc) (pull_request) Canceled after 0s

This commit is contained in:
2026-08-03 17:39:07 +00:00
7 changed files with 660 additions and 128 deletions
+469 -108
View File
@@ -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<u32> {
.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<i64>,
mean: Option<i64>,
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<u32>,
/// 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 3060 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<u32>,
@@ -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<Instant>,
/// 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"
);
@@ -331,9 +416,34 @@ 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;
// 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;
}
@@ -343,14 +453,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 +515,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 +525,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 +558,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 +578,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 +611,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 +641,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 +701,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 +746,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 +1265,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 +1282,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 +1305,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 +1319,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 +1591,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 +1602,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 +1621,130 @@ 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_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)
// 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 +1926,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 +2002,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 +2019,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 +2108,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 +2124,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 +2161,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);
+27 -12
View File
@@ -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.
@@ -492,7 +499,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();
@@ -500,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);
@@ -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(
+12
View File
@@ -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),
+13
View File
@@ -1087,6 +1087,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::<Reconfigured>();
// 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::<u32>();
// 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
@@ -1147,6 +1158,7 @@ async fn serve_session(
probe_tx,
probe_result_rx,
reconfig_result_rx,
retarget_rx,
cursor_shape_rx,
cursor_client_draws,
clip_enabled,
@@ -1598,6 +1610,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,
@@ -40,6 +40,9 @@ pub(super) async fn run(
probe_tx: std::sync::mpsc::Sender<ProbeRequest>,
mut probe_result_rx: tokio::sync::mpsc::UnboundedReceiver<ProbeResult>,
mut reconfig_result_rx: tokio::sync::mpsc::UnboundedReceiver<Reconfigured>,
// 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<u32>,
mut cursor_shape_rx: tokio::sync::mpsc::UnboundedReceiver<punktfunk_core::quic::CursorShape>,
cursor_client_draws: Arc<AtomicBool>,
clip_enabled: Arc<AtomicBool>,
@@ -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
+108 -8
View File
@@ -1214,6 +1214,9 @@ pub(super) struct SessionContext {
/// `Reconfigured { accepted: true, mode: <actually live> }` 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<Reconfigured>,
/// 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<u32>,
/// Adaptive-FEC target the control task updates from the client's loss reports.
pub(super) fec_target: Arc<AtomicU8>,
/// 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<PreparedDispl
probe_rx,
probe_result_tx,
reconfig_result_tx,
retarget_tx,
fec_target,
conn,
timing_conn,
@@ -1597,7 +1601,12 @@ pub(super) fn virtual_stream(ctx: SessionContext, prepared: Option<PreparedDispl
) = pipe;
// The encoder may have opened at a re-resolved rate (a mirrored head delivering a size this
// session never negotiated). Adopt it before anything downstream reads `bitrate_kbps`.
adopt_built_bitrate(&mut bitrate_kbps, built_bitrate, &live_bitrate);
adopt_built_bitrate(
&mut bitrate_kbps,
built_bitrate,
&live_bitrate,
&retarget_tx,
);
// Capture is live — launch the requested title so it renders onto the streamed output and
// grabs focus. Windows spawns the library id into the interactive user session; Linux spawns
@@ -2037,7 +2046,12 @@ pub(super) fn virtual_stream(ctx: SessionContext, prepared: Option<PreparedDispl
// The new compositor may deliver a different size than the old one did (a
// Game→Desktop switch onto a mirrored 4K panel is exactly that), so adopt
// the rate the rebuilt encoder actually opened at.
adopt_built_bitrate(&mut bitrate_kbps, new_bitrate, &live_bitrate);
adopt_built_bitrate(
&mut bitrate_kbps,
new_bitrate,
&live_bitrate,
&retarget_tx,
);
vd = new_vd;
compositor = sw.compositor;
next = std::time::Instant::now();
@@ -2177,7 +2191,12 @@ pub(super) fn virtual_stream(ctx: SessionContext, prepared: Option<PreparedDispl
}
};
if rebuilt {
adopt_built_bitrate(&mut bitrate_kbps, built_bitrate, &live_bitrate);
adopt_built_bitrate(
&mut bitrate_kbps,
built_bitrate,
&live_bitrate,
&retarget_tx,
);
cur_mode = new_mode;
next = std::time::Instant::now();
// H2/H3: the backend may have honored a different mode than requested — KWin caps
@@ -2306,6 +2325,11 @@ pub(super) fn virtual_stream(ctx: SessionContext, prepared: Option<PreparedDispl
);
if applied_kbps < new_kbps {
encoder_ceiling_kbps.store(applied_kbps, Ordering::Relaxed);
// The control task already acked the client with its own resolve, which was
// higher than what the encoder took. Correct it, or the controller climbs
// from a rate the encoder never ran at until its NEXT request happens to be
// pre-clamped by the ceiling we just stored.
let _ = retarget_tx.send(applied_kbps);
}
if applied_kbps < bitrate_kbps {
// Down-step: the behind-cadence backlog was scored against the old,
@@ -2356,6 +2380,9 @@ pub(super) fn virtual_stream(ctx: SessionContext, prepared: Option<PreparedDispl
enc = new_enc;
if applied_kbps < new_kbps {
encoder_ceiling_kbps.store(applied_kbps, Ordering::Relaxed);
// As in the in-place arm: the ack the client already has promises
// more than the fresh encoder accepted — correct it.
let _ = retarget_tx.send(applied_kbps);
}
bitrate_kbps = applied_kbps;
live_bitrate.store(applied_kbps, Ordering::Relaxed);
@@ -2797,7 +2824,7 @@ pub(super) fn virtual_stream(ctx: SessionContext, prepared: Option<PreparedDispl
// A capture-loss rebuild can land on a different source than it lost (this loop
// re-detects the session every cycle, precisely so it can follow a switch), so the
// delivered size — and with it an Automatic rate — may have changed under us.
adopt_built_bitrate(&mut bitrate_kbps, new_bitrate, &live_bitrate);
adopt_built_bitrate(&mut bitrate_kbps, new_bitrate, &live_bitrate, &retarget_tx);
enc.request_keyframe(); // belt-and-suspenders; a fresh encoder opens on an IDR anyway
last_forced_idr = Some(std::time::Instant::now()); // anchor the IDR cooldown from the rebuild
next = std::time::Instant::now();
@@ -3397,11 +3424,22 @@ pub(super) fn virtual_stream(ctx: SessionContext, prepared: Option<PreparedDispl
};
let escalated = cur_depth > 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<AtomicU32>) {
///
/// 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<AtomicU32>,
retarget: &tokio::sync::mpsc::UnboundedSender<u32>,
) {
if built == *current {
return;
}
@@ -4085,6 +4152,7 @@ fn adopt_built_bitrate(current: &mut u32, built: u32, live: &Arc<AtomicU32>) {
);
*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::<u32>();
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.