fix: run_slave 支持主节点下线后自动重新连接
This commit is contained in:
@@ -7,7 +7,7 @@ use crate::stereo_core::jitter_buffer::JitterBuffer;
|
|||||||
use crate::stereo_core::network::SlaveNetwork;
|
use crate::stereo_core::network::SlaveNetwork;
|
||||||
use crate::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket};
|
use crate::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket};
|
||||||
use crate::stereo_core::sync::{ClockSync, now_us};
|
use crate::stereo_core::sync::{ClockSync, now_us};
|
||||||
use anyhow::Result;
|
use anyhow::{Result, anyhow};
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||||
@@ -17,15 +17,30 @@ use tokio::sync::{Mutex, mpsc};
|
|||||||
pub async fn run_slave(role: ChannelRole) -> Result<()> {
|
pub async fn run_slave(role: ChannelRole) -> Result<()> {
|
||||||
println!("--- 从节点模式 ({}) ---", role.to_string());
|
println!("--- 从节点模式 ({}) ---", role.to_string());
|
||||||
|
|
||||||
|
loop {
|
||||||
|
match handle_connection(role.clone()).await {
|
||||||
|
Ok(_) => println!("连接正常结束,正在尝试重新连接..."),
|
||||||
|
Err(e) => {
|
||||||
|
eprintln!("连接异常断开 3 秒后重试... \nError: {:?}.", e);
|
||||||
|
tokio::time::sleep(Duration::from_secs(3)).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// 处理单次完整的连接生命周期
|
||||||
|
async fn handle_connection(role: ChannelRole) -> Result<()> {
|
||||||
println!("正在扫描主节点...");
|
println!("正在扫描主节点...");
|
||||||
|
// 1. 发现主节点
|
||||||
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);
|
println!("在 {} 发现主节点,正在连接...", master_tcp_addr);
|
||||||
|
|
||||||
|
// 2. 建立 TCP 连接
|
||||||
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();
|
||||||
|
|
||||||
// 身份认证
|
// 3. 身份认证
|
||||||
control
|
control
|
||||||
.send_packet(&ControlPacket::ClientIdentify { role: role.clone() })
|
.send_packet(&ControlPacket::ClientIdentify { role: role.clone() })
|
||||||
.await?;
|
.await?;
|
||||||
@@ -34,15 +49,15 @@ pub async fn run_slave(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::anyhow!("应答应为 ServerHello")),
|
_ => return Err(anyhow!("应答应为 ServerHello")),
|
||||||
};
|
};
|
||||||
|
|
||||||
// UDP 打洞以接收音频流
|
// 4. UDP 打洞
|
||||||
audio
|
audio
|
||||||
.punch(format!("{}:{}", master_ip, server_udp_port).parse()?)
|
.punch(format!("{}:{}", master_ip, server_udp_port).parse()?)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
// 音频配置
|
// 5. 初始化音频与同步组件
|
||||||
let config = AudioConfig {
|
let config = AudioConfig {
|
||||||
sample_rate: 48000,
|
sample_rate: 48000,
|
||||||
channels: 1,
|
channels: 1,
|
||||||
@@ -50,28 +65,29 @@ pub async fn run_slave(role: ChannelRole) -> Result<()> {
|
|||||||
bitrate: 64000,
|
bitrate: 64000,
|
||||||
..AudioConfig::default()
|
..AudioConfig::default()
|
||||||
};
|
};
|
||||||
|
|
||||||
let player = AudioPlayer::new(&config)?;
|
let player = AudioPlayer::new(&config)?;
|
||||||
let mut codec = OpusCodec::new(&config)?;
|
let mut codec = OpusCodec::new(&config)?;
|
||||||
let mut jitter = JitterBuffer::new(50_000);
|
let mut jitter = JitterBuffer::new(50_000);
|
||||||
let clock = ClockSync::new(100);
|
let clock = Arc::new(Mutex::new(ClockSync::new(100)));
|
||||||
|
|
||||||
// 分离读写
|
// 用于通知主循环 TCP 已断开的消息通道
|
||||||
|
let (disconnect_tx, mut disconnect_rx) = mpsc::channel::<()>(1);
|
||||||
|
|
||||||
|
// 6. 分离 TCP 读写
|
||||||
let (mut tcp_rx, mut tcp_tx) = control.split();
|
let (mut tcp_rx, mut tcp_tx) = control.split();
|
||||||
let audio_socket = audio.clone_inner();
|
let clock_updater = clock.clone();
|
||||||
let udp_rx = audio_socket.clone();
|
let d_tx_ping = disconnect_tx.clone();
|
||||||
|
let d_tx_pong = disconnect_tx.clone();
|
||||||
|
|
||||||
// 定时发送 Ping 进行时间同步
|
// 定时发送 Ping (心跳 & 时间同步)
|
||||||
let _sync_handle = tokio::spawn(async move {
|
let _sync_handle = tokio::spawn(async move {
|
||||||
let mut seq = 0;
|
let mut seq = 0;
|
||||||
loop {
|
loop {
|
||||||
let t1 = now_us();
|
let t1 = now_us();
|
||||||
let msg = ControlPacket::Ping { client_ts: t1, seq };
|
let msg = ControlPacket::Ping { client_ts: t1, seq };
|
||||||
if tcp_tx
|
let data = postcard::to_allocvec(&msg).unwrap();
|
||||||
.write_all(&postcard::to_allocvec(&msg).unwrap())
|
if tcp_tx.write_all(&data).await.is_err() {
|
||||||
.await
|
let _ = d_tx_ping.send(()).await; // 通知主线程 TCP 失败
|
||||||
.is_err()
|
|
||||||
{
|
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
tokio::time::sleep(Duration::from_secs(1)).await;
|
tokio::time::sleep(Duration::from_secs(1)).await;
|
||||||
@@ -79,10 +95,7 @@ pub async fn run_slave(role: ChannelRole) -> Result<()> {
|
|||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
let clock = Arc::new(Mutex::new(clock));
|
// 接收 Pong
|
||||||
let clock_updater = clock.clone();
|
|
||||||
|
|
||||||
// 接收 Pong 并更新时钟偏移
|
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
let mut buf = [0u8; 1024];
|
let mut buf = [0u8; 1024];
|
||||||
loop {
|
loop {
|
||||||
@@ -98,48 +111,53 @@ pub async fn run_slave(role: ChannelRole) -> Result<()> {
|
|||||||
clock_updater.lock().await.update(client_ts, server_ts, t4);
|
clock_updater.lock().await.update(client_ts, server_ts, t4);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
_ => break,
|
_ => {
|
||||||
}
|
let _ = d_tx_pong.send(()).await; // TCP 断开
|
||||||
}
|
break;
|
||||||
});
|
|
||||||
|
|
||||||
// 接收音频数据包
|
|
||||||
let (audio_tx, mut audio_rx) = mpsc::channel(100);
|
|
||||||
tokio::spawn(async move {
|
|
||||||
let mut buf = [0u8; 2048];
|
|
||||||
loop {
|
|
||||||
if let Ok((len, _)) = udp_rx.recv_from(&mut buf).await {
|
|
||||||
if let Ok(packet) = postcard::from_bytes::<AudioPacket>(&buf[..len]) {
|
|
||||||
let _ = audio_tx.send(packet).await;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
let mut pcm_buf = vec![0i16; config.frame_size];
|
// 7. 接收音频数据包 (UDP)
|
||||||
println!("正在监听音频数据...");
|
let (audio_tx, mut audio_rx) = mpsc::channel(100);
|
||||||
|
let audio_socket = audio.clone_inner();
|
||||||
|
tokio::spawn(async move {
|
||||||
|
let mut buf = [0u8; 2048];
|
||||||
|
loop {
|
||||||
|
if let Ok((len, _)) = audio_socket.recv_from(&mut buf).await {
|
||||||
|
if let Ok(packet) = postcard::from_bytes::<AudioPacket>(&buf[..len]) {
|
||||||
|
if audio_tx.send(packet).await.is_err() {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
// 8. 播放主循环
|
||||||
|
println!("连接成功,正在播放音频...");
|
||||||
|
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 是否已断开
|
||||||
|
if let Ok(_) = disconnect_rx.try_recv() {
|
||||||
|
return Err(anyhow!("主节点连接已断开 (TCP Disconnected)"));
|
||||||
|
}
|
||||||
|
|
||||||
|
// 填充 Jitter Buffer
|
||||||
while let Ok(pkt) = audio_rx.try_recv() {
|
while let Ok(pkt) = audio_rx.try_recv() {
|
||||||
jitter.push(pkt);
|
jitter.push(pkt);
|
||||||
}
|
}
|
||||||
|
|
||||||
let now = now_us();
|
let now = now_us();
|
||||||
let current_server_time = {
|
let current_server_time = clock.lock().await.to_server_time(now);
|
||||||
let c = clock.lock().await;
|
|
||||||
c.to_server_time(now)
|
|
||||||
};
|
|
||||||
|
|
||||||
// 从抖动缓冲区提取待播放帧
|
|
||||||
if let Some((seq, data)) = jitter.pop_frame(current_server_time) {
|
if let Some((seq, data)) = jitter.pop_frame(current_server_time) {
|
||||||
// 丢包处理 (PLC)
|
// 丢失处理 (PLC)
|
||||||
if let Some(last) = last_seq {
|
if let Some(last) = last_seq {
|
||||||
if seq > last + 1 {
|
if seq > last + 1 {
|
||||||
println!("检测到丢包: {} -> {}", last, seq);
|
|
||||||
// 补偿丢失的帧
|
|
||||||
for _ in 0..(seq - last - 1) {
|
for _ in 0..(seq - last - 1) {
|
||||||
if let Ok(len) = codec.decode_loss(&mut pcm_buf) {
|
if let Ok(len) = codec.decode_loss(&mut pcm_buf) {
|
||||||
let _ = player.write(&pcm_buf[..len]);
|
let _ = player.write(&pcm_buf[..len]);
|
||||||
@@ -152,8 +170,7 @@ pub async fn run_slave(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 {
|
||||||
// 稍作等待以减少 CPU 占用
|
tokio::time::sleep(Duration::from_millis(5)).await;
|
||||||
tokio::time::sleep(Duration::from_millis(1)).await;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user