Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
39b5dafd57 | ||
|
|
a78860db15 | ||
|
|
5f52aa1506 |
@@ -2,10 +2,10 @@ 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());
|
||||
}
|
||||
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")
|
||||
|
||||
+1
-1
@@ -3139,7 +3139,7 @@ fn audio_profile_hint(profile: AudioProfile) -> &'static str {
|
||||
}
|
||||
AudioProfile::Balanced => "Default: voice quality with light loss recovery.",
|
||||
AudioProfile::BadNetwork => {
|
||||
"Most resilient on a lossy/congested link: extra loss recovery, lower bitrate."
|
||||
"Most resilient on a lossy/congested link: heavier loss recovery, lower bitrate."
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
+16
-5
@@ -40,7 +40,7 @@ pub fn opus_params(profile: AudioProfile) -> OpusParams {
|
||||
bitrate: 20_000,
|
||||
inband_fec: true,
|
||||
packet_loss_perc: 25,
|
||||
dtx: true,
|
||||
dtx: false,
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -162,6 +162,19 @@ 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)]
|
||||
@@ -180,10 +193,8 @@ mod tests {
|
||||
assert!(bal.inband_fec);
|
||||
assert!(bad.inband_fec);
|
||||
|
||||
// BadNetwork is the only profile that enables DTX, and it expects the
|
||||
// heaviest loss.
|
||||
assert!(bad.dtx);
|
||||
assert!(!low.dtx && !bal.dtx);
|
||||
// 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.
|
||||
|
||||
+1
-1
@@ -105,7 +105,7 @@ pub enum AudioProfile {
|
||||
#[default]
|
||||
Balanced,
|
||||
/// Maximum resilience on a lossy/congested link: in-band FEC tuned for heavy
|
||||
/// loss plus DTX, at a lower bitrate to leave headroom for the redundancy.
|
||||
/// loss, at a lower bitrate to leave headroom for the redundancy.
|
||||
BadNetwork,
|
||||
}
|
||||
|
||||
|
||||
+95
-5
@@ -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();
|
||||
|
||||
+58
-5
@@ -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).
|
||||
@@ -1654,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
|
||||
@@ -2114,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;
|
||||
}
|
||||
|
||||
@@ -3222,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()
|
||||
@@ -3461,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();
|
||||
|
||||
Reference in New Issue
Block a user