From 76260d412c0fd11e6713107e65077dcb13559b3a Mon Sep 17 00:00:00 2001 From: Del Wang Date: Tue, 30 Dec 2025 09:46:01 +0800 Subject: [PATCH] =?UTF-8?q?refactor:=20=E6=8B=86=E5=88=86=20stereo=20?= =?UTF-8?q?=E6=A8=A1=E5=9D=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- apps/hello/src/bin/stereo.rs | 462 +------------------- apps/hello/src/stereo_core/alsa.rs | 70 +++ apps/hello/src/stereo_core/jitter_buffer.rs | 3 + apps/hello/src/stereo_core/master.rs | 241 ++++++++++ apps/hello/src/stereo_core/mod.rs | 3 + apps/hello/src/stereo_core/protocol.rs | 19 +- apps/hello/src/stereo_core/slave.rs | 168 +++++++ apps/hello/src/stereo_core/sync.rs | 17 +- 8 files changed, 511 insertions(+), 472 deletions(-) create mode 100644 apps/hello/src/stereo_core/alsa.rs create mode 100644 apps/hello/src/stereo_core/master.rs create mode 100644 apps/hello/src/stereo_core/slave.rs diff --git a/apps/hello/src/bin/stereo.rs b/apps/hello/src/bin/stereo.rs index b62d82e..939994d 100644 --- a/apps/hello/src/bin/stereo.rs +++ b/apps/hello/src/bin/stereo.rs @@ -1,103 +1,17 @@ #![cfg(target_os = "linux")] -use anyhow::{Context, Result}; -use hello::audio::{AudioPlayer, OpusCodec}; -use hello::config::AudioConfig; -use hello::stereo_core::jitter_buffer::JitterBuffer; -use hello::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket}; -use hello::stereo_core::sync::{ClockSync, now_us}; +use anyhow::Result; +use hello::stereo_core::alsa::AlsaRedirector; +use hello::stereo_core::master::run_master; +use hello::stereo_core::protocol::ChannelRole; +use hello::stereo_core::slave::run_slave; 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::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"; - -// Default ports -const DISCOVERY_PORT: u16 = 53530; -const SERVER_TCP_PORT: u16 = 53531; -// UDP port will be dynamic or fixed - -struct AlsaRedirector; - -impl AlsaRedirector { - fn new() -> Result { - Self::cleanup(); // Ensure clean state - - let original_conf = fs::read_to_string(REAL_ASOUND_CONF).unwrap_or_default(); - - 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 - )); - - fs::write(TEMP_ASOUND_CONF, new_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) - } - - 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."); - } -} - -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!("Usage: {} [master|slave] [left|right]", args[0]); + eprintln!("用法: {} [master|slave] [left|right]", args[0]); return Ok(()); } @@ -108,7 +22,7 @@ async fn main() -> Result<()> { ChannelRole::Right }; - // Explicit cleanup at start just in case + // 启动前先执行清理,确保环境干净 AlsaRedirector::cleanup(); if mode == "master" { @@ -117,365 +31,3 @@ async fn main() -> Result<()> { run_slave(role).await } } - -async fn run_master(role: ChannelRole) -> Result<()> { - println!("--- Master Mode ({:?}) ---", role); - - // 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 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)?; - let playback_config = AudioConfig { - channels: 1, - playback_device: "plug:original_default".into(), - ..config.clone() - }; - let player = AudioPlayer::new(&playback_config)?; - - let mut raw_buf = vec![0u8; config.frame_size * 2 * 2]; - let mut opus_out = vec![0u8; 1500]; - let mut seq = 0u32; - - 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!("Starting session loop..."); - - // 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]; - loop { - 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; - } - } - } - } - } - }); - - 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?, - }; - - 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); - } - - let now = now_us(); - let current_server_time = { - let c = clock.lock().await; - c.to_server_time(now) - }; - - 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]); - } - } - } - } - last_seq = Some(seq); - - 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/stereo_core/alsa.rs b/apps/hello/src/stereo_core/alsa.rs new file mode 100644 index 0000000..694a294 --- /dev/null +++ b/apps/hello/src/stereo_core/alsa.rs @@ -0,0 +1,70 @@ +#![cfg(target_os = "linux")] +use anyhow::{Context, Result}; +use std::fs; +use std::process::Command; + +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"; + +/// ALSA 音频重定向器,用于拦截系统音频输出到 FIFO 管道 +pub struct AlsaRedirector; + +impl AlsaRedirector { + pub fn new() -> Result { + Self::cleanup(); // 确保环境干净 + + let original_conf = fs::read_to_string(REAL_ASOUND_CONF).unwrap_or_default(); + + if !original_conf.contains("pcm.original_default") { + // 重命名原有的 default 逻辑,插入拦截器 + let mut new_conf = original_conf.replace("pcm.!default", "pcm.original_default"); + new_conf.push_str(&format!( + "\npcm.!default {{ type plug slave.pcm \"stereo_interceptor\" }}\n\ + pcm.stereo_interceptor {{ type file slave.pcm \"null\" file \"{}\" format \"raw\" }}\n", + FIFO_PATH + )); + + fs::write(TEMP_ASOUND_CONF, new_conf)?; + + // 挂载覆盖 /etc/asound.conf + let status = Command::new("mount") + .arg("--bind") + .arg(TEMP_ASOUND_CONF) + .arg(REAL_ASOUND_CONF) + .status() + .context("执行 mount 命令失败")?; + + if !status.success() { + return Err(anyhow::anyhow!("挂载 asound.conf 失败")); + } + } + + // 创建 FIFO 管道 + let _ = Command::new("mkfifo").arg(FIFO_PATH).status(); + let _ = Command::new("chmod").arg("666").arg(FIFO_PATH).status(); + + println!("ALSA 输出已重定向至 {}", FIFO_PATH); + Ok(Self) + } + + pub fn cleanup() { + let _ = Command::new("sh") + .arg("-c") + .arg(format!("umount -l {} >/dev/null 2>&1", REAL_ASOUND_CONF)) + .status(); + let _ = fs::remove_file(TEMP_ASOUND_CONF); + let _ = fs::remove_file(FIFO_PATH); + println!("ALSA 配置已恢复。"); + } + + pub fn fifo_path() -> &'static str { + FIFO_PATH + } +} + +impl Drop for AlsaRedirector { + fn drop(&mut self) { + Self::cleanup(); + } +} diff --git a/apps/hello/src/stereo_core/jitter_buffer.rs b/apps/hello/src/stereo_core/jitter_buffer.rs index ad4428a..9e02e14 100644 --- a/apps/hello/src/stereo_core/jitter_buffer.rs +++ b/apps/hello/src/stereo_core/jitter_buffer.rs @@ -22,6 +22,7 @@ impl Ord for OrderedPacket { } } +/// 抖动缓冲区,用于处理网络延迟和乱序 pub struct JitterBuffer { buffer: BinaryHeap, last_played_seq: Option, @@ -38,6 +39,7 @@ impl JitterBuffer { } pub fn push(&mut self, packet: AudioPacket) { + // 丢弃已经播放过的数据包 if let Some(last) = self.last_played_seq { if packet.seq <= last { return; @@ -46,6 +48,7 @@ impl JitterBuffer { 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 { diff --git a/apps/hello/src/stereo_core/master.rs b/apps/hello/src/stereo_core/master.rs new file mode 100644 index 0000000..d6fcfb3 --- /dev/null +++ b/apps/hello/src/stereo_core/master.rs @@ -0,0 +1,241 @@ +#![cfg(target_os = "linux")] +use anyhow::Result; +use crate::audio::{AudioPlayer, OpusCodec}; +use crate::config::AudioConfig; +use std::net::SocketAddr; +use std::sync::Arc; +use std::time::Duration; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::{TcpListener, TcpStream, UdpSocket}; +use crate::stereo_core::alsa::AlsaRedirector; +use crate::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket}; +use crate::stereo_core::sync::now_us; + +pub const SERVER_TCP_PORT: u16 = 53531; +pub const DISCOVERY_PORT: u16 = 53530; + +/// 运行主节点模式 +pub async fn run_master(role: ChannelRole) -> Result<()> { + println!("--- 主节点模式 ({:?}) ---", role); + + // 0. 设置 ALSA 重定向 (在主节点生命周期内持续有效) + let _alsa_guard = AlsaRedirector::new()?; + + // 1. 设置网络 (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. 发现服务 + 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!("正在端口 {} 等待从节点连接...", SERVER_TCP_PORT); + + loop { + let (socket, addr) = listener.accept().await?; + println!("从节点已连接: {}", addr); + + let udp_socket = udp_socket.clone(); + let role = role.clone(); + + // 启动会话处理句柄 + tokio::spawn(async move { + if let Err(e) = handle_master_session(socket, udp_socket, udp_port, role).await { + eprintln!("会话结束: {:?}", 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!("无效的握手协议")), + }; + + let hello = ControlPacket::ServerHello { + udp_port: udp.local_addr()?.port(), + }; + tcp.write_all(&postcard::to_allocvec(&hello)?).await?; + + // 等待 UDP 打洞/确认 + let mut buf = [0u8; 128]; + let (_, client_udp_addr) = udp.recv_from(&mut buf).await?; + println!("从节点 UDP 地址已确认: {}", 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)?; + let playback_config = AudioConfig { + channels: 1, + playback_device: "plug:original_default".into(), + ..config.clone() + }; + let player = AudioPlayer::new(&playback_config)?; + + let mut raw_buf = vec![0u8; config.frame_size * 2 * 2]; + let mut opus_out = vec![0u8; 1500]; + let mut seq = 0u32; + + 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!("开始会话循环..."); + + // 分离 TCP 读写,以便在不同任务中使用 + let (mut tcp_rx, mut tcp_tx) = tcp.into_split(); + let (stop_tx, mut stop_rx) = tokio::sync::mpsc::channel::<()>(1); + + // 处理来自从节点的控制消息(如 Ping/Pong 进行时间同步) + tokio::spawn(async move { + let mut buf = [0u8; 1024]; + loop { + 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; + } + } + } + } + } + }); + + loop { + // 打开 FIFO。这会阻塞直到有写入者打开它。 + // 使用 select 以便在等待 FIFO 时如果从节点断开连接可以退出。 + let mut fifo = tokio::select! { + _ = stop_rx.recv() => { + println!("等待 FIFO 时从节点断开连接。"); + return Ok(()); + } + f = tokio::fs::File::open(AlsaRedirector::fifo_path()) => f?, + }; + + println!("音频流已启动..."); + let mut stream_start_ts = 0; + let stream_start_seq = seq; + + loop { + // 从 FIFO 读取,带超时/select 以检查断开连接 + let read_res = tokio::select! { + _ = stop_rx.recv() => { + println!("串流过程中从节点断开连接。"); + return Ok(()); + } + res = fifo.read_exact(&mut raw_buf) => res, + }; + + if let Err(_) = read_res { + println!("音频流结束 (FIFO 读取完毕)。"); + 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 发送错误: {:?}", 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; + } + } +} + diff --git a/apps/hello/src/stereo_core/mod.rs b/apps/hello/src/stereo_core/mod.rs index e2f7d3c..90aa49d 100644 --- a/apps/hello/src/stereo_core/mod.rs +++ b/apps/hello/src/stereo_core/mod.rs @@ -1,3 +1,6 @@ +pub mod alsa; pub mod jitter_buffer; +pub mod master; pub mod protocol; +pub mod slave; pub mod sync; diff --git a/apps/hello/src/stereo_core/protocol.rs b/apps/hello/src/stereo_core/protocol.rs index a713b2a..afff181 100644 --- a/apps/hello/src/stereo_core/protocol.rs +++ b/apps/hello/src/stereo_core/protocol.rs @@ -8,15 +8,15 @@ pub enum ChannelRole { #[derive(Serialize, Deserialize, Debug, Clone)] pub enum ControlPacket { - // Discovery + // 发现协议 ServerHello { - udp_port: u16, // Port for UDP audio stream + udp_port: u16, // UDP 音频流端口 }, - // Handshake + // 握手协议 ClientIdentify { role: ChannelRole, }, - // Time Sync (Continuous) + // 时间同步 (持续进行) Ping { client_ts: u128, seq: u32, @@ -26,14 +26,13 @@ pub enum ControlPacket { server_ts: u128, seq: u32, }, - // Control - Volume(u8), // 0-100 + // 控制协议 + 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 + pub seq: u32, // 序列号,用于丢包检测 + pub timestamp: u128, // 目标播放时间 (主节点时间) + pub data: Vec, // Opus 编码数据 } - diff --git a/apps/hello/src/stereo_core/slave.rs b/apps/hello/src/stereo_core/slave.rs new file mode 100644 index 0000000..020c4d9 --- /dev/null +++ b/apps/hello/src/stereo_core/slave.rs @@ -0,0 +1,168 @@ +#![cfg(target_os = "linux")] +use anyhow::Result; +use crate::audio::{AudioPlayer, OpusCodec}; +use crate::config::AudioConfig; +use std::sync::Arc; +use std::time::Duration; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::{TcpStream, UdpSocket}; +use tokio::sync::{Mutex, mpsc}; +use crate::stereo_core::jitter_buffer::JitterBuffer; +use crate::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket}; +use crate::stereo_core::sync::{ClockSync, now_us}; +use crate::stereo_core::master::{DISCOVERY_PORT}; + +/// 运行从节点模式 +pub async fn run_slave(role: ChannelRole) -> Result<()> { + println!("--- 从节点模式 ({:?}) ---", role); + + println!("正在扫描主节点..."); + 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!("无效的发现数据包")), + }; + + let master_ip = master_addr.ip(); + let master_tcp_addr = format!("{}:{}", master_ip, master_tcp_port); + println!("在 {} 发现主节点", 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!("应答应为 ServerHello")), + }; + + // UDP 打洞以接收音频流 + 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); + + // 分离读写 + let (mut tcp_rx, mut tcp_tx) = tcp.into_split(); + let udp = Arc::new(udp); + let udp_rx = udp.clone(); + + // 定时发送 Ping 进行时间同步 + 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(); + + // 接收 Pong 并更新时钟偏移 + 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!("正在监听音频数据..."); + + let mut last_seq: Option = None; + + loop { + // 将收到的数据包放入抖动缓冲区 + while let Ok(pkt) = audio_rx.try_recv() { + jitter.push(pkt); + } + + let now = now_us(); + let current_server_time = { + let c = clock.lock().await; + c.to_server_time(now) + }; + + // 从抖动缓冲区提取待播放帧 + if let Some((seq, data)) = jitter.pop_frame(current_server_time) { + // 丢包处理 (PLC) + if let Some(last) = last_seq { + if seq > last + 1 { + println!("检测到丢包: {} -> {}", last, seq); + // 补偿丢失的帧 + for _ in 0..(seq - last - 1) { + if let Ok(len) = codec.decode_loss(&mut pcm_buf) { + let _ = player.write(&pcm_buf[..len]); + } + } + } + } + last_seq = Some(seq); + + let len = codec.decode(&data, &mut pcm_buf)?; + player.write(&pcm_buf[..len])?; + } else { + // 稍作等待以减少 CPU 占用 + tokio::time::sleep(Duration::from_millis(1)).await; + } + } +} + diff --git a/apps/hello/src/stereo_core/sync.rs b/apps/hello/src/stereo_core/sync.rs index a553be7..beac2f3 100644 --- a/apps/hello/src/stereo_core/sync.rs +++ b/apps/hello/src/stereo_core/sync.rs @@ -1,13 +1,15 @@ 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") + .expect("时间倒流") .as_micros() } +/// 时钟同步管理器,用于计算主从节点间的时钟偏移 pub struct ClockSync { offsets: VecDeque, pub current_offset: i128, @@ -23,15 +25,16 @@ impl ClockSync { } } + /// 更新时钟偏移估计 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) + // 基础过滤:如果 RTT 过大则忽略 (例如局域网内 > 100ms) if rtt > 100_000 { return; } - // Clock offset = server_time - client_time - // server_ts is at time (send + recv)/2 + // 时钟偏移 = 主节点时间 - 从节点时间 + // 假设主节点收到 Ping 的时间点在 (发送时间 + 接收时间) / 2 let estimated_server_time = server_ts as i128 + rtt / 2; let offset = estimated_server_time - client_recv_ts as i128; @@ -40,8 +43,7 @@ impl ClockSync { 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() { @@ -49,12 +51,13 @@ impl ClockSync { } } + /// 将本地时间转换为服务器(主节点)时间 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 } } -