perf(recording): move single-file WAV disk I/O off the mixer path (A17)

Recorder::write_frame ran on the playout mixer path and did a blocking write_all
to disk per 20ms frame; slow/contended storage could stall the mixer and cause
local playback underruns. Now the mixer thread only does the cheap mic-sum
(extracted as the pure mix_with_mic helper) and try_sends the frame to a
dedicated writer thread over a bounded sync_channel(256). The writer thread owns
the WavWriter, writes queued frames, records the first write error then drains
without writing, and patches the WAV size fields on channel close. A full queue
DROPS the recording frame (counted + logged at 1 and every 256) rather than
blocking call audio; a disconnected writer surfaces BrokenPipe. finalize() closes
the channel, joins the thread, and returns the first write error or the finalize
result (thread panic handled).

Scope: single-file Recorder only; WavWriter unchanged so the multitrack recorder
is untouched (its writer-thread offload is deferred as A17b). Public method
signatures preserved -> no core/mod.rs changes. New end-to-end threaded WAV
readback test + mix_with_mic helper tests; existing FIFO/mic-sum intent kept.
No new deps, no wire change. Codex-implemented (gpt-5.5 xhigh), senior-reviewed.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-06-29 02:10:29 -04:00
co-authored by Claude Opus 4.8
parent e8a894be49
commit e0325d4590
+151 -43
View File
@@ -2,21 +2,26 @@
//! //!
//! Records the **full call as you experienced it**: the mixed incoming audio //! Records the **full call as you experienced it**: the mixed incoming audio
//! (everyone you hear) summed with your own transmitted mic, into a single mono //! (everyone you hear) summed with your own transmitted mic, into a single mono
//! WAV. Writing is driven by the playout mixer (one [`Recorder::write_frame`] //! WAV. Mixing/enqueue is driven by the playout mixer (one
//! per produced 20ms frame, paced by the hardware clock); your mic arrives //! [`Recorder::write_frame`] per produced 20ms frame, paced by the hardware
//! separately from the capture thread via [`Recorder::push_mic`] and is buffered //! clock), while disk writes happen on a dedicated writer thread; your mic
//! in a small FIFO so the two independently-clocked streams stay roughly aligned. //! arrives separately from the capture thread via [`Recorder::push_mic`] and is
//! buffered in a small FIFO so the two independently-clocked streams stay
//! roughly aligned.
//! Minor clock drift just slowly grows/shrinks that FIFO (capped, so the lag //! Minor clock drift just slowly grows/shrinks that FIFO (capped, so the lag
//! between your voice and the recording is bounded) — harmless for a voice //! between your voice and the recording is bounded) — harmless for a voice
//! recording, no realtime crackle concern. //! recording, no realtime crackle concern.
//! //!
//! No external crates: the WAV writer emits the 44-byte canonical header itself //! No external crates: the WAV writer emits the 44-byte canonical header itself
//! and patches the two size fields on [`Recorder::finalize`]. //! and patches the two size fields on the writer thread during
//! [`Recorder::finalize`].
use std::collections::VecDeque; use std::collections::VecDeque;
use std::fs::{File, OpenOptions}; use std::fs::{File, OpenOptions};
use std::io::{self, Seek, SeekFrom, Write}; use std::io::{self, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
use std::sync::mpsc::{self, SyncSender, TrySendError};
use std::thread::{self, JoinHandle};
/// Capture sample rate (mono, 48kHz, matching the rest of the audio path). /// Capture sample rate (mono, 48kHz, matching the rest of the audio path).
const SAMPLE_RATE: u32 = 48_000; const SAMPLE_RATE: u32 = 48_000;
@@ -25,6 +30,8 @@ const CHANNELS: u16 = 1;
const RIFF_DATA_OVERHEAD: u64 = 36; const RIFF_DATA_OVERHEAD: u64 = 36;
const MAX_RIFF_DATA_BYTES: u64 = u32::MAX as u64 - RIFF_DATA_OVERHEAD; const MAX_RIFF_DATA_BYTES: u64 = u32::MAX as u64 - RIFF_DATA_OVERHEAD;
const MAX_NAME_ATTEMPTS: usize = 1_000; const MAX_NAME_ATTEMPTS: usize = 1_000;
const WRITER_QUEUE_FRAMES: usize = 256;
const DROP_LOG_INTERVAL_FRAMES: u64 = 256;
/// Cap on buffered mic samples (~200ms). Bounds how far recording lag can drift /// Cap on buffered mic samples (~200ms). Bounds how far recording lag can drift
/// if the capture clock runs persistently faster than playout — past this we drop /// if the capture clock runs persistently faster than playout — past this we drop
@@ -118,13 +125,15 @@ impl WavWriter {
} }
} }
/// A live call recorder: a [`WavWriter`] plus a small mic FIFO that aligns your /// A live call recorder: a writer-thread queue plus a small mic FIFO that aligns
/// transmitted mic with the playout mixer's incoming-mix frames. /// your transmitted mic with the playout mixer's incoming-mix frames.
pub struct Recorder { pub struct Recorder {
writer: WavWriter, frame_tx: SyncSender<Vec<i16>>,
writer_thread: JoinHandle<io::Result<()>>,
/// Your transmitted mic samples, awaiting alignment with the next mix frame. /// Your transmitted mic samples, awaiting alignment with the next mix frame.
mic_fifo: VecDeque<i16>, mic_fifo: VecDeque<i16>,
path: PathBuf, path: PathBuf,
dropped_frames: u64,
} }
impl Recorder { impl Recorder {
@@ -142,10 +151,15 @@ impl Recorder {
let path = dir.join(name); let path = dir.join(name);
match OpenOptions::new().write(true).create_new(true).open(&path) { match OpenOptions::new().write(true).create_new(true).open(&path) {
Ok(file) => { Ok(file) => {
let writer = WavWriter::from_file(file)?;
let (frame_tx, frame_rx) = mpsc::sync_channel(WRITER_QUEUE_FRAMES);
let writer_thread = thread::spawn(move || writer_thread_main(writer, frame_rx));
return Ok(Self { return Ok(Self {
writer: WavWriter::from_file(file)?, frame_tx,
writer_thread,
mic_fifo: VecDeque::new(), mic_fifo: VecDeque::new(),
path, path,
dropped_frames: 0,
}); });
} }
Err(e) if e.kind() == io::ErrorKind::AlreadyExists => continue, Err(e) if e.kind() == io::ErrorKind::AlreadyExists => continue,
@@ -179,21 +193,73 @@ impl Recorder {
/// treated as silence (you weren't transmitting), so quiet stretches record /// treated as silence (you weren't transmitting), so quiet stretches record
/// the incoming mix alone. /// the incoming mix alone.
pub fn write_frame(&mut self, mixed: &[i16]) -> io::Result<()> { pub fn write_frame(&mut self, mixed: &[i16]) -> io::Result<()> {
let mut out = Vec::with_capacity(mixed.len()); let out = mix_with_mic(mixed, &mut self.mic_fifo);
for &m in mixed { match self.frame_tx.try_send(out) {
let mic = self.mic_fifo.pop_front().unwrap_or(0); Ok(()) => Ok(()),
let sum = (m as i32 + mic as i32).clamp(i16::MIN as i32, i16::MAX as i32); Err(TrySendError::Full(_)) => {
out.push(sum as i16); self.dropped_frames = self.dropped_frames.saturating_add(1);
if self.dropped_frames == 1
|| self.dropped_frames.is_multiple_of(DROP_LOG_INTERVAL_FRAMES)
{
crate::log_msg(&format!(
"recording: writer queue full; dropped {} frame(s)",
self.dropped_frames
));
}
Ok(())
}
Err(TrySendError::Disconnected(_)) => Err(io::Error::new(
io::ErrorKind::BrokenPipe,
"recording writer thread stopped",
)),
} }
self.writer.write_samples(&out)
} }
/// Finish the file, patching its size fields. Consumes the recorder. /// Finish the file, patching its size fields. Consumes the recorder.
pub fn finalize(self) -> io::Result<()> { pub fn finalize(self) -> io::Result<()> {
self.writer.finalize() let Self {
frame_tx,
writer_thread,
mic_fifo: _,
path: _,
dropped_frames: _,
} = self;
drop(frame_tx);
writer_thread
.join()
.unwrap_or_else(|_| Err(io::Error::other("recording writer thread panicked")))
} }
} }
fn writer_thread_main(mut writer: WavWriter, frame_rx: mpsc::Receiver<Vec<i16>>) -> io::Result<()> {
let mut first_write_error = None;
for frame in frame_rx {
if first_write_error.is_none()
&& let Err(e) = writer.write_samples(&frame)
{
first_write_error = Some(e);
}
}
let finalize_result = writer.finalize();
if let Some(e) = first_write_error {
Err(e)
} else {
finalize_result
}
}
fn mix_with_mic(mixed: &[i16], mic_fifo: &mut VecDeque<i16>) -> Vec<i16> {
let mut out = Vec::with_capacity(mixed.len());
for &m in mixed {
let mic = mic_fifo.pop_front().unwrap_or(0);
let sum = (m as i32 + mic as i32).clamp(i16::MIN as i32, i16::MAX as i32);
out.push(sum as i16);
}
out
}
/// Civil date (year, month, day) from a count of days since the Unix epoch. /// Civil date (year, month, day) from a count of days since the Unix epoch.
/// Howard Hinnant's `civil_from_days`; valid across the whole practical range. /// Howard Hinnant's `civil_from_days`; valid across the whole practical range.
fn civil_from_days(z: i64) -> (i64, u32, u32) { fn civil_from_days(z: i64) -> (i64, u32, u32) {
@@ -222,6 +288,23 @@ pub fn timestamp_filename(unix_secs: u64) -> String {
#[cfg(test)] #[cfg(test)]
mod tests { mod tests {
use super::*; use super::*;
use std::sync::atomic::{AtomicU64, Ordering};
static NEXT_TEMP_ID: AtomicU64 = AtomicU64::new(0);
fn unique_temp_dir(prefix: &str) -> PathBuf {
let id = NEXT_TEMP_ID.fetch_add(1, Ordering::Relaxed);
std::env::temp_dir().join(format!("{prefix}-{}-{id}", std::process::id()))
}
fn read_wav_samples(path: &Path) -> (Vec<u8>, Vec<i16>) {
let bytes = std::fs::read(path).unwrap();
let samples = bytes[44..]
.chunks_exact(2)
.map(|sample| i16::from_le_bytes([sample[0], sample[1]]))
.collect();
(bytes, samples)
}
#[test] #[test]
fn timestamp_filename_is_utc_and_padded() { fn timestamp_filename_is_utc_and_padded() {
@@ -236,10 +319,7 @@ mod tests {
#[test] #[test]
fn same_second_recordings_get_unique_files_without_truncation() { fn same_second_recordings_get_unique_files_without_truncation() {
let dir = std::env::temp_dir().join(format!( let dir = unique_temp_dir("peerspeak-collision");
"peerspeak-collision-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir); let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap(); std::fs::create_dir_all(&dir).unwrap();
@@ -258,6 +338,40 @@ mod tests {
let _ = std::fs::remove_dir_all(&dir); let _ = std::fs::remove_dir_all(&dir);
} }
#[test]
fn recorder_thread_writes_mixed_samples_and_header_on_finalize() {
let dir = unique_temp_dir("peerspeak-recorder-thread");
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let mut recorder = Recorder::create(&dir, 1_700_000_123).unwrap();
let path = recorder.path().to_path_buf();
recorder.push_mic(&[1000, i16::MAX, -1000, i16::MIN, 2222]);
recorder.write_frame(&[10, 20, -32700]).unwrap();
recorder.push_mic(&[300, -300]);
recorder
.write_frame(&[0, 1000, i16::MAX, i16::MIN])
.unwrap();
recorder.finalize().unwrap();
let expected = vec![1010, i16::MAX, i16::MIN, i16::MIN, 3222, i16::MAX, i16::MIN];
let expected_data_bytes = u32::try_from(expected.len() * 2).unwrap();
let (bytes, samples) = read_wav_samples(&path);
assert_eq!(&bytes[0..4], b"RIFF");
assert_eq!(&bytes[8..12], b"WAVE");
assert_eq!(&bytes[36..40], b"data");
let riff = u32::from_le_bytes([bytes[4], bytes[5], bytes[6], bytes[7]]);
let data = u32::from_le_bytes([bytes[40], bytes[41], bytes[42], bytes[43]]);
assert_eq!(data, expected_data_bytes);
assert_eq!(riff, RIFF_DATA_OVERHEAD as u32 + expected_data_bytes);
assert_eq!(bytes.len(), 44 + expected.len() * 2);
assert_eq!(samples, expected);
let _ = std::fs::remove_dir_all(&dir);
}
#[test] #[test]
fn wav_header_round_trips_sizes() { fn wav_header_round_trips_sizes() {
let dir = std::env::temp_dir(); let dir = std::env::temp_dir();
@@ -300,38 +414,32 @@ mod tests {
#[test] #[test]
fn mic_is_summed_with_mix_when_present() { fn mic_is_summed_with_mix_when_present() {
let dir = std::env::temp_dir(); let mut mic_fifo = VecDeque::from([1000, 2000, 3000]);
let mut r = Recorder {
writer: WavWriter::new(&dir.join(format!("ps-sum-{}.wav", std::process::id()))) let first = mix_with_mic(&[10, 20], &mut mic_fifo);
.unwrap(), assert_eq!(first, vec![1010, 2020]);
mic_fifo: VecDeque::new(), assert_eq!(mic_fifo.len(), 1, "two samples consumed, one mic left");
path: PathBuf::new(),
}; let second = mix_with_mic(&[0, 0], &mut mic_fifo);
r.push_mic(&[1000, 2000, 3000]); assert_eq!(second, vec![3000, 0]);
// write_frame pops mic per-sample and sums; we can't read the file mid-stream,
// so assert the FIFO drains exactly by frame length.
r.write_frame(&[10, 20]).unwrap();
assert_eq!(r.mic_fifo.len(), 1, "two samples consumed, one mic left");
r.write_frame(&[0, 0]).unwrap();
assert_eq!( assert_eq!(
r.mic_fifo.len(), mic_fifo.len(),
0, 0,
"remaining mic sample consumed; rest is silence" "remaining mic sample consumed; rest is silence"
); );
let _ = r.finalize();
} }
#[test] #[test]
fn mic_fifo_is_capped() { fn mic_fifo_is_capped() {
let dir = std::env::temp_dir(); let dir = unique_temp_dir("peerspeak-cap");
let mut r = Recorder { let _ = std::fs::remove_dir_all(&dir);
writer: WavWriter::new(&dir.join(format!("ps-cap-{}.wav", std::process::id()))) std::fs::create_dir_all(&dir).unwrap();
.unwrap(),
mic_fifo: VecDeque::new(), let mut r = Recorder::create(&dir, 1_700_000_001).unwrap();
path: PathBuf::new(),
};
r.push_mic(&vec![5i16; MAX_MIC_FIFO * 2]); r.push_mic(&vec![5i16; MAX_MIC_FIFO * 2]);
assert_eq!(r.mic_fifo.len(), MAX_MIC_FIFO, "FIFO is bounded to the cap"); assert_eq!(r.mic_fifo.len(), MAX_MIC_FIFO, "FIFO is bounded to the cap");
let _ = r.finalize(); r.finalize().unwrap();
let _ = std::fs::remove_dir_all(&dir);
} }
} }