From 28d714a2efa74389173a1684f5ccc4d4d16ee385 Mon Sep 17 00:00:00 2001 From: Del Wang Date: Tue, 30 Dec 2025 10:16:27 +0800 Subject: [PATCH] =?UTF-8?q?refactor:=20=E6=8B=86=E5=88=86=E7=BD=91?= =?UTF-8?q?=E7=BB=9C=E9=80=9A=E4=BF=A1=E6=A8=A1=E5=9D=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- apps/hello/src/bin/stereo.rs | 4 - apps/hello/src/stereo_core/discovery.rs | 45 +++++++++ apps/hello/src/stereo_core/master.rs | 76 ++++++-------- apps/hello/src/stereo_core/mod.rs | 2 + apps/hello/src/stereo_core/network.rs | 126 ++++++++++++++++++++++++ apps/hello/src/stereo_core/protocol.rs | 11 ++- apps/hello/src/stereo_core/slave.rs | 53 +++++----- 7 files changed, 233 insertions(+), 84 deletions(-) create mode 100644 apps/hello/src/stereo_core/discovery.rs create mode 100644 apps/hello/src/stereo_core/network.rs diff --git a/apps/hello/src/bin/stereo.rs b/apps/hello/src/bin/stereo.rs index 939994d..4c80c51 100644 --- a/apps/hello/src/bin/stereo.rs +++ b/apps/hello/src/bin/stereo.rs @@ -1,7 +1,6 @@ #![cfg(target_os = "linux")] use anyhow::Result; -use hello::stereo_core::alsa::AlsaRedirector; use hello::stereo_core::master::run_master; use hello::stereo_core::protocol::ChannelRole; use hello::stereo_core::slave::run_slave; @@ -22,9 +21,6 @@ async fn main() -> Result<()> { ChannelRole::Right }; - // 启动前先执行清理,确保环境干净 - AlsaRedirector::cleanup(); - if mode == "master" { run_master(role).await } else { diff --git a/apps/hello/src/stereo_core/discovery.rs b/apps/hello/src/stereo_core/discovery.rs new file mode 100644 index 0000000..2f7b1d5 --- /dev/null +++ b/apps/hello/src/stereo_core/discovery.rs @@ -0,0 +1,45 @@ +use crate::stereo_core::protocol::ControlPacket; +use anyhow::Result; +use std::net::{IpAddr, SocketAddr}; +use std::time::Duration; +use tokio::net::UdpSocket; + +pub const DISCOVERY_PORT: u16 = 53530; + +/// 服务发现模块,用于主从节点的自动发现 +pub struct Discovery; + +impl Discovery { + /// 主节点:启动广播,告知从节点自己的 TCP 端口 + pub async fn start_broadcast(tcp_port: u16) -> Result<()> { + let socket = UdpSocket::bind("0.0.0.0:0").await?; + socket.set_broadcast(true)?; + + let target: SocketAddr = format!("255.255.255.255:{}", DISCOVERY_PORT).parse()?; + let msg = postcard::to_allocvec(&ControlPacket::ServerHello { udp_port: tcp_port })?; + + tokio::spawn(async move { + loop { + let _ = socket.send_to(&msg, target).await; + tokio::time::sleep(Duration::from_secs(1)).await; + } + }); + + Ok(()) + } + + /// 从节点:监听广播,发现主节点的 IP 和 TCP 端口 + pub async fn discover_master() -> Result<(IpAddr, u16)> { + let socket = UdpSocket::bind(format!("0.0.0.0:{}", DISCOVERY_PORT)).await?; + let mut buf = [0u8; 1024]; + + loop { + let (len, addr) = socket.recv_from(&mut buf).await?; + if let Ok(ControlPacket::ServerHello { udp_port }) = + postcard::from_bytes::(&buf[..len]) + { + return Ok((addr.ip(), udp_port)); + } + } + } +} diff --git a/apps/hello/src/stereo_core/master.rs b/apps/hello/src/stereo_core/master.rs index d6fcfb3..47759f1 100644 --- a/apps/hello/src/stereo_core/master.rs +++ b/apps/hello/src/stereo_core/master.rs @@ -1,62 +1,45 @@ #![cfg(target_os = "linux")] -use anyhow::Result; + use crate::audio::{AudioPlayer, OpusCodec}; use crate::config::AudioConfig; -use std::net::SocketAddr; +use crate::stereo_core::alsa::AlsaRedirector; +use crate::stereo_core::discovery::Discovery; +use crate::stereo_core::network::{ControlConnection, MasterNetwork}; +use crate::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket}; +use crate::stereo_core::sync::now_us; +use anyhow::Result; use std::sync::Arc; use std::time::Duration; use tokio::io::{AsyncReadExt, AsyncWriteExt}; -use tokio::net::{TcpListener, TcpStream, UdpSocket}; -use crate::stereo_core::alsa::AlsaRedirector; -use crate::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket}; -use crate::stereo_core::sync::now_us; pub const SERVER_TCP_PORT: u16 = 53531; -pub const DISCOVERY_PORT: u16 = 53530; /// 运行主节点模式 pub async fn run_master(role: ChannelRole) -> Result<()> { - println!("--- 主节点模式 ({:?}) ---", role); + println!("--- 主节点模式 ({}) ---", role.to_string()); - // 0. 设置 ALSA 重定向 (在主节点生命周期内持续有效) + // 0. 设置 ALSA 重定向 let _alsa_guard = AlsaRedirector::new()?; // 1. 设置网络 (UDP + TCP) - let udp_socket = UdpSocket::bind("0.0.0.0:0").await?; - let udp_port = udp_socket.local_addr()?.port(); - let udp_socket = Arc::new(udp_socket); + let network = MasterNetwork::setup(SERVER_TCP_PORT).await?; + let audio_socket = network.audio_socket().clone_inner(); - let listener = TcpListener::bind(format!("0.0.0.0:{}", SERVER_TCP_PORT)).await?; - - // 2. 发现服务 - let _discovery = tokio::spawn(async move { - let socket = UdpSocket::bind("0.0.0.0:0").await.unwrap(); - socket.set_broadcast(true).unwrap(); - let target: SocketAddr = format!("255.255.255.255:{}", DISCOVERY_PORT) - .parse() - .unwrap(); - let msg = postcard::to_allocvec(&ControlPacket::ServerHello { - udp_port: SERVER_TCP_PORT, - }) - .unwrap(); - loop { - let _ = socket.send_to(&msg, target).await; - tokio::time::sleep(Duration::from_secs(1)).await; - } - }); + // 2. 启动服务发现广播 + Discovery::start_broadcast(SERVER_TCP_PORT).await?; println!("正在端口 {} 等待从节点连接...", SERVER_TCP_PORT); loop { - let (socket, addr) = listener.accept().await?; + let (control_conn, addr) = network.accept().await?; println!("从节点已连接: {}", addr); - let udp_socket = udp_socket.clone(); + let audio_socket = audio_socket.clone(); let role = role.clone(); // 启动会话处理句柄 tokio::spawn(async move { - if let Err(e) = handle_master_session(socket, udp_socket, udp_port, role).await { + if let Err(e) = handle_master_session(control_conn, audio_socket, role).await { eprintln!("会话结束: {:?}", e); } }); @@ -65,28 +48,27 @@ pub async fn run_master(role: ChannelRole) -> Result<()> { /// 处理主节点与从节点的会话 async fn handle_master_session( - mut tcp: TcpStream, - udp: Arc, - _local_udp_port: u16, + mut control: ControlConnection, + audio_socket: Arc, role: ChannelRole, ) -> Result<()> { let mut buf = [0u8; 1024]; // 握手 - let len = tcp.read(&mut buf).await?; - match postcard::from_bytes::(&buf[..len])? { + let pkt = control.recv_packet(&mut buf).await?; + match pkt { ControlPacket::ClientIdentify { role: _r } => {} _ => return Err(anyhow::anyhow!("无效的握手协议")), }; let hello = ControlPacket::ServerHello { - udp_port: udp.local_addr()?.port(), + udp_port: audio_socket.local_addr()?.port(), }; - tcp.write_all(&postcard::to_allocvec(&hello)?).await?; + control.send_packet(&hello).await?; // 等待 UDP 打洞/确认 let mut buf = [0u8; 128]; - let (_, client_udp_addr) = udp.recv_from(&mut buf).await?; + let (_, client_udp_addr) = audio_socket.recv_from(&mut buf).await?; println!("从节点 UDP 地址已确认: {}", client_udp_addr); // 配置音频 @@ -121,10 +103,10 @@ async fn handle_master_session( println!("开始会话循环..."); // 分离 TCP 读写,以便在不同任务中使用 - let (mut tcp_rx, mut tcp_tx) = tcp.into_split(); + let (mut tcp_rx, mut tcp_tx) = control.split(); let (stop_tx, mut stop_rx) = tokio::sync::mpsc::channel::<()>(1); - // 处理来自从节点的控制消息(如 Ping/Pong 进行时间同步) + // 处理来自从节点的控制消息 tokio::spawn(async move { let mut buf = [0u8; 1024]; loop { @@ -157,8 +139,7 @@ async fn handle_master_session( }); loop { - // 打开 FIFO。这会阻塞直到有写入者打开它。 - // 使用 select 以便在等待 FIFO 时如果从节点断开连接可以退出。 + // 打开 FIFO let mut fifo = tokio::select! { _ = stop_rx.recv() => { println!("等待 FIFO 时从节点断开连接。"); @@ -172,7 +153,7 @@ async fn handle_master_session( let stream_start_seq = seq; loop { - // 从 FIFO 读取,带超时/select 以检查断开连接 + // 从 FIFO 读取 let read_res = tokio::select! { _ = stop_rx.recv() => { println!("串流过程中从节点断开连接。"); @@ -219,7 +200,7 @@ async fn handle_master_session( }; let bytes = postcard::to_allocvec(&packet)?; - if let Err(e) = udp.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(e.into()); } @@ -238,4 +219,3 @@ async fn handle_master_session( } } } - diff --git a/apps/hello/src/stereo_core/mod.rs b/apps/hello/src/stereo_core/mod.rs index 90aa49d..768f74d 100644 --- a/apps/hello/src/stereo_core/mod.rs +++ b/apps/hello/src/stereo_core/mod.rs @@ -1,6 +1,8 @@ pub mod alsa; +pub mod discovery; pub mod jitter_buffer; pub mod master; +pub mod network; pub mod protocol; pub mod slave; pub mod sync; diff --git a/apps/hello/src/stereo_core/network.rs b/apps/hello/src/stereo_core/network.rs new file mode 100644 index 0000000..133e441 --- /dev/null +++ b/apps/hello/src/stereo_core/network.rs @@ -0,0 +1,126 @@ +use crate::stereo_core::protocol::{AudioPacket, ControlPacket}; +use anyhow::{Context, Result}; +use std::net::SocketAddr; +use std::sync::Arc; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::{TcpListener, TcpStream, UdpSocket}; + +/// UDP 音频传输 +pub struct AudioSocket { + socket: Arc, +} + +impl AudioSocket { + pub async fn bind() -> Result { + let socket = UdpSocket::bind("0.0.0.0:0").await?; + Ok(Self { + socket: Arc::new(socket), + }) + } + + pub fn local_port(&self) -> Result { + Ok(self.socket.local_addr()?.port()) + } + + pub async fn send_packet(&self, packet: &AudioPacket, target: SocketAddr) -> Result<()> { + let bytes = postcard::to_allocvec(packet)?; + self.socket.send_to(&bytes, target).await?; + Ok(()) + } + + pub async fn recv_packet(&self, buf: &mut [u8]) -> Result<(AudioPacket, SocketAddr)> { + let (len, addr) = self.socket.recv_from(buf).await?; + let packet = postcard::from_bytes(&buf[..len])?; + Ok((packet, addr)) + } + + pub async fn punch(&self, target: SocketAddr) -> Result<()> { + self.socket.send_to(&[0u8; 1], target).await?; + Ok(()) + } + + pub fn clone_inner(&self) -> Arc { + self.socket.clone() + } +} + +/// TCP 控制连接 +pub struct ControlConnection { + stream: TcpStream, +} + +impl ControlConnection { + pub fn new(stream: TcpStream) -> Self { + Self { stream } + } + + pub async fn send_packet(&mut self, packet: &ControlPacket) -> Result<()> { + let bytes = postcard::to_allocvec(packet)?; + self.stream.write_all(&bytes).await?; + Ok(()) + } + + pub async fn recv_packet(&mut self, buf: &mut [u8]) -> Result { + let len = self.stream.read(buf).await?; + if len == 0 { + return Err(anyhow::anyhow!("连接已关闭")); + } + let packet = postcard::from_bytes(&buf[..len])?; + Ok(packet) + } + + pub fn split( + self, + ) -> ( + tokio::net::tcp::OwnedReadHalf, + tokio::net::tcp::OwnedWriteHalf, + ) { + self.stream.into_split() + } +} + +/// 主节点网络管理器 +pub struct MasterNetwork { + listener: TcpListener, + audio: AudioSocket, +} + +impl MasterNetwork { + pub async fn setup(port: u16) -> Result { + let listener = TcpListener::bind(format!("0.0.0.0:{}", port)).await?; + let audio = AudioSocket::bind().await?; + Ok(Self { listener, audio }) + } + + pub async fn accept(&self) -> Result<(ControlConnection, SocketAddr)> { + let (stream, addr) = self.listener.accept().await?; + Ok((ControlConnection::new(stream), addr)) + } + + pub fn audio_socket(&self) -> &AudioSocket { + &self.audio + } +} + +/// 从节点网络管理器 +pub struct SlaveNetwork { + control: ControlConnection, + audio: AudioSocket, +} + +impl SlaveNetwork { + pub async fn connect(master_addr: SocketAddr) -> Result { + let stream = TcpStream::connect(master_addr) + .await + .context(format!("无法连接到主节点 TCP 地址: {}", master_addr))?; + let audio = AudioSocket::bind().await?; + Ok(Self { + control: ControlConnection::new(stream), + audio, + }) + } + + pub fn split(self) -> (ControlConnection, AudioSocket) { + (self.control, self.audio) + } +} diff --git a/apps/hello/src/stereo_core/protocol.rs b/apps/hello/src/stereo_core/protocol.rs index afff181..35c352a 100644 --- a/apps/hello/src/stereo_core/protocol.rs +++ b/apps/hello/src/stereo_core/protocol.rs @@ -6,6 +6,15 @@ pub enum ChannelRole { Right, } +impl ChannelRole { + pub fn to_string(&self) -> String { + match self { + ChannelRole::Left => "左声道".to_string(), + ChannelRole::Right => "右声道".to_string(), + } + } +} + #[derive(Serialize, Deserialize, Debug, Clone)] pub enum ControlPacket { // 发现协议 @@ -26,7 +35,7 @@ pub enum ControlPacket { server_ts: u128, seq: u32, }, - // 控制协议 + // 音量同步 Volume(u8), // 音量 0-100 } diff --git a/apps/hello/src/stereo_core/slave.rs b/apps/hello/src/stereo_core/slave.rs index 020c4d9..812083c 100644 --- a/apps/hello/src/stereo_core/slave.rs +++ b/apps/hello/src/stereo_core/slave.rs @@ -1,53 +1,45 @@ #![cfg(target_os = "linux")] -use anyhow::Result; + use crate::audio::{AudioPlayer, OpusCodec}; use crate::config::AudioConfig; +use crate::stereo_core::discovery::Discovery; +use crate::stereo_core::jitter_buffer::JitterBuffer; +use crate::stereo_core::network::SlaveNetwork; +use crate::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket}; +use crate::stereo_core::sync::{ClockSync, now_us}; +use anyhow::Result; use std::sync::Arc; use std::time::Duration; use tokio::io::{AsyncReadExt, AsyncWriteExt}; -use tokio::net::{TcpStream, UdpSocket}; use tokio::sync::{Mutex, mpsc}; -use crate::stereo_core::jitter_buffer::JitterBuffer; -use crate::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket}; -use crate::stereo_core::sync::{ClockSync, now_us}; -use crate::stereo_core::master::{DISCOVERY_PORT}; /// 运行从节点模式 pub async fn run_slave(role: ChannelRole) -> Result<()> { - println!("--- 从节点模式 ({:?}) ---", role); + println!("--- 从节点模式 ({}) ---", role.to_string()); println!("正在扫描主节点..."); - let udp_disc = UdpSocket::bind(format!("0.0.0.0:{}", DISCOVERY_PORT)).await?; - let mut buf = [0u8; 1024]; - let (len, master_addr) = udp_disc.recv_from(&mut buf).await?; - - let master_tcp_port = match postcard::from_bytes::(&buf[..len])? { - ControlPacket::ServerHello { udp_port: p } => p, - _ => return Err(anyhow::anyhow!("无效的发现数据包")), - }; - - let master_ip = master_addr.ip(); + let (master_ip, master_tcp_port) = Discovery::discover_master().await?; let master_tcp_addr = format!("{}:{}", master_ip, master_tcp_port); println!("在 {} 发现主节点", master_tcp_addr); - let mut tcp = TcpStream::connect(master_tcp_addr).await?; + let network = SlaveNetwork::connect(master_tcp_addr.parse()?).await?; + let (mut control, audio) = network.split(); // 身份认证 - tcp.write_all(&postcard::to_allocvec(&ControlPacket::ClientIdentify { - role: role.clone(), - })?) - .await?; + control + .send_packet(&ControlPacket::ClientIdentify { role: role.clone() }) + .await?; - let len = tcp.read(&mut buf).await?; - let server_udp_port = match postcard::from_bytes::(&buf[..len])? { + let mut buf = [0u8; 1024]; + let pkt = control.recv_packet(&mut buf).await?; + let server_udp_port = match pkt { ControlPacket::ServerHello { udp_port } => udp_port, _ => return Err(anyhow::anyhow!("应答应为 ServerHello")), }; // UDP 打洞以接收音频流 - let udp = UdpSocket::bind("0.0.0.0:0").await?; - let punch_packet = vec![0u8; 1]; - udp.send_to(&punch_packet, format!("{}:{}", master_ip, server_udp_port)) + audio + .punch(format!("{}:{}", master_ip, server_udp_port).parse()?) .await?; // 音频配置 @@ -65,9 +57,9 @@ pub async fn run_slave(role: ChannelRole) -> Result<()> { let clock = ClockSync::new(100); // 分离读写 - let (mut tcp_rx, mut tcp_tx) = tcp.into_split(); - let udp = Arc::new(udp); - let udp_rx = udp.clone(); + let (mut tcp_rx, mut tcp_tx) = control.split(); + let audio_socket = audio.clone_inner(); + let udp_rx = audio_socket.clone(); // 定时发送 Ping 进行时间同步 let _sync_handle = tokio::spawn(async move { @@ -165,4 +157,3 @@ pub async fn run_slave(role: ChannelRole) -> Result<()> { } } } -