12 Commits
Author SHA1 Message Date
mollusk 551767f9f5 Show build version on launch screen
CI / check (push) Failing after 37s
cargo-deny / cargo-deny (push) Has been cancelled
windows-build / windows-build (push) Has been cancelled
2026-06-29 16:54:26 -04:00
mollusk fa90cd3ce9 Update Arch package version 2026-06-29 16:53:03 -04:00
molluskandClaude Opus 4.8 660261a9a5 deps: bump memmap2 0.9.10 -> 0.9.11 (clears RUSTSEC-2026-0186)
CI / check (push) Successful in 2m6s
cargo-deny / cargo-deny (push) Has been cancelled
windows-build / windows-build (push) Has been cancelled
cargo-audit flagged memmap2 0.9.10 as unsound (RUSTSEC-2026-0186, unchecked
pointer offset); 0.9.11 is the patched release. Warning-level only (audit/deny
don't fail on it), but cheap to clear. Audit now down to the two deliberately
-accepted unmaintained warnings (audiopus_sys, paste; ignored in deny.toml).
Lockfile-only.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-29 16:06:19 -04:00
molluskandClaude Opus 4.8 3a74fd0230 ci: drop concurrency block (Gitea 1.26 dropped runs with it set)
CI / check (push) Failing after 12m11s
cargo-deny / cargo-deny (push) Has been cancelled
windows-build / windows-build (push) Has been cancelled
A push that only changed Cargo.lock failed to create any Actions run while the
concurrency group was present; removing it restores reliable push triggering.
Single-dev CI doesn't need run-cancellation.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-29 16:02:15 -04:00
molluskandClaude Opus 4.8 2dbb1ea316 deps: bump anyhow 1.0.102 -> 1.0.103 (fixes RUSTSEC-2026-0190)
CI / check (push) Failing after 12m46s
cargo-deny / cargo-deny (push) Has been cancelled
windows-build / windows-build (push) Has been cancelled
CI's cargo-deny flagged RUSTSEC-2026-0190: unsoundness in anyhow's
Error::downcast_mut() (UB via borrow-rule violation after Error::context),
reached transitively (n0-error / iroh + the image/rav1e chain). 1.0.103 is the
patched release; lockfile-only, no API change. cargo deny check now fully clean.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-29 15:56:35 -04:00
molluskandClaude Opus 4.8 c902db2e90 style: rustfmt the 0.6.2 additions (A17b + version-in-UI)
CI / check (push) Failing after 2m34s
cargo-deny / cargo-deny (push) Has been cancelled
windows-build / windows-build (push) Has been cancelled
CI's fmt --check caught that Codex's hand-written additions in these two files
weren't rustfmt-formatted (the senior gate ran clippy + tests but not
fmt --check). Pure line-wrapping, no logic change. Keeps the crate fmt-clean.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-29 15:52:31 -04:00
molluskandClaude Opus 4.8 83e5881768 ci: cancel superseded in-progress runs (concurrency group)
CI / check (push) Failing after 7s
cargo-deny / cargo-deny (push) Has been cancelled
windows-build / windows-build (push) Has been cancelled
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-29 15:51:31 -04:00
molluskandClaude Opus 4.8 8424b44dec ci: add Gitea Actions workflow (self-hosted host-mode runner)
CI / check (push) Failing after 13s
cargo-deny / cargo-deny (push) Has been cancelled
windows-build / windows-build (push) Has been cancelled
CI runs on a self-hosted host-mode gitea-runner on the desktop (label `arch`),
so the cheap gitbutter VPS only queues jobs while all compile/test compute runs
locally. Pipeline on push-to-main / PR / manual dispatch: cargo fmt --check,
clippy --all-targets -D warnings, cargo test --all-targets + doc tests, cargo
deny check, cargo audit.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-29 15:43:35 -04:00
molluskandClaude Opus 4.8 33e49a8ca7 Merge 0.6.2 refinements (Track B): A17b multitrack offload + build-version-in-UI
Track B code body for the 0.6.2 patch release. No wire change (GOSSIP_PROTO stays
5, interoperable with 0.6.0/0.6.1). Two code commits + two investigation closeouts
(A6 root-caused -> deferred to W5; A3 palette audit -> accepted as-is).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-29 15:28:48 -04:00
molluskandClaude Opus 4.8 393c1c7f09 feat(ui): surface the build version in Settings + at startup
A field build is now self-identifying. `run_gui` logs `PeerSpeak v<version>
starting` (from env!("CARGO_PKG_VERSION")) on launch, and the Settings panel
shows a muted `PeerSpeak v<version>` footer — pinned to the bottom of the
220px category sidebar (wide layout) and appended under the body in the narrow
(<820px) layout. Compile-time string, no new test, no deps, local-only.

