From 32dced2fa21bf6d44c7c80baf9cb5c8bb700b9bf Mon Sep 17 00:00:00 2001 From: Del Wang Date: Wed, 31 Dec 2025 10:01:55 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=94=AF=E6=8C=81=E5=A4=9A=E6=92=AD?= =?UTF-8?q?=E6=94=BE=E6=BA=90=E5=B9=B6=E5=8F=91=E6=B7=B7=E9=9F=B3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- apps/hello/src/bin/stereo.rs | 12 +- apps/hello/src/stereo_core/alsa.rs | 41 ++-- apps/hello/src/stereo_core/injector.rs | 34 ++++ apps/hello/src/stereo_core/master.rs | 254 +++++++++++-------------- apps/hello/src/stereo_core/mixer.rs | 124 ++++++++++++ apps/hello/src/stereo_core/mod.rs | 2 + 6 files changed, 312 insertions(+), 155 deletions(-) create mode 100644 apps/hello/src/stereo_core/injector.rs create mode 100644 apps/hello/src/stereo_core/mixer.rs diff --git a/apps/hello/src/bin/stereo.rs b/apps/hello/src/bin/stereo.rs index 4c80c51..85095ca 100644 --- a/apps/hello/src/bin/stereo.rs +++ b/apps/hello/src/bin/stereo.rs @@ -4,13 +4,23 @@ use anyhow::Result; use hello::stereo_core::master::run_master; use hello::stereo_core::protocol::ChannelRole; use hello::stereo_core::slave::run_slave; +use hello::stereo_core::injector::run_injector; use std::env; #[tokio::main] async fn main() -> Result<()> { let args: Vec = env::args().collect(); + + // 优先处理 --inject 模式 + if args.iter().any(|arg| arg == "--inject") { + return run_injector().await; + } + if args.len() < 3 { - eprintln!("用法: {} [master|slave] [left|right]", args[0]); + eprintln!("用法:"); + eprintln!(" 主节点: {} master [left|right]", args[0]); + eprintln!(" 从节点: {} slave [left|right]", args[0]); + eprintln!(" 注入器: {} --inject (由 ALSA 自动拉起)", args[0]); return Ok(()); } diff --git a/apps/hello/src/stereo_core/alsa.rs b/apps/hello/src/stereo_core/alsa.rs index 3e6ad3f..d12bef2 100644 --- a/apps/hello/src/stereo_core/alsa.rs +++ b/apps/hello/src/stereo_core/alsa.rs @@ -1,13 +1,13 @@ #![cfg(target_os = "linux")] use anyhow::{Context, Result}; +use std::env; 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 管道 +/// ALSA 音频重定向器,通过 Pipe 模式将音频流注入到主进程的混音器中 pub struct AlsaRedirector; impl AlsaRedirector { @@ -17,12 +17,32 @@ impl AlsaRedirector { let original_conf = fs::read_to_string(REAL_ASOUND_CONF).unwrap_or_default(); if !original_conf.contains("pcm.original_default") { + let current_exe = env::current_exe() + .context("无法获取当前可执行文件路径")? + .to_string_lossy() + .to_string(); + // 重命名原有的 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\" format S16_LE rate 48000 channels 2 }} }}\n\ - pcm.stereo_interceptor {{ type file slave.pcm \"null\" file \"{}\" format \"raw\" }}\n", - FIFO_PATH + r#" +pcm.!default {{ + type plug + slave {{ + pcm "stereo_interceptor" + format S16_LE + rate 48000 + channels 2 + }} +}} +pcm.stereo_interceptor {{ + type file + slave.pcm "null" + file "| {} --inject" + format "raw" +}} +"#, + current_exe )); fs::write(TEMP_ASOUND_CONF, new_conf)?; @@ -40,11 +60,6 @@ impl AlsaRedirector { } } - // 创建 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) } @@ -54,12 +69,6 @@ impl AlsaRedirector { .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 } } diff --git a/apps/hello/src/stereo_core/injector.rs b/apps/hello/src/stereo_core/injector.rs new file mode 100644 index 0000000..c0e6486 --- /dev/null +++ b/apps/hello/src/stereo_core/injector.rs @@ -0,0 +1,34 @@ +use tokio::net::UnixStream; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use anyhow::Result; +use crate::stereo_core::mixer::MIXER_SOCKET_PATH; + +/// Injector 模式的逻辑:将 stdin 数据转发到 Unix Socket +pub async fn run_injector() -> Result<()> { + // 1. 连接到主进程的 Mixer Socket + let mut socket = match UnixStream::connect(MIXER_SOCKET_PATH).await { + Ok(s) => s, + Err(e) => { + // 如果连接失败,可能是主进程还没启动或已经退出 + // 在 Injector 模式下,静默退出即可 + return Err(e.into()); + } + }; + + let mut stdin = tokio::io::stdin(); + let mut buf = [0u8; 4096]; + + // 2. 数据搬运:stdin -> socket + loop { + let n = stdin.read(&mut buf).await?; + if n == 0 { + break; // stdin 关闭 + } + if let Err(_) = socket.write_all(&buf[..n]).await { + break; // Socket 断开 + } + } + + Ok(()) +} + diff --git a/apps/hello/src/stereo_core/master.rs b/apps/hello/src/stereo_core/master.rs index bf2ef9b..db38d0b 100644 --- a/apps/hello/src/stereo_core/master.rs +++ b/apps/hello/src/stereo_core/master.rs @@ -4,6 +4,7 @@ use crate::audio::{AudioPlayer, OpusCodec}; use crate::config::AudioConfig; use crate::stereo_core::alsa::AlsaRedirector; use crate::stereo_core::discovery::Discovery; +use crate::stereo_core::mixer::Mixer; use crate::stereo_core::network::{ControlConnection, MasterNetwork}; use crate::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket}; use crate::stereo_core::sync::now_us; @@ -30,6 +31,10 @@ pub async fn run_master(master_role: ChannelRole) -> Result<()> { // 0. 设置 ALSA 重定向 let _alsa_guard = AlsaRedirector::new()?; + // 0.1 启动混音器服务 + let mixer = Arc::new(Mixer::new()); + mixer.start().await?; + // 1. 设置网络 (UDP + TCP) let network = MasterNetwork::setup(SERVER_TCP_PORT).await?; let audio_socket = network.audio_socket().clone_inner(); @@ -81,8 +86,6 @@ pub async fn run_master(master_role: ChannelRole) -> Result<()> { let mut current_player_channels = 0; let mut player: Option = None; - let mut raw_buf = vec![0u8; config.frame_size * 2 * 2]; - let mut pcm_out = vec![0i16; config.frame_size * 2]; let mut left_pcm = vec![0i16; config.frame_size]; let mut right_pcm = vec![0i16; config.frame_size]; let mut opus_out = vec![0u8; 1500]; @@ -97,165 +100,140 @@ pub async fn run_master(master_role: ChannelRole) -> Result<()> { let shutdown_flag_clone = shutdown_flag.clone(); let audio_loop = async move { + // 每个新流开始时,重置编码器状态以避免残留音频导致爆音 + let mut left_encoder = OpusCodec::new(&encode_config)?; + let mut right_encoder = OpusCodec::new(&encode_config)?; + loop { if shutdown_flag_clone.load(Ordering::Relaxed) { break; } - // 打开 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 mixer.client_count().await == 0 { + stream_start_ts = 0; + } + + // 从混音器读取一帧 (48kHz, 2ch, 20ms) + let mixed_pcm = mixer + .read_mixed_frame(config.frame_size, config.channels as usize) + .await; + + let active_slaves = { + let s = slaves.lock().await; + if s.is_empty() { None } else { Some(s.clone()) } }; - // 每个新流开始时,重置编码器状态以避免残留音频导致爆音 - let mut left_encoder = OpusCodec::new(&encode_config)?; - let mut right_encoder = OpusCodec::new(&encode_config)?; + // 检查是否需要切换播放器模式 + let target_channels = if active_slaves.is_none() { 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; + } - loop { - if shutdown_flag_clone.load(Ordering::Relaxed) { - break; + let now = now_us(); + if stream_start_ts == 0 { + stream_start_ts = now; + stream_start_seq = seq; + } + + // 计算该帧应当播放的基准时间(相对于流开始) + let target_ts = + stream_start_ts + ((seq - stream_start_seq) as u128 * frame_duration_us) + delay_us; + + if let Some(slaves_list) = active_slaves { + // 情况 1: 有从节点,主从同步 + for i in 0..config.frame_size { + left_pcm[i] = mixed_pcm[i * 2]; + right_pcm[i] = mixed_pcm[i * 2 + 1]; } - // 从 FIFO 读取 - if let Err(_) = fifo.read_exact(&mut raw_buf).await { - break; // FIFO 关闭,重新打开 + // 1. 检查各声道是否有从节点需要 + let needs_left = slaves_list.iter().any(|s| s.role == ChannelRole::Left); + let needs_right = slaves_list.iter().any(|s| s.role == ChannelRole::Right); + + // 2. 编码需要的声道 + let mut left_bytes = None; + let mut right_bytes = None; + + if needs_left { + let len = left_encoder.encode(&left_pcm, &mut opus_out)?; + let packet = AudioPacket { + seq, + timestamp: target_ts, + data: opus_out[..len].to_vec(), + }; + left_bytes = Some(postcard::to_allocvec(&packet)?); } - let active_slaves = { - let s = slaves.lock().await; - if s.is_empty() { None } else { Some(s.clone()) } + if needs_right { + let len = right_encoder.encode(&right_pcm, &mut opus_out)?; + let packet = AudioPacket { + seq, + timestamp: target_ts, + data: opus_out[..len].to_vec(), + }; + right_bytes = Some(postcard::to_allocvec(&packet)?); + } + + // 3. 发送给对应的从节点 + for slave in &slaves_list { + let bytes = match slave.role { + ChannelRole::Left => left_bytes.as_ref(), + ChannelRole::Right => right_bytes.as_ref(), + }; + if let Some(b) = bytes { + let _ = audio_socket.send_to(b, slave.udp_addr).await; + } + } + + // 4. 本地回放 + let master_pcm = match master_role { + ChannelRole::Left => &left_pcm, + ChannelRole::Right => &right_pcm, }; - // 检查是否需要切换播放器模式 - let target_channels = if active_slaves.is_none() { 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 now < target_ts { + let wait = target_ts - now; + if wait > 1000 { + tokio::time::sleep(Duration::from_micros(wait as u64)).await; + } } - - // 计算该帧应当播放的基准时间(相对于流开始) - let target_ts = stream_start_ts - + ((seq - stream_start_seq) as u128 * frame_duration_us) - + delay_us; - - if let Some(slaves_list) = active_slaves { - // 情况 1: 有从节点,主从同步 - for i in 0..config.frame_size { - left_pcm[i] = i16::from_le_bytes([raw_buf[i * 4], raw_buf[i * 4 + 1]]); - right_pcm[i] = i16::from_le_bytes([raw_buf[i * 4 + 2], raw_buf[i * 4 + 3]]); - } - - // 1. 检查各声道是否有从节点需要 - let needs_left = slaves_list.iter().any(|s| s.role == ChannelRole::Left); - let needs_right = slaves_list.iter().any(|s| s.role == ChannelRole::Right); - - // 2. 编码需要的声道 - let mut left_bytes = None; - let mut right_bytes = None; - - if needs_left { - let len = left_encoder.encode(&left_pcm, &mut opus_out)?; - let packet = AudioPacket { - seq, - timestamp: target_ts, - data: opus_out[..len].to_vec(), - }; - left_bytes = Some(postcard::to_allocvec(&packet)?); - } - - if needs_right { - let len = right_encoder.encode(&right_pcm, &mut opus_out)?; - let packet = AudioPacket { - seq, - timestamp: target_ts, - data: opus_out[..len].to_vec(), - }; - right_bytes = Some(postcard::to_allocvec(&packet)?); - } - - // 3. 发送给对应的从节点 - for slave in &slaves_list { - let bytes = match slave.role { - ChannelRole::Left => left_bytes.as_ref(), - ChannelRole::Right => right_bytes.as_ref(), - }; - if let Some(b) = bytes { - let _ = audio_socket.send_to(b, slave.udp_addr).await; - } - } - - // 4. 本地回放 - let master_pcm = match master_role { - ChannelRole::Left => &left_pcm, - ChannelRole::Right => &right_pcm, - }; - - 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 { - if let Err(_) = p.write(master_pcm) { - // 如果写入失败且正在退出,直接跳出循环 - if shutdown_flag_clone.load(Ordering::Relaxed) { - break; - } + if let Some(p) = &player { + if let Err(_) = p.write(master_pcm) { + // 如果写入失败且正在退出,直接跳出循环 + if shutdown_flag_clone.load(Ordering::Relaxed) { + break; } - } - } else { - // 情况 2: 没有从节点,本地立体声播放 - for i in 0..config.frame_size { - pcm_out[i * 2] = i16::from_le_bytes([raw_buf[i * 4], raw_buf[i * 4 + 1]]); - pcm_out[i * 2 + 1] = - i16::from_le_bytes([raw_buf[i * 4 + 2], raw_buf[i * 4 + 3]]); } - if let Some(p) = &player { - if let Err(_) = p.write(&pcm_out) { - // 如果写入失败且正在退出,直接跳出循环 - if shutdown_flag_clone.load(Ordering::Relaxed) { - break; - } + } + } else { + // 情况 2: 没有从节点,本地立体声播放 + if let Some(p) = &player { + if let Err(_) = p.write(&mixed_pcm) { + // 如果写入失败且正在退出,直接跳出循环 + if shutdown_flag_clone.load(Ordering::Relaxed) { + break; } } } - - seq += 1; } - // 重置流计时 - stream_start_ts = 0; - - // 如果是因为退出信号而跳出内层循环,也要跳出外层循环 - if shutdown_flag_clone.load(Ordering::Relaxed) { - break; - } + seq += 1; } Ok::<(), anyhow::Error>(()) diff --git a/apps/hello/src/stereo_core/mixer.rs b/apps/hello/src/stereo_core/mixer.rs new file mode 100644 index 0000000..180d81d --- /dev/null +++ b/apps/hello/src/stereo_core/mixer.rs @@ -0,0 +1,124 @@ +use anyhow::Result; +use std::fs; +use std::path::Path; +use std::sync::Arc; +use tokio::io::AsyncReadExt; +use tokio::net::{UnixListener, UnixStream}; +use tokio::sync::Mutex; + +pub const MIXER_SOCKET_PATH: &str = "/tmp/audio_mixer.sock"; + +/// 混音器服务,负责接收来自多个 Injector 的音频流并进行混音 +pub struct Mixer { + clients: Arc>>, +} + +impl Mixer { + pub fn new() -> Self { + Self { + clients: Arc::new(Mutex::new(Vec::new())), + } + } + + /// 启动 Unix Socket 监听服务 + pub async fn start(&self) -> Result<()> { + if Path::new(MIXER_SOCKET_PATH).exists() { + let _ = fs::remove_file(MIXER_SOCKET_PATH); + } + + let listener = UnixListener::bind(MIXER_SOCKET_PATH)?; + // 允许任何用户写入,确保 ALSA 进程有权连接 + let _ = std::process::Command::new("chmod") + .arg("666") + .arg(MIXER_SOCKET_PATH) + .status(); + + let clients = self.clients.clone(); + + tokio::spawn(async move { + loop { + match listener.accept().await { + Ok((stream, _)) => { + let mut c = clients.lock().await; + c.push(stream); + } + Err(e) => { + eprintln!("❌ Mixer accept error: {:?}", e); + } + } + } + }); + + Ok(()) + } + + /// 获取当前连接的客户端数量 + pub async fn client_count(&self) -> usize { + self.clients.lock().await.len() + } + + /// 读取并混合一帧音频数据 + /// frame_size: 每声道的采样点数 + /// channels: 声道数 + pub async fn read_mixed_frame(&self, frame_size: usize, channels: usize) -> Vec { + let total_samples = frame_size * channels; + let bytes_per_frame = total_samples * 2; + let mut raw_buf = vec![0u8; bytes_per_frame]; + + loop { + let mut clients = self.clients.lock().await; + if clients.is_empty() { + drop(clients); + // 没有客户端时,稍微等待,避免空转 + tokio::time::sleep(tokio::time::Duration::from_millis(1)).await; + continue; + } + + let mut mixed_frame = vec![0i32; total_samples]; + let mut active_clients = Vec::new(); + let mut read_any = false; + + // 遍历所有客户端,各读取一帧 + for mut client in clients.drain(..) { + // 注意:这里 read_exact 是异步的,会按顺序等待每个客户端的数据 + // 在主从同步场景下,所有客户端通常来自本地 ALSA,速率是同步的 + match client.read_exact(&mut raw_buf).await { + Ok(_) => { + for i in 0..total_samples { + let sample = + i16::from_le_bytes([raw_buf[i * 2], raw_buf[i * 2 + 1]]) as i32; + mixed_frame[i] += sample; + } + active_clients.push(client); + read_any = true; + } + Err(_) => { + // 客户端断开连接 + } + } + } + + *clients = active_clients; + + if read_any { + // 饱和截断 (Clipping/Saturating) 将 i32 转回 i16 + return mixed_frame + .into_iter() + .map(|x| { + if x > 32767 { + 32767 + } else if x < -32768 { + -32768 + } else { + x as i16 + } + }) + .collect(); + } else { + // 如果所有客户端都断开了,继续等待新连接 + drop(clients); + tokio::time::sleep(tokio::time::Duration::from_millis(1)).await; + } + } + } +} diff --git a/apps/hello/src/stereo_core/mod.rs b/apps/hello/src/stereo_core/mod.rs index 768f74d..c60cb97 100644 --- a/apps/hello/src/stereo_core/mod.rs +++ b/apps/hello/src/stereo_core/mod.rs @@ -1,7 +1,9 @@ pub mod alsa; pub mod discovery; +pub mod injector; pub mod jitter_buffer; pub mod master; +pub mod mixer; pub mod network; pub mod protocol; pub mod slave;