feat: 支持多播放源并发混音
This commit is contained in:
@@ -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<String> = 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(());
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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(())
|
||||
}
|
||||
|
||||
@@ -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<AudioPlayer> = 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,21 +100,6 @@ pub async fn run_master(master_role: ChannelRole) -> Result<()> {
|
||||
|
||||
let shutdown_flag_clone = shutdown_flag.clone();
|
||||
let audio_loop = async move {
|
||||
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;
|
||||
}
|
||||
};
|
||||
|
||||
// 每个新流开始时,重置编码器状态以避免残留音频导致爆音
|
||||
let mut left_encoder = OpusCodec::new(&encode_config)?;
|
||||
let mut right_encoder = OpusCodec::new(&encode_config)?;
|
||||
@@ -121,11 +109,16 @@ pub async fn run_master(master_role: ChannelRole) -> Result<()> {
|
||||
break;
|
||||
}
|
||||
|
||||
// 从 FIFO 读取
|
||||
if let Err(_) = fifo.read_exact(&mut raw_buf).await {
|
||||
break; // FIFO 关闭,重新打开
|
||||
// 如果当前没有客户端,重置流计时,以便新客户端连接时重新同步
|
||||
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()) }
|
||||
@@ -158,15 +151,14 @@ pub async fn run_master(master_role: ChannelRole) -> Result<()> {
|
||||
}
|
||||
|
||||
// 计算该帧应当播放的基准时间(相对于流开始)
|
||||
let target_ts = stream_start_ts
|
||||
+ ((seq - stream_start_seq) as u128 * frame_duration_us)
|
||||
+ delay_us;
|
||||
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]]);
|
||||
left_pcm[i] = mixed_pcm[i * 2];
|
||||
right_pcm[i] = mixed_pcm[i * 2 + 1];
|
||||
}
|
||||
|
||||
// 1. 检查各声道是否有从节点需要
|
||||
@@ -231,13 +223,8 @@ pub async fn run_master(master_role: ChannelRole) -> Result<()> {
|
||||
}
|
||||
} 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 let Err(_) = p.write(&mixed_pcm) {
|
||||
// 如果写入失败且正在退出,直接跳出循环
|
||||
if shutdown_flag_clone.load(Ordering::Relaxed) {
|
||||
break;
|
||||
@@ -249,15 +236,6 @@ pub async fn run_master(master_role: ChannelRole) -> Result<()> {
|
||||
seq += 1;
|
||||
}
|
||||
|
||||
// 重置流计时
|
||||
stream_start_ts = 0;
|
||||
|
||||
// 如果是因为退出信号而跳出内层循环,也要跳出外层循环
|
||||
if shutdown_flag_clone.load(Ordering::Relaxed) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
Ok::<(), anyhow::Error>(())
|
||||
};
|
||||
|
||||
|
||||
@@ -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<Mutex<Vec<UnixStream>>>,
|
||||
}
|
||||
|
||||
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<i16> {
|
||||
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;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user