From b277fb72344fd9abb6894dae4cfdf5c5d66a7773 Mon Sep 17 00:00:00 2001 From: Del Wang Date: Tue, 30 Dec 2025 11:49:20 +0800 Subject: [PATCH] =?UTF-8?q?chore:=20=E4=BC=98=E5=8C=96=E8=BF=9E=E6=8E=A5?= =?UTF-8?q?=E6=8F=90=E7=A4=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- apps/hello/src/stereo_core/alsa.rs | 4 +- apps/hello/src/stereo_core/master.rs | 60 ++++++++++++++++------------ apps/hello/src/stereo_core/slave.rs | 17 ++++---- 3 files changed, 44 insertions(+), 37 deletions(-) diff --git a/apps/hello/src/stereo_core/alsa.rs b/apps/hello/src/stereo_core/alsa.rs index 694a294..76adc1c 100644 --- a/apps/hello/src/stereo_core/alsa.rs +++ b/apps/hello/src/stereo_core/alsa.rs @@ -44,7 +44,7 @@ impl AlsaRedirector { let _ = Command::new("mkfifo").arg(FIFO_PATH).status(); let _ = Command::new("chmod").arg("666").arg(FIFO_PATH).status(); - println!("ALSA 输出已重定向至 {}", FIFO_PATH); + // println!("ALSA 输出已重定向至 {}", FIFO_PATH); Ok(Self) } @@ -55,7 +55,7 @@ impl AlsaRedirector { .status(); let _ = fs::remove_file(TEMP_ASOUND_CONF); let _ = fs::remove_file(FIFO_PATH); - println!("ALSA 配置已恢复。"); + // println!("ALSA 配置已恢复。"); } pub fn fifo_path() -> &'static str { diff --git a/apps/hello/src/stereo_core/master.rs b/apps/hello/src/stereo_core/master.rs index 47759f1..21cba64 100644 --- a/apps/hello/src/stereo_core/master.rs +++ b/apps/hello/src/stereo_core/master.rs @@ -14,7 +14,6 @@ use tokio::io::{AsyncReadExt, AsyncWriteExt}; pub const SERVER_TCP_PORT: u16 = 53531; -/// 运行主节点模式 pub async fn run_master(role: ChannelRole) -> Result<()> { println!("--- 主节点模式 ({}) ---", role.to_string()); @@ -28,19 +27,21 @@ pub async fn run_master(role: ChannelRole) -> Result<()> { // 2. 启动服务发现广播 Discovery::start_broadcast(SERVER_TCP_PORT).await?; - println!("正在端口 {} 等待从节点连接...", SERVER_TCP_PORT); + println!("✅ 服务已启动,等待连接..."); loop { - let (control_conn, addr) = network.accept().await?; - println!("从节点已连接: {}", addr); + let (control_conn, client_addr) = network.accept().await?; let audio_socket = audio_socket.clone(); let role = role.clone(); // 启动会话处理句柄 tokio::spawn(async move { - if let Err(e) = handle_master_session(control_conn, audio_socket, role).await { - eprintln!("会话结束: {:?}", e); + if let Err(e) = + handle_master_session(control_conn, audio_socket, role, client_addr.to_string()) + .await + { + eprintln!("❌ {:?}", e); } }); } @@ -50,14 +51,18 @@ pub async fn run_master(role: ChannelRole) -> Result<()> { async fn handle_master_session( mut control: ControlConnection, audio_socket: Arc, - role: ChannelRole, + master_role: ChannelRole, + client_tcp_addr: String, ) -> Result<()> { let mut buf = [0u8; 1024]; + #[allow(unused_assignments)] + let mut slave_role = ChannelRole::Left; + // 握手 let pkt = control.recv_packet(&mut buf).await?; match pkt { - ControlPacket::ClientIdentify { role: _r } => {} + ControlPacket::ClientIdentify { role } => slave_role = role, _ => return Err(anyhow::anyhow!("无效的握手协议")), }; @@ -69,7 +74,12 @@ async fn handle_master_session( // 等待 UDP 打洞/确认 let mut buf = [0u8; 128]; let (_, client_udp_addr) = audio_socket.recv_from(&mut buf).await?; - println!("从节点 UDP 地址已确认: {}", client_udp_addr); + + println!( + "✅ 从节点已连接: {} {}", + client_tcp_addr, + slave_role.to_string(), + ); // 配置音频 let config = AudioConfig { @@ -96,11 +106,9 @@ async fn handle_master_session( let mut opus_out = vec![0u8; 1500]; let mut seq = 0u32; + let delay_us = 200_000; // 200ms 延迟,确保主从节点音频同步 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) = control.split(); @@ -113,6 +121,11 @@ async fn handle_master_session( match tcp_rx.read(&mut buf).await { Ok(0) | Err(_) => { let _ = stop_tx.send(()).await; + println!( + "❌ 从节点已断开: {} {}", + client_tcp_addr, + slave_role.to_string(), + ); break; } Ok(n) => { @@ -141,29 +154,22 @@ async fn handle_master_session( loop { // 打开 FIFO let mut fifo = tokio::select! { - _ = stop_rx.recv() => { - println!("等待 FIFO 时从节点断开连接。"); - return Ok(()); - } + _ = stop_rx.recv() => 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 读取 let read_res = tokio::select! { - _ = stop_rx.recv() => { - println!("串流过程中从节点断开连接。"); - return Ok(()); - } + _ = stop_rx.recv() => return Ok(()), res = fifo.read_exact(&mut raw_buf) => res, }; if let Err(_) = read_res { - println!("音频流结束 (FIFO 读取完毕)。"); + // todo 继续等待音频流继续播放 break; } @@ -179,12 +185,15 @@ async fn handle_master_session( 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 { + if master_role == ChannelRole::Left { local_pcm.push(l); - remote_pcm.push(r); } else { local_pcm.push(r); + } + if slave_role == ChannelRole::Left { remote_pcm.push(l); + } else { + remote_pcm.push(r); } } @@ -201,8 +210,7 @@ async fn handle_master_session( let bytes = postcard::to_allocvec(&packet)?; if let Err(e) = audio_socket.send_to(&bytes, client_udp_addr).await { - eprintln!("UDP 发送错误: {:?}", e); - return Err(e.into()); + return Err(anyhow::anyhow!("UDP 发送错误: {:?}", e)); } // 本地回放同步 diff --git a/apps/hello/src/stereo_core/slave.rs b/apps/hello/src/stereo_core/slave.rs index 4aa6ec3..ef168ff 100644 --- a/apps/hello/src/stereo_core/slave.rs +++ b/apps/hello/src/stereo_core/slave.rs @@ -19,24 +19,23 @@ pub async fn run_slave(role: ChannelRole) -> Result<()> { loop { match handle_connection(role.clone()).await { - Ok(_) => println!("连接正常结束,正在尝试重新连接..."), Err(e) => { - eprintln!("连接异常断开 3 秒后重试... \nError: {:?}.", e); + eprintln!("❌ {:?}", e); tokio::time::sleep(Duration::from_secs(3)).await; } + Ok(_) => {} } } } -/// 处理单次完整的连接生命周期 async fn handle_connection(role: ChannelRole) -> Result<()> { - println!("正在扫描主节点..."); // 1. 发现主节点 + println!("🔍 正在扫描主节点..."); let (master_ip, master_tcp_port) = Discovery::discover_master().await?; let master_tcp_addr = format!("{}:{}", master_ip, master_tcp_port); - println!("在 {} 发现主节点,正在连接...", master_tcp_addr); // 2. 建立 TCP 连接 + println!("🔥 发现主节点: {}", master_tcp_addr); let network = SlaveNetwork::connect(master_tcp_addr.parse()?).await?; let (mut control, audio) = network.split(); @@ -49,7 +48,7 @@ async fn handle_connection(role: ChannelRole) -> Result<()> { let pkt = control.recv_packet(&mut buf).await?; let server_udp_port = match pkt { ControlPacket::ServerHello { udp_port } => udp_port, - _ => return Err(anyhow!("应答应为 ServerHello")), + _ => return Err(anyhow!("身份认证应答异常")), }; // 4. UDP 打洞 @@ -136,14 +135,14 @@ async fn handle_connection(role: ChannelRole) -> Result<()> { }); // 8. 播放主循环 - println!("连接成功,正在播放音频..."); + println!("✅ 主节点已连接,音频串流中..."); let mut pcm_buf = vec![0i16; config.frame_size]; let mut last_seq: Option = None; loop { // 检查 TCP 是否已断开 if let Ok(_) = disconnect_rx.try_recv() { - return Err(anyhow!("主节点连接已断开 (TCP Disconnected)")); + return Err(anyhow!("主节点已断开: {}", master_tcp_addr)); } // 填充 Jitter Buffer @@ -170,7 +169,7 @@ async fn handle_connection(role: ChannelRole) -> Result<()> { let len = codec.decode(&data, &mut pcm_buf)?; player.write(&pcm_buf[..len])?; } else { - tokio::time::sleep(Duration::from_millis(5)).await; + tokio::time::sleep(Duration::from_millis(1)).await; } } }