Findings from the post-implementation review of design/audio-quality-and-latency.md. **The bandwidth gap (highest).** Tier `High` (256 kbps) and the redundant `0xD2` plane were added separately, each costed as "~1 % of the video budget", and nobody added them together: 256 kbps sent twice is 512 kbps — ~2.5 % of a 20 Mbps session but ~10 % of a 5 Mbps one. Audio rides QUIC datagrams, OUTSIDE the ABR loop, so ABR could neither see that nor reclaim it; a constrained link quietly handed a tenth of its bandwidth to audio while ABR carefully managed the rest. `plan_audio_budget` now makes tier and redundancy ONE decision against the session's resolved video bitrate, ordered by preference rather than cost — transparent audio beats redundant audio, since the field report was about quality and redundancy only pays under loss, so `High` alone outranks `Standard`+redundancy even though they cost the same. It can lower what the operator asked for, never raise it, and never goes below `Low`: a stream with unintelligible audio is worse than one spending a few percent more. **The Linux host kept the exact defect fixed on Windows.** `let _ = tx.try_send(samples)` — silent, uncounted data loss, where the encoder concatenates across the hole, so every drop is a click AND a permanent shift of everything after it. WP0.2 turned out to be Windows-only and had not said so. Linux now shares `capture_policy::CaptureStats`: drops counted and warned, plus per-window peak/RMS/delivered%. A Linux audio report was until now exactly as un-triageable as the Windows one was on 2026-08-03. **Apple's WP0.3 was half-done** — `bufferedMS` was added and wired to nothing. The drain thread now logs buffer/target/underruns/sheds like the other three, from one locked snapshot so the numbers in a line describe the same instant. Also: the Linux "audio format negotiated" line now says WHICH mode produced it, because that changes what it is worth — in stream-sink mode the host owns the sink so the mix cannot have been narrowed upstream, but in legacy monitor mode a 16 kHz Bluetooth sink would still be reported as a clean 48 kHz through PipeWire's resampler, the same way WASAPI's autoconvert hid it on Windows. Reading the monitored node's own rate needs a registry lookup this stream does not do; recorded as an open gap rather than implied to be covered. Two stale docs: `audio_wasapi.rs` cited `clients/windows/src/audio.rs` (deleted) and still described the pre-shared-policy "prime to ~3 quanta" behaviour. And the Apple ring's `prefill:` parameter, dead since the depth moved into the ring, is gone. Verified: clippy --all-targets -D warnings on Linux (docker) AND Windows (runner .133, forced clean rebuild of punktfunk-host + pf-client-core); core 167 tests; host 57 audio tests on Windows; Android clippy count identical to pristine (6, all documented arm64 artifacts); Apple ring re-simulated. The host suite's `gamestream::stream::tests::sender_delivers_batches` fails under qemu — the recorded environmental flake, unrelated to audio, green on the earlier less-loaded run. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
511 lines
22 KiB
Rust
511 lines
22 KiB
Rust
//! Audio: playback (decoded PCM → a WASAPI shared-mode render stream) and the microphone
|
||
//! uplink (WASAPI capture → Opus → 0xCB datagrams, the inverse of the host's virtual mic).
|
||
//!
|
||
//! The WASAPI twin of `audio.rs` (PipeWire) — same public surface (`AudioPlayer::spawn`/
|
||
//! `take_buffer`/`push`, `MicStreamer::spawn`), swapped in by lib.rs's `#[path]` so the
|
||
//! session pump compiles against one `crate::audio` on both OSes. It began as a copy of the
|
||
//! WinUI shell's own audio path; that shell's built-in streaming path has since been deleted,
|
||
//! so this is now the only WASAPI client ring.
|
||
//!
|
||
//! Playback: the session pump pushes 5 ms Opus-decoded chunks on the network clock; the WASAPI
|
||
//! render thread pulls whole event-driven quanta on the device clock. The depth policy between
|
||
//! them is the SHARED `punktfunk_core::audio::JitterPolicy` (`JitterTuning::WASAPI`) — target in
|
||
//! milliseconds, crossfaded drift correction, de-prime hysteresis — so all four clients behave
|
||
//! the same way and none of them can ratchet latency upward.
|
||
//!
|
||
//! WASAPI objects are COM-apartment-bound and not `Send`, so they live on a dedicated
|
||
//! thread (the same discipline as the host's `wasapi_cap`); only the channels + stop flag
|
||
//! + join handle cross the boundary.
|
||
|
||
use anyhow::{anyhow, Context, Result};
|
||
use punktfunk_core::client::NativeClient;
|
||
use std::collections::VecDeque;
|
||
use std::sync::atomic::{AtomicBool, Ordering};
|
||
use std::sync::mpsc::{Receiver, SyncSender, TrySendError};
|
||
use std::sync::Arc;
|
||
use std::time::Duration;
|
||
use wasapi::{
|
||
AudioClientProperties, DeviceEnumerator, Direction, SampleType, StreamCategory, StreamMode,
|
||
WaveFormat,
|
||
};
|
||
|
||
const SAMPLE_RATE: usize = 48_000;
|
||
/// Mic capture requests STEREO from WASAPI (autoconvert matrixes any endpoint layout down to
|
||
/// it — the proven path; `read_from_device_to_deque` then delivers our requested format) and
|
||
/// downmixes to MONO in code before the encoder: voice is mono at the source, the host accepts
|
||
/// any Opus channel layout (its stereo decoder upmixes), and half the samples halve the
|
||
/// encode + wire cost. The render path is multichannel — its channel count + block align are
|
||
/// runtime, driven by the host-resolved layout.
|
||
const CAPT_CHANNELS: usize = 2;
|
||
/// Mic frames are 10 ms (480 mono samples) — any size ≤ 120 ms is fine host-side; 10 ms
|
||
/// halves the frame-fill share of mouth-to-ear latency vs the old 20 ms.
|
||
const MIC_FRAME: usize = 480;
|
||
|
||
/// A selectable WASAPI endpoint for the settings pickers.
|
||
#[derive(Clone, Debug)]
|
||
pub struct AudioDevice {
|
||
/// The `IMMDevice` endpoint id (`{0.0.0.00000000}.{…}`) — the stable key the render and
|
||
/// capture threads resolve via [`DeviceEnumerator::get_device`]. (The PipeWire twin
|
||
/// stores `node.name` here; both are "the stable key", so the Settings fields and env
|
||
/// contract stay OS-agnostic.)
|
||
pub name: String,
|
||
/// The endpoint's friendly name ("Speakers (Realtek …)") — what the picker shows.
|
||
pub description: String,
|
||
}
|
||
|
||
/// Enumerate active audio endpoints: `(sinks, sources)` — the WASAPI twin of the PipeWire
|
||
/// probe (same tuple shape; no devices → the caller simply shows no pickers). Runs on its
|
||
/// own short-lived MTA thread: the caller is typically a UI thread whose COM apartment is
|
||
/// STA, where a direct `CoInitializeEx(MTA)` would fail with `RPC_E_CHANGED_MODE`.
|
||
pub fn devices() -> Result<(Vec<AudioDevice>, Vec<AudioDevice>)> {
|
||
std::thread::Builder::new()
|
||
.name("pf-audio-enum".into())
|
||
.spawn(|| -> Result<(Vec<AudioDevice>, Vec<AudioDevice>)> {
|
||
wasapi::initialize_mta()
|
||
.ok()
|
||
.context("CoInitializeEx (MTA)")?;
|
||
let enumerator = DeviceEnumerator::new().context("DeviceEnumerator")?;
|
||
let mut out = (Vec::new(), Vec::new());
|
||
for (direction, list) in [
|
||
(Direction::Render, &mut out.0),
|
||
(Direction::Capture, &mut out.1),
|
||
] {
|
||
let coll = enumerator
|
||
.get_device_collection(&direction)
|
||
.context("device collection")?;
|
||
for i in 0..coll.get_nbr_devices().context("device count")? {
|
||
// One broken endpoint (driver limbo) must not hide the rest.
|
||
let Ok(dev) = coll.get_device_at_index(i) else {
|
||
continue;
|
||
};
|
||
let (Ok(id), Ok(name)) = (dev.get_id(), dev.get_friendlyname()) else {
|
||
continue;
|
||
};
|
||
list.push(AudioDevice {
|
||
name: id,
|
||
description: name,
|
||
});
|
||
}
|
||
}
|
||
Ok(out)
|
||
})
|
||
.context("spawn audio enumeration thread")?
|
||
.join()
|
||
.map_err(|_| anyhow!("audio enumeration thread panicked"))?
|
||
}
|
||
|
||
/// The endpoint an env pick names (`PUNKTFUNK_AUDIO_SINK`/`SOURCE` — endpoint ids, the
|
||
/// Settings device pickers via session main), or the OS default. A picked device that's
|
||
/// gone (unplugged USB DAC, remote session) falls back to the default with a warning —
|
||
/// audio keeps working, like the PipeWire twin's `target.object` behavior.
|
||
fn pick_device(
|
||
enumerator: &DeviceEnumerator,
|
||
direction: &Direction,
|
||
var: &str,
|
||
) -> Result<wasapi::Device> {
|
||
if let Some(id) = std::env::var(var).ok().filter(|v| !v.is_empty()) {
|
||
match enumerator.get_device(&id) {
|
||
Ok(d) => {
|
||
tracing::info!(
|
||
var,
|
||
endpoint = %d.get_friendlyname().unwrap_or_else(|_| id.clone()),
|
||
"using the picked audio endpoint"
|
||
);
|
||
return Ok(d);
|
||
}
|
||
Err(e) => tracing::warn!(
|
||
var,
|
||
endpoint_id = %id,
|
||
error = %e,
|
||
"picked audio endpoint not found — using the default"
|
||
),
|
||
}
|
||
}
|
||
enumerator
|
||
.get_default_device(direction)
|
||
.context("default endpoint")
|
||
}
|
||
|
||
pub struct AudioPlayer {
|
||
pcm_tx: SyncSender<Vec<f32>>,
|
||
/// Drained chunk Vecs coming back from the render thread for reuse (the pool half of
|
||
/// the pcm channel — see [`AudioPlayer::take_buffer`]).
|
||
recycle_rx: Receiver<Vec<f32>>,
|
||
stop: Arc<AtomicBool>,
|
||
thread: Option<std::thread::JoinHandle<()>>,
|
||
}
|
||
|
||
impl AudioPlayer {
|
||
/// Spawn the WASAPI render thread for `channels` (2/6/8, canonical wire order
|
||
/// FL FR FC LFE RL RR SL SR). Failure (no render endpoint on this box) is survivable — the
|
||
/// caller streams video-only.
|
||
pub fn spawn(channels: u32) -> Result<AudioPlayer> {
|
||
// 64 × 5 ms = 320 ms of slack between the pump and the WASAPI loop.
|
||
let (pcm_tx, pcm_rx) = std::sync::mpsc::sync_channel::<Vec<f32>>(64);
|
||
// Return path: the render thread sends each drained Vec back for reuse, so
|
||
// steady-state playback stops allocating (~200 chunks/s otherwise). Same capacity
|
||
// as the data channel; a full pool just drops the Vec (plain deallocation).
|
||
let (recycle_tx, recycle_rx) = std::sync::mpsc::sync_channel::<Vec<f32>>(64);
|
||
let stop = Arc::new(AtomicBool::new(false));
|
||
let (ready_tx, ready_rx) = std::sync::mpsc::sync_channel::<Result<()>>(1);
|
||
let stop_t = stop.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, channels as u8)
|
||
{
|
||
tracing::warn!(error = %format!("{e:#}"), "audio playback thread ended");
|
||
}
|
||
})
|
||
.context("spawn audio thread")?;
|
||
match ready_rx.recv_timeout(Duration::from_secs(3)) {
|
||
Ok(Ok(())) => {
|
||
// Default endpoint unless PUNKTFUNK_AUDIO_SINK picked one (logged there).
|
||
tracing::info!(channels, "WASAPI render: 48 kHz f32");
|
||
Ok(AudioPlayer {
|
||
pcm_tx,
|
||
recycle_rx,
|
||
stop,
|
||
thread: Some(thread),
|
||
})
|
||
}
|
||
Ok(Err(e)) => Err(e),
|
||
Err(_) => Err(anyhow!(
|
||
"wasapi render init timed out (no render endpoint?)"
|
||
)),
|
||
}
|
||
}
|
||
|
||
/// 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 WASAPI side dropped chunks).
|
||
pub fn take_buffer(&self) -> Vec<f32> {
|
||
self.recycle_rx.try_recv().unwrap_or_default()
|
||
}
|
||
|
||
/// 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<f32>) {
|
||
if let Err(TrySendError::Disconnected(_)) = self.pcm_tx.try_send(pcm) {
|
||
// Thread already dead — Drop will reap it; nothing to do per-chunk.
|
||
}
|
||
}
|
||
}
|
||
|
||
impl Drop for AudioPlayer {
|
||
fn drop(&mut self) {
|
||
self.stop.store(true, Ordering::SeqCst);
|
||
if let Some(t) = self.thread.take() {
|
||
let _ = t.join();
|
||
}
|
||
}
|
||
}
|
||
|
||
fn render_thread(
|
||
pcm_rx: Receiver<Vec<f32>>,
|
||
recycle_tx: SyncSender<Vec<f32>>,
|
||
stop: Arc<AtomicBool>,
|
||
ready: SyncSender<Result<()>>,
|
||
channels: u8,
|
||
) -> Result<()> {
|
||
if let Err(e) = wasapi::initialize_mta()
|
||
.ok()
|
||
.context("CoInitializeEx (MTA)")
|
||
{
|
||
let _ = ready.send(Err(e));
|
||
return Ok(());
|
||
}
|
||
let res = (|| -> Result<()> {
|
||
// F32LE interleaved: channels × 4 bytes/sample. Stereo (channels == 2) is byte-identical
|
||
// to the old fixed path (mask 0x3, block align 8).
|
||
let block_align = channels as usize * 4;
|
||
let enumerator = DeviceEnumerator::new().context("DeviceEnumerator")?;
|
||
let device = pick_device(&enumerator, &Direction::Render, "PUNKTFUNK_AUDIO_SINK")
|
||
.context("render endpoint")?;
|
||
let mut audio_client = device.get_iaudioclient().context("IAudioClient")?;
|
||
// The explicit dwChannelMask is the wire order (FL FR FC LFE RL RR SL SR); 5.1 = 0x3F,
|
||
// 7.1 = 0x63F. WASAPI delivers channels in ascending mask-bit order, which equals the wire
|
||
// order, so the render mapping is the identity — no permute. `autoconvert` (below) lets the
|
||
// audio engine downmix when the endpoint has fewer speakers.
|
||
let desired = WaveFormat::new(
|
||
32,
|
||
32,
|
||
&SampleType::Float,
|
||
SAMPLE_RATE,
|
||
channels as usize,
|
||
Some(punktfunk_core::audio::wasapi_channel_mask(channels)),
|
||
);
|
||
let (default_period, _min_period) =
|
||
audio_client.get_device_period().context("device period")?;
|
||
let mode = StreamMode::EventsShared {
|
||
autoconvert: true,
|
||
buffer_duration_hns: default_period,
|
||
};
|
||
audio_client
|
||
.initialize_client(&desired, &Direction::Render, &mode)
|
||
.context("initialize render client")?;
|
||
let h_event = audio_client.set_get_eventhandle().context("event handle")?;
|
||
let render_client = audio_client
|
||
.get_audiorenderclient()
|
||
.context("IAudioRenderClient")?;
|
||
audio_client.start_stream().context("start render stream")?;
|
||
let _ = ready.send(Ok(()));
|
||
|
||
// De-jitter ring, in interleaved f32 SAMPLES (it used to be raw bytes, which made the
|
||
// depth arithmetic byte-vs-sample and kept it from sharing the policy and the crossfade
|
||
// helper with the other three clients).
|
||
let mut ring: VecDeque<f32> = VecDeque::new();
|
||
// Shared ms-denominated policy: prime depth, crossfaded drift correction so latency
|
||
// returns to target instead of ratcheting, and de-prime hysteresis — the last replacing
|
||
// the old `if ring.is_empty()`, where a single transient drain manufactured a whole
|
||
// target's worth of fresh silence.
|
||
let mut policy = punktfunk_core::audio::JitterPolicy::new(
|
||
punktfunk_core::audio::JitterTuning::WASAPI,
|
||
channels,
|
||
);
|
||
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() {
|
||
continue;
|
||
}
|
||
// Drain everything the pump has queued into the ring, returning each drained
|
||
// Vec to the pool (a full/closed pool drops it).
|
||
while let Ok(mut chunk) = pcm_rx.try_recv() {
|
||
ring.extend(chunk.iter().copied());
|
||
chunk.clear();
|
||
let _ = recycle_tx.try_send(chunk);
|
||
}
|
||
let avail_frames = audio_client
|
||
.get_available_space_in_frames()
|
||
.context("available space")? as usize;
|
||
if avail_frames == 0 {
|
||
continue;
|
||
}
|
||
let want = avail_frames * channels as usize;
|
||
|
||
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);
|
||
}
|
||
|
||
out.clear();
|
||
out.resize(avail_frames * block_align, 0);
|
||
let mut ran_short = false;
|
||
if !step.silence {
|
||
// `out` is exactly `want` f32s wide (avail_frames × channels × 4 bytes).
|
||
for dst in out.chunks_exact_mut(4) {
|
||
let s = ring.pop_front().unwrap_or_else(|| {
|
||
ran_short = true;
|
||
0.0
|
||
});
|
||
dst.copy_from_slice(&s.to_le_bytes());
|
||
}
|
||
}
|
||
// 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,
|
||
"audio playback"
|
||
);
|
||
}
|
||
render_client
|
||
.write_to_device(avail_frames, &out, None)
|
||
.context("write_to_device")?;
|
||
}
|
||
audio_client.stop_stream().ok();
|
||
Ok(())
|
||
})();
|
||
if let Err(ref e) = res {
|
||
let _ = ready.send(Err(anyhow!("{e:#}")));
|
||
}
|
||
res
|
||
}
|
||
|
||
/// The microphone uplink: capture the default input device, Opus-encode 10 ms mono chunks,
|
||
/// ship them as 0xCB datagrams into the host's virtual mic source.
|
||
pub struct MicStreamer {
|
||
stop: Arc<AtomicBool>,
|
||
thread: Option<std::thread::JoinHandle<()>>,
|
||
}
|
||
|
||
impl MicStreamer {
|
||
/// `muted` is the in-stream mute (B4), shared live with the capture loop: set, the loop
|
||
/// keeps reading the endpoint and discarding whole frames but sends nothing. Muting by
|
||
/// STOPPING the client was rejected — an `IAudioClient` stop/start re-primes the endpoint
|
||
/// buffers and re-runs the category negotiation below on every unmute.
|
||
///
|
||
/// `echo_cancel` is the Settings toggle; `PUNKTFUNK_NO_AEC=1` overrides it off.
|
||
pub fn spawn(
|
||
connector: Arc<NativeClient>,
|
||
muted: Arc<AtomicBool>,
|
||
echo_cancel: bool,
|
||
) -> Result<MicStreamer> {
|
||
let stop = Arc::new(AtomicBool::new(false));
|
||
let stop_t = stop.clone();
|
||
let thread = std::thread::Builder::new()
|
||
.name("punktfunk-mic".into())
|
||
.spawn(move || {
|
||
if let Err(e) = mic_thread(&connector, stop_t, muted, echo_cancel) {
|
||
tracing::warn!(error = %format!("{e:#}"), "mic uplink thread ended");
|
||
}
|
||
})
|
||
.context("spawn mic thread")?;
|
||
Ok(MicStreamer {
|
||
stop,
|
||
thread: Some(thread),
|
||
})
|
||
}
|
||
}
|
||
|
||
impl Drop for MicStreamer {
|
||
fn drop(&mut self) {
|
||
self.stop.store(true, Ordering::SeqCst);
|
||
if let Some(t) = self.thread.take() {
|
||
let _ = t.join();
|
||
}
|
||
}
|
||
}
|
||
|
||
/// Whether the mic echo-cancellation hooks run this session: the `echo_cancel` setting, with
|
||
/// `PUNKTFUNK_NO_AEC=1` as a one-way override OFF. The env var wins — it is the escape hatch
|
||
/// for a box whose canceller misbehaves, and it predates the setting; nothing turns AEC back
|
||
/// on once it is set. Here the hook is the Communications stream category below; the PipeWire
|
||
/// twin gates its echo-cancelled-source preference the same way.
|
||
fn aec_enabled(echo_cancel: bool) -> bool {
|
||
echo_cancel && !std::env::var("PUNKTFUNK_NO_AEC").is_ok_and(|v| !v.is_empty() && v != "0")
|
||
}
|
||
|
||
fn mic_thread(
|
||
connector: &Arc<NativeClient>,
|
||
stop: Arc<AtomicBool>,
|
||
muted: Arc<AtomicBool>,
|
||
echo_cancel: bool,
|
||
) -> Result<()> {
|
||
wasapi::initialize_mta()
|
||
.ok()
|
||
.context("CoInitializeEx (MTA)")?;
|
||
|
||
let mut encoder = opus::Encoder::new(
|
||
SAMPLE_RATE as u32,
|
||
opus::Channels::Mono,
|
||
opus::Application::Voip,
|
||
)
|
||
.map_err(|e| anyhow!("opus encoder: {e}"))?;
|
||
// Voice tuning: 48 kbps mono is transparent for speech; in-band FEC + an assumed 10 %
|
||
// loss let the host's decoder rebuild a lost 0xCB datagram from its successor instead
|
||
// of concealing (datagrams are fire-and-forget — this FEC is the only redundancy).
|
||
let _ = encoder.set_bitrate(opus::Bitrate::Bits(48_000));
|
||
let _ = encoder.set_inband_fec(true);
|
||
let _ = encoder.set_packet_loss_perc(10);
|
||
|
||
let enumerator = DeviceEnumerator::new().context("DeviceEnumerator")?;
|
||
let device = pick_device(&enumerator, &Direction::Capture, "PUNKTFUNK_AUDIO_SOURCE")
|
||
.context("capture endpoint (no microphone?)")?;
|
||
let mut audio_client = device.get_iaudioclient().context("IAudioClient")?;
|
||
// Communications category → the endpoint's communications signal-processing chain. A
|
||
// driver/APO stack with an echo canceller only engages it for communications-category
|
||
// streams; the default (Other) category never did, so the downlink audio playing on
|
||
// this box fed straight back into the host's virtual mic. Must precede Initialize
|
||
// (SetClientProperties is a pre-init call; the wasapi crate QIs IAudioClient2 inside).
|
||
// Best-effort: an endpoint without IAudioClient2 just keeps the default category.
|
||
// The "Echo cancellation" setting opts out, and PUNKTFUNK_NO_AEC=1 overrides that off
|
||
// (same lever as the Linux echo-cancel-source preference) — see `aec_enabled`.
|
||
if aec_enabled(echo_cancel) {
|
||
if let Err(e) = audio_client.set_properties(
|
||
AudioClientProperties::new().set_category(StreamCategory::Communications),
|
||
) {
|
||
tracing::debug!(error = %e, "mic capture: Communications category not set");
|
||
}
|
||
}
|
||
let desired = WaveFormat::new(32, 32, &SampleType::Float, SAMPLE_RATE, CAPT_CHANNELS, None);
|
||
let (default_period, _min_period) =
|
||
audio_client.get_device_period().context("device period")?;
|
||
let mode = StreamMode::EventsShared {
|
||
autoconvert: true,
|
||
buffer_duration_hns: default_period,
|
||
};
|
||
audio_client
|
||
.initialize_client(&desired, &Direction::Capture, &mode)
|
||
.context("initialize capture client")?;
|
||
let h_event = audio_client.set_get_eventhandle().context("event handle")?;
|
||
let capture_client = audio_client
|
||
.get_audiocaptureclient()
|
||
.context("IAudioCaptureClient")?;
|
||
audio_client
|
||
.start_stream()
|
||
.context("start capture stream")?;
|
||
|
||
let mut bytes: VecDeque<u8> = VecDeque::new();
|
||
let mut ring: VecDeque<f32> = VecDeque::new();
|
||
let mut out = vec![0u8; 4000];
|
||
let mut seq = 0u32;
|
||
|
||
while !stop.load(Ordering::Relaxed) {
|
||
if h_event.wait_for_event(100).is_err() {
|
||
continue;
|
||
}
|
||
loop {
|
||
match capture_client.get_next_packet_size() {
|
||
Ok(Some(0)) | Ok(None) => break,
|
||
Ok(Some(_n)) => {
|
||
capture_client
|
||
.read_from_device_to_deque(&mut bytes)
|
||
.context("read capture")?;
|
||
}
|
||
Err(e) => return Err(anyhow!("get_next_packet_size: {e}")),
|
||
}
|
||
}
|
||
// One stereo capture frame (8 bytes) → one mono sample: average L/R. Autoconvert
|
||
// already matrixed the endpoint's real layout (mono/stereo/array mic) into the
|
||
// stereo stream we initialized, so this is the only downmix left to do.
|
||
let stereo_frame = 4 * CAPT_CHANNELS;
|
||
let whole = (bytes.len() / stereo_frame) * stereo_frame;
|
||
for c in bytes
|
||
.drain(..whole)
|
||
.collect::<Vec<u8>>()
|
||
.chunks_exact(stereo_frame)
|
||
{
|
||
let l = f32::from_le_bytes([c[0], c[1], c[2], c[3]]);
|
||
let r = f32::from_le_bytes([c[4], c[5], c[6], c[7]]);
|
||
ring.push_back((l + r) * 0.5);
|
||
}
|
||
// Muted (B4): the capture client stays started and keeps its primed buffers — only
|
||
// the sending stops. Whole frames are discarded so the ring can't grow, and `seq`
|
||
// deliberately does NOT advance: the host sees one continuous sequence with a silent
|
||
// pause in the middle rather than a gap the size of the mute, which its de-jitter
|
||
// would try to conceal frame by frame.
|
||
if muted.load(Ordering::Relaxed) {
|
||
let drop_n = (ring.len() / MIC_FRAME) * MIC_FRAME;
|
||
ring.drain(..drop_n);
|
||
continue;
|
||
}
|
||
// Ship every complete 10 ms mono frame.
|
||
while ring.len() >= MIC_FRAME {
|
||
let pcm: Vec<f32> = ring.drain(..MIC_FRAME).collect();
|
||
match encoder.encode_float(&pcm, &mut out) {
|
||
Ok(len) => {
|
||
let pts = std::time::SystemTime::now()
|
||
.duration_since(std::time::UNIX_EPOCH)
|
||
.map(|d| d.as_nanos() as u64)
|
||
.unwrap_or(0);
|
||
let _ = connector.send_mic(seq, pts, out[..len].to_vec());
|
||
seq = seq.wrapping_add(1);
|
||
}
|
||
Err(e) => tracing::debug!(error = %e, "opus mic encode"),
|
||
}
|
||
}
|
||
}
|
||
audio_client.stop_stream().ok();
|
||
Ok(())
|
||
}
|