refactor: 拆分 stereo 模块
This commit is contained in:
@@ -1,103 +1,17 @@
|
|||||||
#![cfg(target_os = "linux")]
|
#![cfg(target_os = "linux")]
|
||||||
|
|
||||||
use anyhow::{Context, Result};
|
use anyhow::Result;
|
||||||
use hello::audio::{AudioPlayer, OpusCodec};
|
use hello::stereo_core::alsa::AlsaRedirector;
|
||||||
use hello::config::AudioConfig;
|
use hello::stereo_core::master::run_master;
|
||||||
use hello::stereo_core::jitter_buffer::JitterBuffer;
|
use hello::stereo_core::protocol::ChannelRole;
|
||||||
use hello::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket};
|
use hello::stereo_core::slave::run_slave;
|
||||||
use hello::stereo_core::sync::{ClockSync, now_us};
|
|
||||||
use std::env;
|
use std::env;
|
||||||
use std::fs;
|
|
||||||
use std::net::SocketAddr;
|
|
||||||
use std::process::Command;
|
|
||||||
use std::sync::Arc;
|
|
||||||
use std::time::Duration;
|
|
||||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
|
||||||
use tokio::net::{TcpListener, TcpStream, UdpSocket};
|
|
||||||
use tokio::sync::{Mutex, mpsc};
|
|
||||||
|
|
||||||
const FIFO_PATH: &str = "/tmp/stereo_out.fifo";
|
|
||||||
const REAL_ASOUND_CONF: &str = "/etc/asound.conf";
|
|
||||||
const TEMP_ASOUND_CONF: &str = "/tmp/asound.stereo.conf";
|
|
||||||
|
|
||||||
// Default ports
|
|
||||||
const DISCOVERY_PORT: u16 = 53530;
|
|
||||||
const SERVER_TCP_PORT: u16 = 53531;
|
|
||||||
// UDP port will be dynamic or fixed
|
|
||||||
|
|
||||||
struct AlsaRedirector;
|
|
||||||
|
|
||||||
impl AlsaRedirector {
|
|
||||||
fn new() -> Result<Self> {
|
|
||||||
Self::cleanup(); // Ensure clean state
|
|
||||||
|
|
||||||
let original_conf = fs::read_to_string(REAL_ASOUND_CONF).unwrap_or_default();
|
|
||||||
|
|
||||||
if !original_conf.contains("pcm.original_default") {
|
|
||||||
// 重命名原有的 default 逻辑,插入 interceptor
|
|
||||||
let mut new_conf = original_conf.replace("pcm.!default", "pcm.original_default");
|
|
||||||
new_conf.push_str(&format!(
|
|
||||||
r#"
|
|
||||||
pcm.!default {{
|
|
||||||
type plug
|
|
||||||
slave.pcm "stereo_interceptor"
|
|
||||||
}}
|
|
||||||
|
|
||||||
pcm.stereo_interceptor {{
|
|
||||||
type file
|
|
||||||
slave.pcm "null"
|
|
||||||
file "{}"
|
|
||||||
format "raw"
|
|
||||||
}}
|
|
||||||
"#,
|
|
||||||
FIFO_PATH
|
|
||||||
));
|
|
||||||
|
|
||||||
fs::write(TEMP_ASOUND_CONF, new_conf)?;
|
|
||||||
|
|
||||||
// 挂载覆盖 /etc/asound.conf
|
|
||||||
let status = Command::new("mount")
|
|
||||||
.arg("--bind")
|
|
||||||
.arg(TEMP_ASOUND_CONF)
|
|
||||||
.arg(REAL_ASOUND_CONF)
|
|
||||||
.status()
|
|
||||||
.context("Failed to execute mount command")?;
|
|
||||||
|
|
||||||
if !status.success() {
|
|
||||||
return Err(anyhow::anyhow!("Failed to mount asound.conf"));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Create FIFO
|
|
||||||
let _ = Command::new("mkfifo").arg(FIFO_PATH).status();
|
|
||||||
let _ = Command::new("chmod").arg("666").arg(FIFO_PATH).status();
|
|
||||||
|
|
||||||
println!("ALSA output redirected to {}", FIFO_PATH);
|
|
||||||
Ok(Self)
|
|
||||||
}
|
|
||||||
|
|
||||||
fn cleanup() {
|
|
||||||
let _ = Command::new("umount")
|
|
||||||
.arg("-l")
|
|
||||||
.arg(REAL_ASOUND_CONF)
|
|
||||||
.status();
|
|
||||||
let _ = fs::remove_file(TEMP_ASOUND_CONF);
|
|
||||||
let _ = fs::remove_file(FIFO_PATH);
|
|
||||||
println!("ALSA configuration restored.");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl Drop for AlsaRedirector {
|
|
||||||
fn drop(&mut self) {
|
|
||||||
Self::cleanup();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tokio::main]
|
#[tokio::main]
|
||||||
async fn main() -> Result<()> {
|
async fn main() -> Result<()> {
|
||||||
let args: Vec<String> = env::args().collect();
|
let args: Vec<String> = env::args().collect();
|
||||||
if args.len() < 3 {
|
if args.len() < 3 {
|
||||||
eprintln!("Usage: {} [master|slave] [left|right]", args[0]);
|
eprintln!("用法: {} [master|slave] [left|right]", args[0]);
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -108,7 +22,7 @@ async fn main() -> Result<()> {
|
|||||||
ChannelRole::Right
|
ChannelRole::Right
|
||||||
};
|
};
|
||||||
|
|
||||||
// Explicit cleanup at start just in case
|
// 启动前先执行清理,确保环境干净
|
||||||
AlsaRedirector::cleanup();
|
AlsaRedirector::cleanup();
|
||||||
|
|
||||||
if mode == "master" {
|
if mode == "master" {
|
||||||
@@ -117,365 +31,3 @@ async fn main() -> Result<()> {
|
|||||||
run_slave(role).await
|
run_slave(role).await
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn run_master(role: ChannelRole) -> Result<()> {
|
|
||||||
println!("--- Master Mode ({:?}) ---", role);
|
|
||||||
|
|
||||||
// 0. Setup ALSA redirection (persistent for the life of the master)
|
|
||||||
let _alsa_guard = AlsaRedirector::new()?;
|
|
||||||
|
|
||||||
// 1. Setup Network (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 listener = TcpListener::bind(format!("0.0.0.0:{}", SERVER_TCP_PORT)).await?;
|
|
||||||
|
|
||||||
// 2. Discovery Service
|
|
||||||
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!("Waiting for slave on port {}...", SERVER_TCP_PORT);
|
|
||||||
|
|
||||||
loop {
|
|
||||||
let (socket, addr) = listener.accept().await?;
|
|
||||||
println!("Slave connected: {}", addr);
|
|
||||||
|
|
||||||
let udp_socket = udp_socket.clone();
|
|
||||||
let role = role.clone();
|
|
||||||
|
|
||||||
// Spawn connection handler
|
|
||||||
tokio::spawn(async move {
|
|
||||||
if let Err(e) = handle_master_session(socket, udp_socket, udp_port, role).await {
|
|
||||||
eprintln!("Session ended: {:?}", e);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn handle_master_session(
|
|
||||||
mut tcp: TcpStream,
|
|
||||||
udp: Arc<UdpSocket>,
|
|
||||||
_local_udp_port: u16,
|
|
||||||
role: ChannelRole,
|
|
||||||
) -> Result<()> {
|
|
||||||
let mut buf = [0u8; 1024];
|
|
||||||
|
|
||||||
let len = tcp.read(&mut buf).await?;
|
|
||||||
match postcard::from_bytes::<ControlPacket>(&buf[..len])? {
|
|
||||||
ControlPacket::ClientIdentify { role: _r } => {}
|
|
||||||
_ => return Err(anyhow::anyhow!("Invalid handshake")),
|
|
||||||
};
|
|
||||||
|
|
||||||
let hello = ControlPacket::ServerHello {
|
|
||||||
udp_port: udp.local_addr()?.port(),
|
|
||||||
};
|
|
||||||
tcp.write_all(&postcard::to_allocvec(&hello)?).await?;
|
|
||||||
|
|
||||||
let mut buf = [0u8; 128];
|
|
||||||
let (_, client_udp_addr) = udp.recv_from(&mut buf).await?;
|
|
||||||
println!("Client UDP address confirmed: {}", client_udp_addr);
|
|
||||||
|
|
||||||
let config = AudioConfig {
|
|
||||||
sample_rate: 48000,
|
|
||||||
channels: 2,
|
|
||||||
frame_size: 960,
|
|
||||||
bitrate: 64000,
|
|
||||||
..AudioConfig::default()
|
|
||||||
};
|
|
||||||
|
|
||||||
let mono_config = AudioConfig {
|
|
||||||
channels: 1,
|
|
||||||
..config.clone()
|
|
||||||
};
|
|
||||||
let mut codec = OpusCodec::new(&mono_config)?;
|
|
||||||
let playback_config = AudioConfig {
|
|
||||||
channels: 1,
|
|
||||||
playback_device: "plug:original_default".into(),
|
|
||||||
..config.clone()
|
|
||||||
};
|
|
||||||
let player = AudioPlayer::new(&playback_config)?;
|
|
||||||
|
|
||||||
let mut raw_buf = vec![0u8; config.frame_size * 2 * 2];
|
|
||||||
let mut opus_out = vec![0u8; 1500];
|
|
||||||
let mut seq = 0u32;
|
|
||||||
|
|
||||||
let frame_duration_us =
|
|
||||||
(config.frame_size as f64 / config.sample_rate as f64 * 1_000_000.0) as u128;
|
|
||||||
let delay_us = 200_000;
|
|
||||||
|
|
||||||
println!("Starting session loop...");
|
|
||||||
|
|
||||||
// Use into_split for owned ReadHalf/WriteHalf to move into tasks
|
|
||||||
let (mut tcp_rx, mut tcp_tx) = tcp.into_split();
|
|
||||||
let (stop_tx, mut stop_rx) = tokio::sync::mpsc::channel::<()>(1);
|
|
||||||
|
|
||||||
tokio::spawn(async move {
|
|
||||||
let mut buf = [0u8; 1024];
|
|
||||||
loop {
|
|
||||||
match tcp_rx.read(&mut buf).await {
|
|
||||||
Ok(0) | Err(_) => {
|
|
||||||
let _ = stop_tx.send(()).await;
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
Ok(n) => {
|
|
||||||
if let Ok(ControlPacket::Ping { client_ts, seq }) =
|
|
||||||
postcard::from_bytes(&buf[..n])
|
|
||||||
{
|
|
||||||
let pong = ControlPacket::Pong {
|
|
||||||
client_ts,
|
|
||||||
server_ts: now_us(),
|
|
||||||
seq,
|
|
||||||
};
|
|
||||||
if tcp_tx
|
|
||||||
.write_all(&postcard::to_allocvec(&pong).unwrap())
|
|
||||||
.await
|
|
||||||
.is_err()
|
|
||||||
{
|
|
||||||
let _ = stop_tx.send(()).await;
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
loop {
|
|
||||||
// Open FIFO. This will block until a writer opens it.
|
|
||||||
// We use select to allow exiting if the slave disconnects while we wait.
|
|
||||||
let mut fifo = tokio::select! {
|
|
||||||
_ = stop_rx.recv() => {
|
|
||||||
println!("Slave disconnected while waiting for FIFO.");
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
f = tokio::fs::File::open(FIFO_PATH) => f?,
|
|
||||||
};
|
|
||||||
|
|
||||||
println!("Audio stream started...");
|
|
||||||
let mut stream_start_ts = 0;
|
|
||||||
let stream_start_seq = seq;
|
|
||||||
|
|
||||||
loop {
|
|
||||||
// Read from FIFO with timeout/select to check for disconnection
|
|
||||||
let read_res = tokio::select! {
|
|
||||||
_ = stop_rx.recv() => {
|
|
||||||
println!("Slave disconnected during streaming.");
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
res = fifo.read_exact(&mut raw_buf) => res,
|
|
||||||
};
|
|
||||||
|
|
||||||
if let Err(_) = read_res {
|
|
||||||
println!("Audio stream ended (FIFO EOF).");
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
|
|
||||||
let now = now_us();
|
|
||||||
if stream_start_ts == 0 {
|
|
||||||
stream_start_ts = now;
|
|
||||||
}
|
|
||||||
|
|
||||||
let mut local_pcm = Vec::with_capacity(config.frame_size);
|
|
||||||
let mut remote_pcm = Vec::with_capacity(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 r = i16::from_le_bytes([raw_buf[i * 4 + 2], raw_buf[i * 4 + 3]]);
|
|
||||||
if role == ChannelRole::Left {
|
|
||||||
local_pcm.push(l);
|
|
||||||
remote_pcm.push(r);
|
|
||||||
} else {
|
|
||||||
local_pcm.push(r);
|
|
||||||
remote_pcm.push(l);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
let len = 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)?;
|
|
||||||
if let Err(e) = udp.send_to(&bytes, client_udp_addr).await {
|
|
||||||
eprintln!("UDP send error: {:?}", e);
|
|
||||||
return Err(e.into());
|
|
||||||
}
|
|
||||||
|
|
||||||
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;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
player.write(&local_pcm)?;
|
|
||||||
|
|
||||||
seq += 1;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn run_slave(role: ChannelRole) -> Result<()> {
|
|
||||||
println!("--- Slave Mode ({:?}) ---", role);
|
|
||||||
|
|
||||||
println!("Scanning for Master...");
|
|
||||||
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::<ControlPacket>(&buf[..len])? {
|
|
||||||
ControlPacket::ServerHello { udp_port: p } => p,
|
|
||||||
_ => return Err(anyhow::anyhow!("Invalid discovery packet")),
|
|
||||||
};
|
|
||||||
|
|
||||||
let master_ip = master_addr.ip();
|
|
||||||
let master_tcp_addr = format!("{}:{}", master_ip, master_tcp_port);
|
|
||||||
println!("Found Master at {}", master_tcp_addr);
|
|
||||||
|
|
||||||
let mut tcp = TcpStream::connect(master_tcp_addr).await?;
|
|
||||||
|
|
||||||
tcp.write_all(&postcard::to_allocvec(&ControlPacket::ClientIdentify {
|
|
||||||
role: role.clone(),
|
|
||||||
})?)
|
|
||||||
.await?;
|
|
||||||
|
|
||||||
let len = tcp.read(&mut buf).await?;
|
|
||||||
let server_udp_port = match postcard::from_bytes::<ControlPacket>(&buf[..len])? {
|
|
||||||
ControlPacket::ServerHello { udp_port } => udp_port,
|
|
||||||
_ => return Err(anyhow::anyhow!("Expected ServerHello")),
|
|
||||||
};
|
|
||||||
|
|
||||||
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))
|
|
||||||
.await?;
|
|
||||||
|
|
||||||
let config = AudioConfig {
|
|
||||||
sample_rate: 48000,
|
|
||||||
channels: 1,
|
|
||||||
frame_size: 960,
|
|
||||||
bitrate: 64000,
|
|
||||||
..AudioConfig::default()
|
|
||||||
};
|
|
||||||
|
|
||||||
let player = AudioPlayer::new(&config)?;
|
|
||||||
let mut codec = OpusCodec::new(&config)?;
|
|
||||||
let mut jitter = JitterBuffer::new(50_000);
|
|
||||||
let clock = ClockSync::new(100);
|
|
||||||
|
|
||||||
// Use into_split for owned ReadHalf/WriteHalf to move into tasks
|
|
||||||
let (mut tcp_rx, mut tcp_tx) = tcp.into_split();
|
|
||||||
let udp = Arc::new(udp);
|
|
||||||
let udp_rx = udp.clone();
|
|
||||||
|
|
||||||
let _sync_handle = tokio::spawn(async move {
|
|
||||||
let mut seq = 0;
|
|
||||||
loop {
|
|
||||||
let t1 = now_us();
|
|
||||||
let msg = ControlPacket::Ping { client_ts: t1, seq };
|
|
||||||
if tcp_tx
|
|
||||||
.write_all(&postcard::to_allocvec(&msg).unwrap())
|
|
||||||
.await
|
|
||||||
.is_err()
|
|
||||||
{
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
tokio::time::sleep(Duration::from_secs(1)).await;
|
|
||||||
seq += 1;
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
let clock = Arc::new(Mutex::new(clock));
|
|
||||||
let clock_updater = clock.clone();
|
|
||||||
|
|
||||||
tokio::spawn(async move {
|
|
||||||
let mut buf = [0u8; 1024];
|
|
||||||
loop {
|
|
||||||
match tcp_rx.read(&mut buf).await {
|
|
||||||
Ok(n) if n > 0 => {
|
|
||||||
if let Ok(ControlPacket::Pong {
|
|
||||||
client_ts,
|
|
||||||
server_ts,
|
|
||||||
..
|
|
||||||
}) = postcard::from_bytes(&buf[..n])
|
|
||||||
{
|
|
||||||
let t4 = now_us();
|
|
||||||
clock_updater.lock().await.update(client_ts, server_ts, t4);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
_ => 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];
|
|
||||||
println!("Listening for audio...");
|
|
||||||
|
|
||||||
let mut last_seq: Option<u32> = None;
|
|
||||||
|
|
||||||
loop {
|
|
||||||
while let Ok(pkt) = audio_rx.try_recv() {
|
|
||||||
jitter.push(pkt);
|
|
||||||
}
|
|
||||||
|
|
||||||
let now = now_us();
|
|
||||||
let current_server_time = {
|
|
||||||
let c = clock.lock().await;
|
|
||||||
c.to_server_time(now)
|
|
||||||
};
|
|
||||||
|
|
||||||
if let Some((seq, data)) = jitter.pop_frame(current_server_time) {
|
|
||||||
// PLC logic
|
|
||||||
let mut _lost = false;
|
|
||||||
if let Some(last) = last_seq {
|
|
||||||
if seq > last + 1 {
|
|
||||||
println!("Packet loss detected: {} -> {}", last, seq);
|
|
||||||
_lost = true;
|
|
||||||
// Conceal missing frames
|
|
||||||
for _ in 0..(seq - last - 1) {
|
|
||||||
if let Ok(len) = codec.decode_loss(&mut pcm_buf) {
|
|
||||||
let _ = player.write(&pcm_buf[..len]);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
last_seq = Some(seq);
|
|
||||||
|
|
||||||
let len = codec.decode(&data, &mut pcm_buf)?;
|
|
||||||
player.write(&pcm_buf[..len])?;
|
|
||||||
} else {
|
|
||||||
tokio::time::sleep(Duration::from_millis(1)).await;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -0,0 +1,70 @@
|
|||||||
|
#![cfg(target_os = "linux")]
|
||||||
|
use anyhow::{Context, Result};
|
||||||
|
use std::fs;
|
||||||
|
use std::process::Command;
|
||||||
|
|
||||||
|
const FIFO_PATH: &str = "/tmp/stereo_out.fifo";
|
||||||
|
const REAL_ASOUND_CONF: &str = "/etc/asound.conf";
|
||||||
|
const TEMP_ASOUND_CONF: &str = "/tmp/asound.stereo.conf";
|
||||||
|
|
||||||
|
/// ALSA 音频重定向器,用于拦截系统音频输出到 FIFO 管道
|
||||||
|
pub struct AlsaRedirector;
|
||||||
|
|
||||||
|
impl AlsaRedirector {
|
||||||
|
pub fn new() -> Result<Self> {
|
||||||
|
Self::cleanup(); // 确保环境干净
|
||||||
|
|
||||||
|
let original_conf = fs::read_to_string(REAL_ASOUND_CONF).unwrap_or_default();
|
||||||
|
|
||||||
|
if !original_conf.contains("pcm.original_default") {
|
||||||
|
// 重命名原有的 default 逻辑,插入拦截器
|
||||||
|
let mut new_conf = original_conf.replace("pcm.!default", "pcm.original_default");
|
||||||
|
new_conf.push_str(&format!(
|
||||||
|
"\npcm.!default {{ type plug slave.pcm \"stereo_interceptor\" }}\n\
|
||||||
|
pcm.stereo_interceptor {{ type file slave.pcm \"null\" file \"{}\" format \"raw\" }}\n",
|
||||||
|
FIFO_PATH
|
||||||
|
));
|
||||||
|
|
||||||
|
fs::write(TEMP_ASOUND_CONF, new_conf)?;
|
||||||
|
|
||||||
|
// 挂载覆盖 /etc/asound.conf
|
||||||
|
let status = Command::new("mount")
|
||||||
|
.arg("--bind")
|
||||||
|
.arg(TEMP_ASOUND_CONF)
|
||||||
|
.arg(REAL_ASOUND_CONF)
|
||||||
|
.status()
|
||||||
|
.context("执行 mount 命令失败")?;
|
||||||
|
|
||||||
|
if !status.success() {
|
||||||
|
return Err(anyhow::anyhow!("挂载 asound.conf 失败"));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 创建 FIFO 管道
|
||||||
|
let _ = Command::new("mkfifo").arg(FIFO_PATH).status();
|
||||||
|
let _ = Command::new("chmod").arg("666").arg(FIFO_PATH).status();
|
||||||
|
|
||||||
|
println!("ALSA 输出已重定向至 {}", FIFO_PATH);
|
||||||
|
Ok(Self)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn cleanup() {
|
||||||
|
let _ = Command::new("sh")
|
||||||
|
.arg("-c")
|
||||||
|
.arg(format!("umount -l {} >/dev/null 2>&1", REAL_ASOUND_CONF))
|
||||||
|
.status();
|
||||||
|
let _ = fs::remove_file(TEMP_ASOUND_CONF);
|
||||||
|
let _ = fs::remove_file(FIFO_PATH);
|
||||||
|
println!("ALSA 配置已恢复。");
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn fifo_path() -> &'static str {
|
||||||
|
FIFO_PATH
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Drop for AlsaRedirector {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
Self::cleanup();
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -22,6 +22,7 @@ impl Ord for OrderedPacket {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// 抖动缓冲区,用于处理网络延迟和乱序
|
||||||
pub struct JitterBuffer {
|
pub struct JitterBuffer {
|
||||||
buffer: BinaryHeap<OrderedPacket>,
|
buffer: BinaryHeap<OrderedPacket>,
|
||||||
last_played_seq: Option<u32>,
|
last_played_seq: Option<u32>,
|
||||||
@@ -38,6 +39,7 @@ impl JitterBuffer {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub fn push(&mut self, packet: AudioPacket) {
|
pub fn push(&mut self, packet: AudioPacket) {
|
||||||
|
// 丢弃已经播放过的数据包
|
||||||
if let Some(last) = self.last_played_seq {
|
if let Some(last) = self.last_played_seq {
|
||||||
if packet.seq <= last {
|
if packet.seq <= last {
|
||||||
return;
|
return;
|
||||||
@@ -46,6 +48,7 @@ impl JitterBuffer {
|
|||||||
self.buffer.push(OrderedPacket(packet));
|
self.buffer.push(OrderedPacket(packet));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// 根据当前时间弹出待播放的帧
|
||||||
pub fn pop_frame(&mut self, current_time: u128) -> Option<(u32, Vec<u8>)> {
|
pub fn pop_frame(&mut self, current_time: u128) -> Option<(u32, Vec<u8>)> {
|
||||||
if let Some(OrderedPacket(pkt)) = self.buffer.peek() {
|
if let Some(OrderedPacket(pkt)) = self.buffer.peek() {
|
||||||
if current_time >= pkt.timestamp {
|
if current_time >= pkt.timestamp {
|
||||||
|
|||||||
@@ -0,0 +1,241 @@
|
|||||||
|
#![cfg(target_os = "linux")]
|
||||||
|
use anyhow::Result;
|
||||||
|
use crate::audio::{AudioPlayer, OpusCodec};
|
||||||
|
use crate::config::AudioConfig;
|
||||||
|
use std::net::SocketAddr;
|
||||||
|
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);
|
||||||
|
|
||||||
|
// 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 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;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
println!("正在端口 {} 等待从节点连接...", SERVER_TCP_PORT);
|
||||||
|
|
||||||
|
loop {
|
||||||
|
let (socket, addr) = listener.accept().await?;
|
||||||
|
println!("从节点已连接: {}", addr);
|
||||||
|
|
||||||
|
let udp_socket = udp_socket.clone();
|
||||||
|
let role = role.clone();
|
||||||
|
|
||||||
|
// 启动会话处理句柄
|
||||||
|
tokio::spawn(async move {
|
||||||
|
if let Err(e) = handle_master_session(socket, udp_socket, udp_port, role).await {
|
||||||
|
eprintln!("会话结束: {:?}", e);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// 处理主节点与从节点的会话
|
||||||
|
async fn handle_master_session(
|
||||||
|
mut tcp: TcpStream,
|
||||||
|
udp: Arc<UdpSocket>,
|
||||||
|
_local_udp_port: u16,
|
||||||
|
role: ChannelRole,
|
||||||
|
) -> Result<()> {
|
||||||
|
let mut buf = [0u8; 1024];
|
||||||
|
|
||||||
|
// 握手
|
||||||
|
let len = tcp.read(&mut buf).await?;
|
||||||
|
match postcard::from_bytes::<ControlPacket>(&buf[..len])? {
|
||||||
|
ControlPacket::ClientIdentify { role: _r } => {}
|
||||||
|
_ => return Err(anyhow::anyhow!("无效的握手协议")),
|
||||||
|
};
|
||||||
|
|
||||||
|
let hello = ControlPacket::ServerHello {
|
||||||
|
udp_port: udp.local_addr()?.port(),
|
||||||
|
};
|
||||||
|
tcp.write_all(&postcard::to_allocvec(&hello)?).await?;
|
||||||
|
|
||||||
|
// 等待 UDP 打洞/确认
|
||||||
|
let mut buf = [0u8; 128];
|
||||||
|
let (_, client_udp_addr) = udp.recv_from(&mut buf).await?;
|
||||||
|
println!("从节点 UDP 地址已确认: {}", client_udp_addr);
|
||||||
|
|
||||||
|
// 配置音频
|
||||||
|
let config = AudioConfig {
|
||||||
|
sample_rate: 48000,
|
||||||
|
channels: 2,
|
||||||
|
frame_size: 960,
|
||||||
|
bitrate: 64000,
|
||||||
|
..AudioConfig::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let mono_config = AudioConfig {
|
||||||
|
channels: 1,
|
||||||
|
..config.clone()
|
||||||
|
};
|
||||||
|
let mut codec = OpusCodec::new(&mono_config)?;
|
||||||
|
let playback_config = AudioConfig {
|
||||||
|
channels: 1,
|
||||||
|
playback_device: "plug:original_default".into(),
|
||||||
|
..config.clone()
|
||||||
|
};
|
||||||
|
let player = AudioPlayer::new(&playback_config)?;
|
||||||
|
|
||||||
|
let mut raw_buf = vec![0u8; config.frame_size * 2 * 2];
|
||||||
|
let mut opus_out = vec![0u8; 1500];
|
||||||
|
let mut seq = 0u32;
|
||||||
|
|
||||||
|
let frame_duration_us =
|
||||||
|
(config.frame_size as f64 / config.sample_rate as f64 * 1_000_000.0) as u128;
|
||||||
|
let delay_us = 200_000;
|
||||||
|
|
||||||
|
println!("开始会话循环...");
|
||||||
|
|
||||||
|
// 分离 TCP 读写,以便在不同任务中使用
|
||||||
|
let (mut tcp_rx, mut tcp_tx) = tcp.into_split();
|
||||||
|
let (stop_tx, mut stop_rx) = tokio::sync::mpsc::channel::<()>(1);
|
||||||
|
|
||||||
|
// 处理来自从节点的控制消息(如 Ping/Pong 进行时间同步)
|
||||||
|
tokio::spawn(async move {
|
||||||
|
let mut buf = [0u8; 1024];
|
||||||
|
loop {
|
||||||
|
match tcp_rx.read(&mut buf).await {
|
||||||
|
Ok(0) | Err(_) => {
|
||||||
|
let _ = stop_tx.send(()).await;
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
Ok(n) => {
|
||||||
|
if let Ok(ControlPacket::Ping { client_ts, seq }) =
|
||||||
|
postcard::from_bytes(&buf[..n])
|
||||||
|
{
|
||||||
|
let pong = ControlPacket::Pong {
|
||||||
|
client_ts,
|
||||||
|
server_ts: now_us(),
|
||||||
|
seq,
|
||||||
|
};
|
||||||
|
if tcp_tx
|
||||||
|
.write_all(&postcard::to_allocvec(&pong).unwrap())
|
||||||
|
.await
|
||||||
|
.is_err()
|
||||||
|
{
|
||||||
|
let _ = stop_tx.send(()).await;
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
loop {
|
||||||
|
// 打开 FIFO。这会阻塞直到有写入者打开它。
|
||||||
|
// 使用 select 以便在等待 FIFO 时如果从节点断开连接可以退出。
|
||||||
|
let mut fifo = tokio::select! {
|
||||||
|
_ = stop_rx.recv() => {
|
||||||
|
println!("等待 FIFO 时从节点断开连接。");
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
f = tokio::fs::File::open(AlsaRedirector::fifo_path()) => f?,
|
||||||
|
};
|
||||||
|
|
||||||
|
println!("音频流已启动...");
|
||||||
|
let mut stream_start_ts = 0;
|
||||||
|
let stream_start_seq = seq;
|
||||||
|
|
||||||
|
loop {
|
||||||
|
// 从 FIFO 读取,带超时/select 以检查断开连接
|
||||||
|
let read_res = tokio::select! {
|
||||||
|
_ = stop_rx.recv() => {
|
||||||
|
println!("串流过程中从节点断开连接。");
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
res = fifo.read_exact(&mut raw_buf) => res,
|
||||||
|
};
|
||||||
|
|
||||||
|
if let Err(_) = read_res {
|
||||||
|
println!("音频流结束 (FIFO 读取完毕)。");
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
|
let now = now_us();
|
||||||
|
if stream_start_ts == 0 {
|
||||||
|
stream_start_ts = now;
|
||||||
|
}
|
||||||
|
|
||||||
|
let mut local_pcm = Vec::with_capacity(config.frame_size);
|
||||||
|
let mut remote_pcm = Vec::with_capacity(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 r = i16::from_le_bytes([raw_buf[i * 4 + 2], raw_buf[i * 4 + 3]]);
|
||||||
|
if role == ChannelRole::Left {
|
||||||
|
local_pcm.push(l);
|
||||||
|
remote_pcm.push(r);
|
||||||
|
} else {
|
||||||
|
local_pcm.push(r);
|
||||||
|
remote_pcm.push(l);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 编码并发送给从节点
|
||||||
|
let len = 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)?;
|
||||||
|
if let Err(e) = udp.send_to(&bytes, client_udp_addr).await {
|
||||||
|
eprintln!("UDP 发送错误: {:?}", e);
|
||||||
|
return Err(e.into());
|
||||||
|
}
|
||||||
|
|
||||||
|
// 本地回放同步
|
||||||
|
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;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
player.write(&local_pcm)?;
|
||||||
|
|
||||||
|
seq += 1;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@@ -1,3 +1,6 @@
|
|||||||
|
pub mod alsa;
|
||||||
pub mod jitter_buffer;
|
pub mod jitter_buffer;
|
||||||
|
pub mod master;
|
||||||
pub mod protocol;
|
pub mod protocol;
|
||||||
|
pub mod slave;
|
||||||
pub mod sync;
|
pub mod sync;
|
||||||
|
|||||||
@@ -8,15 +8,15 @@ pub enum ChannelRole {
|
|||||||
|
|
||||||
#[derive(Serialize, Deserialize, Debug, Clone)]
|
#[derive(Serialize, Deserialize, Debug, Clone)]
|
||||||
pub enum ControlPacket {
|
pub enum ControlPacket {
|
||||||
// Discovery
|
// 发现协议
|
||||||
ServerHello {
|
ServerHello {
|
||||||
udp_port: u16, // Port for UDP audio stream
|
udp_port: u16, // UDP 音频流端口
|
||||||
},
|
},
|
||||||
// Handshake
|
// 握手协议
|
||||||
ClientIdentify {
|
ClientIdentify {
|
||||||
role: ChannelRole,
|
role: ChannelRole,
|
||||||
},
|
},
|
||||||
// Time Sync (Continuous)
|
// 时间同步 (持续进行)
|
||||||
Ping {
|
Ping {
|
||||||
client_ts: u128,
|
client_ts: u128,
|
||||||
seq: u32,
|
seq: u32,
|
||||||
@@ -26,14 +26,13 @@ pub enum ControlPacket {
|
|||||||
server_ts: u128,
|
server_ts: u128,
|
||||||
seq: u32,
|
seq: u32,
|
||||||
},
|
},
|
||||||
// Control
|
// 控制协议
|
||||||
Volume(u8), // 0-100
|
Volume(u8), // 音量 0-100
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Serialize, Deserialize, Debug, Clone)]
|
#[derive(Serialize, Deserialize, Debug, Clone)]
|
||||||
pub struct AudioPacket {
|
pub struct AudioPacket {
|
||||||
pub seq: u32, // Sequence number for packet loss detection
|
pub seq: u32, // 序列号,用于丢包检测
|
||||||
pub timestamp: u128, // Target playback time (server time)
|
pub timestamp: u128, // 目标播放时间 (主节点时间)
|
||||||
pub data: Vec<u8>, // Opus encoded data
|
pub data: Vec<u8>, // Opus 编码数据
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,168 @@
|
|||||||
|
#![cfg(target_os = "linux")]
|
||||||
|
use anyhow::Result;
|
||||||
|
use crate::audio::{AudioPlayer, OpusCodec};
|
||||||
|
use crate::config::AudioConfig;
|
||||||
|
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!("正在扫描主节点...");
|
||||||
|
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::<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);
|
||||||
|
println!("在 {} 发现主节点", master_tcp_addr);
|
||||||
|
|
||||||
|
let mut tcp = TcpStream::connect(master_tcp_addr).await?;
|
||||||
|
|
||||||
|
// 身份认证
|
||||||
|
tcp.write_all(&postcard::to_allocvec(&ControlPacket::ClientIdentify {
|
||||||
|
role: role.clone(),
|
||||||
|
})?)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
let len = tcp.read(&mut buf).await?;
|
||||||
|
let server_udp_port = match postcard::from_bytes::<ControlPacket>(&buf[..len])? {
|
||||||
|
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))
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
// 音频配置
|
||||||
|
let config = AudioConfig {
|
||||||
|
sample_rate: 48000,
|
||||||
|
channels: 1,
|
||||||
|
frame_size: 960,
|
||||||
|
bitrate: 64000,
|
||||||
|
..AudioConfig::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let player = AudioPlayer::new(&config)?;
|
||||||
|
let mut codec = OpusCodec::new(&config)?;
|
||||||
|
let mut jitter = JitterBuffer::new(50_000);
|
||||||
|
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();
|
||||||
|
|
||||||
|
// 定时发送 Ping 进行时间同步
|
||||||
|
let _sync_handle = tokio::spawn(async move {
|
||||||
|
let mut seq = 0;
|
||||||
|
loop {
|
||||||
|
let t1 = now_us();
|
||||||
|
let msg = ControlPacket::Ping { client_ts: t1, seq };
|
||||||
|
if tcp_tx
|
||||||
|
.write_all(&postcard::to_allocvec(&msg).unwrap())
|
||||||
|
.await
|
||||||
|
.is_err()
|
||||||
|
{
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
tokio::time::sleep(Duration::from_secs(1)).await;
|
||||||
|
seq += 1;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
let clock = Arc::new(Mutex::new(clock));
|
||||||
|
let clock_updater = clock.clone();
|
||||||
|
|
||||||
|
// 接收 Pong 并更新时钟偏移
|
||||||
|
tokio::spawn(async move {
|
||||||
|
let mut buf = [0u8; 1024];
|
||||||
|
loop {
|
||||||
|
match tcp_rx.read(&mut buf).await {
|
||||||
|
Ok(n) if n > 0 => {
|
||||||
|
if let Ok(ControlPacket::Pong {
|
||||||
|
client_ts,
|
||||||
|
server_ts,
|
||||||
|
..
|
||||||
|
}) = postcard::from_bytes(&buf[..n])
|
||||||
|
{
|
||||||
|
let t4 = now_us();
|
||||||
|
clock_updater.lock().await.update(client_ts, server_ts, t4);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
_ => 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];
|
||||||
|
println!("正在监听音频数据...");
|
||||||
|
|
||||||
|
let mut last_seq: Option<u32> = None;
|
||||||
|
|
||||||
|
loop {
|
||||||
|
// 将收到的数据包放入抖动缓冲区
|
||||||
|
while let Ok(pkt) = audio_rx.try_recv() {
|
||||||
|
jitter.push(pkt);
|
||||||
|
}
|
||||||
|
|
||||||
|
let now = now_us();
|
||||||
|
let current_server_time = {
|
||||||
|
let c = clock.lock().await;
|
||||||
|
c.to_server_time(now)
|
||||||
|
};
|
||||||
|
|
||||||
|
// 从抖动缓冲区提取待播放帧
|
||||||
|
if let Some((seq, data)) = jitter.pop_frame(current_server_time) {
|
||||||
|
// 丢包处理 (PLC)
|
||||||
|
if let Some(last) = last_seq {
|
||||||
|
if seq > last + 1 {
|
||||||
|
println!("检测到丢包: {} -> {}", last, seq);
|
||||||
|
// 补偿丢失的帧
|
||||||
|
for _ in 0..(seq - last - 1) {
|
||||||
|
if let Ok(len) = codec.decode_loss(&mut pcm_buf) {
|
||||||
|
let _ = player.write(&pcm_buf[..len]);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
last_seq = Some(seq);
|
||||||
|
|
||||||
|
let len = codec.decode(&data, &mut pcm_buf)?;
|
||||||
|
player.write(&pcm_buf[..len])?;
|
||||||
|
} else {
|
||||||
|
// 稍作等待以减少 CPU 占用
|
||||||
|
tokio::time::sleep(Duration::from_millis(1)).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@@ -1,13 +1,15 @@
|
|||||||
use std::collections::VecDeque;
|
use std::collections::VecDeque;
|
||||||
use std::time::{SystemTime, UNIX_EPOCH};
|
use std::time::{SystemTime, UNIX_EPOCH};
|
||||||
|
|
||||||
|
/// 获取当前微秒级时间戳
|
||||||
pub fn now_us() -> u128 {
|
pub fn now_us() -> u128 {
|
||||||
SystemTime::now()
|
SystemTime::now()
|
||||||
.duration_since(UNIX_EPOCH)
|
.duration_since(UNIX_EPOCH)
|
||||||
.expect("Time went backwards")
|
.expect("时间倒流")
|
||||||
.as_micros()
|
.as_micros()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// 时钟同步管理器,用于计算主从节点间的时钟偏移
|
||||||
pub struct ClockSync {
|
pub struct ClockSync {
|
||||||
offsets: VecDeque<i128>,
|
offsets: VecDeque<i128>,
|
||||||
pub current_offset: i128,
|
pub current_offset: i128,
|
||||||
@@ -23,15 +25,16 @@ impl ClockSync {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// 更新时钟偏移估计
|
||||||
pub fn update(&mut self, client_send_ts: u128, server_ts: u128, client_recv_ts: u128) {
|
pub fn update(&mut self, client_send_ts: u128, server_ts: u128, client_recv_ts: u128) {
|
||||||
let rtt = (client_recv_ts - client_send_ts) as i128;
|
let rtt = (client_recv_ts - client_send_ts) as i128;
|
||||||
// Basic filter: ignore if RTT is absurdly large (e.g. > 100ms on LAN)
|
// 基础过滤:如果 RTT 过大则忽略 (例如局域网内 > 100ms)
|
||||||
if rtt > 100_000 {
|
if rtt > 100_000 {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
// Clock offset = server_time - client_time
|
// 时钟偏移 = 主节点时间 - 从节点时间
|
||||||
// server_ts is at time (send + recv)/2
|
// 假设主节点收到 Ping 的时间点在 (发送时间 + 接收时间) / 2
|
||||||
let estimated_server_time = server_ts as i128 + rtt / 2;
|
let estimated_server_time = server_ts as i128 + rtt / 2;
|
||||||
let offset = estimated_server_time - client_recv_ts as i128;
|
let offset = estimated_server_time - client_recv_ts as i128;
|
||||||
|
|
||||||
@@ -40,8 +43,7 @@ impl ClockSync {
|
|||||||
self.offsets.pop_front();
|
self.offsets.pop_front();
|
||||||
}
|
}
|
||||||
|
|
||||||
// Calculate median or average offset
|
// 计算中位数偏移,以增强抗干扰能力
|
||||||
// Median is more robust to outliers
|
|
||||||
let mut sorted: Vec<i128> = self.offsets.iter().cloned().collect();
|
let mut sorted: Vec<i128> = self.offsets.iter().cloned().collect();
|
||||||
sorted.sort_unstable();
|
sorted.sort_unstable();
|
||||||
if !sorted.is_empty() {
|
if !sorted.is_empty() {
|
||||||
@@ -49,12 +51,13 @@ impl ClockSync {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// 将本地时间转换为服务器(主节点)时间
|
||||||
pub fn to_server_time(&self, client_time: u128) -> u128 {
|
pub fn to_server_time(&self, client_time: u128) -> u128 {
|
||||||
(client_time as i128 + self.current_offset) as u128
|
(client_time as i128 + self.current_offset) as u128
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// 将服务器(主节点)时间转换为本地时间
|
||||||
pub fn to_client_time(&self, server_time: u128) -> u128 {
|
pub fn to_client_time(&self, server_time: u128) -> u128 {
|
||||||
(server_time as i128 - self.current_offset) as u128
|
(server_time as i128 - self.current_offset) as u128
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user