diff --git a/apps/hello/src/stereo_core/master.rs b/apps/hello/src/stereo_core/master.rs index 9d5ab3f..581c08c 100644 --- a/apps/hello/src/stereo_core/master.rs +++ b/apps/hello/src/stereo_core/master.rs @@ -12,6 +12,7 @@ use std::net::SocketAddr; use std::sync::Arc; use std::time::Duration; use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::signal::unix::{SignalKind, signal}; use tokio::sync::Mutex; pub const SERVER_TCP_PORT: u16 = 53531; @@ -96,121 +97,152 @@ pub async fn run_master(master_role: ChannelRole) -> Result<()> { let mut stream_start_ts = 0; let mut stream_start_seq = 0; - loop { - // 打开 FIFO - let mut fifo = match tokio::fs::File::open(AlsaRedirector::fifo_path()).await { - Ok(f) => f, - Err(e) => { - eprintln!("❌ 无法打开 FIFO: {:?}, 重试...", e); - tokio::time::sleep(Duration::from_secs(1)).await; - continue; - } - }; - + let audio_loop = async { loop { - // 从 FIFO 读取 - if let Err(_) = fifo.read_exact(&mut raw_buf).await { - break; // FIFO 关闭,重新打开 - } - - let active_slaves = slaves.lock().await.clone(); - - // 检查是否需要切换播放器模式 - let target_channels = if active_slaves.is_empty() { 2 } else { 1 }; - if player.is_none() || current_player_channels != target_channels { - println!( - "🔄 切换播放模式: {}", - if target_channels == 2 { - "本地立体声" - } else { - "主从同步 (单声道)" - } - ); - let playback_config = AudioConfig { - channels: target_channels, - playback_device: "plug:original_default".into(), - ..config.clone() - }; - player = Some(AudioPlayer::new(&playback_config)?); - current_player_channels = target_channels; - } - - let now = now_us(); - if stream_start_ts == 0 { - stream_start_ts = now; - stream_start_seq = seq; - } - - if active_slaves.is_empty() { - // 情况 1: 没有从节点,本地立体声播放 - let mut pcm = Vec::with_capacity(config.frame_size * 2); - 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]]); - pcm.push(l); - pcm.push(r); + // 打开 FIFO + let mut fifo = match tokio::fs::File::open(AlsaRedirector::fifo_path()).await { + Ok(f) => f, + Err(e) => { + eprintln!("❌ 无法打开 FIFO: {:?}, 重试...", e); + tokio::time::sleep(Duration::from_secs(1)).await; + continue; } - if let Some(p) = &player { - p.write(&pcm)?; - } - } else { - // 情况 2: 有从节点,主从同步 - let mut local_pcm = Vec::with_capacity(config.frame_size); - let mut remote_pcm = Vec::with_capacity(config.frame_size); + }; - // 提取左右声道 (假设当前逻辑只处理一个从节点的情况,或所有从节点角色一致) - // 如果有多个从节点角色不同,这里需要更复杂的逻辑 - let slave_role = active_slaves[0].role; - - 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 master_role == ChannelRole::Left { - local_pcm.push(l); - } else { - local_pcm.push(r); - } - if slave_role == ChannelRole::Left { - remote_pcm.push(l); - } else { - remote_pcm.push(r); - } + loop { + // 从 FIFO 读取 + if let Err(_) = fifo.read_exact(&mut raw_buf).await { + break; // FIFO 关闭,重新打开 } - // 编码并发送给所有从节点 - let len = mono_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 active_slaves = slaves.lock().await.clone(); - let packet = AudioPacket { - seq, - timestamp: target_ts, - data: opus_out[..len].to_vec(), - }; - - let bytes = postcard::to_allocvec(&packet)?; - for slave in &active_slaves { - let _ = audio_socket.send_to(&bytes, slave.udp_addr).await; + // 检查是否需要切换播放器模式 + let target_channels = if active_slaves.is_empty() { 2 } else { 1 }; + if player.is_none() || current_player_channels != target_channels { + println!( + "🔄 切换播放模式: {}", + if target_channels == 2 { + "本地立体声" + } else { + "主从同步 (单声道)" + } + ); + let playback_config = AudioConfig { + channels: target_channels, + playback_device: "plug:original_default".into(), + ..config.clone() + }; + player = Some(AudioPlayer::new(&playback_config)?); + current_player_channels = target_channels; } - // 本地回放同步 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; + if stream_start_ts == 0 { + stream_start_ts = now; + stream_start_seq = seq; + } + + if active_slaves.is_empty() { + // 情况 1: 没有从节点,本地立体声播放 + let mut pcm = Vec::with_capacity(config.frame_size * 2); + 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]]); + pcm.push(l); + pcm.push(r); + } + if let Some(p) = &player { + p.write(&pcm)?; + } + } else { + // 情况 2: 有从节点,主从同步 + let mut local_pcm = Vec::with_capacity(config.frame_size); + let mut remote_pcm = Vec::with_capacity(config.frame_size); + + // 提取左右声道 (假设当前逻辑只处理一个从节点的情况,或所有从节点角色一致) + let slave_role = active_slaves[0].role; + + 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 master_role == ChannelRole::Left { + local_pcm.push(l); + } else { + local_pcm.push(r); + } + if slave_role == ChannelRole::Left { + remote_pcm.push(l); + } else { + remote_pcm.push(r); + } + } + + // 编码并发送给所有从节点 + let len = mono_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)?; + for slave in &active_slaves { + let _ = audio_socket.send_to(&bytes, slave.udp_addr).await; + } + + // 本地回放同步 + 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; + } + } + if let Some(p) = &player { + p.write(&local_pcm)?; } } - if let Some(p) = &player { - p.write(&local_pcm)?; - } - } - seq += 1; + seq += 1; + } + // 重置流计时 + stream_start_ts = 0; } - // 重置流计时 - stream_start_ts = 0; + #[allow(unreachable_code)] + Ok::<(), anyhow::Error>(()) + }; + + tokio::select! { + res = audio_loop => { + if let Err(e) = res { + eprintln!("❌ 音频循环错误: {:?}", e); + } + }, + _ = shutdown_signal() => {}, + } + + // 显式清理 + AlsaRedirector::cleanup(); + + // 强制退出 + std::process::exit(0); +} + +/// 监听系统退出信号 (SIGINT, SIGTERM, SIGQUIT) +async fn shutdown_signal() { + let mut sigint = signal(SignalKind::interrupt()).expect("无法注册 SIGINT 处理器"); + let mut sigterm = signal(SignalKind::terminate()).expect("无法注册 SIGTERM 处理器"); + let mut sigquit = signal(SignalKind::quit()).expect("无法注册 SIGQUIT 处理器"); + + tokio::select! { + _ = sigint.recv() => {}, + _ = sigterm.recv() => {}, + _ = sigquit.recv() => {}, } }