refactor: 拆分网络通信模块
This commit is contained in:
@@ -1,7 +1,6 @@
|
|||||||
#![cfg(target_os = "linux")]
|
#![cfg(target_os = "linux")]
|
||||||
|
|
||||||
use anyhow::Result;
|
use anyhow::Result;
|
||||||
use hello::stereo_core::alsa::AlsaRedirector;
|
|
||||||
use hello::stereo_core::master::run_master;
|
use hello::stereo_core::master::run_master;
|
||||||
use hello::stereo_core::protocol::ChannelRole;
|
use hello::stereo_core::protocol::ChannelRole;
|
||||||
use hello::stereo_core::slave::run_slave;
|
use hello::stereo_core::slave::run_slave;
|
||||||
@@ -22,9 +21,6 @@ async fn main() -> Result<()> {
|
|||||||
ChannelRole::Right
|
ChannelRole::Right
|
||||||
};
|
};
|
||||||
|
|
||||||
// 启动前先执行清理,确保环境干净
|
|
||||||
AlsaRedirector::cleanup();
|
|
||||||
|
|
||||||
if mode == "master" {
|
if mode == "master" {
|
||||||
run_master(role).await
|
run_master(role).await
|
||||||
} else {
|
} else {
|
||||||
|
|||||||
@@ -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::<ControlPacket>(&buf[..len])
|
||||||
|
{
|
||||||
|
return Ok((addr.ip(), udp_port));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,62 +1,45 @@
|
|||||||
#![cfg(target_os = "linux")]
|
#![cfg(target_os = "linux")]
|
||||||
use anyhow::Result;
|
|
||||||
use crate::audio::{AudioPlayer, OpusCodec};
|
use crate::audio::{AudioPlayer, OpusCodec};
|
||||||
use crate::config::AudioConfig;
|
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::sync::Arc;
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
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 SERVER_TCP_PORT: u16 = 53531;
|
||||||
pub const DISCOVERY_PORT: u16 = 53530;
|
|
||||||
|
|
||||||
/// 运行主节点模式
|
/// 运行主节点模式
|
||||||
pub async fn run_master(role: ChannelRole) -> Result<()> {
|
pub async fn run_master(role: ChannelRole) -> Result<()> {
|
||||||
println!("--- 主节点模式 ({:?}) ---", role);
|
println!("--- 主节点模式 ({}) ---", role.to_string());
|
||||||
|
|
||||||
// 0. 设置 ALSA 重定向 (在主节点生命周期内持续有效)
|
// 0. 设置 ALSA 重定向
|
||||||
let _alsa_guard = AlsaRedirector::new()?;
|
let _alsa_guard = AlsaRedirector::new()?;
|
||||||
|
|
||||||
// 1. 设置网络 (UDP + TCP)
|
// 1. 设置网络 (UDP + TCP)
|
||||||
let udp_socket = UdpSocket::bind("0.0.0.0:0").await?;
|
let network = MasterNetwork::setup(SERVER_TCP_PORT).await?;
|
||||||
let udp_port = udp_socket.local_addr()?.port();
|
let audio_socket = network.audio_socket().clone_inner();
|
||||||
let udp_socket = Arc::new(udp_socket);
|
|
||||||
|
|
||||||
let listener = TcpListener::bind(format!("0.0.0.0:{}", SERVER_TCP_PORT)).await?;
|
// 2. 启动服务发现广播
|
||||||
|
Discovery::start_broadcast(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;
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
println!("正在端口 {} 等待从节点连接...", SERVER_TCP_PORT);
|
println!("正在端口 {} 等待从节点连接...", SERVER_TCP_PORT);
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
let (socket, addr) = listener.accept().await?;
|
let (control_conn, addr) = network.accept().await?;
|
||||||
println!("从节点已连接: {}", addr);
|
println!("从节点已连接: {}", addr);
|
||||||
|
|
||||||
let udp_socket = udp_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(socket, udp_socket, udp_port, role).await {
|
if let Err(e) = handle_master_session(control_conn, audio_socket, role).await {
|
||||||
eprintln!("会话结束: {:?}", e);
|
eprintln!("会话结束: {:?}", e);
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
@@ -65,28 +48,27 @@ pub async fn run_master(role: ChannelRole) -> Result<()> {
|
|||||||
|
|
||||||
/// 处理主节点与从节点的会话
|
/// 处理主节点与从节点的会话
|
||||||
async fn handle_master_session(
|
async fn handle_master_session(
|
||||||
mut tcp: TcpStream,
|
mut control: ControlConnection,
|
||||||
udp: Arc<UdpSocket>,
|
audio_socket: Arc<tokio::net::UdpSocket>,
|
||||||
_local_udp_port: u16,
|
|
||||||
role: ChannelRole,
|
role: ChannelRole,
|
||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
let mut buf = [0u8; 1024];
|
let mut buf = [0u8; 1024];
|
||||||
|
|
||||||
// 握手
|
// 握手
|
||||||
let len = tcp.read(&mut buf).await?;
|
let pkt = control.recv_packet(&mut buf).await?;
|
||||||
match postcard::from_bytes::<ControlPacket>(&buf[..len])? {
|
match pkt {
|
||||||
ControlPacket::ClientIdentify { role: _r } => {}
|
ControlPacket::ClientIdentify { role: _r } => {}
|
||||||
_ => return Err(anyhow::anyhow!("无效的握手协议")),
|
_ => return Err(anyhow::anyhow!("无效的握手协议")),
|
||||||
};
|
};
|
||||||
|
|
||||||
let hello = ControlPacket::ServerHello {
|
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 打洞/确认
|
// 等待 UDP 打洞/确认
|
||||||
let mut buf = [0u8; 128];
|
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);
|
println!("从节点 UDP 地址已确认: {}", client_udp_addr);
|
||||||
|
|
||||||
// 配置音频
|
// 配置音频
|
||||||
@@ -121,10 +103,10 @@ async fn handle_master_session(
|
|||||||
println!("开始会话循环...");
|
println!("开始会话循环...");
|
||||||
|
|
||||||
// 分离 TCP 读写,以便在不同任务中使用
|
// 分离 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);
|
let (stop_tx, mut stop_rx) = tokio::sync::mpsc::channel::<()>(1);
|
||||||
|
|
||||||
// 处理来自从节点的控制消息(如 Ping/Pong 进行时间同步)
|
// 处理来自从节点的控制消息
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
let mut buf = [0u8; 1024];
|
let mut buf = [0u8; 1024];
|
||||||
loop {
|
loop {
|
||||||
@@ -157,8 +139,7 @@ async fn handle_master_session(
|
|||||||
});
|
});
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
// 打开 FIFO。这会阻塞直到有写入者打开它。
|
// 打开 FIFO
|
||||||
// 使用 select 以便在等待 FIFO 时如果从节点断开连接可以退出。
|
|
||||||
let mut fifo = tokio::select! {
|
let mut fifo = tokio::select! {
|
||||||
_ = stop_rx.recv() => {
|
_ = stop_rx.recv() => {
|
||||||
println!("等待 FIFO 时从节点断开连接。");
|
println!("等待 FIFO 时从节点断开连接。");
|
||||||
@@ -172,7 +153,7 @@ async fn handle_master_session(
|
|||||||
let stream_start_seq = seq;
|
let stream_start_seq = seq;
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
// 从 FIFO 读取,带超时/select 以检查断开连接
|
// 从 FIFO 读取
|
||||||
let read_res = tokio::select! {
|
let read_res = tokio::select! {
|
||||||
_ = stop_rx.recv() => {
|
_ = stop_rx.recv() => {
|
||||||
println!("串流过程中从节点断开连接。");
|
println!("串流过程中从节点断开连接。");
|
||||||
@@ -219,7 +200,7 @@ async fn handle_master_session(
|
|||||||
};
|
};
|
||||||
|
|
||||||
let bytes = postcard::to_allocvec(&packet)?;
|
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);
|
eprintln!("UDP 发送错误: {:?}", e);
|
||||||
return Err(e.into());
|
return Err(e.into());
|
||||||
}
|
}
|
||||||
@@ -238,4 +219,3 @@ async fn handle_master_session(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,8 @@
|
|||||||
pub mod alsa;
|
pub mod alsa;
|
||||||
|
pub mod discovery;
|
||||||
pub mod jitter_buffer;
|
pub mod jitter_buffer;
|
||||||
pub mod master;
|
pub mod master;
|
||||||
|
pub mod network;
|
||||||
pub mod protocol;
|
pub mod protocol;
|
||||||
pub mod slave;
|
pub mod slave;
|
||||||
pub mod sync;
|
pub mod sync;
|
||||||
|
|||||||
@@ -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<UdpSocket>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl AudioSocket {
|
||||||
|
pub async fn bind() -> Result<Self> {
|
||||||
|
let socket = UdpSocket::bind("0.0.0.0:0").await?;
|
||||||
|
Ok(Self {
|
||||||
|
socket: Arc::new(socket),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn local_port(&self) -> Result<u16> {
|
||||||
|
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<UdpSocket> {
|
||||||
|
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<ControlPacket> {
|
||||||
|
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<Self> {
|
||||||
|
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<Self> {
|
||||||
|
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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -6,6 +6,15 @@ pub enum ChannelRole {
|
|||||||
Right,
|
Right,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl ChannelRole {
|
||||||
|
pub fn to_string(&self) -> String {
|
||||||
|
match self {
|
||||||
|
ChannelRole::Left => "左声道".to_string(),
|
||||||
|
ChannelRole::Right => "右声道".to_string(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Serialize, Deserialize, Debug, Clone)]
|
#[derive(Serialize, Deserialize, Debug, Clone)]
|
||||||
pub enum ControlPacket {
|
pub enum ControlPacket {
|
||||||
// 发现协议
|
// 发现协议
|
||||||
@@ -26,7 +35,7 @@ pub enum ControlPacket {
|
|||||||
server_ts: u128,
|
server_ts: u128,
|
||||||
seq: u32,
|
seq: u32,
|
||||||
},
|
},
|
||||||
// 控制协议
|
// 音量同步
|
||||||
Volume(u8), // 音量 0-100
|
Volume(u8), // 音量 0-100
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,53 +1,45 @@
|
|||||||
#![cfg(target_os = "linux")]
|
#![cfg(target_os = "linux")]
|
||||||
use anyhow::Result;
|
|
||||||
use crate::audio::{AudioPlayer, OpusCodec};
|
use crate::audio::{AudioPlayer, OpusCodec};
|
||||||
use crate::config::AudioConfig;
|
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::sync::Arc;
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||||
use tokio::net::{TcpStream, UdpSocket};
|
|
||||||
use tokio::sync::{Mutex, mpsc};
|
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<()> {
|
pub async fn run_slave(role: ChannelRole) -> Result<()> {
|
||||||
println!("--- 从节点模式 ({:?}) ---", role);
|
println!("--- 从节点模式 ({}) ---", role.to_string());
|
||||||
|
|
||||||
println!("正在扫描主节点...");
|
println!("正在扫描主节点...");
|
||||||
let udp_disc = UdpSocket::bind(format!("0.0.0.0:{}", DISCOVERY_PORT)).await?;
|
let (master_ip, master_tcp_port) = Discovery::discover_master().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::<ControlPacket>(&buf[..len])? {
|
|
||||||
ControlPacket::ServerHello { udp_port: p } => p,
|
|
||||||
_ => return Err(anyhow::anyhow!("无效的发现数据包")),
|
|
||||||
};
|
|
||||||
|
|
||||||
let master_ip = master_addr.ip();
|
|
||||||
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);
|
||||||
|
|
||||||
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 {
|
control
|
||||||
role: role.clone(),
|
.send_packet(&ControlPacket::ClientIdentify { role: role.clone() })
|
||||||
})?)
|
.await?;
|
||||||
.await?;
|
|
||||||
|
|
||||||
let len = tcp.read(&mut buf).await?;
|
let mut buf = [0u8; 1024];
|
||||||
let server_udp_port = match postcard::from_bytes::<ControlPacket>(&buf[..len])? {
|
let pkt = control.recv_packet(&mut buf).await?;
|
||||||
|
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::anyhow!("应答应为 ServerHello")),
|
||||||
};
|
};
|
||||||
|
|
||||||
// UDP 打洞以接收音频流
|
// UDP 打洞以接收音频流
|
||||||
let udp = UdpSocket::bind("0.0.0.0:0").await?;
|
audio
|
||||||
let punch_packet = vec![0u8; 1];
|
.punch(format!("{}:{}", master_ip, server_udp_port).parse()?)
|
||||||
udp.send_to(&punch_packet, format!("{}:{}", master_ip, server_udp_port))
|
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
// 音频配置
|
// 音频配置
|
||||||
@@ -65,9 +57,9 @@ pub async fn run_slave(role: ChannelRole) -> Result<()> {
|
|||||||
let clock = ClockSync::new(100);
|
let clock = ClockSync::new(100);
|
||||||
|
|
||||||
// 分离读写
|
// 分离读写
|
||||||
let (mut tcp_rx, mut tcp_tx) = tcp.into_split();
|
let (mut tcp_rx, mut tcp_tx) = control.split();
|
||||||
let udp = Arc::new(udp);
|
let audio_socket = audio.clone_inner();
|
||||||
let udp_rx = udp.clone();
|
let udp_rx = audio_socket.clone();
|
||||||
|
|
||||||
// 定时发送 Ping 进行时间同步
|
// 定时发送 Ping 进行时间同步
|
||||||
let _sync_handle = tokio::spawn(async move {
|
let _sync_handle = tokio::spawn(async move {
|
||||||
@@ -165,4 +157,3 @@ pub async fn run_slave(role: ChannelRole) -> Result<()> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user