diff --git a/apps/hello/src/audio/codec.rs b/apps/hello/src/audio/codec.rs index 6ac2dd2..ea59d17 100644 --- a/apps/hello/src/audio/codec.rs +++ b/apps/hello/src/audio/codec.rs @@ -1,6 +1,6 @@ -use opus::{Encoder, Decoder, Application, Bitrate, Channels}; -use anyhow::{Result, Context}; use crate::config::AudioConfig; +use anyhow::{Context, Result}; +use opus::{Application, Bitrate, Channels, Decoder, Encoder}; pub struct OpusCodec { encoder: Encoder, @@ -9,23 +9,37 @@ pub struct OpusCodec { impl OpusCodec { pub fn new(config: &AudioConfig) -> Result { - let channels = if config.channels == 1 { Channels::Mono } else { Channels::Stereo }; + let channels = if config.channels == 1 { + Channels::Mono + } else { + Channels::Stereo + }; let mut encoder = Encoder::new(config.sample_rate, channels, Application::Audio) .context("Failed to create Opus encoder")?; encoder.set_bitrate(Bitrate::Bits(config.bitrate))?; - let decoder = Decoder::new(config.sample_rate, channels) - .context("Failed to create Opus decoder")?; + let decoder = + Decoder::new(config.sample_rate, channels).context("Failed to create Opus decoder")?; Ok(Self { encoder, decoder }) } pub fn encode(&mut self, pcm: &[i16], out: &mut [u8]) -> Result { - self.encoder.encode(pcm, out).context("Opus encoding failed") + self.encoder + .encode(pcm, out) + .context("Opus encoding failed") } pub fn decode(&mut self, opus: &[u8], out: &mut [i16]) -> Result { - self.decoder.decode(opus, out, false).context("Opus decoding failed") + self.decoder + .decode(opus, out, false) + .context("Opus decoding failed") + } + + // Packet Loss Concealment + pub fn decode_loss(&mut self, out: &mut [i16]) -> Result { + // Fallback to silence as opus crate 0.3 doesn't expose safe PLC + out.fill(0); + Ok(out.len()) } } - diff --git a/apps/hello/src/bin/stereo.rs b/apps/hello/src/bin/stereo.rs index bd28c5c..b62d82e 100644 --- a/apps/hello/src/bin/stereo.rs +++ b/apps/hello/src/bin/stereo.rs @@ -3,88 +3,101 @@ use anyhow::{Context, Result}; use hello::audio::{AudioPlayer, OpusCodec}; use hello::config::AudioConfig; -use hello::net::{AudioPacket, ChannelRole, ControlPacket, DISCOVERY_PORT, SERVER_PORT}; -use hello::sync::{ClockSync, now_us}; -use std::collections::VecDeque; +use hello::stereo_core::jitter_buffer::JitterBuffer; +use hello::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket}; +use hello::stereo_core::sync::{ClockSync, now_us}; use std::env; use std::fs; use std::net::SocketAddr; use std::process::Command; +use std::sync::Arc; use std::time::Duration; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::{TcpListener, TcpStream, UdpSocket}; -use tokio::signal; +use tokio::sync::{Mutex, mpsc}; const FIFO_PATH: &str = "/tmp/stereo_out.fifo"; const REAL_ASOUND_CONF: &str = "/etc/asound.conf"; const TEMP_ASOUND_CONF: &str = "/tmp/asound.stereo.conf"; -fn setup_alsa_config() -> Result<()> { - cleanup_alsa_config(); +// Default ports +const DISCOVERY_PORT: u16 = 53530; +const SERVER_TCP_PORT: u16 = 53531; +// UDP port will be dynamic or fixed - let original_conf = fs::read_to_string(REAL_ASOUND_CONF)?; +struct AlsaRedirector; - if !original_conf.contains("pcm.original_default") { - // 重命名原有的 default 逻辑,插入 interceptor - let mut new_conf = original_conf.replace("pcm.!default", "pcm.original_default"); - new_conf.push_str(&format!( - r#" -pcm.!default {{ - type plug - slave.pcm "stereo_interceptor" -}} +impl AlsaRedirector { + fn new() -> Result { + Self::cleanup(); // Ensure clean state -pcm.stereo_interceptor {{ - type file - slave.pcm "null" - file "{}" - format "raw" -}} -"#, - FIFO_PATH - )); + let original_conf = fs::read_to_string(REAL_ASOUND_CONF).unwrap_or_default(); - fs::write(TEMP_ASOUND_CONF, new_conf)?; + if !original_conf.contains("pcm.original_default") { + // 重命名原有的 default 逻辑,插入 interceptor + let mut new_conf = original_conf.replace("pcm.!default", "pcm.original_default"); + new_conf.push_str(&format!( + r#" + pcm.!default {{ + type plug + slave.pcm "stereo_interceptor" + }} + + pcm.stereo_interceptor {{ + type file + slave.pcm "null" + file "{}" + format "raw" + }} + "#, + FIFO_PATH + )); - // 挂载覆盖 /etc/asound.conf - let status = Command::new("mount") - .arg("--bind") - .arg(TEMP_ASOUND_CONF) - .arg(REAL_ASOUND_CONF) - .status() - .context("Failed to execute mount command")?; + fs::write(TEMP_ASOUND_CONF, new_conf)?; - if !status.success() { - return Err(anyhow::anyhow!("Failed to mount asound.conf")); + // 挂载覆盖 /etc/asound.conf + let status = Command::new("mount") + .arg("--bind") + .arg(TEMP_ASOUND_CONF) + .arg(REAL_ASOUND_CONF) + .status() + .context("Failed to execute mount command")?; + + if !status.success() { + return Err(anyhow::anyhow!("Failed to mount asound.conf")); + } } + + // Create FIFO + let _ = Command::new("mkfifo").arg(FIFO_PATH).status(); + let _ = Command::new("chmod").arg("666").arg(FIFO_PATH).status(); + + println!("ALSA output redirected to {}", FIFO_PATH); + Ok(Self) } - // 创建 FIFO - let _ = Command::new("mkfifo").arg(FIFO_PATH).status(); - let _ = Command::new("chmod").arg("666").arg(FIFO_PATH).status(); - - println!("Successfully redirected ALSA output to {}", FIFO_PATH); - Ok(()) + fn cleanup() { + let _ = Command::new("umount") + .arg("-l") + .arg(REAL_ASOUND_CONF) + .status(); + let _ = fs::remove_file(TEMP_ASOUND_CONF); + let _ = fs::remove_file(FIFO_PATH); + println!("ALSA configuration restored."); + } } -fn cleanup_alsa_config() { - println!("Cleaning up ALSA configurations..."); - let _ = Command::new("umount") - .arg("-l") - .arg(REAL_ASOUND_CONF) - .status(); - let _ = fs::remove_file(TEMP_ASOUND_CONF); - let _ = fs::remove_file(FIFO_PATH); +impl Drop for AlsaRedirector { + fn drop(&mut self) { + Self::cleanup(); + } } #[tokio::main] async fn main() -> Result<()> { let args: Vec = env::args().collect(); if args.len() < 3 { - eprintln!("用法: {} [master|slave] [left|right]", args[0]); - eprintln!("示例:"); - eprintln!(" 主设备: {} master left", args[0]); - eprintln!(" 从设备: {} slave right", args[0]); + eprintln!("Usage: {} [master|slave] [left|right]", args[0]); return Ok(()); } @@ -95,379 +108,374 @@ async fn main() -> Result<()> { ChannelRole::Right }; - let result = tokio::select! { - res = async { - if mode == "master" { - run_master(role).await - } else { - run_slave(role).await - } - } => res, - _ = signal::ctrl_c() => { - println!("\n收到退出信号"); - Ok(()) - } - }; + // Explicit cleanup at start just in case + AlsaRedirector::cleanup(); - cleanup_alsa_config(); - return result; -} - -async fn run_master(role: ChannelRole) -> Result<()> { - let config = AudioConfig { - sample_rate: 48000, - channels: 2, // 拦截的是立体声 - frame_size: 960, - bitrate: 48000, - ..AudioConfig::default() - }; - - println!("--- 主设备模式 ---"); - println!("本地声道: {:?}", role); - - loop { - println!("正在启动发现服务 (UDP Broadcast)..."); - // 1. 启动发现广播 - let socket = UdpSocket::bind("0.0.0.0:0").await?; - socket.set_broadcast(true)?; - let target_addr: SocketAddr = format!("255.255.255.255:{}", DISCOVERY_PORT).parse()?; - let hello = ControlPacket::ServerHello { port: SERVER_PORT }; - let msg = postcard::to_allocvec(&hello)?; - - let broadcast_handle = tokio::spawn(async move { - loop { - let _ = socket.send_to(&msg, target_addr).await; - tokio::time::sleep(Duration::from_secs(1)).await; - } - }); - - // 2. 等待从设备连接 - let listener = TcpListener::bind(format!("0.0.0.0:{}", SERVER_PORT)).await?; - println!("等待从设备连接于端口 {}...", SERVER_PORT); - - let (mut slave_stream, addr) = tokio::select! { - res = listener.accept() => res?, - _ = signal::ctrl_c() => { - broadcast_handle.abort(); - return Ok(()); - } - }; - - println!("从设备已连接: {}", addr); - broadcast_handle.abort(); - - // 3. 握手与时间同步 - let mut buf = [0u8; 1024]; - let n = slave_stream.read(&mut buf).await?; - if let Ok(ControlPacket::ClientIdentify { role: slave_role }) = - postcard::from_bytes::(&buf[..n]) - { - println!("从设备识别为: {:?}", slave_role); - } - - let n = slave_stream.read(&mut buf).await?; - if let Ok(ControlPacket::Ping { client_ts }) = - postcard::from_bytes::(&buf[..n]) - { - let pong = ControlPacket::Pong { - client_ts, - server_ts: now_us(), - }; - slave_stream - .write_all(&postcard::to_allocvec(&pong)?) - .await?; - println!("时钟同步完成"); - } - - // 4. 初始化音频并开始拦截 - setup_alsa_config()?; - - let res = run_master_audio_loop(&mut slave_stream, role.clone(), &config).await; - - cleanup_alsa_config(); - - if let Err(e) = res { - eprintln!("主设备音频循环出错: {:?}, 准备重连...", e); - } else { - println!("从设备正常断开,准备下一次连接..."); - } - - tokio::time::sleep(Duration::from_secs(1)).await; + if mode == "master" { + run_master(role).await + } else { + run_slave(role).await } } -async fn run_master_audio_loop( - slave_stream: &mut TcpStream, - role: ChannelRole, - config: &AudioConfig, -) -> Result<()> { - // 打开 FIFO (使用 tokio::fs 以支持异步读取) - let mut fifo = tokio::fs::File::open(FIFO_PATH) - .await - .context("Failed to open FIFO for reading")?; +async fn run_master(role: ChannelRole) -> Result<()> { + println!("--- Master Mode ({:?}) ---", role); - // 使用 plug:original_default 以支持系统主音量控制 - let playback_config = AudioConfig { - channels: 1, - playback_device: "plug:original_default".to_string(), - ..config.clone() + // 0. Setup ALSA redirection (persistent for the life of the master) + let _alsa_guard = AlsaRedirector::new()?; + + // 1. Setup Network (UDP + TCP) + let udp_socket = UdpSocket::bind("0.0.0.0:0").await?; + let udp_port = udp_socket.local_addr()?.port(); + let udp_socket = Arc::new(udp_socket); + + let listener = TcpListener::bind(format!("0.0.0.0:{}", SERVER_TCP_PORT)).await?; + + // 2. Discovery Service + let _discovery = tokio::spawn(async move { + let socket = UdpSocket::bind("0.0.0.0:0").await.unwrap(); + socket.set_broadcast(true).unwrap(); + let target: SocketAddr = format!("255.255.255.255:{}", DISCOVERY_PORT) + .parse() + .unwrap(); + let msg = postcard::to_allocvec(&ControlPacket::ServerHello { + udp_port: SERVER_TCP_PORT, + }) + .unwrap(); + loop { + let _ = socket.send_to(&msg, target).await; + tokio::time::sleep(Duration::from_secs(1)).await; + } + }); + + println!("Waiting for slave on port {}...", SERVER_TCP_PORT); + + loop { + let (socket, addr) = listener.accept().await?; + println!("Slave connected: {}", addr); + + let udp_socket = udp_socket.clone(); + let role = role.clone(); + + // Spawn connection handler + tokio::spawn(async move { + if let Err(e) = handle_master_session(socket, udp_socket, udp_port, role).await { + eprintln!("Session ended: {:?}", e); + } + }); + } +} + +async fn handle_master_session( + mut tcp: TcpStream, + udp: Arc, + _local_udp_port: u16, + role: ChannelRole, +) -> Result<()> { + let mut buf = [0u8; 1024]; + + let len = tcp.read(&mut buf).await?; + match postcard::from_bytes::(&buf[..len])? { + ControlPacket::ClientIdentify { role: _r } => {} + _ => return Err(anyhow::anyhow!("Invalid handshake")), }; - let player = AudioPlayer::new(&playback_config).context("无法打开本地播放设备")?; + + let hello = ControlPacket::ServerHello { + udp_port: udp.local_addr()?.port(), + }; + tcp.write_all(&postcard::to_allocvec(&hello)?).await?; + + let mut buf = [0u8; 128]; + let (_, client_udp_addr) = udp.recv_from(&mut buf).await?; + println!("Client UDP address confirmed: {}", client_udp_addr); + + let config = AudioConfig { + sample_rate: 48000, + channels: 2, + frame_size: 960, + bitrate: 64000, + ..AudioConfig::default() + }; + let mono_config = AudioConfig { channels: 1, ..config.clone() }; - let mut codec = OpusCodec::new(&mono_config).context("无法初始化编码器")?; - - let mut stereo_raw_buf = vec![0i16; config.frame_size * 2]; - let mut opus_buf = vec![0u8; 2048]; - - let delay_us = 100_000; // 100ms 缓冲 - let mut current_ts = 0; - let frame_duration = - (config.frame_size as f64 / config.sample_rate as f64 * 1_000_000.0) as u128; - - println!("开始拦截并分发立体声音频..."); - - let mut byte_buf = vec![0u8; config.frame_size * 2 * 2]; - - loop { - // 从 FIFO 读取 (这是同步源) - if let Err(e) = fifo.read_exact(&mut byte_buf).await { - if e.kind() == std::io::ErrorKind::UnexpectedEof { - tokio::time::sleep(Duration::from_millis(10)).await; - continue; - } - return Err(e.into()); - } - - let now = now_us(); - if current_ts == 0 { - current_ts = now + delay_us; - } - - for i in 0..stereo_raw_buf.len() { - stereo_raw_buf[i] = i16::from_le_bytes([byte_buf[i * 2], byte_buf[i * 2 + 1]]); - } - - let mut local_pcm = Vec::with_capacity(config.frame_size); - let mut remote_pcm = Vec::with_capacity(config.frame_size); - - for i in 0..config.frame_size { - let (l, r) = (stereo_raw_buf[i * 2], stereo_raw_buf[i * 2 + 1]); - if role == ChannelRole::Left { - local_pcm.push(l); - remote_pcm.push(r); - } else { - local_pcm.push(r); - remote_pcm.push(l); - } - } - - // 先发送网络包 - let opus_len = codec.encode(&remote_pcm, &mut opus_buf)?; - let packet = AudioPacket { - timestamp: current_ts, - data: opus_buf[..opus_len].to_vec(), - }; - let packet_data = postcard::to_allocvec(&packet)?; - - if slave_stream - .write_all(&(packet_data.len() as u32).to_be_bytes()) - .await - .is_err() - || slave_stream.write_all(&packet_data).await.is_err() - { - return Err(anyhow::anyhow!("从设备断开连接")); - } - - // 本地播放也需要同步 - let now = now_us(); - if now < current_ts { - let wait = current_ts - now; - if wait > 1000 { - tokio::time::sleep(Duration::from_micros(wait as u64)).await; - } - } - player.write(&local_pcm)?; - - current_ts += frame_duration; - - // 漂移校正 - let now = now_us(); - if now > current_ts + 200_000 { - current_ts = now + 50_000; - } - } -} - -async fn run_slave(role: ChannelRole) -> Result<()> { - let config = AudioConfig { - sample_rate: 48000, + let mut codec = OpusCodec::new(&mono_config)?; + let playback_config = AudioConfig { channels: 1, - frame_size: 960, - bitrate: 32000, - ..AudioConfig::default() + playback_device: "plug:original_default".into(), + ..config.clone() }; + let player = AudioPlayer::new(&playback_config)?; - println!("--- 从设备模式 ---"); - println!("本地声道: {:?}", role); + let mut raw_buf = vec![0u8; config.frame_size * 2 * 2]; + let mut opus_out = vec![0u8; 1500]; + let mut seq = 0u32; - loop { - // 1. 发现主设备 - println!("正在搜索主设备..."); - let master_addr = match tokio::select! { - addr = hello::net::Discovery::client_discover_server() => addr, - _ = signal::ctrl_c() => return Ok(()), - } { - Ok(addr) => addr, - Err(e) => { - eprintln!("搜索主设备失败: {:?}, 1秒后重试...", e); - tokio::time::sleep(Duration::from_secs(1)).await; - continue; - } - }; + let frame_duration_us = + (config.frame_size as f64 / config.sample_rate as f64 * 1_000_000.0) as u128; + let delay_us = 200_000; - println!("发现主设备: {}", master_addr); - let mut stream = match TcpStream::connect(master_addr).await { - Ok(s) => s, - Err(e) => { - eprintln!("连接主设备失败: {:?}, 1秒后重试...", e); - tokio::time::sleep(Duration::from_secs(1)).await; - continue; - } - }; + println!("Starting session loop..."); - // 2. 身份识别 - if let Err(e) = stream - .write_all(&postcard::to_allocvec(&ControlPacket::ClientIdentify { - role: role.clone(), - })?) - .await - { - eprintln!("发送身份识别失败: {:?}, 准备重连...", e); - continue; - } - - // 3. 时间同步 - let mut clock = ClockSync::new(); - let t1 = now_us(); - if let Err(e) = stream - .write_all(&postcard::to_allocvec(&ControlPacket::Ping { - client_ts: t1, - })?) - .await - { - eprintln!("发送 Ping 失败: {:?}, 准备重连...", e); - continue; - } + // Use into_split for owned ReadHalf/WriteHalf to move into tasks + let (mut tcp_rx, mut tcp_tx) = tcp.into_split(); + let (stop_tx, mut stop_rx) = tokio::sync::mpsc::channel::<()>(1); + tokio::spawn(async move { let mut buf = [0u8; 1024]; - let n = match stream.read(&mut buf).await { - Ok(n) => n, - Err(e) => { - eprintln!("读取 Pong 失败: {:?}, 准备重连...", e); - continue; - } - }; - if let Ok(ControlPacket::Pong { - client_ts, - server_ts, - }) = postcard::from_bytes::(&buf[..n]) - { - let t4 = now_us(); - clock.update(client_ts, server_ts, t4); - println!("时钟同步完成. 偏移: {}us, RTT: {}us", clock.offset, t4 - t1); - } else { - eprintln!("收到无效的 Pong 响应, 准备重连..."); - continue; - } - - let res = run_slave_audio_loop(stream, &config, clock).await; - - if let Err(e) = res { - eprintln!("从设备音频循环出错: {:?}, 准备重连...", e); - } else { - println!("主设备已断开,准备重连..."); - } - - tokio::time::sleep(Duration::from_secs(1)).await; - } -} - -async fn run_slave_audio_loop( - stream: TcpStream, - config: &AudioConfig, - clock: ClockSync, -) -> Result<()> { - let player = AudioPlayer::new(config).context("无法打开播放设备")?; - let mut codec = OpusCodec::new(config).context("无法初始化解码器")?; - let mut jitter_buffer: VecDeque = VecDeque::new(); - let mut pcm_buf = vec![0i16; config.frame_size]; - - let (tx, mut rx) = tokio::sync::mpsc::channel(100); - - // 网络接收线程 - let mut stream_read = stream; - let receive_handle = tokio::spawn(async move { loop { - let mut s_buf = [0u8; 4]; - if stream_read.read_exact(&mut s_buf).await.is_err() { - break; - } - let size = u32::from_be_bytes(s_buf) as usize; - let mut data = vec![0u8; size]; - if stream_read.read_exact(&mut data).await.is_err() { - break; - } - if let Ok(p) = postcard::from_bytes::(&data) { - if tx.send(p).await.is_err() { + match tcp_rx.read(&mut buf).await { + Ok(0) | Err(_) => { + let _ = stop_tx.send(()).await; break; } + Ok(n) => { + if let Ok(ControlPacket::Ping { client_ts, seq }) = + postcard::from_bytes(&buf[..n]) + { + let pong = ControlPacket::Pong { + client_ts, + server_ts: now_us(), + seq, + }; + if tcp_tx + .write_all(&postcard::to_allocvec(&pong).unwrap()) + .await + .is_err() + { + let _ = stop_tx.send(()).await; + break; + } + } + } } } }); - println!("开始接收并播放音频..."); + loop { + // Open FIFO. This will block until a writer opens it. + // We use select to allow exiting if the slave disconnects while we wait. + let mut fifo = tokio::select! { + _ = stop_rx.recv() => { + println!("Slave disconnected while waiting for FIFO."); + return Ok(()); + } + f = tokio::fs::File::open(FIFO_PATH) => f?, + }; - let res = loop { - while let Ok(p) = rx.try_recv() { - jitter_buffer.push_back(p); + println!("Audio stream started..."); + let mut stream_start_ts = 0; + let stream_start_seq = seq; + + loop { + // Read from FIFO with timeout/select to check for disconnection + let read_res = tokio::select! { + _ = stop_rx.recv() => { + println!("Slave disconnected during streaming."); + return Ok(()); + } + res = fifo.read_exact(&mut raw_buf) => res, + }; + + if let Err(_) = read_res { + println!("Audio stream ended (FIFO EOF)."); + break; + } + + let now = now_us(); + if stream_start_ts == 0 { + stream_start_ts = now; + } + + let mut local_pcm = Vec::with_capacity(config.frame_size); + let mut remote_pcm = Vec::with_capacity(config.frame_size); + + for i in 0..config.frame_size { + let l = i16::from_le_bytes([raw_buf[i * 4], raw_buf[i * 4 + 1]]); + let r = i16::from_le_bytes([raw_buf[i * 4 + 2], raw_buf[i * 4 + 3]]); + if role == ChannelRole::Left { + local_pcm.push(l); + remote_pcm.push(r); + } else { + local_pcm.push(r); + remote_pcm.push(l); + } + } + + let len = codec.encode(&remote_pcm, &mut opus_out)?; + let target_ts = + stream_start_ts + ((seq - stream_start_seq) as u128 * frame_duration_us) + delay_us; + + let packet = AudioPacket { + seq, + timestamp: target_ts, + data: opus_out[..len].to_vec(), + }; + + let bytes = postcard::to_allocvec(&packet)?; + if let Err(e) = udp.send_to(&bytes, client_udp_addr).await { + eprintln!("UDP send error: {:?}", e); + return Err(e.into()); + } + + let now = now_us(); + if now < target_ts { + let wait = target_ts - now; + if wait > 1000 { + tokio::time::sleep(Duration::from_micros(wait as u64)).await; + } + } + player.write(&local_pcm)?; + + seq += 1; + } + } +} + +async fn run_slave(role: ChannelRole) -> Result<()> { + println!("--- Slave Mode ({:?}) ---", role); + + println!("Scanning for Master..."); + let udp_disc = UdpSocket::bind(format!("0.0.0.0:{}", DISCOVERY_PORT)).await?; + let mut buf = [0u8; 1024]; + let (len, master_addr) = udp_disc.recv_from(&mut buf).await?; + + let master_tcp_port = match postcard::from_bytes::(&buf[..len])? { + ControlPacket::ServerHello { udp_port: p } => p, + _ => return Err(anyhow::anyhow!("Invalid discovery packet")), + }; + + let master_ip = master_addr.ip(); + let master_tcp_addr = format!("{}:{}", master_ip, master_tcp_port); + println!("Found Master at {}", master_tcp_addr); + + let mut tcp = TcpStream::connect(master_tcp_addr).await?; + + tcp.write_all(&postcard::to_allocvec(&ControlPacket::ClientIdentify { + role: role.clone(), + })?) + .await?; + + let len = tcp.read(&mut buf).await?; + let server_udp_port = match postcard::from_bytes::(&buf[..len])? { + ControlPacket::ServerHello { udp_port } => udp_port, + _ => return Err(anyhow::anyhow!("Expected ServerHello")), + }; + + let udp = UdpSocket::bind("0.0.0.0:0").await?; + let punch_packet = vec![0u8; 1]; + udp.send_to(&punch_packet, format!("{}:{}", master_ip, server_udp_port)) + .await?; + + let config = AudioConfig { + sample_rate: 48000, + channels: 1, + frame_size: 960, + bitrate: 64000, + ..AudioConfig::default() + }; + + let player = AudioPlayer::new(&config)?; + let mut codec = OpusCodec::new(&config)?; + let mut jitter = JitterBuffer::new(50_000); + let clock = ClockSync::new(100); + + // Use into_split for owned ReadHalf/WriteHalf to move into tasks + let (mut tcp_rx, mut tcp_tx) = tcp.into_split(); + let udp = Arc::new(udp); + let udp_rx = udp.clone(); + + let _sync_handle = tokio::spawn(async move { + let mut seq = 0; + loop { + let t1 = now_us(); + let msg = ControlPacket::Ping { client_ts: t1, seq }; + if tcp_tx + .write_all(&postcard::to_allocvec(&msg).unwrap()) + .await + .is_err() + { + break; + } + tokio::time::sleep(Duration::from_secs(1)).await; + seq += 1; + } + }); + + let clock = Arc::new(Mutex::new(clock)); + let clock_updater = clock.clone(); + + tokio::spawn(async move { + let mut buf = [0u8; 1024]; + loop { + match tcp_rx.read(&mut buf).await { + Ok(n) if n > 0 => { + if let Ok(ControlPacket::Pong { + client_ts, + server_ts, + .. + }) = postcard::from_bytes(&buf[..n]) + { + let t4 = now_us(); + clock_updater.lock().await.update(client_ts, server_ts, t4); + } + } + _ => break, + } + } + }); + + let (audio_tx, mut audio_rx) = mpsc::channel(100); + tokio::spawn(async move { + let mut buf = [0u8; 2048]; + loop { + if let Ok((len, _)) = udp_rx.recv_from(&mut buf).await { + if let Ok(packet) = postcard::from_bytes::(&buf[..len]) { + let _ = audio_tx.send(packet).await; + } + } + } + }); + + let mut pcm_buf = vec![0i16; config.frame_size]; + println!("Listening for audio..."); + + let mut last_seq: Option = None; + + loop { + while let Ok(pkt) = audio_rx.try_recv() { + jitter.push(pkt); } - if let Some(p) = jitter_buffer.front() { - let target_client_time = clock.to_client_time(p.timestamp); - let now = now_us(); + let now = now_us(); + let current_server_time = { + let c = clock.lock().await; + c.to_server_time(now) + }; - if now >= target_client_time { - let packet = jitter_buffer.pop_front().unwrap(); - - // 如果包太旧了(延迟超过150ms),跳过以赶上进度,防止累积卡顿 - if now > target_client_time + 150_000 { - continue; - } - - let len = codec.decode(&packet.data, &mut pcm_buf)?; - player.write(&pcm_buf[..len])?; - } else { - let wait = (target_client_time - now) as u64; - if wait > 1000 { - // 如果时间差太大,可能是时钟跳变,清空缓冲重新同步 - if wait > 1_000_000 { - jitter_buffer.clear(); - } else { - tokio::time::sleep(Duration::from_micros(wait)).await; + if let Some((seq, data)) = jitter.pop_frame(current_server_time) { + // PLC logic + let mut _lost = false; + if let Some(last) = last_seq { + if seq > last + 1 { + println!("Packet loss detected: {} -> {}", last, seq); + _lost = true; + // Conceal missing frames + for _ in 0..(seq - last - 1) { + if let Ok(len) = codec.decode_loss(&mut pcm_buf) { + let _ = player.write(&pcm_buf[..len]); + } } } } - } else { - if receive_handle.is_finished() { - break Ok(()); - } - tokio::time::sleep(Duration::from_millis(5)).await; - } - }; + last_seq = Some(seq); - receive_handle.abort(); - res + let len = codec.decode(&data, &mut pcm_buf)?; + player.write(&pcm_buf[..len])?; + } else { + tokio::time::sleep(Duration::from_millis(1)).await; + } + } } diff --git a/apps/hello/src/lib.rs b/apps/hello/src/lib.rs index 8ffc885..147dfb3 100644 --- a/apps/hello/src/lib.rs +++ b/apps/hello/src/lib.rs @@ -1,4 +1,5 @@ pub mod audio; pub mod config; pub mod net; +pub mod stereo_core; pub mod sync; diff --git a/apps/hello/src/stereo_core/jitter_buffer.rs b/apps/hello/src/stereo_core/jitter_buffer.rs new file mode 100644 index 0000000..ad4428a --- /dev/null +++ b/apps/hello/src/stereo_core/jitter_buffer.rs @@ -0,0 +1,64 @@ +use super::protocol::AudioPacket; +use std::collections::BinaryHeap; +use std::cmp::Ordering; + +#[derive(Debug)] +struct OrderedPacket(AudioPacket); + +impl PartialEq for OrderedPacket { + fn eq(&self, other: &Self) -> bool { + self.0.seq == other.0.seq + } +} +impl Eq for OrderedPacket {} +impl PartialOrd for OrderedPacket { + fn partial_cmp(&self, other: &Self) -> Option { + Some(other.0.seq.cmp(&self.0.seq)) + } +} +impl Ord for OrderedPacket { + fn cmp(&self, other: &Self) -> Ordering { + other.0.seq.cmp(&self.0.seq) + } +} + +pub struct JitterBuffer { + buffer: BinaryHeap, + last_played_seq: Option, + pub target_delay_us: u128, +} + +impl JitterBuffer { + pub fn new(target_delay_us: u128) -> Self { + Self { + buffer: BinaryHeap::new(), + last_played_seq: None, + target_delay_us, + } + } + + pub fn push(&mut self, packet: AudioPacket) { + if let Some(last) = self.last_played_seq { + if packet.seq <= last { + return; + } + } + self.buffer.push(OrderedPacket(packet)); + } + + pub fn pop_frame(&mut self, current_time: u128) -> Option<(u32, Vec)> { + if let Some(OrderedPacket(pkt)) = self.buffer.peek() { + if current_time >= pkt.timestamp { + let pkt = self.buffer.pop().unwrap().0; + self.last_played_seq = Some(pkt.seq); + return Some((pkt.seq, pkt.data)); + } + } + None + } + + pub fn clear(&mut self) { + self.buffer.clear(); + self.last_played_seq = None; + } +} diff --git a/apps/hello/src/stereo_core/mod.rs b/apps/hello/src/stereo_core/mod.rs new file mode 100644 index 0000000..e2f7d3c --- /dev/null +++ b/apps/hello/src/stereo_core/mod.rs @@ -0,0 +1,3 @@ +pub mod jitter_buffer; +pub mod protocol; +pub mod sync; diff --git a/apps/hello/src/stereo_core/protocol.rs b/apps/hello/src/stereo_core/protocol.rs new file mode 100644 index 0000000..a713b2a --- /dev/null +++ b/apps/hello/src/stereo_core/protocol.rs @@ -0,0 +1,39 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq, Copy)] +pub enum ChannelRole { + Left, + Right, +} + +#[derive(Serialize, Deserialize, Debug, Clone)] +pub enum ControlPacket { + // Discovery + ServerHello { + udp_port: u16, // Port for UDP audio stream + }, + // Handshake + ClientIdentify { + role: ChannelRole, + }, + // Time Sync (Continuous) + Ping { + client_ts: u128, + seq: u32, + }, + Pong { + client_ts: u128, + server_ts: u128, + seq: u32, + }, + // Control + Volume(u8), // 0-100 +} + +#[derive(Serialize, Deserialize, Debug, Clone)] +pub struct AudioPacket { + pub seq: u32, // Sequence number for packet loss detection + pub timestamp: u128, // Target playback time (server time) + pub data: Vec, // Opus encoded data +} + diff --git a/apps/hello/src/stereo_core/sync.rs b/apps/hello/src/stereo_core/sync.rs new file mode 100644 index 0000000..a553be7 --- /dev/null +++ b/apps/hello/src/stereo_core/sync.rs @@ -0,0 +1,60 @@ +use std::collections::VecDeque; +use std::time::{SystemTime, UNIX_EPOCH}; + +pub fn now_us() -> u128 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("Time went backwards") + .as_micros() +} + +pub struct ClockSync { + offsets: VecDeque, + pub current_offset: i128, + window_size: usize, +} + +impl ClockSync { + pub fn new(window_size: usize) -> Self { + Self { + offsets: VecDeque::with_capacity(window_size), + current_offset: 0, + window_size, + } + } + + pub fn update(&mut self, client_send_ts: u128, server_ts: u128, client_recv_ts: u128) { + let rtt = (client_recv_ts - client_send_ts) as i128; + // Basic filter: ignore if RTT is absurdly large (e.g. > 100ms on LAN) + if rtt > 100_000 { + return; + } + + // Clock offset = server_time - client_time + // server_ts is at time (send + recv)/2 + let estimated_server_time = server_ts as i128 + rtt / 2; + let offset = estimated_server_time - client_recv_ts as i128; + + self.offsets.push_back(offset); + if self.offsets.len() > self.window_size { + self.offsets.pop_front(); + } + + // Calculate median or average offset + // Median is more robust to outliers + let mut sorted: Vec = self.offsets.iter().cloned().collect(); + sorted.sort_unstable(); + if !sorted.is_empty() { + self.current_offset = sorted[sorted.len() / 2]; + } + } + + pub fn to_server_time(&self, client_time: u128) -> u128 { + (client_time as i128 + self.current_offset) as u128 + } + + pub fn to_client_time(&self, server_time: u128) -> u128 { + (server_time as i128 - self.current_offset) as u128 + } +} +