diff --git a/Cargo.lock b/Cargo.lock index f69e911f..1a600c63 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3532,12 +3532,14 @@ dependencies = [ name = "punktfunk-client-session" version = "0.30.0" dependencies = [ + "log", "pf-client-core", "pf-console-ui", "pf-presenter", "punktfunk-core", "serde_json", "tracing", + "tracing-log", "tracing-subscriber", "winresource", ] diff --git a/clients/session/Cargo.toml b/clients/session/Cargo.toml index 5be2aef9..abee381b 100644 --- a/clients/session/Cargo.toml +++ b/clients/session/Cargo.toml @@ -39,6 +39,11 @@ punktfunk-core = { path = "../../crates/punktfunk-core", features = ["quic"] } serde_json = { version = "1", optional = true } tracing = "0.1" tracing-subscriber = { version = "0.3", features = ["env-filter"] } +# The log ring normalizes `log`-crate events off the bridge's "log" shim target +# (`ring_layer.rs`) so the vendored decoder's per-frame DEBUG chatter can be gated by its real +# module path; both are already in the graph through tracing-subscriber's default features. +tracing-log = "0.2" +log = "0.4" # This crate carries NO toolkit, deliberately: it is the renderer the shells spawn, and the # `--no-default-features` build is what a minimal/embedded image installs. GTK4/libadwaita/relm4 diff --git a/clients/session/src/ring_layer.rs b/clients/session/src/ring_layer.rs index dae3e7c5..58f3d0be 100644 --- a/clients/session/src/ring_layer.rs +++ b/clients/session/src/ring_layer.rs @@ -2,15 +2,51 @@ //! "Send logs to host" action. Captures at DEBUG+ regardless of `RUST_LOG` (its own filter is //! applied at install), mirroring the host's `log_capture::RingLayer`: the whole point is that //! a field report carries the diagnostics nobody thought to enable beforehand. +//! +//! …which is exactly why it also has to keep OUT the chatter that would evict them. The ring +//! holds 4096 lines. The vendored H.265 parser (`cros_codecs`, behind `pf-bitstream`) DEBUG-logs +//! its DPB bookkeeping — "Retaining pic POC", "Stored picture", "Set reference", "Bumping POC", +//! one `find_short_term_ref_by_poc` per reference — a dozen lines PER FRAME, so at 120 fps the +//! ring turns over in about three seconds. The 2026-08-17 field bundle from a Steam Deck read +//! `… 2037456 older lines evicted from the ring …` followed by 3.5 s of DPB chatter: the whole +//! 27-minute session, including the 10 s `audio playback buffer_ms= underruns=` line three +//! investigation rounds had been waiting for, was gone. A field ring that a healthy decoder can +//! flush is worse than no ring, because it looks like diagnostics and carries none. use std::fmt::Write as _; use tracing::field::{Field, Visit}; use tracing_subscriber::layer::Context; +/// Targets whose DEBUG/TRACE output is steady-state per-frame chatter, not diagnostics. The ring +/// keeps their INFO-and-up. Prefix-matched on module-path boundaries, so `cros_codecs::codec::…` +/// is gated and a hypothetical `cros_codecs_probe` is not. Same shape as the host's +/// `log_capture::NOISY_DEBUG_TARGETS`. +const NOISY_DEBUG_TARGETS: &[&str] = &["cros_codecs"]; + +fn is_noisy_debug(target: &str) -> bool { + NOISY_DEBUG_TARGETS.iter().any(|t| { + target + .strip_prefix(t) + .is_some_and(|rest| rest.is_empty() || rest.starts_with("::")) + }) +} + pub(crate) struct RingLayer; impl tracing_subscriber::Layer for RingLayer { fn on_event(&self, event: &tracing::Event<'_>, _ctx: Context<'_, S>) { + // Events from `log`-crate dependencies (the vendored decoder among them) arrive through + // the tracing-log bridge under the shim target "log", with the record's real module path + // tucked into `log.target=`. Normalize back to the real metadata so the noise gate below + // and the target column both see `cros_codecs::…` — under the shim target every bridged + // event is indistinguishable from every other, and the field bundle's target column read + // `log` for two million lines. + use tracing_log::NormalizeEvent; + let normalized = event.normalized_metadata(); + let meta = normalized.as_ref().unwrap_or_else(|| event.metadata()); + if *meta.level() > tracing::Level::INFO && is_noisy_debug(meta.target()) { + return; + } struct V(String); impl Visit for V { fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) { @@ -20,14 +56,16 @@ impl tracing_subscriber::Layer for RingLayer { let rest = std::mem::take(&mut self.0); let _ = write!(self.0, "{value:?}"); self.0.push_str(&rest); - } else { + } else if !field.name().starts_with("log.") { + // `log.target`/`log.module_path`/`log.file`/`log.line` are the bridge's own + // bookkeeping — already surfaced through the normalized target above, and + // 150 bytes of repeated path per line otherwise. let _ = write!(self.0, " {}={:?}", field.name(), value); } } } let mut v = V(String::new()); event.record(&mut v); - let meta = event.metadata(); pf_client_core::logring::note(format!( "{} {:5} {} {}", wallclock(), @@ -65,3 +103,70 @@ fn wallclock() -> String { ms % 1000 ) } + +#[cfg(test)] +mod tests { + use super::*; + + /// The gate is a prefix match on module-path boundaries, nothing looser. + #[test] + fn noisy_gate_matches_the_crate_and_its_modules_only() { + assert!(is_noisy_debug("cros_codecs")); + assert!(is_noisy_debug("cros_codecs::codec::h265::dpb")); + assert!(!is_noisy_debug("cros_codecs_probe")); + assert!(!is_noisy_debug("pf_bitstream::h265")); + assert!(!is_noisy_debug("pf_client_core::audio")); + } + + /// End to end through the bridge: a `log::debug!` from the vendored decoder's module path + /// must NOT reach the ring, its `warn!` must (under its real target, without the bridge's + /// bookkeeping fields), and a DEBUG event from our own audio module — the very line the gate + /// exists to protect — must land. + /// + /// The ring is process-global, so the assertions look for lines this test wrote (unique + /// markers) rather than at the ring's size, and the subscriber is installed only for the + /// duration of the test. + #[test] + fn bridged_decoder_debug_is_dropped_and_the_audio_line_survives() { + use tracing_subscriber::layer::SubscriberExt; + use tracing_subscriber::Layer; + let sub = tracing_subscriber::registry() + .with(RingLayer.with_filter(tracing_subscriber::filter::LevelFilter::DEBUG)); + // The bridge may already be installed by another test in this binary; either way the + // `log` max level has to admit DEBUG for the planted records to be dispatched at all. + let _ = tracing_log::LogTracer::builder() + .with_max_level(log::LevelFilter::Debug) + .init(); + log::set_max_level(log::LevelFilter::Debug); + let _guard = tracing::subscriber::set_default(sub); + + let marker = format!("ringgate-{}", std::process::id()); + log::debug!(target: "cros_codecs::codec::h265::dpb", "Retaining pic POC {marker}-dpb: true"); + log::warn!(target: "cros_codecs::codec::h265::parser", "{marker}-parser-warn"); + tracing::debug!(target: "pf_client_core::audio", buffer_ms = 15u32, "audio playback {marker}-audio"); + + let text = pf_client_core::logring::render("test"); + assert!( + !text.contains(&format!("{marker}-dpb")), + "decoder DPB DEBUG chatter must not reach the ring" + ); + let warn_line = text + .lines() + .find(|l| l.contains(&format!("{marker}-parser-warn"))) + .expect("decoder WARN must be kept"); + assert!( + warn_line.contains("cros_codecs::codec::h265::parser"), + "bridged events must carry their real target, not the `log` shim: {warn_line}" + ); + assert!( + !warn_line.contains("log.target="), + "bridge bookkeeping fields must be dropped: {warn_line}" + ); + let audio_line = text + .lines() + .find(|l| l.contains(&format!("{marker}-audio"))) + .expect("our own DEBUG audio line must survive"); + assert!(audio_line.contains("pf_client_core::audio")); + assert!(audio_line.contains("buffer_ms=15")); + } +} diff --git a/crates/pf-client-core/src/audio.rs b/crates/pf-client-core/src/audio.rs index 591c5bb1..057d564a 100644 --- a/crates/pf-client-core/src/audio.rs +++ b/crates/pf-client-core/src/audio.rs @@ -135,11 +135,6 @@ impl PlaybackFormat { fn quantum_frames(&self) -> u32 { ((self.rate_hz as u64 * self.frame_us as u64 / 1_000_000) as u32).max(1) } - - /// Frames (per channel) per millisecond — 48 at the protocol default, 96 at 96 kHz. - fn frames_per_ms(&self) -> usize { - (self.rate_hz / 1000).max(1) as usize - } } pub struct AudioPlayer { @@ -152,6 +147,9 @@ pub struct AudioPlayer { /// A/V sync hand-off with the PipeWire callback: it publishes the ring depth, the decode /// thread posts the depth the sync loop wants. See [`punktfunk_core::audio::AudioSyncCell`]. sync: Arc, + /// The callback's vitals, published as atomics and logged by the decode thread — the + /// callback runs on the graph's realtime loop and formats nothing (see [`crate::audio_vitals`]). + vitals: Arc, } impl AudioPlayer { @@ -171,10 +169,12 @@ impl AudioPlayer { let (quit_tx, quit_rx) = pipewire::channel::channel::(); let sync: Arc = Arc::default(); let sync_cb = sync.clone(); + let vitals: Arc = Arc::default(); + let vitals_cb = vitals.clone(); let thread = std::thread::Builder::new() .name("punktfunk-audio".into()) .spawn(move || { - if let Err(e) = pw_thread(pcm_rx, recycle_tx, quit_rx, fmt, sync_cb) { + if let Err(e) = pw_thread(pcm_rx, recycle_tx, quit_rx, fmt, sync_cb, vitals_cb) { tracing::warn!(error = %e, "audio playback thread ended"); } }) @@ -185,6 +185,7 @@ impl AudioPlayer { quit_tx, thread: Some(thread), sync, + vitals, }) } @@ -194,6 +195,11 @@ impl AudioPlayer { self.sync.clone() } + /// The callback's vitals — the decode thread logs them (see [`crate::audio_vitals`]). + pub fn vitals(&self) -> Arc { + self.vitals.clone() + } + /// A recycled chunk Vec from the pool, empty but with its capacity intact — fill it /// and hand it back through [`push`](Self::push). Allocates only when the pool is dry /// (startup, or after the PipeWire side dropped chunks). @@ -263,13 +269,13 @@ struct PlayerData { /// What `param_changed` last saw, so a graph RESUME (which re-announces the same format) is /// not logged as a format change — the host's virtual sink learned the same lesson. negotiated: Option<(u32, u32)>, - /// Diagnostics (WP0.3), logged ~every 10 s: the audio plane used to be entirely silent in a - /// client log, so a latency or dropout report had nothing to go on. - underruns: u64, - sheds: u64, - callbacks: u64, /// A/V sync hand-off with the decode thread (depth out, target in). sync: Arc, + /// Diagnostics (WP0.3): the audio plane used to be entirely silent in a client log, so a + /// latency or dropout report had nothing to go on. Published from the callback as atomics + /// and LOGGED by the decode thread — see [`crate::audio_vitals`] for why the callback itself + /// no longer formats a line. + vitals: Arc, } fn pw_thread( @@ -278,6 +284,7 @@ fn pw_thread( quit_rx: pipewire::channel::Receiver, fmt: PlaybackFormat, sync: Arc, + vitals: Arc, ) -> Result<()> { use pipewire as pw; use pw::{properties::properties, spa}; @@ -328,10 +335,25 @@ fn pw_thread( let stream = pw::stream::StreamBox::new(&core, "punktfunk-client", props).context("pw Stream")?; + // Pre-reserved so `extend` never reallocates on the realtime loop (the same shape as the + // Android callback's ring): the policy trims the ring back under its hard cap on every + // callback, so the most it can ever hold on entry is that cap plus everything the pump could + // have queued since the last callback — the whole 64-chunk channel of the plane's own frame. + // Sized from the resolved format on all three axes (rate, frame, channels), so a 96 kHz + // 7.1 lossless session reserves what it needs and a stereo Opus one does not over-reserve. + let ring_capacity = { + let per_ms = fmt.rate_hz as usize * channels / 1000; + let frame = punktfunk_core::audio::pcm::samples_per_frame( + fmt.rate_hz, + fmt.frame_us, + fmt.channels as u8, + ); + per_ms * TUNING.hard_cap_ms as usize + 64 * frame + }; let ud = PlayerData { rx: pcm_rx, recycle: recycle_tx, - ring: VecDeque::new(), + ring: VecDeque::with_capacity(ring_capacity), policy: { // Both at the RESOLVED format: `new_at_rate` denominates every depth/target/shed // figure — and the `buffer_ms`/`target_ms` this client reports — in the right @@ -350,10 +372,8 @@ fn pw_thread( channels, fmt, negotiated: None, - underruns: 0, - sheds: 0, - callbacks: 0, sync, + vitals, }; let _listener = stream @@ -400,6 +420,21 @@ fn pw_thread( ); } }) + // ⚠ REALTIME. With `RT_PROCESS` (the connect below) this closure runs on libpipewire's + // data loop — the thread the graph drives, scheduled realtime wherever the client's + // PipeWire has rtkit (`module-rt` in `client.conf`; SteamOS does). Everything in it has + // to be what a realtime callback is allowed to do: no allocation (the ring is + // pre-reserved, drained chunks keep their capacity), no locks that a lower-priority + // thread can hold (both channels are lock-free for one producer/one consumer), and no + // logging — the vitals go out as atomics and the decode thread prints them. + // + // Before this the stream connected WITHOUT `RT_PROCESS`, so `process` ran on our own + // main-loop thread at ordinary priority: PipeWire's data loop signalled it, and if this + // thread was not scheduled within the cycle the graph rendered silence for our node and + // moved on — an underrun neither our counters nor the ring saw, because by the time we + // ran the ring was full and the callback drained normally. On a Steam Deck decoding + // 1440p120 alongside, that is a real and invisible source of clicks. The host's own + // PipeWire stream nodes have run `RT_PROCESS` since they were written. .process(|stream, ud| { let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { let Some(mut buffer) = stream.dequeue_buffer() else { @@ -443,19 +478,11 @@ fn pw_thread( let want = want_frames * ud.channels; // Once per stream, in the shape of the host's per-capture-open quantum log: // whether the graph's request or the buffer ceiling is sizing our writes is - // exactly what an on-glass latency report needs to say. - if ud.callbacks == 0 { - tracing::info!( - requested_frames = requested, - capacity_frames = max_frames, - write_frames = want_frames, - // From the session's rate, not from 48: a 96 kHz quantum divided by 48 - // reads as twice the latency it is, in the one line an on-glass latency - // report is triaged from. - write_ms = want_frames / ud.fmt.frames_per_ms(), - rate_hz = ud.fmt.rate_hz, - "audio playback quantum" - ); + // exactly what an on-glass latency report needs to say. Published, not logged + // (realtime, see above); the decode thread prints it the first time it sees it. + if !ud.vitals.quantum_known() { + ud.vitals + .note_quantum(requested as u32, max_frames as u32, want_frames as u32); } // A/V sync: take whatever depth the decode thread's sync loop last asked for, and @@ -470,7 +497,6 @@ fn pw_thread( // and a hard cap as the backstop. let step = ud.policy.step(ud.ring.len(), want); if step.drop_front > 0 { - ud.sheds += 1; punktfunk_core::audio::crossfade_drop( &mut ud.ring, step.drop_front, @@ -499,23 +525,13 @@ fn pw_thread( // No-op while un-primed (the policy ignores it), so a deliberate priming silence // is never miscounted as an underrun. ud.policy.note_read(ran_short); - ud.underruns += u64::from(ran_short); - ud.callbacks += 1; - // ~10 s at a 5 ms quantum; the exact cadence does not matter, only that the - // plane stops being invisible. - if ud.callbacks % 2_000 == 0 { - tracing::debug!( - buffer_ms = ud.policy.avg_depth_ms(), - target_ms = ud.policy.target_ms(), - underruns = ud.underruns, - drift_sheds = ud.sheds, - // Concealment must be visible next to the underruns it prevented: a - // healthy `underruns` bought with a climbing `plc_ms` is a link in - // trouble, not a link that is fine. - plc_ms = ud.sync.plc_ms(), - "audio playback" - ); - } + // The 10 s `audio playback` line is printed by the decode thread from these. + ud.vitals.note_callback( + ran_short, + step.drop_front > 0, + ud.policy.avg_depth_ms(), + ud.policy.target_ms(), + ); let chunk = data.chunk_mut(); *chunk.offset_mut() = 0; *chunk.stride_mut() = stride as _; @@ -552,11 +568,16 @@ fn pw_thread( .into_inner(); let mut params = [Pod::from_bytes(&values).context("pod from bytes")?]; + // `RT_PROCESS`: run `process` on libpipewire's realtime data loop rather than on this + // main-loop thread — see the callback's own note for what that buys and what it forbids. + // Same flag the host's virtual mic and stream sink connect with. stream .connect( spa::utils::Direction::Output, None, - pw::stream::StreamFlags::AUTOCONNECT | pw::stream::StreamFlags::MAP_BUFFERS, + pw::stream::StreamFlags::AUTOCONNECT + | pw::stream::StreamFlags::MAP_BUFFERS + | pw::stream::StreamFlags::RT_PROCESS, &mut params, ) .context("pw stream connect")?; @@ -861,14 +882,10 @@ mod tests { } } - /// The one-shot quantum log divides by this, and reading it off a constant is how a 96 kHz - /// session reports twice the latency it actually has in the line a report is triaged from. + /// A nonsense format must not ask the graph for a zero quantum. #[test] - fn frames_per_ms_follows_the_negotiated_rate() { - assert_eq!(fmt(48_000, 5_000).frames_per_ms(), 48); - assert_eq!(fmt(96_000, 2_000).frames_per_ms(), 96); - // A nonsense rate must not divide by zero in a log line. - assert_eq!(fmt(0, 5_000).frames_per_ms(), 1); + fn a_nonsense_format_still_asks_for_a_frame() { assert_eq!(fmt(0, 0).quantum_frames(), 1); + assert_eq!(fmt(0, 5_000).quantum_frames(), 1); } } diff --git a/crates/pf-client-core/src/audio_vitals.rs b/crates/pf-client-core/src/audio_vitals.rs new file mode 100644 index 00000000..674bfc3f --- /dev/null +++ b/crates/pf-client-core/src/audio_vitals.rs @@ -0,0 +1,120 @@ +//! Playback vitals: what the device callback knows and the decode thread logs. +//! +//! The PipeWire callback (`audio.rs`) runs on the graph's realtime data loop now +//! (`RT_PROCESS`), where formatting a log line — a `String` allocation, the subscriber's mutex, +//! a write to stderr and the log ring — is exactly the class of thing a realtime thread must +//! not do: at best it is a priority inversion against whatever holds the lock, at worst it is a +//! missed graph cycle, which is a click. So the callback publishes numbers into these atomics +//! and the decode thread, an ordinary thread that already wakes every frame, prints them at the +//! old cadence with the old field names (`audio playback buffer_ms= target_ms= underruns= +//! drift_sheds= plc_ms=`), so a field-log grep keeps working. The WASAPI twin runs its render +//! loop on a plain thread and could log in place, but publishes here too: one logging site, +//! one line shape, on both platforms. + +use std::sync::atomic::{AtomicU32, AtomicU64, Ordering}; + +/// Counters and gauges the device callback publishes. All `Relaxed`: every field is a +/// self-contained reading, nothing here orders anything else. +#[derive(Debug, Default)] +pub struct PlaybackVitals { + /// Device callbacks served (primed or not) — proof of life for the pull side. + pub callbacks: AtomicU64, + /// Callbacks the ring could not fill (a genuine underrun) — see `JitterPolicy::note_read`. + pub underruns: AtomicU64, + /// Drops the policy asked for: drift sheds and hard trims together. + pub sheds: AtomicU64, + /// The policy's smoothed ring depth, ms — what drift correction reacts to. + pub buffer_ms: AtomicU32, + /// The policy's LIVE target depth, ms (grows under underrun pressure, follows A/V sync). + pub target_ms: AtomicU32, + /// The device quantum as first seen: frames the graph/engine asked for per callback, the + /// mapped buffer's capacity, and what we actually write. `write_frames == 0` = not seen yet. + pub requested_frames: AtomicU32, + pub capacity_frames: AtomicU32, + pub write_frames: AtomicU32, +} + +impl PlaybackVitals { + /// Callback side: one callback done. `ran_short` = it could not be filled from the ring; + /// `shed` = the policy dropped something this callback. + pub fn note_callback(&self, ran_short: bool, shed: bool, buffer_ms: u32, target_ms: u32) { + self.callbacks.fetch_add(1, Ordering::Relaxed); + if ran_short { + self.underruns.fetch_add(1, Ordering::Relaxed); + } + if shed { + self.sheds.fetch_add(1, Ordering::Relaxed); + } + self.buffer_ms.store(buffer_ms, Ordering::Relaxed); + self.target_ms.store(target_ms, Ordering::Relaxed); + } + + /// Callback side: the quantum, published once (the first callback that has one). + pub fn note_quantum(&self, requested: u32, capacity: u32, write: u32) { + self.requested_frames.store(requested, Ordering::Relaxed); + self.capacity_frames.store(capacity, Ordering::Relaxed); + self.write_frames.store(write, Ordering::Relaxed); + } + + /// Whether [`note_quantum`](Self::note_quantum) has been called. + pub fn quantum_known(&self) -> bool { + self.write_frames.load(Ordering::Relaxed) > 0 + } + + /// A consistent-enough snapshot for a log line (each field is read once; a callback landing + /// between two reads skews a counter by one, which a 10 s log line does not care about). + pub fn snapshot(&self) -> Snapshot { + Snapshot { + callbacks: self.callbacks.load(Ordering::Relaxed), + underruns: self.underruns.load(Ordering::Relaxed), + sheds: self.sheds.load(Ordering::Relaxed), + buffer_ms: self.buffer_ms.load(Ordering::Relaxed), + target_ms: self.target_ms.load(Ordering::Relaxed), + requested_frames: self.requested_frames.load(Ordering::Relaxed), + capacity_frames: self.capacity_frames.load(Ordering::Relaxed), + write_frames: self.write_frames.load(Ordering::Relaxed), + } + } +} + +/// One reading of [`PlaybackVitals`]. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct Snapshot { + pub callbacks: u64, + pub underruns: u64, + pub sheds: u64, + pub buffer_ms: u32, + pub target_ms: u32, + pub requested_frames: u32, + pub capacity_frames: u32, + pub write_frames: u32, +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn counters_accumulate_and_gauges_overwrite() { + let v = PlaybackVitals::default(); + assert!(!v.quantum_known()); + v.note_callback(false, false, 15, 15); + v.note_callback(true, true, 9, 25); + v.note_callback(true, false, 12, 25); + v.note_quantum(240, 8192, 240); + let s = v.snapshot(); + assert_eq!(s.callbacks, 3); + assert_eq!(s.underruns, 2); + assert_eq!(s.sheds, 1); + assert_eq!( + (s.buffer_ms, s.target_ms), + (12, 25), + "gauges hold the latest reading" + ); + assert_eq!( + (s.requested_frames, s.capacity_frames, s.write_frames), + (240, 8192, 240) + ); + assert!(v.quantum_known()); + } +} diff --git a/crates/pf-client-core/src/audio_wasapi.rs b/crates/pf-client-core/src/audio_wasapi.rs index b545cb0e..6d12395e 100644 --- a/crates/pf-client-core/src/audio_wasapi.rs +++ b/crates/pf-client-core/src/audio_wasapi.rs @@ -283,6 +283,10 @@ pub struct AudioPlayer { /// A/V sync hand-off with the render thread: it publishes the ring depth, the decode thread /// posts the depth the sync loop wants. See [`punktfunk_core::audio::AudioSyncCell`]. sync: Arc, + /// The render loop's vitals, logged by the decode thread — the same surface the PipeWire + /// twin exposes, so `session.rs` prints one line shape on both platforms + /// (see [`crate::audio_vitals`]). + vitals: Arc, } impl AudioPlayer { @@ -306,10 +310,14 @@ impl AudioPlayer { let stop_t = stop.clone(); let sync: Arc = Arc::default(); let sync_t = sync.clone(); + let vitals: Arc = Arc::default(); + let vitals_t = vitals.clone(); let thread = std::thread::Builder::new() .name("punktfunk-audio".into()) .spawn(move || { - if let Err(e) = render_thread(pcm_rx, recycle_tx, stop_t, ready_tx, fmt, sync_t) { + if let Err(e) = + render_thread(pcm_rx, recycle_tx, stop_t, ready_tx, fmt, sync_t, vitals_t) + { tracing::warn!(error = %format!("{e:#}"), "audio playback thread ended"); } }) @@ -333,6 +341,7 @@ impl AudioPlayer { stop, thread: Some(thread), sync, + vitals, }) } Ok(Err(e)) => Err(e), @@ -355,6 +364,11 @@ impl AudioPlayer { self.sync.clone() } + /// The render loop's vitals — the decode thread logs them (see [`crate::audio_vitals`]). + pub fn vitals(&self) -> Arc { + self.vitals.clone() + } + /// Queue one interleaved f32 chunk (in the session's channel layout). Drops the chunk if the /// WASAPI side is wedged (the renderer conceals the gap; never block the session pump). pub fn push(&self, pcm: Vec) { @@ -380,6 +394,7 @@ fn render_thread( ready: SyncSender>>, fmt: PlaybackFormat, sync: Arc, + vitals: Arc, ) -> Result<()> { if let Err(e) = wasapi::initialize_mta() .ok() @@ -484,7 +499,6 @@ fn render_thread( punktfunk_core::audio::JitterPolicy::new_at_rate(TUNING, channels, fmt.rate_hz); policy.set_frame_us(fmt.frame_us); let mut out = Vec::new(); // per-quantum scratch, reused across iterations - let (mut underruns, mut sheds, mut callbacks) = (0u64, 0u64, 0u64); while !stop.load(Ordering::Relaxed) { if h_event.wait_for_event(100).is_err() { @@ -504,6 +518,15 @@ fn render_thread( continue; } let want = avail_frames * channels as usize; + // Once per stream: the engine's period as we first see it — same field meanings as + // the PipeWire twin's line, printed by the decode thread. + if !vitals.quantum_known() { + vitals.note_quantum( + avail_frames as u32, + avail_frames as u32, + avail_frames as u32, + ); + } // A/V sync: same contract as the PipeWire ring — take the decode thread's request, // publish where the ring actually is. The policy clamps the request against its own @@ -513,7 +536,6 @@ fn render_thread( let step = policy.step(ring.len(), want); if step.drop_front > 0 { - sheds += 1; punktfunk_core::audio::crossfade_drop(&mut ring, step.drop_front, step.crossfade); } @@ -533,21 +555,13 @@ fn render_thread( // No-op while un-primed (the policy ignores it), so a deliberate priming silence is // never miscounted as an underrun. policy.note_read(ran_short); - underruns += u64::from(ran_short); - callbacks += 1; - if callbacks % 1_000 == 0 { - tracing::debug!( - buffer_ms = policy.avg_depth_ms(), - target_ms = policy.target_ms(), - underruns, - drift_sheds = sheds, - // Concealment must be visible next to the underruns it prevented: a healthy - // `underruns` bought with a climbing `plc_ms` is a link in trouble, not a - // link that is fine. - plc_ms = sync.plc_ms(), - "audio playback" - ); - } + // The 10 s `audio playback` line is printed by the decode thread from these. + vitals.note_callback( + ran_short, + step.drop_front > 0, + policy.avg_depth_ms(), + policy.target_ms(), + ); render_client .write_to_device(avail_frames, &out, None) .context("write_to_device")?; diff --git a/crates/pf-client-core/src/lib.rs b/crates/pf-client-core/src/lib.rs index 19dd000e..234e8235 100644 --- a/crates/pf-client-core/src/lib.rs +++ b/crates/pf-client-core/src/lib.rs @@ -26,6 +26,10 @@ pub mod audio; #[cfg(windows)] #[path = "audio_wasapi.rs"] pub mod audio; +// The playback vitals both twins publish from their device callback and the decode thread logs +// — atomics only, because the PipeWire callback runs on the graph's realtime loop. +#[cfg(any(target_os = "linux", windows))] +pub mod audio_vitals; #[cfg(any(target_os = "linux", windows))] pub mod discovery; #[cfg(any(target_os = "linux", windows))] diff --git a/crates/pf-client-core/src/session.rs b/crates/pf-client-core/src/session.rs index cd86bacc..44019d01 100644 --- a/crates/pf-client-core/src/session.rs +++ b/crates/pf-client-core/src/session.rs @@ -2099,6 +2099,9 @@ fn spawn_audio( Ok("1") | Ok("true") ); let sync_cell = player.sync_cell(); + // The device callback's counters. Logged from THIS thread, on wall clock — the PipeWire + // callback runs on the graph's realtime loop and formats nothing (`crate::audio_vitals`). + let vitals = player.vitals(); let video_e2e = connector.video_e2e_shared(); let av_offset_out = connector.audio_av_offset_shared(); let buffer_ms_out = connector.audio_buffer_ms_shared(); @@ -2131,6 +2134,20 @@ fn spawn_audio( std::thread::Builder::new() .name("punktfunk-audio-rx".into()) .spawn(move || { + // Best-effort priority for the decode leg. This thread's lateness is absorbed by + // the ring (target 15 ms and up), so it is not the callback's problem in kind — but + // on a Steam Deck the same four cores decode 1440p120 and present it, and a decode + // thread descheduled past the ring depth is a drought the callback then has to + // conceal. A plain `setpriority` is honoured wherever RLIMIT_NICE allows (rtkit is + // the sanctioned unprivileged path and a follow-up); where it is refused this is a + // no-op, which is exactly what it was before. + #[cfg(target_os = "linux")] + { + // SAFETY: three by-value integers, no pointers; `PRIO_PROCESS` with `who == 0` + // targets the calling thread on Linux and only adjusts its nice value. + let rc = unsafe { libc::setpriority(libc::PRIO_PROCESS, 0, -10) }; + tracing::debug!(raised = rc == 0, "audio decode thread priority"); + } let mut pcm = vec![0f32; scratch]; let mut gaps = punktfunk_core::audio::AudioGapTracker::new(); // Interleaved samples in the last decoded frame — the unit concealment is produced in. @@ -2152,7 +2169,44 @@ fn spawn_audio( frame_us, ); let mut last_packet = std::time::Instant::now(); + // The playback vitals line, ~every 10 s on wall clock (it used to be every 2 000 + // device callbacks from inside the callback — same fields, same name, so a field-log + // grep keeps working), plus the one-shot quantum line the first time the callback + // has published one. + let mut last_vitals = std::time::Instant::now(); + let mut quantum_logged = false; while !stop.load(Ordering::SeqCst) { + if !quantum_logged && vitals.quantum_known() { + quantum_logged = true; + let v = vitals.snapshot(); + tracing::info!( + requested_frames = v.requested_frames, + capacity_frames = v.capacity_frames, + write_frames = v.write_frames, + // From the session's rate, not from 48: a 96 kHz quantum divided by 48 + // reads as twice the latency it is, in the one line an on-glass latency + // report is triaged from. + write_ms = v.write_frames / (rate_hz / 1000).max(1), + rate_hz, + "audio playback quantum" + ); + } + if last_vitals.elapsed() >= Duration::from_secs(10) { + last_vitals = std::time::Instant::now(); + let v = vitals.snapshot(); + tracing::debug!( + buffer_ms = v.buffer_ms, + target_ms = v.target_ms, + underruns = v.underruns, + drift_sheds = v.sheds, + callbacks = v.callbacks, + // Concealment must be visible next to the underruns it prevented: a + // healthy `underruns` bought with a climbing `plc_ms` is a link in + // trouble, not a link that is fine. + plc_ms = sync_cell.plc_ms(), + "audio playback" + ); + } // Wait at most one frame WHILE there is a stream to protect: the drought decision // has to be made on the wire's schedule, not whenever the next packet happens to // turn up. Before anything has decoded there is no state to conceal from and diff --git a/crates/punktfunk-core/src/audio/pcm.rs b/crates/punktfunk-core/src/audio/pcm.rs index 5c8675ac..f4f6427e 100644 --- a/crates/punktfunk-core/src/audio/pcm.rs +++ b/crates/punktfunk-core/src/audio/pcm.rs @@ -350,7 +350,13 @@ impl PcmConceal { /// Apply a raised-cosine fade to the final `n` samples in place, so a spliced or repeated frame /// meets what follows it without a step. -fn raised_cosine_tail(buf: &mut [f32], n: usize) { +/// +/// `pub` for the host's capture-hole infill (`punktfunk-host::native::audio`), which fades the +/// audio into a hole rather than stepping to digital zero — same curve, same reason. `n` counts +/// interleaved samples: a fade meant to span whole frames of a multi-channel signal passes +/// `frames × channels`, and adjacent channels of one frame then sit one step apart on the curve +/// (a 1/n gain difference — nothing, at any fade longer than a few frames). +pub fn raised_cosine_tail(buf: &mut [f32], n: usize) { let n = n.min(buf.len()); if n == 0 { return; @@ -362,6 +368,20 @@ fn raised_cosine_tail(buf: &mut [f32], n: usize) { } } +/// The mirror of [`raised_cosine_tail`]: fade the FIRST `n` samples in from zero, so audio that +/// resumes mid-waveform after a hole (or after silence) starts without a step. `n` counts +/// interleaved samples, like the tail. +pub fn raised_cosine_head(buf: &mut [f32], n: usize) { + let n = n.min(buf.len()); + if n == 0 { + return; + } + for (i, s) in buf[..n].iter_mut().enumerate() { + let t = (i as f32 + 0.5) / n as f32; + *s *= 0.5 * (1.0 - (std::f32::consts::PI * t).cos()); + } +} + #[cfg(test)] mod tests { use super::*; @@ -776,4 +796,27 @@ mod tests { out.last().unwrap() ); } + + /// The head fade is the tail fade run backwards: it starts at silence, lands at full level, + /// touches nothing past `n`, and the two together are unity — so a tail-faded frame followed + /// by a head-faded one is a crossfade, not a dip. + #[test] + fn the_head_fade_mirrors_the_tail_fade() { + let mut head = vec![1.0f32; 64]; + raised_cosine_head(&mut head, 32); + assert!(head[0] < 0.01, "starts at silence: {}", head[0]); + assert!(head[31] > 0.99, "lands at full level: {}", head[31]); + assert_eq!(head[32], 1.0, "untouched past n"); + let mut tail = vec![1.0f32; 32]; + raised_cosine_tail(&mut tail, 32); + for i in 0..32 { + let sum = head[i] + tail[i]; + assert!((sum - 1.0).abs() < 1e-5, "unity at {i}: {sum}"); + } + // Degenerate lengths must not panic or touch anything. + let mut short = vec![1.0f32; 4]; + raised_cosine_head(&mut short, 0); + raised_cosine_head(&mut short, 100); + assert!(short[3] > 0.9); + } } diff --git a/crates/punktfunk-host/src/audio/capture_policy.rs b/crates/punktfunk-host/src/audio/capture_policy.rs index 5209e385..b4e7571b 100644 --- a/crates/punktfunk-host/src/audio/capture_policy.rs +++ b/crates/punktfunk-host/src/audio/capture_policy.rs @@ -103,6 +103,13 @@ pub(crate) const STATS_EVERY: Duration = Duration::from_secs(30); /// 2.7 ms) from scoring ordinary scheduling noise as a hole. const GAP_FLOOR: Duration = Duration::from_millis(10); +/// Upper edges (exclusive, ms) of the gap-size histogram's first buckets; the last bucket is +/// everything at or above the final edge. Chosen around what the client can hide: libopus PLC +/// covers ~50 ms of a seq gap, the drought fuse twice a de-prime window (80–120 ms by preset), so +/// `<20` is a hiccup nobody hears, `<50` is concealable, `<100` is on the edge, and `≥100` is a +/// dropout however good the concealment. +pub(crate) const GAP_HIST_EDGES_MS: [u64; 3] = [20, 50, 100]; + /// One reporting window's worth of capture vitals. /// /// The point is to make three states that used to look identical in a log tell themselves apart: a @@ -134,6 +141,16 @@ pub(crate) struct CaptureStats { pub(crate) gaps: u64, /// The largest of those, µs. Reported in ms; kept in µs so a sub-ms threshold is expressible. pub(crate) max_gap_us: u64, + /// The SHAPE of those gaps: how many fell in each of [`GAP_HIST_EDGES_MS`]'s buckets, and + /// the audio they cost in total. `gaps` + `max_gap_ms` say "sixty holes, the worst 146 ms" + /// and leave a reader unable to tell sixty 30 ms stalls (a periodic scheduler on the box) + /// from fifty-nine 12 ms hiccups and one outage — a different fault with a different fix. + /// The 2026-08-17 field log had exactly that ambiguity across 53 windows. `missing_us` also + /// closes the arithmetic between `gaps` and `delivered_pct`: when the sum of the gaps + /// accounts for the shortfall the loss is all in counted holes, and when it does not, the + /// remainder is sub-threshold losses the counter cannot see (see [`GAP_FLOOR`]). + pub(crate) gap_hist: [u64; GAP_HIST_EDGES_MS.len() + 1], + pub(crate) missing_us: u64, /// Callbacks that ran but carried nothing — no buffer to dequeue, no `datas`, no mapped /// 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. @@ -187,22 +204,51 @@ impl CaptureStats { pub(crate) fn observe_callback(&mut self, since_last: Option, quantum: Duration) { let Some(delta) = since_last else { return }; if delta > (quantum * 2).max(GAP_FLOOR) { - self.gaps += 1; // The MISSING audio, not the callback delta: one quantum of that delta is the buffer // we were legitimately handed. Reporting the delta would inflate every gap by the // quantum and — worse — mean something different from the Windows feed, which sizes // its holes from the device position and so reports missing audio by construction. - self.max_gap_us = self - .max_gap_us - .max(delta.saturating_sub(quantum).as_micros() as u64); + self.observe_gap(delta.saturating_sub(quantum)); } } + /// Score one hole of `missing` audio — the shared accounting behind both feeds: the Linux + /// callback cadence above, and the Windows discontinuity flag, which measures the hole from + /// the device position and calls this directly. + pub(crate) fn observe_gap(&mut self, missing: Duration) { + self.gaps += 1; + let us = missing.as_micros() as u64; + self.max_gap_us = self.max_gap_us.max(us); + self.missing_us = self.missing_us.saturating_add(us); + let ms = us / 1_000; + let bucket = GAP_HIST_EDGES_MS + .iter() + .position(|&edge| ms < edge) + .unwrap_or(GAP_HIST_EDGES_MS.len()); + self.gap_hist[bucket] += 1; + } + /// The window's worst gap in whole ms — the unit the log line and the field reports speak. pub(crate) fn max_gap_ms(&self) -> u64 { self.max_gap_us / 1_000 } + /// Audio the window's counted gaps cost in total, whole ms. + pub(crate) fn missing_ms(&self) -> u64 { + self.missing_us / 1_000 + } + + /// The gap-size histogram as one log field, `a/b/c/d` = counts under 20 / 50 / 100 ms and + /// at-or-above 100 ms ([`GAP_HIST_EDGES_MS`]). One field rather than four so the line stays + /// greppable and the buckets read as a shape. + pub(crate) fn gap_hist(&self) -> String { + self.gap_hist + .iter() + .map(u64::to_string) + .collect::>() + .join("/") + } + /// 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 @@ -400,6 +446,9 @@ pub(crate) struct InfillPolicy { filled: Duration, /// Latched once a hole outlives the budget and the wire falls silent. broke: bool, + /// The largest capture chunk seen, as a duration — see [`Self::after`]. Zero until the + /// caller reports one, which leaves every figure here at its frame-denominated default. + quantum: Duration, } impl InfillPolicy { @@ -419,15 +468,31 @@ impl InfillPolicy { frame: Duration::from_micros(frame_us.max(1) as u64), filled: Duration::ZERO, broke: false, + quantum: Duration::ZERO, } } + /// Tell the policy how big the capture chunks actually are. Called with each chunk's + /// duration; only the largest is kept, so a graph that hands us 21 ms buffers lifts + /// [`after`](Self::after) once and a single short buffer never lowers it. + pub(crate) fn note_quantum(&mut self, chunk: Duration) { + self.quantum = self.quantum.max(chunk); + } + /// 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. Two frames of THIS session — the client's ring is sized in /// its own frames, so that is what "before it notices" is measured in. + /// + /// …at a quantum of a frame or two. A graph that clamps us to a 21 ms buffer (a VM's + /// `min-quantum = 1024`, see `audio::linux`) legitimately goes 21 ms between chunks, and + /// two frames after the last one is the MIDDLE of a normal cycle: a chunk merely a couple + /// of milliseconds late would read as a hole, be covered, and leave one frame of surplus + /// behind for good. So the threshold is never less than one chunk plus one frame — the + /// jitter budget a chunk is allowed before it counts as missing. Identical to two frames + /// wherever the quantum is a frame, which is every host that asked for one and got it. pub(crate) fn after(&self) -> Duration { - self.frame * 2 + (self.frame * 2).max(self.quantum + self.frame) } /// Decide the slot that is due now. Call EXACTLY once per due frame — it consumes budget. @@ -451,6 +516,19 @@ impl InfillPolicy { self.filled >= INFILL_MAX } + /// Silence sent for the hole currently open. Read after [`decide`](Self::decide) says + /// `Silence`: equal to one frame means that frame is the FIRST of the hole — the one that + /// carries the fade out of the last real audio (see the caller) — and anything larger is + /// deep enough into the hole to be plain silence. + pub(crate) fn covered(&self) -> Duration { + self.filled + } + + /// One frame of this session, the unit [`covered`](Self::covered) counts in. + pub(crate) fn frame(&self) -> Duration { + self.frame + } + /// A real chunk arrived. Returns whether the hole it closed BROKE continuity — the wire went /// silent across it, so the redundancy predecessor and any partial frame straddling the hole /// both describe audio from before a discontinuity, and neither may be spliced onto what @@ -598,6 +676,44 @@ mod tests { ); } + /// The histogram is what tells sixty 30 ms stalls from fifty-nine 12 ms hiccups and one + /// outage — `gaps=60` either way, and the 2026-08-17 field log had exactly that ambiguity. + /// Buckets sit on the client-concealment edges, the total is the audio the holes cost, and + /// both feeds (Linux cadence, Windows discontinuity) land in the same accounting. + #[test] + fn gap_histogram_gives_the_shape_and_the_cost() { + let mut stalls = CaptureStats::default(); + for _ in 0..60 { + // A 30 ms hole arrives as a 35 ms callback delta at the 5 ms quantum. + stalls.observe_callback(Some(Q + Duration::from_millis(30)), Q); + } + assert_eq!(stalls.gaps, 60); + assert_eq!(stalls.gap_hist(), "0/60/0/0", "all in the <50 ms bucket"); + assert_eq!(stalls.missing_ms(), 60 * 30); + + let mut mixed = CaptureStats::default(); + for _ in 0..59 { + mixed.observe_callback(Some(Q + Duration::from_millis(12)), Q); + } + // The Windows feed measures the hole from the device position and calls this directly. + mixed.observe_gap(Duration::from_millis(1_092)); + assert_eq!(mixed.gaps, 60, "same count as the stalls above…"); + assert_eq!( + mixed.gap_hist(), + "59/0/0/1", + "…and nothing like the same shape" + ); + assert_eq!(mixed.missing_ms(), 59 * 12 + 1_092); + assert_eq!(mixed.max_gap_ms(), 1_092); + + // Edges are exclusive on the upper side: exactly 20 ms is the second bucket. + let mut edge = CaptureStats::default(); + edge.observe_gap(Duration::from_millis(20)); + edge.observe_gap(Duration::from_micros(19_999)); + edge.observe_gap(Duration::from_millis(100)); + assert_eq!(edge.gap_hist(), "1/1/0/1"); + } + /// A stream delivering exactly what it negotiated is never a gap — including the clamped /// 21.3 ms quantum a VM's `default.clock.min-quantum` forces, which would otherwise score a /// gap on every single callback and bury the real ones. @@ -669,6 +785,43 @@ mod tests { assert!(p.exhausted(), "…and then stop asking"); } + /// A clamped quantum lifts the hole threshold to one chunk plus one frame, so a 21 ms VM + /// buffer arriving a couple of milliseconds late is jitter, not a hole to be covered — and + /// at a frame-sized quantum nothing moves. `covered()`/`frame()` let the caller pick out the + /// first frame of a hole, which is the one that carries the fade. + #[test] + fn infill_threshold_follows_a_clamped_quantum() { + let frame = Duration::from_millis(FRAME_MS as u64); + let mut tight = InfillPolicy::new(OPUS_FRAME_US); + tight.note_quantum(Duration::from_micros(2_667)); // 128 frames at 48 kHz + tight.note_quantum(frame); // 240 frames — what we ask for + assert_eq!( + tight.after(), + frame * 2, + "a frame-sized quantum changes nothing" + ); + + let mut vm = InfillPolicy::new(OPUS_FRAME_US); + vm.note_quantum(Duration::from_micros(21_333)); // 1024 frames at 48 kHz + vm.note_quantum(Duration::from_micros(2_667)); // one short buffer never lowers it + assert_eq!( + vm.after(), + Duration::from_micros(26_333), + "one chunk plus one frame" + ); + assert_eq!( + vm.decide(Duration::from_millis(23)), + Infill::Wait, + "a chunk 2 ms late on a 21 ms quantum is not a hole" + ); + assert_eq!(vm.decide(Duration::from_millis(27)), Infill::Silence); + assert_eq!(vm.covered(), vm.frame(), "the first frame of the hole"); + assert_eq!(vm.decide(Duration::from_millis(32)), Infill::Silence); + assert_eq!(vm.covered(), vm.frame() * 2, "…and no longer the first"); + vm.chunk_arrived(); + assert_eq!(vm.covered(), Duration::ZERO, "a chunk closes the hole"); + } + /// **The gate on re-deriving these policies from `audio_frame_us`.** A 5 ms Opus session must /// behave exactly as it did when all three figures were written against `FRAME_MS` — same /// 10 ms infill threshold, same 100 frames of cover, same 5 ms slip threshold — or a change diff --git a/crates/punktfunk-host/src/audio/linux/mod.rs b/crates/punktfunk-host/src/audio/linux/mod.rs index 5f921c64..2717f736 100644 --- a/crates/punktfunk-host/src/audio/linux/mod.rs +++ b/crates/punktfunk-host/src/audio/linux/mod.rs @@ -1199,6 +1199,11 @@ fn pw_thread( // percentage and mean entirely different things. gaps = ud.stats.gaps, max_gap_ms = ud.stats.max_gap_ms(), + // …and their SHAPE: bucket counts under 20/50/100 ms and ≥ 100 ms, + // plus the audio they cost. Sixty 30 ms stalls and fifty-nine + // hiccups plus one outage share `gaps=60`; they do not share this. + gap_hist = %ud.stats.gap_hist(), + missing_ms = ud.stats.missing_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 diff --git a/crates/punktfunk-host/src/audio/windows/wasapi_cap.rs b/crates/punktfunk-host/src/audio/windows/wasapi_cap.rs index 998271a9..42cfde44 100644 --- a/crates/punktfunk-host/src/audio/windows/wasapi_cap.rs +++ b/crates/punktfunk-host/src/audio/windows/wasapi_cap.rs @@ -249,6 +249,14 @@ fn capture_thread( let _ = ready.send(Err(e)); return Ok(()); } + // This is the thread that produces the chunks the paced sender consumes — the one that + // has to wake on the engine's event every 10 ms and read the loopback packets before the + // engine's buffer wraps. The SENDER has carried `THREAD_PRIORITY_HIGHEST` + MMCSS since the + // data-plane QoS work; this reader ran at normal priority beside it, so a CPU-saturating + // game could hold it off long enough for `DATA_DISCONTINUITY` to fire — a hole the sender + // then dutifully covered. Same boost, same helper: it registers the MMCSS task and raises + // the class, and is a no-op where either is refused. + pf_frame::thread_qos::boost_thread_priority(true); // Self-heal for the capturer's whole life: each `capture_once` is one endpoint open + inner // capture loop; it returns to reopen (default-device change) or errors (device invalidated, // engine restart). The FIRST open gets [`FIRST_OPEN_ATTEMPTS`] tries (session-start endpoint @@ -777,11 +785,10 @@ fn capture_once( stats.missed_dequeues += 1; } else { if info.flags.data_discontinuity && flowing { - stats.gaps += 1; let lost = info.index.saturating_sub(next_index); - stats.max_gap_us = stats - .max_gap_us - .max(lost.saturating_mul(1_000_000) / open_hz.max(1) as u64); + stats.observe_gap(Duration::from_micros( + lost.saturating_mul(1_000_000) / open_hz.max(1) as u64, + )); } next_index = info.index.saturating_add(frames); last_packet = Some(now); @@ -846,6 +853,10 @@ fn capture_once( // resumes is not among them — see [`LOOPBACK_IDLE_AFTER`]. gaps = stats.gaps, max_gap_ms = stats.max_gap_ms(), + // Their shape (bucket counts under 20/50/100 ms and ≥ 100 ms) and their total + // cost — same fields, same meaning as the Linux line. + gap_hist = %stats.gap_hist(), + missing_ms = stats.missing_ms(), missed_dequeues = stats.missed_dequeues, dropped_chunks = 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 d7dfee27..2d7d20be 100644 --- a/crates/punktfunk-host/src/native/audio.rs +++ b/crates/punktfunk-host/src/native/audio.rs @@ -195,6 +195,10 @@ pub(super) fn audio_thread( /// re-anchors. Chasing an old schedule after a stall would send a burst — the exact thing /// pacing exists to prevent — so past this point the debt is forgiven, not repaid. const PACE_REANCHOR: std::time::Duration = std::time::Duration::from_millis(100); + /// How much audio is faded at each edge of a capture hole — out into the hole, in out of + /// it — in µs. One millisecond: long enough to be a slope rather than an edge, far too short + /// to read as a swell. See `last_real` / `resume_fade` below. + const EDGE_FADE_US: u64 = 1_000; let want = punktfunk_core::audio::normalize_channels(channels); // The three session values that used to be compile-time constants. Every one of them is now // a property of the negotiated plane, and everything downstream — the pacer, the sample @@ -352,6 +356,26 @@ pub(super) fn audio_thread( // Nothing may be synthesized before the first real frame: there is no continuity to protect // yet, and the wire clock has no anchor to continue from. let mut sent_any = false; + // The two EDGES of a hole. The infill used to be digital zero from its first sample, and + // zero is not the absence of sound — it is a step from whatever level the last real sample + // sat at down to nothing, which the listener hears as a click at the front of every hole + // (and which the codec faithfully encodes); the audio then resumes mid-waveform, a second + // step. The Skynet field host opens ~2 such holes a second. So the first synthesized frame + // fades the audio OUT over `EDGE_FADE_US` and the first real frame after the hole fades it + // back IN over the same span, both on the same raised cosine — a crossfade through silence + // rather than a step at each end. `last_real` is what the fade-out is built from when a hole + // opens on an empty partial (nothing of the pre-hole audio left in hand to fade). + let mut last_real: Vec = Vec::with_capacity(frame_len); + let mut resume_fade = false; + // In interleaved samples: `frames × channels`, so the curve spans whole frames. (Adjacent + // channels of one frame land one step apart on the curve — a 1/n gain difference, nothing at + // any fade longer than a few frames.) + let edge_fade_samples = (rate_hz as u64 * EDGE_FADE_US / 1_000_000) as usize * want as usize; + // Backlog above which a slot sends TWO frames — see the bonus arm in the loop below. Sized + // from the largest capture chunk this session has actually seen (the graph's quantum, or a + // VM's clamped 1024 frames), plus one full protocol frame: anything up to that is a chunk + // being paced out, anything past it is audio that arrived faster than the schedule sends. + let mut max_chunk_len: usize = 0; // Reopen-with-backoff: hold the capturer in an Option so a mid-session capture-thread death // (device unplug, daemon restart) — or a first open lost to session-start churn above — // reopens instead of muting the rest of a multi-hour session. A quiet sink is NOT a death — @@ -571,6 +595,15 @@ pub(super) fn audio_thread( // long session. let arrival_ns = now_ns(); acc.extend_from_slice(&chunk); + max_chunk_len = max_chunk_len.max(chunk.len()); + // How big the graph's buffers really are, so the infill threshold is one chunk plus + // one frame and never the middle of a legitimately long cycle — see + // `InfillPolicy::after`. + infill.note_quantum(std::time::Duration::from_nanos(pcm::frame_duration_ns( + chunk.len(), + rate_hz, + want, + ))); let queued_frames = (acc.len() / want as usize) as u64; // The session's rate, not the module constant: at 96 kHz a 48 000 divisor would put // the anchor twice as far into the past as the queued audio really is, so every pts @@ -581,9 +614,42 @@ pub(super) fn audio_thread( clock.reanchor(anchor); } // Everything the wire owes for the slots that have come due — real or synthesized, one - // schedule, one encoder, one `seq`. A schedule that has fallen more than one frame behind - // is re-anchored rather than chased, so a scheduling hiccup cannot turn into a permanent - // send-time debt. + // schedule, one encoder, one `seq`. A schedule that has fallen more than + // `PACE_REANCHOR` behind is re-anchored rather than chased, so a scheduling hiccup cannot + // turn into a permanent send-time debt. + // + // **The schedule is wall clock; the source is not, and the two disagree in both + // directions.** What this loop owes the client is a wire that carries one frame per + // frame-time of WALL clock — the client's ring is drained on its own device clock, so + // any frame this schedule fails to send is a frame the ring goes without, and any frame + // it sends over is one the ring holds forever: + // + // - **Source behind** (capture lost audio, or the graph clock is slow): a slot comes due + // with no complete frame in `acc`. That is a hole. The ≥ 10 ms kind was already covered + // by the infill policy, measured from the last chunk; the SHORT kind — a missed + // 2.7 ms graph cycle or two, invisible to the capture gap counter — was not, and it + // left the schedule permanently behind by the loss. The 2026-08-17 field log shows the + // shape: 33–72 % of departures "late", `max_late_ms` climbing to 99, `reanchors` ≈ 0 — + // a lag that only ever accumulated, and was only ever repaid when the next long hole's + // infill burst out (lag + 10 ms) / 5 ms silence frames back to back. On the client + // that is a ring drained a frame at a time and then refilled with a burst it has to + // trim. So the infill decision now looks at the LAG as well as the time since the + // last chunk: a schedule `after()` behind with no audio to send is owed a frame of + // cover exactly as a hole that long is, and the lag stays bounded at that (a couple + // of frames at a frame-sized quantum, one chunk plus a frame on a clamped one). + // - **Source ahead** (the graph clock is fast, or a chunk landed after a stall): `acc` + // holds more than a chunk's worth beyond the frame being sent. Left alone the surplus + // is never sent — one frame per slot, forever — so a source a hundred ppm fast grows + // the backlog, and the latency, by five milliseconds every fifty seconds for the life + // of the session, and only a hole ever takes it back. So a slot whose remaining + // backlog exceeds one chunk plus one frame sends a second frame in the same slot + // (`bonus`): at most two per slot, so a burst never exceeds what a 10 ms quantum + // sends every cycle anyway, and the surplus drains at twice real time until it is gone. + // + // Both together bound what `late` can accumulate: a lag can no longer grow past + // `after()` without being covered, so a `max_late_ms` beyond that band is this thread's + // own scheduling — the reading WP-C built the counter for — and not the source's. + let mut bonus_taken = false; loop { let now = std::time::Instant::now(); // How far past its slot this frame is leaving. Measured before the re-anchor arm can @@ -605,23 +671,64 @@ pub(super) fn audio_thread( } else if !sent_any { break; } else { - match infill.decide(last_chunk_at.elapsed()) { + // The hole is the LATER of "since the last chunk" and "how far the schedule is + // behind": the first is a graph that stopped feeding us, the second is a graph + // that fed us less than wall clock — see the loop comment. + match infill.decide(last_chunk_at.elapsed().max(late)) { 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 - // half of it was captured. + // half of it was captured. (And it is not dropped either: those samples + // are audio the wire owes, and every sample dropped here would come + // straight back as schedule lag — see the loop comment.) + let partial = acc.len(); frame_buf.append(&mut acc); + if infill.covered() == infill.frame() { + // First frame of the hole: the audio does not stop at a step. Fade + // the tail of what we have over up to `EDGE_FADE_US` — the partial + // frame if there is one, else the same slice of the last real frame + // again, so the fade is a slope from the level the listener was at + // down to nothing rather than an edge there. + if partial == 0 && !last_real.is_empty() { + let n = edge_fade_samples.min(last_real.len()); + frame_buf.extend_from_slice(&last_real[..n]); + } + let n = frame_buf.len().min(edge_fade_samples); + pcm::raised_cosine_tail(&mut frame_buf, n); + } frame_buf.resize(frame_len, 0.0); + // Whatever follows this hole starts mid-waveform. + resume_fade = true; } crate::audio::capture_policy::Infill::Wait | crate::audio::capture_policy::Infill::Quiet => break, } } - pace_due = Some(pace_due.unwrap_or_else(std::time::Instant::now) + frame_interval); - if gain != 1.0 { - punktfunk_core::audio::apply_gain(&mut frame_buf, gain); + // A slot whose remaining backlog exceeds one chunk plus one whole frame sends a + // second frame in the same slot — see the loop comment. Decided on what is left + // AFTER this frame, so the frame that merely completes a chunk never triggers it, + // and never twice in a row, so a big surplus drains at 2× and not in one burst. + let bonus = !infilled && !bonus_taken && acc.len() >= max_chunk_len + frame_len; + pace_due = match pace_due { + Some(due) if bonus => Some(due), // same slot again for the next frame + other => Some(other.unwrap_or_else(std::time::Instant::now) + frame_interval), + }; + bonus_taken = bonus; + if !infilled { + if gain != 1.0 { + punktfunk_core::audio::apply_gain(&mut frame_buf, gain); + } + if std::mem::take(&mut resume_fade) { + // First real frame after a hole: fade it in from silence, the mirror of the + // fade the hole started with. + pcm::raised_cosine_head(&mut frame_buf, edge_fade_samples); + } + // What the next hole's fade is built from when it opens on an empty partial — + // recorded post-gain, so the fade sits at the level the listener was hearing. + last_real.clear(); + last_real.extend_from_slice(&frame_buf); } // W1.1 — the wire clock. ⚠ Charged the frame's REAL sample count, never the negotiated // `frame_us`, which is a label on the 44.1 kHz family and would run this 2.3 ms/s fast