From 6286f91f89ab9804263285a8b5c0f4517c1c4e07 Mon Sep 17 00:00:00 2001 From: enricobuehler Date: Fri, 31 Jul 2026 21:50:29 +0200 Subject: [PATCH] feat(host/audio): the mic uplink keeps its sequence, and a backlog heals itself MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The 0xCB datagrams always carried a seq + pts, and the ingest threw both away one line after decoding them — the pump saw an anonymous byte pile. Frames now travel as MicFrame {seq, pts_ns, opus} so the de-jitter that follows can reorder, conceal and measure. The shared queue also stops being a latency reservoir: cap 64 → 12, and a pump that wakes to a backlog deeper than 6 frames jumps to the newest 4 instead of replaying the pile — before this, a scheduling stall could park up to 1.28 s of standing mic delay that only a >600 ms silence gap ever flushed. Co-Authored-By: Claude Fable 5 --- crates/punktfunk-host/src/audio.rs | 2 +- crates/punktfunk-host/src/audio/mic_pump.rs | 127 ++++++++++++++------ crates/punktfunk-host/src/native.rs | 12 +- 3 files changed, 101 insertions(+), 40 deletions(-) diff --git a/crates/punktfunk-host/src/audio.rs b/crates/punktfunk-host/src/audio.rs index f1997202..ceef8618 100644 --- a/crates/punktfunk-host/src/audio.rs +++ b/crates/punktfunk-host/src/audio.rs @@ -153,4 +153,4 @@ mod wasapi_mic; pub(crate) mod wiring_plan; mod mic_pump; -pub use mic_pump::MicPump; +pub use mic_pump::{MicFrame, MicPump}; diff --git a/crates/punktfunk-host/src/audio/mic_pump.rs b/crates/punktfunk-host/src/audio/mic_pump.rs index e09b9e12..f65c49c0 100644 --- a/crates/punktfunk-host/src/audio/mic_pump.rs +++ b/crates/punktfunk-host/src/audio/mic_pump.rs @@ -10,10 +10,29 @@ use anyhow::Result; /// Mic is 48 kHz stereo — matches the Opus stereo decoder and the host→client audio layout. pub const MIC_CHANNELS: u32 = 2; -/// Bound for the shared mic frame queue (drop-newest when full): the host-lifetime queue is -/// shared across all concurrent sessions and must not grow without limit under a near-line-rate -/// flood (security-review 2026-06-28 S6). 64 × 5–20 ms frames ≈ 0.3–1.3 s of slack. -const MIC_QUEUE_CAP: usize = 64; + +/// One client mic frame off the wire (`0xCB` — `punktfunk_core::quic::decode_mic_datagram`). +/// The datagram's seq + pts used to be decoded at ingest and thrown away; the pump's de-jitter +/// needs them (reorder window, loss concealment, cadence/drift), so they ride the queue with +/// the Opus payload. +pub struct MicFrame { + pub seq: u32, + pub pts_ns: u64, + pub opus: Vec, +} + +/// Bound for the shared mic frame queue (drop-newest at the producer when full): the +/// host-lifetime queue is shared across all concurrent sessions and must not grow without limit +/// under a near-line-rate flood (security-review 2026-06-28 S6). 12 × 5–20 ms frames ≈ +/// 60–240 ms of slack — enough for a scheduling hiccup, and the consumer heals the rest (see +/// [`DRAIN_ABOVE`]). The old cap of 64 could park 1.28 s of standing latency here that only a +/// >600 ms silence gap ever flushed. +const MIC_QUEUE_CAP: usize = 12; +/// Consumer self-healing: a pump that wakes to a backlog deeper than this drops down to the +/// newest [`DRAIN_KEEP`] frames before decoding — the oldest frames are already late, and +/// playing them would turn a one-off stall into permanent mic delay. +const DRAIN_ABOVE: usize = 6; +const DRAIN_KEEP: usize = 4; /// Tuning for [`MicPump`]'s open/reopen/flush behaviour — parameterized so the tests can run the /// real pump loop in milliseconds instead of seconds. @@ -62,14 +81,14 @@ const PUMP_TUNING: PumpTuning = PumpTuning { /// (security-review 2026-06-28 S2). The thread exits when every sender is dropped (host /// shutdown), tearing the backend down. pub struct MicPump { - tx: std::sync::mpsc::SyncSender>, + tx: std::sync::mpsc::SyncSender, } impl MicPump { /// Start the host-lifetime pump (Linux/Windows). On platforms without a virtual-mic backend /// the thread just drains and drops frames (sessions still count the datagrams). pub fn start() -> MicPump { - let (tx, rx) = std::sync::mpsc::sync_channel::>(MIC_QUEUE_CAP); + let (tx, rx) = std::sync::mpsc::sync_channel::(MIC_QUEUE_CAP); let spawned = std::thread::Builder::new() .name("punktfunk-mic-pump".into()) .spawn(move || { @@ -90,7 +109,7 @@ impl MicPump { /// A sender a session forwards the client's Opus mic frames to (`try_send` — never block a /// datagram loop). Cloned per session; dropping a clone does NOT stop the pump (it holds /// the original sender for the host life). - pub fn sender(&self) -> std::sync::mpsc::SyncSender> { + pub fn sender(&self) -> std::sync::mpsc::SyncSender { self.tx.clone() } } @@ -99,7 +118,7 @@ impl MicPump { /// never accumulates a stale backlog and senders never see a wedged queue. Returns `false` when /// every sender is gone (host shutdown). #[cfg_attr(not(any(target_os = "linux", target_os = "windows")), allow(dead_code))] -fn drain_sleep(rx: &std::sync::mpsc::Receiver>, dur: std::time::Duration) -> bool { +fn drain_sleep(rx: &std::sync::mpsc::Receiver, dur: std::time::Duration) -> bool { use std::sync::mpsc::RecvTimeoutError; let deadline = std::time::Instant::now() + dur; loop { @@ -118,7 +137,7 @@ fn drain_sleep(rx: &std::sync::mpsc::Receiver>, dur: std::time::Duration /// The pump loop. `opener` is injected so the tests can run the REAL loop against a mock /// backend; production passes [`open_virtual_mic`](super::open_virtual_mic). #[cfg_attr(not(any(target_os = "linux", target_os = "windows")), allow(dead_code))] -fn pump_thread(rx: std::sync::mpsc::Receiver>, opener: O, tuning: PumpTuning) +fn pump_thread(rx: std::sync::mpsc::Receiver, opener: O, tuning: PumpTuning) where O: Fn() -> Result>, { @@ -160,34 +179,61 @@ where // Pump phase — runs until the backend dies (break) or the host shuts down (return). let mut decode_fails: u64 = 0; + let mut drain_drops: u64 = 0; let mut pcm = vec![0f32; 5760 * MIC_CHANNELS as usize]; // up to 120 ms scratch let mut last_push = Instant::now(); - loop { + let mut batch: Vec = Vec::new(); + 'pump: loop { match rx.recv_timeout(tuning.heartbeat) { - Ok(frame) => { - if frame.is_empty() { - continue; // DTX silence — the source underruns to silence on its own + Ok(first) => { + // Take everything already queued in one gulp: normally that's just `first`, + // but after a stall the backlog IS standing mic latency — heal by jumping + // to the newest few frames instead of replaying the pile. + batch.clear(); + batch.push(first); + while batch.len() <= MIC_QUEUE_CAP { + match rx.try_recv() { + Ok(f) => batch.push(f), + Err(_) => break, + } + } + if batch.len() > DRAIN_ABOVE { + let drop_n = batch.len() - DRAIN_KEEP; + drain_drops += drop_n as u64; + if drain_drops.is_power_of_two() { + tracing::debug!( + dropped = drop_n, + total = drain_drops, + "mic pump backlog — dropped oldest frames to recover latency" + ); + } + batch.drain(..drop_n); } if last_push.elapsed() > tuning.stale_gap { mic.discard(); } - match decoder.decode_float(&frame, &mut pcm, false) { - Ok(samples_per_ch) => { - let total = (samples_per_ch * MIC_CHANNELS as usize).min(pcm.len()); - if !mic.push(&pcm[..total]) { - tracing::warn!("virtual mic backend died — reopening"); - break; - } - last_push = Instant::now(); - decode_fails = 0; + for frame in batch.drain(..) { + if frame.opus.is_empty() { + continue; // DTX silence — the source underruns to silence on its own } - Err(e) => { - // Malformed/garbage frame: drop it, keep the shared mic + decoder - // (see the struct docs). Throttled log (1, 2, 4, … fails). - decode_fails += 1; - if decode_fails.is_power_of_two() { - tracing::warn!(error = %e, fails = decode_fails, - "mic opus decode failed — dropping frame"); + match decoder.decode_float(&frame.opus, &mut pcm, false) { + Ok(samples_per_ch) => { + let total = (samples_per_ch * MIC_CHANNELS as usize).min(pcm.len()); + if !mic.push(&pcm[..total]) { + tracing::warn!("virtual mic backend died — reopening"); + break 'pump; + } + last_push = Instant::now(); + decode_fails = 0; + } + Err(e) => { + // Malformed/garbage frame: drop it, keep the shared mic + decoder + // (see the struct docs). Throttled log (1, 2, 4, … fails). + decode_fails += 1; + if decode_fails.is_power_of_two() { + tracing::warn!(error = %e, fails = decode_fails, + "mic opus decode failed — dropping frame"); + } } } } @@ -252,7 +298,7 @@ mod pump_tests { } struct Harness { - tx: std::sync::mpsc::SyncSender>, + tx: std::sync::mpsc::SyncSender, opens: Arc, alive: Arc>>>, // latest instance's kill switch pushed: Arc, @@ -265,7 +311,7 @@ mod pump_tests { /// pre-killed (exercises the rapid-death churn guard). `stable_after` mirrors the tuning /// field (ZERO = every death counts as stable → immediate reopen, keeping tests fast). fn start_tuned(fail_first: usize, dead_on_arrival: bool, stable_after: Duration) -> Harness { - let (tx, rx) = std::sync::mpsc::sync_channel::>(MIC_QUEUE_CAP); + let (tx, rx) = std::sync::mpsc::sync_channel::(MIC_QUEUE_CAP); let opens = Arc::new(AtomicUsize::new(0)); let alive = Arc::new(Mutex::new(None::>)); let pushed = Arc::new(AtomicUsize::new(0)); @@ -336,6 +382,15 @@ mod pump_tests { out } + /// A wire-shaped frame: `seq` with the matching 20 ms pts, payload from [`opus_frame`]. + fn mic_frame(seq: u32) -> MicFrame { + MicFrame { + seq, + pts_ns: seq as u64 * 20_000_000, + opus: opus_frame(), + } + } + /// Eager: the backend opens (after transient failures) with NO frame ever sent. #[test] fn opens_eagerly_with_backoff() { @@ -352,7 +407,7 @@ mod pump_tests { fn decodes_and_pushes() { let h = start(0); wait_until("open", || h.alive.lock().unwrap().is_some()); - h.tx.send(opus_frame()).unwrap(); + h.tx.send(mic_frame(0)).unwrap(); wait_until("pcm pushed", || h.pushed.load(Ordering::SeqCst) > 0); drop(h.tx); h.join.join().unwrap(); @@ -388,9 +443,9 @@ mod pump_tests { .as_ref() .unwrap() .store(false, Ordering::Release); - h.tx.send(opus_frame()).unwrap(); // push sees death → reopen + h.tx.send(mic_frame(0)).unwrap(); // push sees death → reopen wait_until("reopen", || h.opens.load(Ordering::SeqCst) >= 2); - h.tx.send(opus_frame()).unwrap(); + h.tx.send(mic_frame(1)).unwrap(); wait_until("pcm after reopen", || h.pushed.load(Ordering::SeqCst) > 0); drop(h.tx); h.join.join().unwrap(); @@ -421,10 +476,10 @@ mod pump_tests { fn discards_after_gap() { let h = start(0); wait_until("instance", || h.alive.lock().unwrap().is_some()); - h.tx.send(opus_frame()).unwrap(); + h.tx.send(mic_frame(0)).unwrap(); wait_until("first push", || h.pushed.load(Ordering::SeqCst) > 0); std::thread::sleep(Duration::from_millis(150)); // > stale_gap - h.tx.send(opus_frame()).unwrap(); + h.tx.send(mic_frame(1)).unwrap(); wait_until("discard on gap", || h.discards.load(Ordering::SeqCst) >= 1); drop(h.tx); h.join.join().unwrap(); diff --git a/crates/punktfunk-host/src/native.rs b/crates/punktfunk-host/src/native.rs index 05807e7c..7dad3025 100644 --- a/crates/punktfunk-host/src/native.rs +++ b/crates/punktfunk-host/src/native.rs @@ -760,7 +760,7 @@ async fn serve_session( opts: &Punktfunk1Options, audio_cap: &AudioCapSlot, inj_tx: std::sync::mpsc::Sender, - mic_tx: std::sync::mpsc::SyncSender>, + mic_tx: std::sync::mpsc::SyncSender, host_fp: &[u8; 32], np: &NativePairing, last_pairing: &std::sync::Mutex>, @@ -1178,11 +1178,17 @@ async fn serve_session( tokio::spawn(async move { let (mut input_count, mut mic_count, mut rich_count) = (0u64, 0u64, 0u64); while let Ok(d) = input_conn.read_datagram().await { - if let Some((_seq, _pts, opus)) = punktfunk_core::quic::decode_mic_datagram(&d) { + if let Some((seq, pts, opus)) = punktfunk_core::quic::decode_mic_datagram(&d) { mic_count += 1; // Host-lifetime mic service (bounded queue): `try_send` drops the frame when the // service is full or gone, never blocking this datagram loop (security-review S6). - let _ = mic_tx.try_send(opus.to_vec()); + // seq + pts ride along — the pump's de-jitter reorders, conceals losses and + // tracks cadence with them (they used to be decoded here and thrown away). + let _ = mic_tx.try_send(crate::audio::MicFrame { + seq, + pts_ns: pts, + opus: opus.to_vec(), + }); } else if let Some(rich) = punktfunk_core::quic::RichInput::decode(&d) { rich_count += 1; if rich_tx.send(ClientInput::Rich(rich)).is_err() {