Compare commits

...
Author SHA1 Message Date
mollusk 39b5dafd57 fix(audio): bound playback handoff queue 2026-07-01 13:39:10 -04:00
mollusk a78860db15 Merge W12 FEC/DTX follow-up (Codex, senior-reviewed)
CI / check (push) Failing after 23s
cargo-deny / cargo-deny (push) Has been cancelled
windows-build / windows-build (push) Has been cancelled
2026-06-30 16:56:19 -04:00
5f52aa1506 W12 follow-up: consume in-band FEC, drop redundant DTX
Fixes the two P2 efficacy findings from the Codex audit of the W12
profiles feature.

FEC was enabled on the encoder but never used: the jitter buffer's
loss path did pure PLC, so the redundancy was wasted bitrate. Now the
gap path reconstructs the lost frame from the next buffered packet via
Opus in-band FEC (new `AudioDecoder::decode_fec`, libopus decode with
fec=true into a one-frame buffer), keeping that packet for its own
normal decode and falling back to PLC if FEC decode fails. This is the
documented libopus FEC pattern; receiver-side only, no wire change.

DTX was enabled on BadNetwork but provided no benefit — the capture
noise gate already suppresses silence transmission, and the broadcast
DTX silence packets only created seq gaps that grew the jitter cushion.
All profiles now set dtx=false (plumbing kept for a future revisit).

Adds a jitter-buffer test proving FEC reconstruction beats pure PLC
(RMS error < 0.75x) and that the FEC source packet stays buffered.
500 lib tests, clippy + fmt clean, release build clean.

Co-Authored-By: Codex (gpt-5.5) <noreply@openai.com>
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-30 16:56:19 -04:00
molluskandClaude Opus 4.8 d92d0f6f6b Add W12 Opus/network quality profiles
CI / check (push) Failing after 25s
cargo-deny / cargo-deny (push) Has been cancelled
windows-build / windows-build (push) Has been cancelled
Add a small, named codec-policy picker (Low latency / Balanced / Bad
network) instead of exposing raw Opus knobs. The profile->params mapping
is a pure function (`codec::opus_impl::opus_params`) for unit testing;
profiles tune bitrate, in-band FEC, expected packet-loss, and DTX.

- config: `AudioProfile` enum (serde + Display + ALL + u8 round-trip),
  persisted `audio_profile` field (default Balanced).
- codec: `OpusParams` + pure `opus_params()` + `OpusEncoder::apply_params`
  / `apply_profile`.
- core: new `SetAudioProfile` command (Reliable, no coalesce); a shared
  `AtomicU8` lets the capture thread re-tune the live encoder on a
  mid-call switch and read it at each new call's encoder creation.
- app: Settings "Connection quality" picker in the Audio tab, startup
  config-sync send, and a one-line hint per profile.

