chore: 优化连接提示
This commit is contained in:
@@ -44,7 +44,7 @@ impl AlsaRedirector {
|
|||||||
let _ = Command::new("mkfifo").arg(FIFO_PATH).status();
|
let _ = Command::new("mkfifo").arg(FIFO_PATH).status();
|
||||||
let _ = Command::new("chmod").arg("666").arg(FIFO_PATH).status();
|
let _ = Command::new("chmod").arg("666").arg(FIFO_PATH).status();
|
||||||
|
|
||||||
println!("ALSA 输出已重定向至 {}", FIFO_PATH);
|
// println!("ALSA 输出已重定向至 {}", FIFO_PATH);
|
||||||
Ok(Self)
|
Ok(Self)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -55,7 +55,7 @@ impl AlsaRedirector {
|
|||||||
.status();
|
.status();
|
||||||
let _ = fs::remove_file(TEMP_ASOUND_CONF);
|
let _ = fs::remove_file(TEMP_ASOUND_CONF);
|
||||||
let _ = fs::remove_file(FIFO_PATH);
|
let _ = fs::remove_file(FIFO_PATH);
|
||||||
println!("ALSA 配置已恢复。");
|
// println!("ALSA 配置已恢复。");
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn fifo_path() -> &'static str {
|
pub fn fifo_path() -> &'static str {
|
||||||
|
|||||||
@@ -14,7 +14,6 @@ use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
|||||||
|
|
||||||
pub const SERVER_TCP_PORT: u16 = 53531;
|
pub const SERVER_TCP_PORT: u16 = 53531;
|
||||||
|
|
||||||
/// 运行主节点模式
|
|
||||||
pub async fn run_master(role: ChannelRole) -> Result<()> {
|
pub async fn run_master(role: ChannelRole) -> Result<()> {
|
||||||
println!("--- 主节点模式 ({}) ---", role.to_string());
|
println!("--- 主节点模式 ({}) ---", role.to_string());
|
||||||
|
|
||||||
@@ -28,19 +27,21 @@ pub async fn run_master(role: ChannelRole) -> Result<()> {
|
|||||||
// 2. 启动服务发现广播
|
// 2. 启动服务发现广播
|
||||||
Discovery::start_broadcast(SERVER_TCP_PORT).await?;
|
Discovery::start_broadcast(SERVER_TCP_PORT).await?;
|
||||||
|
|
||||||
println!("正在端口 {} 等待从节点连接...", SERVER_TCP_PORT);
|
println!("✅ 服务已启动,等待连接...");
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
let (control_conn, addr) = network.accept().await?;
|
let (control_conn, client_addr) = network.accept().await?;
|
||||||
println!("从节点已连接: {}", addr);
|
|
||||||
|
|
||||||
let audio_socket = audio_socket.clone();
|
let audio_socket = audio_socket.clone();
|
||||||
let role = role.clone();
|
let role = role.clone();
|
||||||
|
|
||||||
// 启动会话处理句柄
|
// 启动会话处理句柄
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
if let Err(e) = handle_master_session(control_conn, audio_socket, role).await {
|
if let Err(e) =
|
||||||
eprintln!("会话结束: {:?}", 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(
|
async fn handle_master_session(
|
||||||
mut control: ControlConnection,
|
mut control: ControlConnection,
|
||||||
audio_socket: Arc<tokio::net::UdpSocket>,
|
audio_socket: Arc<tokio::net::UdpSocket>,
|
||||||
role: ChannelRole,
|
master_role: ChannelRole,
|
||||||
|
client_tcp_addr: String,
|
||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
let mut buf = [0u8; 1024];
|
let mut buf = [0u8; 1024];
|
||||||
|
|
||||||
|
#[allow(unused_assignments)]
|
||||||
|
let mut slave_role = ChannelRole::Left;
|
||||||
|
|
||||||
// 握手
|
// 握手
|
||||||
let pkt = control.recv_packet(&mut buf).await?;
|
let pkt = control.recv_packet(&mut buf).await?;
|
||||||
match pkt {
|
match pkt {
|
||||||
ControlPacket::ClientIdentify { role: _r } => {}
|
ControlPacket::ClientIdentify { role } => slave_role = role,
|
||||||
_ => return Err(anyhow::anyhow!("无效的握手协议")),
|
_ => return Err(anyhow::anyhow!("无效的握手协议")),
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -69,7 +74,12 @@ async fn handle_master_session(
|
|||||||
// 等待 UDP 打洞/确认
|
// 等待 UDP 打洞/确认
|
||||||
let mut buf = [0u8; 128];
|
let mut buf = [0u8; 128];
|
||||||
let (_, client_udp_addr) = audio_socket.recv_from(&mut buf).await?;
|
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 {
|
let config = AudioConfig {
|
||||||
@@ -96,11 +106,9 @@ async fn handle_master_session(
|
|||||||
let mut opus_out = vec![0u8; 1500];
|
let mut opus_out = vec![0u8; 1500];
|
||||||
let mut seq = 0u32;
|
let mut seq = 0u32;
|
||||||
|
|
||||||
|
let delay_us = 200_000; // 200ms 延迟,确保主从节点音频同步
|
||||||
let frame_duration_us =
|
let frame_duration_us =
|
||||||
(config.frame_size as f64 / config.sample_rate as f64 * 1_000_000.0) as u128;
|
(config.frame_size as f64 / config.sample_rate as f64 * 1_000_000.0) as u128;
|
||||||
let delay_us = 200_000;
|
|
||||||
|
|
||||||
println!("开始会话循环...");
|
|
||||||
|
|
||||||
// 分离 TCP 读写,以便在不同任务中使用
|
// 分离 TCP 读写,以便在不同任务中使用
|
||||||
let (mut tcp_rx, mut tcp_tx) = control.split();
|
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 {
|
match tcp_rx.read(&mut buf).await {
|
||||||
Ok(0) | Err(_) => {
|
Ok(0) | Err(_) => {
|
||||||
let _ = stop_tx.send(()).await;
|
let _ = stop_tx.send(()).await;
|
||||||
|
println!(
|
||||||
|
"❌ 从节点已断开: {} {}",
|
||||||
|
client_tcp_addr,
|
||||||
|
slave_role.to_string(),
|
||||||
|
);
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
Ok(n) => {
|
Ok(n) => {
|
||||||
@@ -141,29 +154,22 @@ async fn handle_master_session(
|
|||||||
loop {
|
loop {
|
||||||
// 打开 FIFO
|
// 打开 FIFO
|
||||||
let mut fifo = tokio::select! {
|
let mut fifo = tokio::select! {
|
||||||
_ = stop_rx.recv() => {
|
_ = stop_rx.recv() => return Ok(()),
|
||||||
println!("等待 FIFO 时从节点断开连接。");
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
f = tokio::fs::File::open(AlsaRedirector::fifo_path()) => f?,
|
f = tokio::fs::File::open(AlsaRedirector::fifo_path()) => f?,
|
||||||
};
|
};
|
||||||
|
|
||||||
println!("音频流已启动...");
|
|
||||||
let mut stream_start_ts = 0;
|
let mut stream_start_ts = 0;
|
||||||
let stream_start_seq = seq;
|
let stream_start_seq = seq;
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
// 从 FIFO 读取
|
// 从 FIFO 读取
|
||||||
let read_res = tokio::select! {
|
let read_res = tokio::select! {
|
||||||
_ = stop_rx.recv() => {
|
_ = stop_rx.recv() => return Ok(()),
|
||||||
println!("串流过程中从节点断开连接。");
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
res = fifo.read_exact(&mut raw_buf) => res,
|
res = fifo.read_exact(&mut raw_buf) => res,
|
||||||
};
|
};
|
||||||
|
|
||||||
if let Err(_) = read_res {
|
if let Err(_) = read_res {
|
||||||
println!("音频流结束 (FIFO 读取完毕)。");
|
// todo 继续等待音频流继续播放
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -179,12 +185,15 @@ async fn handle_master_session(
|
|||||||
for i in 0..config.frame_size {
|
for i in 0..config.frame_size {
|
||||||
let l = i16::from_le_bytes([raw_buf[i * 4], raw_buf[i * 4 + 1]]);
|
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]]);
|
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);
|
local_pcm.push(l);
|
||||||
remote_pcm.push(r);
|
|
||||||
} else {
|
} else {
|
||||||
local_pcm.push(r);
|
local_pcm.push(r);
|
||||||
|
}
|
||||||
|
if slave_role == ChannelRole::Left {
|
||||||
remote_pcm.push(l);
|
remote_pcm.push(l);
|
||||||
|
} else {
|
||||||
|
remote_pcm.push(r);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -201,8 +210,7 @@ async fn handle_master_session(
|
|||||||
|
|
||||||
let bytes = postcard::to_allocvec(&packet)?;
|
let bytes = postcard::to_allocvec(&packet)?;
|
||||||
if let Err(e) = audio_socket.send_to(&bytes, client_udp_addr).await {
|
if let Err(e) = audio_socket.send_to(&bytes, client_udp_addr).await {
|
||||||
eprintln!("UDP 发送错误: {:?}", e);
|
return Err(anyhow::anyhow!("UDP 发送错误: {:?}", e));
|
||||||
return Err(e.into());
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// 本地回放同步
|
// 本地回放同步
|
||||||
|
|||||||
@@ -19,24 +19,23 @@ pub async fn run_slave(role: ChannelRole) -> Result<()> {
|
|||||||
|
|
||||||
loop {
|
loop {
|
||||||
match handle_connection(role.clone()).await {
|
match handle_connection(role.clone()).await {
|
||||||
Ok(_) => println!("连接正常结束,正在尝试重新连接..."),
|
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
eprintln!("连接异常断开 3 秒后重试... \nError: {:?}.", e);
|
eprintln!("❌ {:?}", e);
|
||||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
tokio::time::sleep(Duration::from_secs(3)).await;
|
||||||
}
|
}
|
||||||
|
Ok(_) => {}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// 处理单次完整的连接生命周期
|
|
||||||
async fn handle_connection(role: ChannelRole) -> Result<()> {
|
async fn handle_connection(role: ChannelRole) -> Result<()> {
|
||||||
println!("正在扫描主节点...");
|
|
||||||
// 1. 发现主节点
|
// 1. 发现主节点
|
||||||
|
println!("🔍 正在扫描主节点...");
|
||||||
let (master_ip, master_tcp_port) = Discovery::discover_master().await?;
|
let (master_ip, master_tcp_port) = Discovery::discover_master().await?;
|
||||||
let master_tcp_addr = format!("{}:{}", master_ip, master_tcp_port);
|
let master_tcp_addr = format!("{}:{}", master_ip, master_tcp_port);
|
||||||
println!("在 {} 发现主节点,正在连接...", master_tcp_addr);
|
|
||||||
|
|
||||||
// 2. 建立 TCP 连接
|
// 2. 建立 TCP 连接
|
||||||
|
println!("🔥 发现主节点: {}", master_tcp_addr);
|
||||||
let network = SlaveNetwork::connect(master_tcp_addr.parse()?).await?;
|
let network = SlaveNetwork::connect(master_tcp_addr.parse()?).await?;
|
||||||
let (mut control, audio) = network.split();
|
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 pkt = control.recv_packet(&mut buf).await?;
|
||||||
let server_udp_port = match pkt {
|
let server_udp_port = match pkt {
|
||||||
ControlPacket::ServerHello { udp_port } => udp_port,
|
ControlPacket::ServerHello { udp_port } => udp_port,
|
||||||
_ => return Err(anyhow!("应答应为 ServerHello")),
|
_ => return Err(anyhow!("身份认证应答异常")),
|
||||||
};
|
};
|
||||||
|
|
||||||
// 4. UDP 打洞
|
// 4. UDP 打洞
|
||||||
@@ -136,14 +135,14 @@ async fn handle_connection(role: ChannelRole) -> Result<()> {
|
|||||||
});
|
});
|
||||||
|
|
||||||
// 8. 播放主循环
|
// 8. 播放主循环
|
||||||
println!("连接成功,正在播放音频...");
|
println!("✅ 主节点已连接,音频串流中...");
|
||||||
let mut pcm_buf = vec![0i16; config.frame_size];
|
let mut pcm_buf = vec![0i16; config.frame_size];
|
||||||
let mut last_seq: Option<u32> = None;
|
let mut last_seq: Option<u32> = None;
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
// 检查 TCP 是否已断开
|
// 检查 TCP 是否已断开
|
||||||
if let Ok(_) = disconnect_rx.try_recv() {
|
if let Ok(_) = disconnect_rx.try_recv() {
|
||||||
return Err(anyhow!("主节点连接已断开 (TCP Disconnected)"));
|
return Err(anyhow!("主节点已断开: {}", master_tcp_addr));
|
||||||
}
|
}
|
||||||
|
|
||||||
// 填充 Jitter Buffer
|
// 填充 Jitter Buffer
|
||||||
@@ -170,7 +169,7 @@ async fn handle_connection(role: ChannelRole) -> Result<()> {
|
|||||||
let len = codec.decode(&data, &mut pcm_buf)?;
|
let len = codec.decode(&data, &mut pcm_buf)?;
|
||||||
player.write(&pcm_buf[..len])?;
|
player.write(&pcm_buf[..len])?;
|
||||||
} else {
|
} else {
|
||||||
tokio::time::sleep(Duration::from_millis(5)).await;
|
tokio::time::sleep(Duration::from_millis(1)).await;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user