Renders the current crate version, so it tracks the Cargo.toml bump at each
release cut (shows v0.6.1 until 0.6.2 is stamped in Track A).

Codex-implemented (gpt-5.5 xhigh), senior-reviewed.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-29 14:58:33 -04:00
molluskandClaude Opus 4.8 1bf14ba08f perf(recording): move multitrack stem disk I/O off the mixer path (A17b)
The multitrack recorder wrote every per-stem WAV frame (and the potentially
large late-joiner back-pad) inline on the caller thread while holding the
recorder mutex, so a slow/contended disk stalled the playout mixer (local
underruns) and the events loop. This is the multitrack counterpart to A17
(e0325d4), which moved the single-file recorder's writes off the mixer path.

Design: the front (MultitrackRecorder) now keeps only cheap in-memory state
(known-peer set, mic FIFO, a pending-cycle builder) and on each end_cycle
assembles ONE whole-cycle batch (new peers + mic frame + optional mix frame +
the map of peer frames written this cycle) and try_sends it over a bounded
sync_channel(256) to a dedicated writer thread. The writer thread owns every
WavWriter, is authoritative for its own cycle count, back-pads a brand-new
peer by cycles_written*frame_samples, fills absent peer/mix frames with
silence, latches the first write/create error then drains, and finalizes all
headers on channel close.

The unit of hand-off is a whole cycle, not a track: the writer appends exactly
frame_samples to every existing track per applied batch, and a full queue
DROPS the entire batch (counted + logged at 1 and every 256). So a dropped
cycle omits the same 20ms from every stem at once and all tracks stay
equal-length and sample-aligned by construction even under disk back-pressure.
On drop the batch's new-peer announcements are rolled back out of the known set
so they re-announce (and correctly re-back-pad) on the next applied cycle.

Public method signatures are unchanged -> zero core/mod.rs edits. The
WAV/file format is unchanged (no wire/on-disk change), no new deps
(std::sync::mpsc + std::thread, as A17). Writer logic is factored behind a
generic SampleWriter seam so the apply-batch alignment invariant is unit-tested
without spawning the thread; new tests cover the back-pad-on-apply invariant,
the dropped-cycle equal-length property, and async create-error surfacing at
finalize. The three existing end-to-end tests pass unchanged (now exercising
the threaded path). 496 lib tests, clippy --all-targets clean, release builds.

