fix: 修复 master 退出环境清理流程

This commit is contained in:
Del Wang
2025-12-31 22:53:33 +08:00
committed by Del
parent 81ac343098
commit f4b70088fe
+134 -102
View File
@@ -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() => {},
}
}