No wire-format change (GOSSIP/audio planes untouched). 499 lib tests
green (config + codec mapping/apply tests added), clippy + fmt clean.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-30 16:13:26 -04:00
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
11 changed files with 494 additions and 29 deletions
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")
&& 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')
+49 -8
View File
@@ -4,7 +4,7 @@ use crate::audio::clip_player::{
};
use crate::audio::eq::{EQ_GAIN_DB_MAX, EQ_GAIN_DB_MIN, EqSettings};
use crate::audio::{AudioDevice, enumerate_audio_devices};
use crate::config::{AppConfig, NetworkMode, RecordingMode, RoomLayout};
use crate::config::{AppConfig, AudioProfile, NetworkMode, RecordingMode, RoomLayout};
use crate::core::{
CoreController,
messages::{CoreCommand, UiEvent},
@@ -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,
@@ -366,6 +377,7 @@ pub enum AppMessage {
/// immediately but does not persist (saved once on release via NoiseGateChanged).
NoiseGateDragging(f32),
NetworkModeSelected(NetworkMode),
AudioProfileSelected(AudioProfile),
RecordingModeSelected(RecordingMode),
/// Choose the friends presence posture (W7): invisible / normal / discoverable.
PresenceModeSelected(PresenceMode),
@@ -908,6 +920,7 @@ impl Default for AppState {
let _ = controller.send(CoreCommand::SetInputVolume(config.input_volume));
let _ = controller.send(CoreCommand::SetOutputVolume(config.output_volume));
let _ = controller.send(CoreCommand::SetNetworkMode(config.network_mode));
let _ = controller.send(CoreCommand::SetAudioProfile(config.audio_profile));
let _ = controller.send(CoreCommand::SetRecordingMode(config.recording_mode));
let _ = controller.send(CoreCommand::SetPixelpassPath(config.pixelpass_path.clone()));
let _ = controller.send(CoreCommand::SetPresenceMode(config.presence_mode));
@@ -1100,7 +1113,7 @@ fn effective_background_bytes(
}
pub fn run_gui() -> iced::Result {
crate::log_msg(&format!("PeerSpeak v{} starting", env!("CARGO_PKG_VERSION")));
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.
@@ -2039,6 +2052,12 @@ fn update(state: &mut AppState, message: AppMessage) -> Task<AppMessage> {
// Applied on the next join, since the endpoint is rebuilt then.
let _ = state.controller.send(CoreCommand::SetNetworkMode(mode));
}
AppMessage::AudioProfileSelected(profile) => {
state.config.audio_profile = profile;
state.config.save();
// Applies live to the running encoder, and to the next call.
let _ = state.controller.send(CoreCommand::SetAudioProfile(profile));
}
AppMessage::RecordingModeSelected(mode) => {
state.config.recording_mode = mode;
state.config.save();
@@ -3112,6 +3131,19 @@ fn network_mode_hint(mode: NetworkMode) -> &'static str {
}
}
/// One-line explanation of an audio/network profile for the settings picker (W12).
fn audio_profile_hint(profile: AudioProfile) -> &'static str {
match profile {
AudioProfile::LowLatency => {
"Lowest delay, no loss recovery. Best on a clean LAN or wired link."
}
AudioProfile::Balanced => "Default: voice quality with light loss recovery.",
AudioProfile::BadNetwork => {
"Most resilient on a lossy/congested link: heavier loss recovery, lower bitrate."
}
}
}
/// One-line explanation of a recording mode for the settings picker.
fn recording_mode_hint(mode: RecordingMode) -> &'static str {
match mode {
@@ -3807,6 +3839,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),
@@ -3862,6 +3895,7 @@ fn connect_card(state: &AppState) -> Element<'_, AppMessage> {
column![
logo,
subtitle,
build_label,
vertical_space(20.0),
nickname_input,
vertical_space(16.0),
@@ -4968,6 +5002,17 @@ fn view(state: &AppState) -> Element<'_, AppMessage> {
control
},
].spacing(8).width(iced::Length::Fill),
vertical_space(section_gap),
section_header("Connection quality"),
column![
pick_list(
&AudioProfile::ALL[..],
Some(state.config.audio_profile),
AppMessage::AudioProfileSelected,
).width(iced::Length::Fill),
text(audio_profile_hint(state.config.audio_profile)).size(11).color(color_subtext),
text("Applies immediately, even mid-call.").size(11).color(color_subtext),
].spacing(4).width(iced::Length::Fill),
]
.spacing(10)
.width(iced::Length::Fill)
@@ -5243,9 +5288,7 @@ fn view(state: &AppState) -> Element<'_, AppMessage> {
}
settings_nav = settings_nav
.push(iced::widget::Space::new().height(iced::Length::Fill))
.push(text(format!("PeerSpeak v{}", env!("CARGO_PKG_VERSION")))
.size(11)
.color(color_subtext));
.push(text(app_build_label()).size(11).color(color_subtext));
let settings_nav = container(settings_nav)
.padding(12)
.width(iced::Length::Fixed(220.0))
@@ -5265,9 +5308,7 @@ fn view(state: &AppState) -> Element<'_, AppMessage> {
vertical_space(10.0),
settings_body,
vertical_space(10.0),
text(format!("PeerSpeak v{}", env!("CARGO_PKG_VERSION")))
.size(11)
.color(color_subtext),
text(app_build_label()).size(11).color(color_subtext),
]
.spacing(8)
.width(iced::Length::Fill),
+13 -3
View File
@@ -380,7 +380,11 @@ impl MultitrackRecorder {
let batch = CycleBatch {
new_peers: pending.new_peers,
mic_frame,
mix_frame: if self.with_mix { pending.mix_frame } else { None },
mix_frame: if self.with_mix {
pending.mix_frame
} else {
None
},
peer_frames: pending.peer_frames,
};
match self.batch_tx.try_send(batch) {
@@ -560,9 +564,15 @@ mod tests {
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.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(&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],
+3
View File
@@ -20,6 +20,9 @@ pub trait AudioDecoder: Send {
/// If `compressed` is `None` (or `Some(&[])`), it indicates packet loss,
/// enabling the decoder to perform packet loss concealment (PLC).
fn decode(&mut self, compressed: Option<&[u8]>) -> Result<Vec<i16>, CodecError>;
/// Reconstructs the previous lost frame from the next packet's in-band FEC.
fn decode_fec(&mut self, next_payload: &[u8]) -> Result<Vec<i16>, CodecError>;
}
pub mod opus_impl;
+120 -1
View File
@@ -1,5 +1,49 @@
use crate::codec::{AudioDecoder, AudioEncoder, CodecError};
use opus::{Application, Channels, Decoder, Encoder};
use crate::config::AudioProfile;
use opus::{Application, Bitrate, Channels, Decoder, Encoder};
/// Concrete libopus encoder settings derived from an [`AudioProfile`]. Plain
/// data, so the profile→params mapping ([`opus_params`]) stays a pure,
/// unit-testable function (W12).
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct OpusParams {
/// Target bitrate in bits/sec.
pub bitrate: i32,
/// Enable in-band forward error correction (loss redundancy in the bitstream).
pub inband_fec: bool,
/// Expected packet-loss percentage (0..=100); tunes how much FEC libopus adds.
pub packet_loss_perc: i32,
/// Discontinuous transmission: stop sending during silence to save bandwidth.
pub dtx: bool,
}
/// Map a named profile to concrete Opus parameters. Pure — the W12 testable seam.
///
/// `BadNetwork` deliberately runs a *lower* bitrate than `Balanced`: in-band FEC
/// redundancy is carried inside the same bitstream, so trimming the base bitrate
/// leaves headroom for the redundancy on a congested link.
pub fn opus_params(profile: AudioProfile) -> OpusParams {
match profile {
AudioProfile::LowLatency => OpusParams {
bitrate: 24_000,
inband_fec: false,
packet_loss_perc: 0,
dtx: false,
},
AudioProfile::Balanced => OpusParams {
bitrate: 32_000,
inband_fec: true,
packet_loss_perc: 10,
dtx: false,
},
AudioProfile::BadNetwork => OpusParams {
bitrate: 20_000,
inband_fec: true,
packet_loss_perc: 25,
dtx: false,
},
}
}
pub struct OpusEncoder {
encoder: Encoder,
@@ -17,6 +61,29 @@ impl OpusEncoder {
.map_err(|e| CodecError::Init(format!("Failed to create Opus encoder: {}", e)))?;
Ok(Self { encoder })
}
/// Apply concrete codec parameters to the live encoder. Safe to call between
/// frames, so the user can switch profile mid-call.
pub fn apply_params(&mut self, params: &OpusParams) -> Result<(), CodecError> {
self.encoder
.set_bitrate(Bitrate::Bits(params.bitrate))
.map_err(|e| CodecError::Init(format!("set_bitrate: {}", e)))?;
self.encoder
.set_inband_fec(params.inband_fec)
.map_err(|e| CodecError::Init(format!("set_inband_fec: {}", e)))?;
self.encoder
.set_packet_loss_perc(params.packet_loss_perc)
.map_err(|e| CodecError::Init(format!("set_packet_loss_perc: {}", e)))?;
self.encoder
.set_dtx(params.dtx)
.map_err(|e| CodecError::Init(format!("set_dtx: {}", e)))?;
Ok(())
}
/// Apply a named [`AudioProfile`] (shorthand for `apply_params(&opus_params(p))`).
pub fn apply_profile(&mut self, profile: AudioProfile) -> Result<(), CodecError> {
self.apply_params(&opus_params(profile))
}
}
impl AudioEncoder for OpusEncoder {
@@ -95,12 +162,64 @@ impl AudioDecoder for OpusDecoder {
pcm.truncate(decoded_per_channel * channels_count);
Ok(pcm)
}
fn decode_fec(&mut self, next_payload: &[u8]) -> Result<Vec<i16>, CodecError> {
let channels_count = self.channels_count();
let mut pcm = vec![0i16; self.frame_samples * channels_count];
let decoded_per_channel = self
.decoder
.decode(next_payload, &mut pcm, true)
.map_err(|e| CodecError::Decode(format!("Opus FEC decoding failed: {}", e)))?;
pcm.truncate(decoded_per_channel * channels_count);
Ok(pcm)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_opus_params_mapping() {
let low = opus_params(AudioProfile::LowLatency);
let bal = opus_params(AudioProfile::Balanced);
let bad = opus_params(AudioProfile::BadNetwork);
// LowLatency has no loss redundancy; the other two do.
assert!(!low.inband_fec);
assert_eq!(low.packet_loss_perc, 0);
assert!(bal.inband_fec);
assert!(bad.inband_fec);
// Capture-side gating suppresses silence; no profile adds Opus DTX.
assert!(!low.dtx && !bal.dtx && !bad.dtx);
assert!(bad.packet_loss_perc > bal.packet_loss_perc);
// BadNetwork trims base bitrate to make room for FEC redundancy.
assert!(bad.bitrate < bal.bitrate);
// All bitrates are sane positive voice rates.
for p in [low, bal, bad] {
assert!(p.bitrate > 0 && p.bitrate <= 64_000);
assert!((0..=100).contains(&p.packet_loss_perc));
}
}
#[test]
fn test_apply_profile_sets_bitrate() {
let mut encoder = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap();
// Every profile applies cleanly to a real encoder...
for profile in AudioProfile::ALL {
encoder.apply_profile(profile).unwrap();
}
// ...and the last-applied bitrate is reflected by the encoder.
encoder.apply_profile(AudioProfile::Balanced).unwrap();
let want = opus_params(AudioProfile::Balanced).bitrate;
assert_eq!(encoder.encoder.get_bitrate().unwrap(), Bitrate::Bits(want));
}
#[test]
fn test_round_trip() {
let mut encoder = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap();
+91
View File
@@ -93,6 +93,60 @@ impl std::fmt::Display for RecordingMode {
}
}
/// Named Opus encoder / network-resilience policy (W12). The user picks a
/// profile instead of raw codec knobs; the concrete libopus parameters live in
/// `codec::opus_impl::opus_params`. Applies live to the running encoder.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
pub enum AudioProfile {
/// Lowest mouth-to-ear delay: modest bitrate, no FEC redundancy. Best on a
/// clean LAN / low-loss link where added latency matters more than loss.
LowLatency,
/// Sensible default: voice bitrate with in-band FEC for light packet loss.
#[default]
Balanced,
/// Maximum resilience on a lossy/congested link: in-band FEC tuned for heavy
/// loss, at a lower bitrate to leave headroom for the redundancy.
BadNetwork,
}
impl AudioProfile {
/// All variants, in picker display order.
pub const ALL: [AudioProfile; 3] = [
AudioProfile::LowLatency,
AudioProfile::Balanced,
AudioProfile::BadNetwork,
];
/// Compact discriminant for handing the profile to the capture thread via an
/// atomic. Pairs with [`AudioProfile::from_u8`].
pub fn as_u8(self) -> u8 {
match self {
AudioProfile::LowLatency => 0,
AudioProfile::Balanced => 1,
AudioProfile::BadNetwork => 2,
}
}
/// Inverse of [`AudioProfile::as_u8`]; unknown values fall back to the default.
pub fn from_u8(v: u8) -> AudioProfile {
match v {
0 => AudioProfile::LowLatency,
2 => AudioProfile::BadNetwork,
_ => AudioProfile::Balanced,
}
}
}
impl std::fmt::Display for AudioProfile {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(match self {
AudioProfile::LowLatency => "Low latency",
AudioProfile::Balanced => "Balanced",
AudioProfile::BadNetwork => "Bad network",
})
}
}
impl std::fmt::Display for RoomLayout {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(match self {
@@ -199,6 +253,10 @@ pub struct AppConfig {
pub clip_volume_universal: bool,
#[serde(default)]
pub network_mode: NetworkMode,
/// Opus encoder / network-resilience profile (W12). Applies live to the
/// running encoder; default `Balanced`.
#[serde(default)]
pub audio_profile: AudioProfile,
/// Presence posture for the friends idle listener (W7): invisible / normal /
/// discoverable. Default `Normal` = answer friends only, no DNS beacon.
#[serde(default)]
@@ -380,6 +438,7 @@ impl Default for AppConfig {
show_player_bar: true,
clip_volume_universal: true,
network_mode: NetworkMode::default(),
audio_profile: AudioProfile::default(),
presence_mode: crate::presence::PresenceMode::default(),
echo_cancellation_enabled: false,
notifications_enabled: true,
@@ -1019,6 +1078,38 @@ mod tests {
assert_ne!(display_0, display_2);
}
#[test]
fn test_audio_profile() {
// Default is Balanced.
assert_eq!(AudioProfile::default(), AudioProfile::Balanced);
// ALL holds the three variants.
assert_eq!(AudioProfile::ALL.len(), 3);
assert!(AudioProfile::ALL.contains(&AudioProfile::LowLatency));
assert!(AudioProfile::ALL.contains(&AudioProfile::Balanced));
assert!(AudioProfile::ALL.contains(&AudioProfile::BadNetwork));
// as_u8 / from_u8 round-trip every variant, and unknown bytes fall back
// to the default rather than panicking.
for p in AudioProfile::ALL {
assert_eq!(AudioProfile::from_u8(p.as_u8()), p);
}
assert_eq!(AudioProfile::from_u8(99), AudioProfile::Balanced);
// serde round-trips, and Display strings are non-empty + distinct.
let mut labels = Vec::new();
for p in AudioProfile::ALL {
let s = serde_json::to_string(&p).unwrap();
assert_eq!(serde_json::from_str::<AudioProfile>(&s).unwrap(), p);
let label = p.to_string();
assert!(!label.is_empty());
labels.push(label);
}
labels.sort();
labels.dedup();
assert_eq!(labels.len(), 3);
}
#[test]
fn test_unknown_field_tolerance() {
// Unknown/extra field tolerance: a config JSON containing an extra unrecognized key should still deserialize.
+95 -5
View File
@@ -202,11 +202,15 @@ impl JitterBuffer {
None
} else {
// Gap with later packets already buffered: a packet was lost
// or reordered out of window. Conceal this frame via Opus PLC
// and grow the cushion — the jitter beat our current delay.
// or reordered out of window. First try Opus in-band FEC from
// the next packet; if unavailable, fall back to plain PLC.
self.next_seq = Some(next.wrapping_add(1));
self.note_disruption();
self.decoder.decode(None).ok()
let next_payload = self.packets.values().next().expect("non-empty");
self.decoder
.decode_fec(next_payload)
.or_else(|_| self.decoder.decode(None))
.ok()
}
}
}
@@ -221,8 +225,8 @@ impl JitterBuffer {
#[cfg(test)]
mod tests {
use super::*;
use crate::codec::AudioEncoder;
use crate::codec::opus_impl::OpusEncoder;
use crate::codec::opus_impl::{OpusDecoder, OpusEncoder, OpusParams};
use crate::codec::{AudioDecoder, AudioEncoder};
use opus::{Application, Channels};
/// A real, decodable Opus packet for one 20ms mono frame at amplitude `amp`.
@@ -233,6 +237,32 @@ mod tests {
enc.encode(&pcm).unwrap()
}
fn tone_frame(enc: &mut OpusEncoder, amp: i16, frame_index: usize) -> Vec<u8> {
let pcm: Vec<i16> = (0..FRAME_SAMPLES)
.map(|i| {
let sample_index = frame_index * FRAME_SAMPLES + i;
let t = sample_index as f32 / 48_000.0;
let fundamental = (t * 220.0 * 2.0 * std::f32::consts::PI).sin();
let harmonic = (t * 440.0 * 2.0 * std::f32::consts::PI).sin();
((fundamental * 0.7 + harmonic * 0.3) * amp as f32) as i16
})
.collect();
enc.encode(&pcm).unwrap()
}
fn rms_error(a: &[i16], b: &[i16]) -> f64 {
assert_eq!(a.len(), b.len());
let sum_sq: f64 = a
.iter()
.zip(b)
.map(|(&left, &right)| {
let diff = left as f64 - right as f64;
diff * diff
})
.sum();
(sum_sq / a.len() as f64).sqrt()
}
#[test]
fn buffers_then_plays_in_order() {
let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap();
@@ -289,6 +319,66 @@ mod tests {
assert!(jb.pop_frame().is_none());
}
#[test]
fn uses_in_band_fec_from_next_packet_for_gap() {
let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap();
enc.apply_params(&OpusParams {
bitrate: 20_000,
inband_fec: true,
packet_loss_perc: 60,
dtx: false,
})
.unwrap();
let dropped_seq = 5usize;
let amps = [1800, 1800, 1800, 1800, 1800, 12_000, 12_000, 12_000];
let packets: Vec<Vec<u8>> = amps
.into_iter()
.enumerate()
.map(|(seq, amp)| tone_frame(&mut enc, amp, seq))
.collect();
let mut expected_decoder = OpusDecoder::new(48000, Channels::Mono, FRAME_SAMPLES).unwrap();
for packet in packets.iter().take(dropped_seq) {
expected_decoder.decode(Some(packet)).unwrap();
}
let expected_lost = expected_decoder
.decode(Some(&packets[dropped_seq]))
.unwrap();
let mut plc_decoder = OpusDecoder::new(48000, Channels::Mono, FRAME_SAMPLES).unwrap();
for packet in packets.iter().take(dropped_seq) {
plc_decoder.decode(Some(packet)).unwrap();
}
let pure_plc = plc_decoder.decode(None).unwrap();
let mut jb = JitterBuffer::new().unwrap();
for (seq, packet) in packets.iter().enumerate() {
if seq != dropped_seq {
jb.insert(seq as u32, packet.clone());
}
}
for _ in 0..dropped_seq {
assert_eq!(jb.pop_frame().map(|frame| frame.len()), Some(FRAME_SAMPLES));
}
let recovered = jb.pop_frame().expect("gap should be reconstructed");
assert_eq!(recovered.len(), FRAME_SAMPLES);
assert!(
jb.packets.contains_key(&(dropped_seq as u32 + 1)),
"FEC source packet must remain buffered for normal decode"
);
assert_eq!(jb.pop_frame().map(|frame| frame.len()), Some(FRAME_SAMPLES));
let fec_error = rms_error(&recovered, &expected_lost);
let plc_error = rms_error(&pure_plc, &expected_lost);
assert!(
fec_error < plc_error * 0.75,
"FEC reconstruction should be materially closer than PLC (fec_error={fec_error}, plc_error={plc_error})"
);
}
#[test]
fn drops_packets_already_played() {
let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap();
+8 -1
View File
@@ -1,4 +1,4 @@
use crate::config::{NetworkMode, RecordingMode};
use crate::config::{AudioProfile, NetworkMode, RecordingMode};
use crate::friends::Friend;
use crate::network::PeerState;
use crate::presence::{FriendPresence, PresenceMode};
@@ -56,6 +56,10 @@ pub enum CoreCommand {
/// Set the relay/discovery posture. Takes effect on the next room join,
/// since the endpoint is (re)built then.
SetNetworkMode(NetworkMode),
/// Set the Opus encoder / network-resilience profile (W12). Applies live to
/// the running capture encoder, and to the next call's encoder. Sent at
/// startup from config and whenever the user changes it.
SetAudioProfile(AudioProfile),
/// Start/stop recording the call to a local WAV (your mic + the incoming
/// mix). No-op start if already recording / not in a call.
SetRecording(bool),
@@ -212,6 +216,7 @@ pub fn delivery_class(cmd: &CoreCommand) -> DeliveryClass {
input_device: _,
}
| CoreCommand::SetNetworkMode(_)
| CoreCommand::SetAudioProfile(_)
| CoreCommand::SetRecording(_)
| CoreCommand::SetRecordingMode(_)
| CoreCommand::SendChat(_)
@@ -292,6 +297,7 @@ pub fn coalesce_key(cmd: &CoreCommand) -> Option<CoalesceKey> {
input_device: _,
}
| CoreCommand::SetNetworkMode(_)
| CoreCommand::SetAudioProfile(_)
| CoreCommand::SetRecording(_)
| CoreCommand::SetRecordingMode(_)
| CoreCommand::SendChat(_)
@@ -574,6 +580,7 @@ mod tests {
},
CoreCommand::SetPeerMuted(peer, true),
CoreCommand::SetPresenceMode(PresenceMode::Normal),
CoreCommand::SetAudioProfile(crate::config::AudioProfile::BadNetwork),
CoreCommand::SendChat("hello".to_string()),
];
+88 -6
View File
@@ -17,7 +17,7 @@ use crate::network::{
};
use crate::audio::multitrack::MultitrackRecorder;
use crate::config::{NetworkMode, RecordingMode};
use crate::config::{AudioProfile, NetworkMode, RecordingMode};
use crate::presence::PresenceMode;
use iroh::{
Endpoint, EndpointAddr, EndpointId, RelayMode, SecretKey, endpoint::presets, protocol::Router,
@@ -31,6 +31,12 @@ use tokio::sync::{Mutex, mpsc};
type CoalesceStore = Arc<StdMutex<HashMap<CoalesceKey, CoreCommand>>>;
// Mixer -> playback-worker handoff. The playback ring itself targets three
// 20ms frames; allow at most two more in flight so worker lag applies
// backpressure before the ring can overshoot to its 200ms cap (A6).
const PLAYBACK_HANDOFF_QUEUE_FRAMES: usize = 2;
const PLAYBACK_HANDOFF_RETRY: Duration = Duration::from_millis(1);
pub struct CoreController {
reliable_tx: mpsc::UnboundedSender<CoreCommand>,
coalesce: CoalesceStore,
@@ -156,6 +162,22 @@ fn audio_datagram_len_ok(len: usize) -> bool {
(4..=4 + MAX_OPUS_PAYLOAD).contains(&len)
}
async fn send_playback_frame(
tx: &std::sync::mpsc::SyncSender<Vec<i16>>,
mut frame: Vec<i16>,
) -> bool {
loop {
match tx.try_send(frame) {
Ok(()) => return true,
Err(std::sync::mpsc::TrySendError::Full(returned)) => {
frame = returned;
tokio::time::sleep(PLAYBACK_HANDOFF_RETRY).await;
}
Err(std::sync::mpsc::TrySendError::Disconnected(_)) => return false,
}
}
}
/// The presence label to broadcast for a detected game: its display name,
/// sanitized + length-capped, or `None` when there's no game or no broadcastable
/// name (a Steam appid without a manifest name, or a label that sanitizes empty).
@@ -1162,6 +1184,12 @@ async fn run_core_loop(
// App-internal capture/playback gains (f32 bits), live-read by the audio loops.
let input_gain = Arc::new(std::sync::atomic::AtomicU32::new(1.0f32.to_bits()));
let output_gain = Arc::new(std::sync::atomic::AtomicU32::new(1.0f32.to_bits()));
// Opus encoder profile (W12) as a discriminant, live-read by the capture
// thread so a mid-call profile switch re-tunes the running encoder. Set from
// config via the GUI's startup `SetAudioProfile`; defaults to Balanced.
let audio_profile = Arc::new(std::sync::atomic::AtomicU8::new(
AudioProfile::default().as_u8(),
));
// Call recording: an optional live recorder (mic FIFO + WAV writer), shared
// by the capture thread (pushes mic) and the mixer task (writes mix frames).
// `is_recording` is a fast-path gate so the audio loops only take the lock
@@ -1648,7 +1676,8 @@ async fn run_core_loop(
// Setup raw audio channels
let (capture_tx, capture_rx) = std::sync::mpsc::channel();
let (playback_tx, playback_rx) = std::sync::mpsc::channel();
let (playback_tx, playback_rx) =
std::sync::mpsc::sync_channel(PLAYBACK_HANDOFF_QUEUE_FRAMES);
// Echo cancellation: if enabled, load PipeWire's echo-cancel module
// bound to the chosen real devices and route capture/playback
@@ -1740,6 +1769,7 @@ async fn run_core_loop(
let is_recording_capture = is_recording.clone();
let multitrack_capture = multitrack.clone();
let is_multitrack_capture = is_multitrack.clone();
let audio_profile_capture = audio_profile.clone();
let capture_thread = std::thread::spawn(move || {
use opus::{Application, Channels};
@@ -1751,6 +1781,13 @@ async fn run_core_loop(
return;
}
};
// Tune the encoder to the configured profile (W12), then track
// the live discriminant so a mid-call switch re-applies it.
let mut current_profile =
AudioProfile::from_u8(audio_profile_capture.load(Ordering::Relaxed));
if let Err(e) = encoder.apply_profile(current_profile) {
crate::log_msg(&format!("Opus profile apply failed: {:?}", e));
}
// Per-sender packet sequence number, prepended to every frame so
// receivers can reorder and conceal loss. Wraps after ~years.
let mut seq: u32 = 0;
@@ -1763,6 +1800,14 @@ async fn run_core_loop(
let mut mic_meter = MicLevelMeter::new();
while let Ok(mut pcm) = capture_rx.recv() {
// Re-tune the encoder if the user switched profile mid-call.
// Cheap atomic load per frame; only reconfigures on change.
let want =
AudioProfile::from_u8(audio_profile_capture.load(Ordering::Relaxed));
if want != current_profile && encoder.apply_profile(want).is_ok() {
current_profile = want;
}
// Apply the input gain first so the meter, gate, and what we
// transmit all reflect the same (gained) signal.
apply_volume(
@@ -2092,7 +2137,7 @@ async fn run_core_loop(
mixed
};
if playback_tx.send(frame_to_send).is_err() {
if !send_playback_frame(&playback_tx, frame_to_send).await {
break;
}
@@ -2682,6 +2727,13 @@ async fn run_core_loop(
}
}
CoreCommand::SetAudioProfile(profile) => {
// Publish the new profile to the capture thread (W12). It picks up
// the change on its next frame and re-tunes the live encoder; a
// call that starts later reads the same atomic at encoder creation.
audio_profile.store(profile.as_u8(), Ordering::Relaxed);
}
CoreCommand::RegenerateIdentity => {
// Mint + persist a fresh identity, discarding the old one. The
// persistent endpoint is rebuilt with the new key (now if idle, else
@@ -3193,12 +3245,15 @@ async fn run_core_loop(
mod tests {
use super::{
KnownPeers, MAX_OPUS_PAYLOAD, MAX_RETAINED_PEERS, MIC_LEVEL_REPORT_SAMPLES, MicLevelMeter,
PeerSpeakTicket, admit_retained, apply_peer_volume, apply_volume, audio_datagram_len_ok,
coalesce_insert, coalesce_pop, frame_level, mix_frames, mix_stereo_frames,
next_game_change, should_auto_fetch, stereo_to_mono,
PLAYBACK_HANDOFF_QUEUE_FRAMES, PeerSpeakTicket, admit_retained, apply_peer_volume,
apply_volume, audio_datagram_len_ok, coalesce_insert, coalesce_pop, frame_level,
mix_frames, mix_stereo_frames, next_game_change, send_playback_frame, should_auto_fetch,
stereo_to_mono,
};
use crate::core::messages::{CoalesceKey, CoreCommand, coalesce_key};
use std::collections::{HashMap, HashSet};
use std::sync::mpsc::sync_channel;
use std::time::Duration;
fn endpoint_id() -> iroh::EndpointId {
iroh::SecretKey::generate().public()
@@ -3432,6 +3487,33 @@ mod tests {
assert!(!audio_datagram_len_ok(5 + MAX_OPUS_PAYLOAD));
}
#[tokio::test]
async fn playback_handoff_waits_for_bounded_queue_space() {
let (tx, rx) = sync_channel::<Vec<i16>>(PLAYBACK_HANDOFF_QUEUE_FRAMES);
for n in 0..PLAYBACK_HANDOFF_QUEUE_FRAMES {
tx.try_send(vec![n as i16]).unwrap();
}
let worker = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(20));
for n in 0..PLAYBACK_HANDOFF_QUEUE_FRAMES {
assert_eq!(rx.recv().unwrap(), vec![n as i16]);
}
assert_eq!(rx.recv().unwrap(), vec![99, 100]);
});
let sent = tokio::time::timeout(
Duration::from_secs(1),
send_playback_frame(&tx, vec![99, 100]),
)
.await
.expect("bounded handoff should unblock after the worker drains a frame");
assert!(sent);
drop(tx);
worker.join().unwrap();
}
#[test]
fn mic_meter_holds_the_peak_across_the_window() {
let mut m = MicLevelMeter::new();