Codex-implemented (gpt-5.5 xhigh), senior-reviewed.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-29 14:53:52 -04:00
molluskandClaude Opus 4.8 6f14d2668d docs(protocol): correct GOSSIP_PROTO version mapping (v5 = 0.6.0, not 0.7.0)
cargo-deny / cargo-deny (push) Has been cancelled
windows-build / windows-build (push) Has been cancelled
The const comment claimed "v5 (0.7.0)" while GOSSIP_PROTO has been 5 since the
v0.6.0 tag (introduced by bca2ccd, "release 0.6.0"). Git confirms the value went
straight 3 -> 5 in that one release and a GOSSIP_PROTO == 4 build never existed.
Merge the two mislabeled v4/v5 bullets into one accurate v4-v5 (0.6.0) entry and
note the 3->5 jump + that this breaking gossip change correctly rode the
0.5.1 -> 0.6.0 MINOR bump per VERSIONING.md (0.6.1 is a wire-compatible PATCH,
still proto 5). Comment-only; no wire/behavior change.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-29 02:24:53 -04:00
7 changed files with 526 additions and 95 deletions
+42
View File
@@ -0,0 +1,42 @@
name: CI
# Runs on the self-hosted host-mode runner on the desktop (label `arch`). The
# gitbutter VPS only queues the job; all compile/test compute happens locally.
on:
push:
branches: [main]
pull_request:
workflow_dispatch:
jobs:
check:
runs-on: arch
steps:
- name: Checkout
uses: actions/checkout@v4
- name: Toolchain versions
run: |
rustc --version
cargo --version
cargo clippy --version
cargo deny --version
cargo audit --version
- name: Format check
run: cargo fmt --all -- --check
- name: Clippy (all targets, warnings as errors)
run: cargo clippy --all-targets -- -D warnings
- name: Tests
run: cargo test --all-targets
- name: Doc tests
run: cargo test --doc
- name: cargo-deny (advisories, bans, licenses, sources)
run: cargo deny check
- name: cargo-audit
run: cargo audit
Generated
+4 -4
View File
@@ -200,9 +200,9 @@ checksum = "940b3a0ca603d1eade50a4846a2afffd5ef57a9feac2c0e2ec2e14f9ead76000"
[[package]]
name = "anyhow"
version = "1.0.102"
version = "1.0.103"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7f202df86484c868dbad7eaa557ef785d5c66295e41b460ef922eca0723b842c"
checksum = "2a4385e2e34eb35d6b3efe798b9eb88096925d87726c0798709bf56d9ed84af3"
[[package]]
name = "arbitrary"
@@ -3682,9 +3682,9 @@ checksum = "6b947ae49db0d222b1dbc6b113ce7248a3fc3a6ca21b696717bfc000ba4484d8"
[[package]]
name = "memmap2"
version = "0.9.10"
version = "0.9.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "714098028fe011992e1c3962653c96b2d578c4b4bce9036e15ff220319b1e0e3"
checksum = "d1219ed1b7f229ee7104d281dd01d6802fe28bb6e95d292942c4daacdeb798c0"
dependencies = [
"libc",
]
+22
View File
@@ -0,0 +1,22 @@
use std::process::Command;
fn main() {
println!("cargo:rerun-if-changed=.git/HEAD");
if let Ok(head) = std::fs::read_to_string(".git/HEAD") {
if let Some(reference) = head.strip_prefix("ref: ") {
println!("cargo:rerun-if-changed=.git/{}", reference.trim());
}
}
let short = Command::new("git")
.args(["rev-parse", "--short=8", "HEAD"])
.output()
.ok()
.filter(|output| output.status.success())
.and_then(|output| String::from_utf8(output.stdout).ok())
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty())
.unwrap_or_else(|| "unknown".to_string());
println!("cargo:rustc-env=PEERSPEAK_GIT_SHORT={short}");
}
+1 -1
View File
@@ -1,7 +1,7 @@
# Maintainer: mollusk <jitty+lc1iz0dc@protonmail.com>
pkgname=peerspeak-git
_pkgname=peerspeak
pkgver=0.5.0.r0.g0000000
pkgver=0.6.1.r310.g660261a
pkgrel=1
pkgdesc="Decentralized peer-to-peer voice chat (Rust/iroh/PipeWire/Opus/iced)"
arch=('x86_64')
+19
View File
@@ -34,6 +34,17 @@ use tokio::sync::Mutex;
static UI_RX: OnceLock<Mutex<Option<tokio::sync::mpsc::Receiver<UiEvent>>>> = OnceLock::new();
const APP_VERSION: &str = env!("CARGO_PKG_VERSION");
const GIT_SHORT: &str = env!("PEERSPEAK_GIT_SHORT");
fn app_build_label() -> String {
if GIT_SHORT == "unknown" {
format!("PeerSpeak v{APP_VERSION}")
} else {
format!("PeerSpeak v{APP_VERSION} ({GIT_SHORT})")
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Screen {
Home,
@@ -1100,6 +1111,7 @@ fn effective_background_bytes(
}
pub fn run_gui() -> iced::Result {
crate::log_msg(&format!("{} starting", app_build_label()));
// Restore the last window size (saved on close). Position is restored too,
// but only on X11 — Wayland's xdg-shell gives clients no way to set their own
// position, so we center there and leave placement to the compositor.
@@ -3806,6 +3818,7 @@ fn connect_card(state: &AppState) -> Element<'_, AppMessage> {
let subtitle = text("NAT-traversing full-mesh voice chat")
.size(16)
.color(color_subtext);
let build_label = text(app_build_label()).size(11).color(color_subtext);
let nickname_input = column![
text("Nickname").size(14).color(color_subtext),
@@ -3861,6 +3874,7 @@ fn connect_card(state: &AppState) -> Element<'_, AppMessage> {
column![
logo,
subtitle,
build_label,
vertical_space(20.0),
nickname_input,
vertical_space(16.0),
@@ -5240,6 +5254,9 @@ fn view(state: &AppState) -> Element<'_, AppMessage> {
for category in SettingsCategory::ALL {
settings_nav = settings_nav.push(category_button(category));
}
settings_nav = settings_nav
.push(iced::widget::Space::new().height(iced::Length::Fill))
.push(text(app_build_label()).size(11).color(color_subtext));
let settings_nav = container(settings_nav)
.padding(12)
.width(iced::Length::Fixed(220.0))
@@ -5258,6 +5275,8 @@ fn view(state: &AppState) -> Element<'_, AppMessage> {
.width(iced::Length::Fill),
vertical_space(10.0),
settings_body,
vertical_space(10.0),
text(app_build_label()).size(11).color(color_subtext),
]
.spacing(8)
.width(iced::Length::Fill),
+428 -84
View File
@@ -12,9 +12,11 @@
//! This module is pure plumbing over [`WavWriter`]: no audio decode, no
//! networking, no realtime work. The mixer (a non-RT task) drives it.
use std::collections::{HashMap, VecDeque};
use std::collections::{HashMap, HashSet, VecDeque};
use std::io;
use std::path::{Path, PathBuf};
use std::sync::mpsc::{self, SyncSender, TrySendError};
use std::thread::{self, JoinHandle};
use iroh::EndpointId;
@@ -24,6 +26,8 @@ use crate::core::jitter::FRAME_SAMPLES;
/// Cap on the silence chunk written at once when pre-padding a late joiner, so a
/// long-running call can't trigger a single multi-hundred-MB allocation.
const SILENCE_CHUNK: usize = FRAME_SAMPLES * 256;
const WRITER_QUEUE_CYCLES: usize = 256;
const DROP_LOG_INTERVAL_CYCLES: u64 = 256;
/// Cap on buffered mic samples (~200ms @ 48kHz). Bounds how far the mic track
/// can drift if the capture clock runs ahead of the mixer cycle; past it the
@@ -56,41 +60,6 @@ pub fn create_session_dir(base: &Path, now_unix_secs: u64) -> io::Result<PathBuf
))
}
/// One output track: its WAV writer plus whether it has been written *this*
/// cycle (so `end_cycle` knows which tracks to pad with silence).
struct Track {
writer: WavWriter,
written_this_cycle: bool,
}
impl Track {
fn create(path: &Path) -> io::Result<Self> {
Ok(Self {
writer: WavWriter::new(path)?,
written_this_cycle: false,
})
}
/// Append `frame` fitted to exactly `frame_samples` (zero-padded if short),
/// and mark the track as written for this cycle.
fn write_frame(&mut self, frame: &[i16], frame_samples: usize) -> io::Result<()> {
self.writer.write_samples(&fit(frame, frame_samples))?;
self.written_this_cycle = true;
Ok(())
}
/// Append `samples` of silence (no cycle-marking — used for padding).
fn write_silence(&mut self, samples: usize) -> io::Result<()> {
let mut remaining = samples;
while remaining > 0 {
let n = remaining.min(SILENCE_CHUNK);
self.writer.write_samples(&vec![0i16; n])?;
remaining -= n;
}
Ok(())
}
}
/// Return `frame` resized to exactly `n` samples: truncated if longer (shouldn't
/// happen — Opus frames are uniform), zero-padded if shorter.
fn fit(frame: &[i16], n: usize) -> Vec<i16> {
@@ -126,41 +95,210 @@ pub fn track_filename(name: &str, id: &EndpointId) -> String {
format!("{slug}-{short}.wav")
}
/// A live multitrack recording: per-peer stems + your mic, plus an optional
/// mixed track, all under one session directory and clocked together.
pub struct MultitrackRecorder {
dir: PathBuf,
frame_samples: usize,
/// Cycles recorded so far = the shared length (in frames) of every track.
cycles: u64,
peers: HashMap<EndpointId, Track>,
/// Your mic track. Fed asynchronously from the capture thread via
/// [`push_mic`](MultitrackRecorder::push_mic) into `mic_fifo`, then drained
/// one frame per `end_cycle` so it aligns with the cycle clock.
mic: WavWriter,
mic_fifo: VecDeque<i16>,
/// Present in "Both" mode (stems + mixed), absent in "stems only".
mix: Option<Track>,
#[derive(Default)]
struct PendingCycle {
new_peers: Vec<NewPeer>,
peer_frames: HashMap<EndpointId, Vec<i16>>,
mix_frame: Option<Vec<i16>>,
}
impl MultitrackRecorder {
/// Create a recording in `dir` (which must already exist). `with_mix` adds
/// the convenience mixed track (`mix.wav`). Your mic is always `me.wav`.
pub fn create(dir: &Path, frame_samples: usize, with_mix: bool) -> io::Result<Self> {
struct NewPeer {
id: EndpointId,
filename: String,
}
struct CycleBatch {
new_peers: Vec<NewPeer>,
mic_frame: Vec<i16>,
mix_frame: Option<Vec<i16>>,
peer_frames: HashMap<EndpointId, Vec<i16>>,
}
trait SampleWriter {
fn write_samples(&mut self, samples: &[i16]) -> io::Result<()>;
fn finalize(self) -> io::Result<()>;
}
impl SampleWriter for WavWriter {
fn write_samples(&mut self, samples: &[i16]) -> io::Result<()> {
WavWriter::write_samples(self, samples)
}
fn finalize(self) -> io::Result<()> {
WavWriter::finalize(self)
}
}
struct WriterState<W> {
dir: PathBuf,
frame_samples: usize,
peers: HashMap<EndpointId, W>,
mic: W,
mix: Option<W>,
cycles_written: u64,
}
impl WriterState<WavWriter> {
fn create(dir: &Path, frame_samples: usize, with_mix: bool) -> io::Result<Self> {
let mic = WavWriter::new(&dir.join("me.wav"))?;
let mix = if with_mix {
Some(Track::create(&dir.join("mix.wav"))?)
Some(WavWriter::new(&dir.join("mix.wav"))?)
} else {
None
};
Ok(Self {
dir: dir.to_path_buf(),
frame_samples,
cycles: 0,
peers: HashMap::new(),
mic,
mic_fifo: VecDeque::new(),
mix,
cycles_written: 0,
})
}
}
impl<W: SampleWriter> WriterState<W> {
fn apply_batch<F>(&mut self, batch: &CycleBatch, mut create_peer: F) -> io::Result<()>
where
F: FnMut(&Path) -> io::Result<W>,
{
for peer in &batch.new_peers {
if !self.peers.contains_key(&peer.id) {
let writer = create_peer(&self.dir.join(&peer.filename))?;
self.peers.insert(peer.id, writer);
let pad = self.back_pad_samples()?;
let writer = self.peers.get_mut(&peer.id).unwrap();
Self::write_silence(writer, pad)?;
}
}
self.mic.write_samples(&batch.mic_frame)?;
if let Some(mix) = self.mix.as_mut() {
if let Some(frame) = batch.mix_frame.as_deref() {
mix.write_samples(frame)?;
} else {
Self::write_silence(mix, self.frame_samples)?;
}
}
let silence = vec![0i16; self.frame_samples];
for (id, writer) in &mut self.peers {
let frame = batch
.peer_frames
.get(id)
.map(Vec::as_slice)
.unwrap_or(&silence);
writer.write_samples(frame)?;
}
self.cycles_written += 1;
Ok(())
}
fn back_pad_samples(&self) -> io::Result<usize> {
let cycles = usize::try_from(self.cycles_written)
.map_err(|_| io::Error::other("multitrack recording too long"))?;
cycles
.checked_mul(self.frame_samples)
.ok_or_else(|| io::Error::other("multitrack recording too long"))
}
fn write_silence(writer: &mut W, samples: usize) -> io::Result<()> {
let mut remaining = samples;
let silence = vec![0i16; remaining.min(SILENCE_CHUNK)];
while remaining > 0 {
let n = remaining.min(silence.len());
writer.write_samples(&silence[..n])?;
remaining -= n;
}
Ok(())
}
fn finalize(self) -> io::Result<()> {
let mut first_finalize_error = None;
record_first_error(&mut first_finalize_error, self.mic.finalize());
if let Some(mix) = self.mix {
record_first_error(&mut first_finalize_error, mix.finalize());
}
for writer in self.peers.into_values() {
record_first_error(&mut first_finalize_error, writer.finalize());
}
if let Some(e) = first_finalize_error {
Err(e)
} else {
Ok(())
}
}
}
fn record_first_error(slot: &mut Option<io::Error>, result: io::Result<()>) {
if slot.is_none()
&& let Err(e) = result
{
*slot = Some(e);
}
}
/// Applies whole-cycle batches on the writer thread. Each applied batch appends
/// exactly `frame_samples` to every existing track, and a dropped batch never
/// reaches this loop for any track, so stem lengths stay equal even when the
/// bounded queue applies back-pressure.
fn writer_thread_main(
mut state: WriterState<WavWriter>,
batch_rx: mpsc::Receiver<CycleBatch>,
) -> io::Result<()> {
let mut first_write_error = None;
for batch in batch_rx {
if first_write_error.is_none()
&& let Err(e) = state.apply_batch(&batch, WavWriter::new)
{
first_write_error = Some(e);
}
}
let finalize_result = state.finalize();
if let Some(e) = first_write_error {
Err(e)
} else {
finalize_result
}
}
/// A live multitrack recording: per-peer stems + your mic, plus an optional
/// mixed track, all under one session directory and clocked together.
pub struct MultitrackRecorder {
dir: PathBuf,
frame_samples: usize,
known_peers: HashSet<EndpointId>,
/// Your mic track. Fed asynchronously from the capture thread via
/// [`push_mic`](MultitrackRecorder::push_mic) into `mic_fifo`, then drained
/// one frame per `end_cycle` so it aligns with the cycle clock.
mic_fifo: VecDeque<i16>,
/// Present in "Both" mode (stems + mixed), absent in "stems only".
with_mix: bool,
batch_tx: SyncSender<CycleBatch>,
writer_thread: JoinHandle<io::Result<()>>,
dropped_cycles: u64,
pending: PendingCycle,
}
impl MultitrackRecorder {
/// Create a recording in `dir` (which must already exist). `with_mix` adds
/// the convenience mixed track (`mix.wav`). Your mic is always `me.wav`.
pub fn create(dir: &Path, frame_samples: usize, with_mix: bool) -> io::Result<Self> {
let writer_state = WriterState::create(dir, frame_samples, with_mix)?;
let (batch_tx, batch_rx) = mpsc::sync_channel(WRITER_QUEUE_CYCLES);
let writer_thread = thread::spawn(move || writer_thread_main(writer_state, batch_rx));
Ok(Self {
dir: dir.to_path_buf(),
frame_samples,
known_peers: HashSet::new(),
mic_fifo: VecDeque::new(),
with_mix,
batch_tx,
writer_thread,
dropped_cycles: 0,
pending: PendingCycle::default(),
})
}
@@ -173,12 +311,14 @@ impl MultitrackRecorder {
/// so it aligns with the others. Idempotent: a peer already tracked is left
/// as-is (re-announce / name change doesn't restart their file).
pub fn add_peer(&mut self, id: EndpointId, name: &str) -> io::Result<()> {
if self.peers.contains_key(&id) {
if self.known_peers.contains(&id) {
return Ok(());
}
let mut track = Track::create(&self.dir.join(track_filename(name, &id)))?;
track.write_silence(self.cycles as usize * self.frame_samples)?;
self.peers.insert(id, track);
self.known_peers.insert(id);
self.pending.new_peers.push(NewPeer {
id,
filename: track_filename(name, &id),
});
Ok(())
}
@@ -186,11 +326,13 @@ impl MultitrackRecorder {
/// registered yet (write raced ahead of the join event), auto-register it
/// with an id-only name so no audio is dropped.
pub fn write_peer(&mut self, id: EndpointId, frame: &[i16]) -> io::Result<()> {
if !self.peers.contains_key(&id) {
if !self.known_peers.contains(&id) {
self.add_peer(id, "")?;
}
let fs = self.frame_samples;
self.peers.get_mut(&id).unwrap().write_frame(frame, fs)
self.pending
.peer_frames
.insert(id, fit(frame, self.frame_samples));
Ok(())
}
/// Buffer a frame of your transmitted mic audio (called from the capture
@@ -215,9 +357,8 @@ impl MultitrackRecorder {
/// Record the finished mixed-bus frame for the current cycle (no-op in
/// stems-only mode).
pub fn write_mix(&mut self, frame: &[i16]) -> io::Result<()> {
let fs = self.frame_samples;
if let Some(mix) = self.mix.as_mut() {
mix.write_frame(frame, fs)?;
if self.with_mix {
self.pending.mix_frame = Some(fit(frame, self.frame_samples));
}
Ok(())
}
@@ -230,28 +371,63 @@ impl MultitrackRecorder {
// Mic: always one frame per cycle, drained from the FIFO (silence on
// underrun), so it tracks the cycle clock like the peer stems.
let mic_frame = self.drain_mic(fs);
self.mic.write_samples(&mic_frame)?;
// Peers + the optional mix track: pad any not written this cycle.
for track in self.peers.values_mut().chain(self.mix.as_mut()) {
if !track.written_this_cycle {
track.write_silence(fs)?;
let mut pending = std::mem::take(&mut self.pending);
pending.new_peers.sort_by(|a, b| {
a.filename
.cmp(&b.filename)
.then_with(|| a.id.to_string().cmp(&b.id.to_string()))
});
let batch = CycleBatch {
new_peers: pending.new_peers,
mic_frame,
mix_frame: if self.with_mix {
pending.mix_frame
} else {
None
},
peer_frames: pending.peer_frames,
};
match self.batch_tx.try_send(batch) {
Ok(()) => Ok(()),
Err(TrySendError::Full(batch)) => {
for peer in &batch.new_peers {
self.known_peers.remove(&peer.id);
}
self.dropped_cycles = self.dropped_cycles.saturating_add(1);
if self.dropped_cycles == 1
|| self.dropped_cycles.is_multiple_of(DROP_LOG_INTERVAL_CYCLES)
{
crate::log_msg(&format!(
"multitrack recording: writer queue full; dropped {} cycle(s)",
self.dropped_cycles
));
}
Ok(())
}
track.written_this_cycle = false;
Err(TrySendError::Disconnected(_)) => Err(io::Error::new(
io::ErrorKind::BrokenPipe,
"multitrack writer thread stopped",
)),
}
self.cycles += 1;
Ok(())
}
/// Finalize every track's WAV header. Consumes the recorder.
pub fn finalize(self) -> io::Result<()> {
self.mic.finalize()?;
if let Some(mix) = self.mix {
mix.writer.finalize()?;
}
for (_, track) in self.peers {
track.writer.finalize()?;
}
Ok(())
let Self {
dir: _,
frame_samples: _,
known_peers: _,
mic_fifo: _,
with_mix: _,
batch_tx,
writer_thread,
dropped_cycles: _,
pending: _,
} = self;
drop(batch_tx);
writer_thread
.join()
.unwrap_or_else(|_| Err(io::Error::other("multitrack writer thread panicked")))
}
}
@@ -277,6 +453,51 @@ mod tests {
d
}
#[derive(Default)]
struct TestWriter {
samples: Vec<i16>,
}
impl SampleWriter for TestWriter {
fn write_samples(&mut self, samples: &[i16]) -> io::Result<()> {
self.samples.extend_from_slice(samples);
Ok(())
}
fn finalize(self) -> io::Result<()> {
Ok(())
}
}
fn test_writer_state(frame_samples: usize, with_mix: bool) -> WriterState<TestWriter> {
WriterState {
dir: PathBuf::new(),
frame_samples,
peers: HashMap::new(),
mic: TestWriter::default(),
mix: if with_mix {
Some(TestWriter::default())
} else {
None
},
cycles_written: 0,
}
}
fn test_batch(
new_peers: Vec<NewPeer>,
mic_frame: Vec<i16>,
mix_frame: Option<Vec<i16>>,
peer_frames: Vec<(EndpointId, Vec<i16>)>,
) -> CycleBatch {
CycleBatch {
new_peers,
mic_frame,
mix_frame,
peer_frames: peer_frames.into_iter().collect(),
}
}
#[test]
fn fit_pads_and_truncates() {
assert_eq!(fit(&[1, 2], 4), vec![1, 2, 0, 0]);
@@ -311,6 +532,110 @@ mod tests {
let _ = std::fs::remove_dir_all(&base);
}
#[test]
fn apply_batch_advances_existing_tracks_and_back_pads_late_peer() {
let frame = 3;
let early = an_id();
let late = an_id();
let mut state = test_writer_state(frame, true);
state.cycles_written = 2;
state.mic.samples = vec![8; 2 * frame];
state.mix.as_mut().unwrap().samples = vec![6; 2 * frame];
state.peers.insert(
early,
TestWriter {
samples: vec![1; 2 * frame],
},
);
let batch = test_batch(
vec![NewPeer {
id: late,
filename: "late.wav".to_string(),
}],
vec![9; frame],
None,
vec![(early, vec![2; frame]), (late, vec![7; frame])],
);
state
.apply_batch(&batch, |_| Ok(TestWriter::default()))
.unwrap();
assert_eq!(state.cycles_written, 3);
assert_eq!(state.mic.samples.len(), 3 * frame);
assert_eq!(state.mix.as_ref().unwrap().samples.len(), 3 * frame);
assert_eq!(
&state.mix.as_ref().unwrap().samples[2 * frame..],
&[0, 0, 0]
);
assert_eq!(state.peers.get(&early).unwrap().samples.len(), 3 * frame);
assert_eq!(
&state.peers.get(&early).unwrap().samples[2 * frame..],
&[2, 2, 2]
);
assert_eq!(
state.peers.get(&late).unwrap().samples,
vec![0, 0, 0, 0, 0, 0, 7, 7, 7],
"late peer is back-padded by completed cycles before this batch"
);
}
#[test]
fn skipped_batches_keep_all_tracks_equal_length() {
let frame = 2;
let p1 = an_id();
let p2 = an_id();
let mut state = test_writer_state(frame, true);
let first = test_batch(
vec![
NewPeer {
id: p1,
filename: "p1.wav".to_string(),
},
NewPeer {
id: p2,
filename: "p2.wav".to_string(),
},
],
vec![1; frame],
Some(vec![5; frame]),
vec![(p1, vec![10; frame]), (p2, vec![20; frame])],
);
state
.apply_batch(&first, |_| Ok(TestWriter::default()))
.unwrap();
let _dropped_cycle = test_batch(
Vec::new(),
vec![2; frame],
Some(vec![6; frame]),
vec![(p1, vec![11; frame])],
);
let after_drop = test_batch(
Vec::new(),
vec![3; frame],
None,
vec![(p1, vec![12; frame])],
);
state
.apply_batch(&after_drop, |_| Ok(TestWriter::default()))
.unwrap();
let expected = 2 * frame;
assert_eq!(state.cycles_written, 2);
assert_eq!(state.mic.samples.len(), expected);
assert_eq!(state.mix.as_ref().unwrap().samples.len(), expected);
assert_eq!(state.peers.get(&p1).unwrap().samples.len(), expected);
assert_eq!(state.peers.get(&p2).unwrap().samples.len(), expected);
assert_eq!(
&state.peers.get(&p2).unwrap().samples[frame..],
&[0, 0],
"peer absent from an applied batch gets silence for that cycle"
);
}
#[test]
fn all_tracks_equal_length_after_n_cycles() {
let dir = tmpdir("equal");
@@ -406,4 +731,23 @@ mod tests {
"no mix track in stems-only mode"
);
}
#[cfg(unix)]
#[test]
fn async_peer_create_error_surfaces_at_finalize() {
use std::os::unix::fs::PermissionsExt;
let dir = tmpdir("asyncerr");
let mut rec = MultitrackRecorder::create(&dir, 4, false).unwrap();
std::fs::set_permissions(&dir, std::fs::Permissions::from_mode(0o500)).unwrap();
rec.add_peer(an_id(), "blocked").unwrap();
rec.end_cycle().unwrap();
let result = rec.finalize();
std::fs::set_permissions(&dir, std::fs::Permissions::from_mode(0o700)).unwrap();
let err = result.unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::PermissionDenied);
let _ = std::fs::remove_dir_all(&dir);
}
}
+10 -6
View File
@@ -33,12 +33,16 @@ pub const FRIENDS_PROTO: u32 = 1;
/// change is isolated into its own topic + signature domain so v2 and v3 peers
/// never share a swarm. Resync everyone, exactly like the W4 avatar bump.
///
/// v4 (0.6.0): `PeerState` gained an optional `music` presence field carrying a
/// current shared-listening track descriptor and playback timeline. Bytes still
/// ride the files plane by id; gossip carries only the descriptor/timeline.
///
/// v5 (0.7.0): `MusicPresence` gained optional prefetch hints for the next
/// track so tuned-in listeners can fetch it before the DJ advances.
/// v4v5 (0.6.0): the W22 shared-listening / music presence work. `PeerState`
/// gained an optional `music` presence field (a current shared-listening track
/// descriptor + playback timeline; the audio bytes still ride the files plane by
/// id, gossip carries only the descriptor/timeline), and `MusicPresence` then
/// gained optional prefetch hints for the next track so tuned-in listeners can
/// fetch it before the DJ advances. Both shipped together in the **0.6.0** release
/// (commit `bca2ccd`), where the const advanced straight `3 → 5`: there was never
/// a `GOSSIP_PROTO == 4` build — 4 is a skipped step. (Per `VERSIONING.md` this
/// breaking gossip change rode the `0.5.1 → 0.6.0` MINOR bump, so the discipline
/// was honoured; 0.6.1 is a wire-compatible PATCH on top, still proto 5.)
pub const GOSSIP_PROTO: u32 = 5;
/// File-transfer plane version (chat attachment request/stream shape). Bump on
/// any change. Mirrored in [`FILES_ALPN`].