From 2563041e8866191c5f7e81f3239f4cd26b94e3bf Mon Sep 17 00:00:00 2001 From: enricobuehler Date: Sat, 15 Aug 2026 20:05:10 +0200 Subject: [PATCH 1/2] feat(host,audio): a 16-second outage and a starved stream reported the same numbers MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Measured on a live host running this commit's parent: our PipeWire capture stream spent 16.2 s of one window in `Paused`, and the line said delivered_pct=63 gaps=0 max_gap_ms=0 Every number is correct. `gaps` scores inter-callback deltas and a stream that is not scheduled fires no callbacks, so there is nothing to score — that is deliberate, and `a_paused_span_is_not_scored` pins it. But the outage did not vanish; it moved into `delivered_pct`, because the reporting window is flushed from the process callback and therefore STRETCHES by exactly the time we were absent. The only explanation lived in the state DEBUG lines, which a field journal at INFO does not carry. So a shortfall had two possible causes and no way to tell them apart: a sink nobody was rendering into (benign — the 2026-08-15 logs show ~40 s of it at every session start, at peak_db=-120.0, i.e. digital silence), or a capture path losing real audio under load (Skynet, 5-8 % through loud gameplay). Those want opposite responses and cost days of investigation to separate by hand. Three changes, one theme — make the line answer the question it invites: - `pauses` / `paused_ms` beside the percentage they explain. `gaps` keeps its narrow meaning (holes inside a running stream); absence is counted separately, because a burst of sub-10 ms holes is a scheduling problem on the box and a multi-second pause is our node not being in the graph at all. - The quantum now tracks what the graph is actually handing us instead of latching the first callback of the open. A graph re-plans whenever anything else on the box asks for a different latency, and the stale value silently corrupted the very threshold gaps are scored against. A new size must survive three callbacks before it is believed, so one short buffer cannot move it. - `audio egress` — the send path had no periodic metric of any kind. Across five field logs it emitted 14 lines, all the same session banner, which made "the host paces audio badly" unfalsifiable and left it on the suspect list forever. It now reports sent/infilled/late/max_late_ms/max_spacing_ms/reanchors on the same 30 s window as capture, so the two read as a pair: holes at the tap with clean departures means the host delivered everything it had. Note `reanchors` in particular — the pacer forgives accumulated debt silently, and that is precisely the event that leaves no trace and then gets blamed on the network. Windows keeps its own field set; its capture model differs (loopback stops delivering while the endpoint idles) and zero-valued fields would imply it had measured something it did not. Gated in the amd64 CI container: fmt, clippy -D warnings, 18 capture-policy tests (4 new, plus the pause pair sitting next to the test that pins the blindness they answer). --- .../src/audio/capture_policy.rs | 234 +++++++++++++++++- crates/punktfunk-host/src/audio/linux/mod.rs | 139 ++++++++--- crates/punktfunk-host/src/native/audio.rs | 44 +++- 3 files changed, 374 insertions(+), 43 deletions(-) diff --git a/crates/punktfunk-host/src/audio/capture_policy.rs b/crates/punktfunk-host/src/audio/capture_policy.rs index 9e88e654..931c8133 100644 --- a/crates/punktfunk-host/src/audio/capture_policy.rs +++ b/crates/punktfunk-host/src/audio/capture_policy.rs @@ -138,6 +138,22 @@ pub(crate) struct CaptureStats { /// memory. Every one of these used to `return` silently, so a stream that fired its callback /// on time and handed us nothing looked identical to a stream nobody was feeding. pub(crate) missed_dequeues: u64, + /// Spans this window spent with the stream NOT in `Streaming`, and how long they totalled. + /// + /// `gaps` deliberately cannot see these (see [`Self::observe_callback`]) — a paused stream + /// fires no callbacks at all, so there is no delta to score and the caller drops its cadence + /// stamp on every transition. The cost of that correct decision was that the outage went + /// somewhere else entirely: into `delivered_pct`, as an unattributed shortfall, because the + /// reporting window is flushed from the callback and therefore STRETCHES by exactly the time + /// we were not being scheduled. + /// + /// Measured on a live host on 2026-08-15: a 16.2 s pause produced + /// `delivered_pct=63 gaps=0 max_gap_ms=0`. Every number was correct and the line still could + /// not say what happened — the explanation existed only in the state DEBUG lines, which a + /// field journal at INFO does not carry. These two fields are that explanation, at INFO, + /// beside the percentage they explain. + pub(crate) pauses: u64, + pub(crate) paused_us: u64, } impl CaptureStats { @@ -157,9 +173,13 @@ impl CaptureStats { /// /// `since_last` is `None` for the first callback of a stream — and, deliberately, for the /// first after a state transition: the caller drops its stamp when the stream pauses, so a - /// legitimately Paused span is not scored as one enormous hole. (The Paused↔Streaming flaps - /// around a format renegotiation stay visible as the state DEBUG lines next to a small - /// post-resume gap, which is the honest reading of what happened.) + /// legitimately Paused span is not scored as one enormous hole. + /// + /// That leaves this counter about ONE thing — holes inside a stream that is running — and + /// pushes the other kind onto [`Self::observe_pause`]. The split matters because the two want + /// opposite answers: a run of sub-10 ms holes is a scheduling problem on the box, whereas a + /// multi-second pause is our node not being in the graph at all. A single "gap" number that + /// mixed them would be worse than either. /// /// `quantum` is the NEGOTIATED buffer duration, not the one we asked for: a graph handing us /// 21.3 ms buffers is not gapping when its callbacks are 21.3 ms apart — it is doing exactly @@ -183,6 +203,21 @@ impl CaptureStats { self.max_gap_us / 1_000 } + /// Record one span the stream spent away from `Streaming`. + /// + /// Called on the transition BACK, so the whole span lands in the window that is flushed after + /// the resume — which is the same window whose `delivered_pct` the span diluted. Keeping the + /// two together is the entire point: apart, neither is interpretable. + pub(crate) fn observe_pause(&mut self, span: Duration) { + self.pauses += 1; + self.paused_us += span.as_micros() as u64; + } + + /// Total time away from `Streaming` this window, in whole ms. + pub(crate) fn paused_ms(&self) -> u64 { + self.paused_us / 1_000 + } + /// `(peak dBFS, rms dBFS, delivered %)` for this window. Silence reports -120 dB rather than /// -inf so the log line stays parseable. pub(crate) fn summary(&self, elapsed: Duration, sample_rate: u32) -> (f64, f64, f64) { @@ -199,6 +234,76 @@ impl CaptureStats { } } +/// A departure this far past its slot is a slip worth counting rather than ordinary jitter: one +/// whole protocol frame, so a frame that merely rounds late never scores. +const LATE_DEPARTURE: Duration = Duration::from_millis(FRAME_MS as u64); + +/// One reporting window of AUDIO EGRESS vitals (WP-C). +/// +/// Capture has been instrumented since WP-A2 and the send path has not, so a field log could show +/// audio arriving at the tap and say nothing whatsoever about how it left. That asymmetry is not +/// neutral: it made "the host paces audio badly" unfalsifiable, and an unfalsifiable suspect stays +/// on the list forever. Across five 2026-08-15 field logs the entire egress path emitted 14 lines, +/// all of them the same session-open banner. +/// +/// The point of these counters is to be *boring*. If departures are clean while capture reports +/// holes, the pacing rework introduced in v0.25 is acquitted permanently and the search moves +/// upstream for good. +#[derive(Default)] +pub(crate) struct SendStats { + pub(crate) sent: u64, + /// Frames synthesized to cover a capture hole. Wire continuity and captured continuity are + /// different claims and a log that conflates them cannot be used to judge either. + pub(crate) infilled: u64, + /// Departures that missed their paced slot by at least [`LATE_DEPARTURE`]. + pub(crate) late: u64, + /// The worst such miss, µs — kept even when the count is zero, because "never late" and + /// "never late by a whole frame" are different statements. + pub(crate) max_late_us: u64, + /// Widest gap between two consecutive departures, µs. The number a client-side starvation + /// complaint is actually about: the wire going quiet, whatever the reason. + pub(crate) max_spacing_us: u64, + /// Times the schedule fell more than `PACE_REANCHOR` behind and was re-anchored instead of + /// chased. Each one silently forgives accumulated debt, which is exactly the kind of event + /// that leaves no trace and then gets blamed on the network. + pub(crate) reanchors: u64, +} + +impl SendStats { + /// Score one frame leaving the host. `late` is how far past its paced slot it went (zero when + /// the schedule is unanchored), `since_prev` the spacing from the previous departure. + pub(crate) fn observe_departure( + &mut self, + late: Duration, + since_prev: Option, + infilled: bool, + ) { + self.sent += 1; + if infilled { + self.infilled += 1; + } + self.max_late_us = self.max_late_us.max(late.as_micros() as u64); + if late >= LATE_DEPARTURE { + self.late += 1; + } + if let Some(gap) = since_prev { + self.max_spacing_us = self.max_spacing_us.max(gap.as_micros() as u64); + } + } + + pub(crate) fn observe_reanchor(&mut self) { + self.reanchors += 1; + } + + pub(crate) fn max_late_ms(&self) -> u64 { + self.max_late_us / 1_000 + } + + pub(crate) fn max_spacing_ms(&self) -> u64 { + self.max_spacing_us / 1_000 + } +} + /// How long a capture hole may run before the wire starts covering it. Two protocol frames: long /// enough that ordinary quantum jitter never trips it, short enough that the client's ring never /// notices the hole. @@ -525,4 +630,127 @@ mod tests { assert_eq!(s.gaps, 0); assert_eq!(s.max_gap_ms(), 0); } + + /// The companion to the test above, and the reason it is safe: a pause stays out of `gaps`, + /// but it does NOT stay out of the log line. Numbers are the ones measured on a live host on + /// 2026-08-15, where a 16.2 s pause reported `delivered_pct=63 gaps=0 max_gap_ms=0` and no + /// field in the line could say why. + #[test] + fn a_paused_span_is_reported_even_though_it_is_not_a_gap() { + let mut s = CaptureStats::default(); + s.observe_callback(Some(Duration::from_millis(5)), Q); + s.observe_pause(Duration::from_millis(16_214)); + s.observe_callback(None, Q); // resumed + s.observe_callback(Some(Duration::from_millis(5)), Q); + + assert_eq!(s.gaps, 0, "a pause is still not a delivery gap"); + assert_eq!(s.max_gap_ms(), 0); + assert_eq!(s.pauses, 1, "…but it is now countable"); + assert_eq!(s.paused_ms(), 16_214); + } + + /// One long outage and a burst of short flaps must not read alike — the same argument that + /// makes `gaps` and `max_gap_ms` two fields instead of one. The triple here is the shape every + /// Skynet and AVALON session start produced: three dwells, no format actually changing. + #[test] + fn pause_spans_accumulate_and_stay_countable() { + let mut long = CaptureStats::default(); + long.observe_pause(Duration::from_millis(38_400)); + + let mut flappy = CaptureStats::default(); + for ms in [12_534, 17_030, 8_765] { + flappy.observe_pause(Duration::from_millis(ms)); + } + + assert_eq!(long.pauses, 1); + assert_eq!(flappy.pauses, 3); + assert_eq!(flappy.paused_ms(), 38_329); + assert!( + long.paused_ms().abs_diff(flappy.paused_ms()) < 100, + "near-identical dead time, and the count is the only thing that separates them" + ); + } + + /// The discriminator the field logs needed. A stream that is running and starved reports gaps + /// and NO pause; a stream that was never scheduled reports the mirror image. Both dilute + /// `delivered_pct` identically, which is exactly why neither can be diagnosed from it alone. + #[test] + fn starvation_and_absence_are_told_apart() { + let mut starved = CaptureStats::default(); + for _ in 0..60 { + starved.observe_callback(Some(Duration::from_millis(30)), Q); + } + + let mut absent = CaptureStats::default(); + absent.observe_pause(Duration::from_millis(1_800)); + + assert_eq!(starved.gaps, 60); + assert_eq!(starved.pauses, 0, "a running stream was never absent"); + assert_eq!(absent.gaps, 0); + assert_eq!(absent.pauses, 1, "an absent stream never got to be slow"); + assert_eq!(absent.paused_ms(), 1_800); + } + + /// The acquittal case, and the whole reason [`SendStats`] exists: a pacer doing its job must + /// produce a line a reader can dismiss at a glance. + #[test] + fn a_healthy_pacer_reports_nothing_alarming() { + let mut s = SendStats::default(); + let frame = Duration::from_millis(FRAME_MS as u64); + for i in 0..200 { + s.observe_departure(Duration::ZERO, (i > 0).then_some(frame), false); + } + assert_eq!(s.sent, 200); + assert_eq!(s.late, 0); + assert_eq!(s.reanchors, 0); + assert_eq!(s.infilled, 0); + assert_eq!(s.max_late_ms(), 0); + assert_eq!(s.max_spacing_ms(), FRAME_MS as u64); + } + + /// Lateness under one frame is jitter, not a slip — but it must still be *visible*, or + /// "never late" and "never late by a whole frame" become the same report. + #[test] + fn sub_frame_lateness_is_measured_without_being_counted() { + let mut s = SendStats::default(); + s.observe_departure(Duration::from_micros(3_400), None, false); + assert_eq!(s.late, 0, "3.4 ms has not slipped a whole 5 ms slot"); + assert_eq!(s.max_late_ms(), 3, "…and it is still on the record"); + } + + /// A slot missed by a whole frame or more is the event the field logs could never show. + #[test] + fn a_slipped_slot_is_counted_and_its_worst_case_kept() { + let mut s = SendStats::default(); + s.observe_departure(Duration::from_millis(6), None, false); + s.observe_departure( + Duration::from_millis(41), + Some(Duration::from_millis(47)), + false, + ); + s.observe_departure(Duration::ZERO, Some(Duration::from_millis(5)), false); + s.observe_reanchor(); + + assert_eq!(s.late, 2); + assert_eq!(s.max_late_ms(), 41); + assert_eq!(s.max_spacing_ms(), 47, "the wire's worst quiet stretch"); + assert_eq!(s.reanchors, 1); + } + + /// Wire continuity is not captured continuity. A window whose frames were all synthesized + /// looks perfect on every other counter, and must not be readable as healthy audio. + #[test] + fn synthesized_frames_stay_distinguishable_from_captured_ones() { + let mut s = SendStats::default(); + let frame = Duration::from_millis(FRAME_MS as u64); + for _ in 0..100 { + s.observe_departure(Duration::ZERO, Some(frame), true); + } + assert_eq!(s.sent, 100); + assert_eq!( + s.infilled, 100, + "every one of these was silence we invented" + ); + assert_eq!(s.late, 0); + } } diff --git a/crates/punktfunk-host/src/audio/linux/mod.rs b/crates/punktfunk-host/src/audio/linux/mod.rs index 473fc502..779cf853 100644 --- a/crates/punktfunk-host/src/audio/linux/mod.rs +++ b/crates/punktfunk-host/src/audio/linux/mod.rs @@ -395,6 +395,11 @@ const MIC_STALE: Duration = Duration::from_secs(1); /// against the same number the ask used. const CAPTURE_QUANTUM_FRAMES: u32 = 240; +/// Callbacks that must agree on a new buffer size before it replaces the one gaps are scored +/// against. Three is enough to reject a boundary artefact and still adopt a genuine re-plan +/// within ~15 ms. +const QUANTUM_CONFIRM: u8 = 3; + fn mic_pw_thread( pcm_rx: Receiver<(std::time::Instant, Vec)>, quit_rx: pipewire::channel::Receiver, @@ -782,12 +787,15 @@ fn pw_thread( channels: u32, stats: crate::audio::capture_policy::CaptureStats, last_stats: std::time::Instant, - /// Whether this OPEN has reported its negotiated buffer size yet. Per-open, not the - /// process-wide `static AtomicBool` this replaces: a host runs for days across many - /// sessions, so the old form reported the very first capture of the process and then - /// never again — the one number that identifies a clamped quantum, invisible on every + /// Frames per callback the graph is currently handing us, `0` until the first is + /// confirmed. Per-open, not a process-wide latch: a host runs for days across many + /// sessions, so a process-wide form reported the very first capture and then never + /// again — the one number that identifies a clamped quantum, invisible on every /// subsequent open (including every reopen after a device change). - reported_quantum: bool, + quantum_frames: usize, + /// A buffer size seen but not yet believed, with how many callbacks in a row have + /// agreed on it. Stops one short buffer from moving the gap threshold. + quantum_candidate: Option<(usize, u8)>, /// When the callback last ran (WP-A2), so its CADENCE can be scored and not just its /// content. Cleared across a state transition — a deliberate Paused span must not /// read as one enormous hole. Lives here rather than in `stats` because the stats @@ -804,19 +812,24 @@ fn pw_thread( /// Shared with the capturer — see [`PwAudioCapturer::active`]. Read on every /// failed hand-off to keep parked-capturer backpressure out of the drop count. active: Arc, + /// When the stream last left `Streaming`, so the span can be charged to the window + /// that the span itself stretched. `None` while streaming. + paused_since: Option, } let ud = CapUd { tx, channels, stats: Default::default(), last_stats: std::time::Instant::now(), - reported_quantum: false, + quantum_frames: 0, + quantum_candidate: None, last_cb: None, quantum: Duration::from_micros( CAPTURE_QUANTUM_FRAMES as u64 * 1_000_000 / SAMPLE_RATE as u64, ), negotiated: None, active, + paused_since: None, }; let _listener = stream .add_local_listener_with_user_data(ud) @@ -829,6 +842,22 @@ fn pw_thread( // existing. Scoring it would report one huge hole per renegotiation and bury // the sub-10 ms ones the field log is actually about (WP-A2). ud.last_cb = None; + // …but it still has to be reported, because the reporting window is flushed + // from the process callback and therefore stretches by the whole span. Charge + // it to the window flushed after the resume — the same window it diluted. + // Without this the line says `delivered_pct=4 gaps=0` and cannot say whether + // that is a dead capture path or a sink nobody was rendering into; the + // 2026-08-15 field logs are 40 s of exactly that ambiguity per session start. + match new { + pw::stream::StreamState::Streaming => { + if let Some(since) = ud.paused_since.take() { + ud.stats.observe_pause(since.elapsed()); + } + } + _ => { + ud.paused_since.get_or_insert_with(std::time::Instant::now); + } + } // A stream error is unrecoverable for this instance — exit so the sessions' // reopen path builds a fresh one (same contract as the core-error path above). if matches!(new, pw::stream::StreamState::Error(_)) { @@ -913,42 +942,70 @@ fn pw_thread( let region = &buf[offset..(offset + size).min(buf.len())]; // Negotiated as F32LE; reinterpret the byte region as interleaved f32. let n = region.len() / 4; - if !ud.reported_quantum { - ud.reported_quantum = true; - // What we ASKED for vs what PipeWire actually handed us. Stating only the - // result ("samples=2048") reads as a fact about the device; stating it - // next to the request is what makes a clamp legible. A VM is the common - // cause — stock `pipewire.conf` raises `default.clock.min-quantum` to - // 1024 whenever `cpu.vm.name` is set, so a 5 ms ask silently becomes - // 21.3 ms and the audio plane starts arriving in bursts. That cost a - // whole field investigation to find; it should cost one log line. - let frames = n / (ud.channels.max(1) as usize); - let want = CAPTURE_QUANTUM_FRAMES as usize; - // What a gap is measured against from here on — see `CapUd::quantum`. - if frames > 0 { + // Track the quantum the graph is ACTUALLY handing us, not merely the first one + // it ever did. The graph re-plans whenever anything else on the box asks for a + // different latency, and latching the first callback of the open left every + // subsequent gap scored against a buffer size that no longer existed — a + // silent corruption of the one metric this whole diagnosis rests on. A new + // size has to survive `QUANTUM_CONFIRM` callbacks before it is believed, + // because one short buffer at a boundary is not a new deal. + let frames = n / (ud.channels.max(1) as usize); + if frames > 0 && frames != ud.quantum_frames { + let streak = match ud.quantum_candidate { + Some((f, c)) if f == frames => c.saturating_add(1), + _ => 1, + }; + if streak < QUANTUM_CONFIRM { + ud.quantum_candidate = Some((frames, streak)); + } else { + let was = ud.quantum_frames; + ud.quantum_frames = frames; + ud.quantum_candidate = None; + // What a gap is measured against from here on — see `CapUd::quantum`. ud.quantum = Duration::from_micros( frames as u64 * 1_000_000 / SAMPLE_RATE as u64, ); + let want = CAPTURE_QUANTUM_FRAMES as usize; + let negotiated_ms = + format!("{:.1}", frames as f32 * 1000.0 / SAMPLE_RATE as f32); + if was != 0 { + // A mid-open change. Rare, and worth a line of its own: it moves + // the gap threshold under a reader who is comparing windows. + tracing::info!( + previous_frames = was, + negotiated_frames = frames, + negotiated_ms, + "the audio graph re-planned our quantum mid-stream" + ); + } else if frames > want { + // What we ASKED for vs what PipeWire actually handed us. Stating + // only the result ("samples=2048") reads as a fact about the + // device; stating it next to the request is what makes a clamp + // legible. A VM is the common cause — stock `pipewire.conf` raises + // `default.clock.min-quantum` to 1024 whenever `cpu.vm.name` is + // set, so a 5 ms ask silently becomes 21.3 ms and the audio plane + // starts arriving in bursts. That cost a whole field + // investigation to find; it should cost one log line. + tracing::warn!( + requested_frames = want, + negotiated_frames = frames, + negotiated_ms, + "the audio graph refused our low-latency quantum — capture \ + arrives in bursts this size, and the client must buffer at \ + least that much to play them smoothly. On a VM this is \ + PipeWire's `default.clock.min-quantum = 1024` rule; check \ + `pw-metadata -n settings`" + ); + } else { + tracing::info!( + requested_frames = want, + negotiated_frames = frames, + "audio capture quantum negotiated" + ); + } } - if frames > want { - tracing::warn!( - requested_frames = want, - negotiated_frames = frames, - negotiated_ms = - format!("{:.1}", frames as f32 * 1000.0 / SAMPLE_RATE as f32), - "the audio graph refused our low-latency quantum — capture arrives \ - in bursts this size, and the client must buffer at least that \ - much to play them smoothly. On a VM this is PipeWire's \ - `default.clock.min-quantum = 1024` rule; check \ - `pw-metadata -n settings`" - ); - } else { - tracing::info!( - requested_frames = want, - negotiated_frames = frames, - "audio capture quantum negotiated" - ); - } + } else if frames == ud.quantum_frames { + ud.quantum_candidate = None; } let mut samples = Vec::with_capacity(n); for i in 0..n { @@ -991,6 +1048,12 @@ fn pw_thread( // percentage and mean entirely different things. gaps = ud.stats.gaps, max_gap_ms = ud.stats.max_gap_ms(), + // The OTHER thing a shortfall can be (see `CaptureStats::pauses`): + // time our node was not in the graph at all. `gaps` deliberately + // cannot see it, so without these two a paused span and a starved + // stream are the same number. + pauses = ud.stats.pauses, + paused_ms = ud.stats.paused_ms(), missed_dequeues = ud.stats.missed_dequeues, dropped_chunks = ud.stats.dropped_chunks, "desktop audio capture" diff --git a/crates/punktfunk-host/src/native/audio.rs b/crates/punktfunk-host/src/native/audio.rs index 356add00..4cc3e096 100644 --- a/crates/punktfunk-host/src/native/audio.rs +++ b/crates/punktfunk-host/src/native/audio.rs @@ -218,6 +218,12 @@ pub(super) fn audio_thread( // frame. `sent_any` is what keeps the seed from ever reaching the wire. let mut next_pts_ns: u64 = 0; let mut pace_due: Option = None; + // WP-C — what the wire actually did, as opposed to what the tap handed us. See [`SendStats`]: + // until this existed the send path was the one stage of the audio pipeline that could not be + // ruled in or out from a field log. + let mut send_stats = crate::audio::capture_policy::SendStats::default(); + let mut last_send_stats = std::time::Instant::now(); + let mut last_departure: Option = None; if capturer.is_some() { tracing::info!( channels = want, @@ -315,12 +321,20 @@ pub(super) fn audio_thread( // send-time debt. loop { let now = std::time::Instant::now(); + // How far past its slot this frame is leaving. Measured before the re-anchor arm can + // erase the evidence — that arm is the one that forgives debt silently (WP-C). + let mut late = std::time::Duration::ZERO; match pace_due { Some(due) if due > now => break, // this frame's slot has not arrived yet - Some(due) if now.duration_since(due) > PACE_REANCHOR => pace_due = None, - _ => {} + Some(due) if now.duration_since(due) > PACE_REANCHOR => { + send_stats.observe_reanchor(); + pace_due = None; + } + Some(due) => late = now.duration_since(due), + None => {} } frame_buf.clear(); + let mut infilled = false; if acc.len() >= frame_len { frame_buf.extend(acc.drain(..frame_len)); } else if !sent_any { @@ -328,6 +342,7 @@ pub(super) fn audio_thread( } else { match infill.decide(last_chunk_at.elapsed()) { crate::audio::capture_policy::Infill::Silence => { + infilled = true; // Pad the partial frame out with silence and send THAT, rather than // leaving it for post-gap samples to complete: one frame carrying audio // from both sides of a hole is a click, and its pts is a lie about when @@ -366,6 +381,15 @@ pub(super) fn audio_thread( prev_frame.extend_from_slice(opus); } seq = seq.wrapping_add(1); + // Score the departure against its slot and against the previous one. `now` is + // from the top of this iteration — microseconds earlier and one clock read + // cheaper, 200 times a second. + send_stats.observe_departure( + late, + last_departure.map(|t| now.duration_since(t)), + infilled, + ); + last_departure = Some(now); // From here there is a continuity worth protecting, and `next_pts_ns` has a // real anchor to continue from — both preconditions for synthesizing anything. sent_any = true; @@ -382,6 +406,22 @@ pub(super) fn audio_thread( } } } + if last_send_stats.elapsed() >= crate::audio::capture_policy::STATS_EVERY { + // Deliberately the same window as the capture line, so the two can be read as a pair: + // holes at the tap with clean departures means the host delivered everything it had, + // and the search belongs upstream of us. + tracing::info!( + sent = send_stats.sent, + infilled = send_stats.infilled, + late = send_stats.late, + max_late_ms = send_stats.max_late_ms(), + max_spacing_ms = send_stats.max_spacing_ms(), + reanchors = send_stats.reanchors, + "audio egress" + ); + send_stats = Default::default(); + last_send_stats = std::time::Instant::now(); + } } // Park the live capturer for the next session (None if it died and never reopened), // releasing its session-scoped routing claim (Linux: the default sink moves back; -- 2.54.0 From 7c964e95c687a08c5d77012a95b21da17e7ea9fe Mon Sep 17 00:00:00 2001 From: enricobuehler Date: Sat, 15 Aug 2026 20:13:14 +0200 Subject: [PATCH 2/2] fix(host,audio): the rtkit boost was applied to a thread that does not run the capture callback MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The comment above it asserted "the stream's `process` callbacks run ON this mainloop thread (we never hand PipeWire a separate data loop)". We do: the stream is created with `RT_PROCESS`, so libpipewire runs `process()` on a data loop it creates and schedules itself. Measured in one live host process on 2026-08-15: punktfunk-pw-au SCHED_OTHER nice 0 <- the thread we boost data-loop.0 SCHED_RR prio 20 <- the thread running process() This is not a stale-comment nit. #232 shipped the rtkit boost to answer a field report of audio stutter, its success line was read as evidence that the capture callback had been prioritised, and the follow-up round concluded priorities were "engaged but insufficient" — when they had never been applied to the thread in question. Whether the capture callback is realtime decides whether a Wine shader storm can deschedule it for tens of milliseconds at a 2.7 ms quantum, which is the exact shape of the one field signature still unexplained (~2 stalls/s of ~30 ms with the node reporting itself continuously Streaming). So the assumption is replaced by a measurement rather than a guess about which thread to boost: - `pf_frame::thread_qos::current_thread_sched()` reports the calling thread's policy, RT priority and nice. Three by-value syscalls, no allocation and no blocking, so it is safe to call from an RT callback. - The capture callback reports its own scheduling once per open. Every future field log now states what the audio path actually runs as, instead of what we asked for somewhere else. The boost itself is kept: this thread still dispatches state and format events, and it IS the capture thread when `PUNKTFUNK_STREAM_SINK=0` selects the legacy monitor path. Deliberately not attempted here: boosting the data loop. rtkit is a blocking D-Bus call and must never run inside an RT audio callback, and on the one host that could be measured PipeWire already gives that thread SCHED_RR/20 — a nice boost would be inert. Ship the instrument first; a host that reports SCHED_OTHER here is the evidence that would justify the plumbing. Gated in the amd64 CI container: fmt, clippy -D warnings on punktfunk-host and pf-frame, pf-frame tests (incl. a non-vacuity test that the introspection returns a policy the kernel could have named), and the full punktfunk-host suite at 560 passed. The one failure, mgmt::tests::local_summary_is_loopback_only_and_ non_sensitive, is the recorded process-global parallel-test race: it passes when run alone, which is the documented discriminator, and it is untouched by this. --- crates/pf-frame/src/thread_qos.rs | 69 ++++++++++++++++++++ crates/punktfunk-host/src/audio/linux/mod.rs | 37 ++++++++++- 2 files changed, 103 insertions(+), 3 deletions(-) diff --git a/crates/pf-frame/src/thread_qos.rs b/crates/pf-frame/src/thread_qos.rs index 9672bd98..a34bebbd 100644 --- a/crates/pf-frame/src/thread_qos.rs +++ b/crates/pf-frame/src/thread_qos.rs @@ -82,6 +82,51 @@ pub fn boost_thread_priority(critical: bool) { } } +/// What the OS is actually giving the CALLING thread: `(policy, rt_priority, nice)`. +/// +/// Exists because a boost we *asked for* and a boost the hot thread *has* turned out to be +/// different questions. Callbacks handed to a library — PipeWire's `RT_PROCESS` streams above all +/// — run on a thread that library created and schedules, so a `boost_thread_priority` call in our +/// own setup path can log a cheerful success about a thread that never touches audio. A +/// 2026-08-15 measurement found exactly that shape: our loop thread at SCHED_OTHER/0 while the +/// data loop actually running the capture callback sat at SCHED_RR/20, both in the same process. +/// +/// Report this from inside the hot callback, where "the calling thread" is the one that matters. +#[cfg(target_os = "linux")] +pub fn current_thread_sched() -> (&'static str, i32, i32) { + // SAFETY: all three calls take by-value integers (plus, for `sched_getparam`, a pointer to a + // fully-initialised local we own and outlive) and return integers. `0` means "the calling + // task" on Linux, so nothing outside this thread is read or written, and no allocation, + // locking or blocking happens — which is what makes this callable from an RT callback. + unsafe { + let policy = libc::sched_getscheduler(0); + let mut param: libc::sched_param = std::mem::zeroed(); + let rt_priority = if libc::sched_getparam(0, &mut param) == 0 { + param.sched_priority + } else { + -1 + }; + // `getpriority` legitimately returns -1, so errno is the only way to tell a nice of -1 + // from a failure. + *libc::__errno_location() = 0; + let nice = libc::getpriority(libc::PRIO_PROCESS, 0); + let nice = if *libc::__errno_location() == 0 { + nice + } else { + 0 + }; + let policy = match policy { + libc::SCHED_FIFO => "SCHED_FIFO", + libc::SCHED_RR => "SCHED_RR", + libc::SCHED_OTHER => "SCHED_OTHER", + libc::SCHED_BATCH => "SCHED_BATCH", + libc::SCHED_IDLE => "SCHED_IDLE", + _ => "unknown", + }; + (policy, rt_priority, nice) + } +} + /// RealtimeKit fallback for [`boost_thread_priority`]: ask the system-bus broker /// (`org.freedesktop.RealtimeKit1`) to renice the calling thread when the direct /// `setpriority` was refused. This is how PulseAudio/PipeWire clients get their boosts on a @@ -121,3 +166,27 @@ mod linux_rtkit { Ok(()) } } + +#[cfg(all(test, target_os = "linux"))] +mod tests { + /// Non-vacuity: the introspection has to come back with something the OS could actually have + /// said. A helper whose whole job is to be quoted in a field log is worthless if it can + /// quietly report a placeholder, and it only ever runs on hosts nobody can attach a debugger + /// to. + #[test] + fn current_thread_sched_reports_a_real_policy() { + let (policy, rt_priority, nice) = super::current_thread_sched(); + assert!( + matches!( + policy, + "SCHED_OTHER" | "SCHED_RR" | "SCHED_FIFO" | "SCHED_BATCH" | "SCHED_IDLE" + ), + "unrecognised policy {policy}" + ); + assert!( + (0..=99).contains(&rt_priority), + "rt priority {rt_priority} outside the kernel's range" + ); + assert!((-20..=19).contains(&nice), "nice {nice} outside PRIO range"); + } +} diff --git a/crates/punktfunk-host/src/audio/linux/mod.rs b/crates/punktfunk-host/src/audio/linux/mod.rs index 779cf853..c26b7ef5 100644 --- a/crates/punktfunk-host/src/audio/linux/mod.rs +++ b/crates/punktfunk-host/src/audio/linux/mod.rs @@ -691,9 +691,19 @@ fn pw_thread( use pw::{properties::properties, spa}; use spa::param::audio::{AudioFormat, AudioInfoRaw}; use spa::pod::Pod; - // The stream's `process` callbacks run ON this mainloop thread (we never hand PipeWire a - // separate data loop), so PipeWire's own client `module-rt` boost of its data loops does not - // cover it — the ~2.7 ms capture quantum lives or dies by this thread's scheduling. + // ⚠ This boosts the MAINLOOP thread, which is NOT where the capture callback runs. + // + // The previous comment here asserted the opposite ("we never hand PipeWire a separate data + // loop"), and it was wrong: we pass `RT_PROCESS` below, so libpipewire runs `process()` on a + // data loop it creates and schedules itself. Measured in one live host process on 2026-08-15 + // — this thread at SCHED_OTHER/nice 0, `data-loop.0` at SCHED_RR/20. That mattered more than + // a stale comment usually does: a field investigation read the boost's success line as + // evidence that the audio callback was prioritised, and spent a round concluding priorities + // were "engaged but insufficient" when they had never been applied to the thread in question. + // + // The boost is kept — this thread still dispatches state and format events, and it IS the + // capture thread when `PUNKTFUNK_STREAM_SINK=0` selects the legacy monitor path. What replaces + // the assumption is a measurement: the callback reports its own scheduling on first entry. pf_frame::thread_qos::boost_thread_priority(true); // Setup errors funnel through the ready handshake (mirrors mic_pw_thread's IIFE). @@ -796,6 +806,8 @@ fn pw_thread( /// A buffer size seen but not yet believed, with how many callbacks in a row have /// agreed on it. Stops one short buffer from moving the gap threshold. quantum_candidate: Option<(usize, u8)>, + /// Whether this open has reported the scheduling of the thread running `process()`. + reported_sched: bool, /// When the callback last ran (WP-A2), so its CADENCE can be scored and not just its /// content. Cleared across a state transition — a deliberate Paused span must not /// read as one enormous hole. Lives here rather than in `stats` because the stats @@ -823,6 +835,7 @@ fn pw_thread( last_stats: std::time::Instant::now(), quantum_frames: 0, quantum_candidate: None, + reported_sched: false, last_cb: None, quantum: Duration::from_micros( CAPTURE_QUANTUM_FRAMES as u64 * 1_000_000 / SAMPLE_RATE as u64, @@ -917,6 +930,24 @@ fn pw_thread( ud.last_cb = Some(now); ud.stats.observe_callback(since_last, ud.quantum); + if !ud.reported_sched { + ud.reported_sched = true; + // Say what the thread that ACTUALLY runs this callback is scheduled as. + // Whether the capture callback is realtime decides whether a Wine shader + // storm can deschedule it for tens of ms at a 2.7 ms quantum, and until + // now no log anywhere carried the answer — only that we had asked for a + // boost, on a different thread. Once per open, off the hot path after + // that. + let (policy, rt_priority, nice) = + pf_frame::thread_qos::current_thread_sched(); + tracing::info!( + policy, + rt_priority, + nice, + "audio capture callback scheduling" + ); + } + let Some(mut buffer) = stream.dequeue_buffer() else { ud.stats.missed_dequeues += 1; return; -- 2.54.0