refactor: 使用 UDP 传输音频流
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
use opus::{Encoder, Decoder, Application, Bitrate, Channels};
|
||||
use anyhow::{Result, Context};
|
||||
use crate::config::AudioConfig;
|
||||
use anyhow::{Context, Result};
|
||||
use opus::{Application, Bitrate, Channels, Decoder, Encoder};
|
||||
|
||||
pub struct OpusCodec {
|
||||
encoder: Encoder,
|
||||
@@ -9,23 +9,37 @@ pub struct OpusCodec {
|
||||
|
||||
impl OpusCodec {
|
||||
pub fn new(config: &AudioConfig) -> Result<Self> {
|
||||
let channels = if config.channels == 1 { Channels::Mono } else { Channels::Stereo };
|
||||
let channels = if config.channels == 1 {
|
||||
Channels::Mono
|
||||
} else {
|
||||
Channels::Stereo
|
||||
};
|
||||
let mut encoder = Encoder::new(config.sample_rate, channels, Application::Audio)
|
||||
.context("Failed to create Opus encoder")?;
|
||||
encoder.set_bitrate(Bitrate::Bits(config.bitrate))?;
|
||||
|
||||
let decoder = Decoder::new(config.sample_rate, channels)
|
||||
.context("Failed to create Opus decoder")?;
|
||||
let decoder =
|
||||
Decoder::new(config.sample_rate, channels).context("Failed to create Opus decoder")?;
|
||||
|
||||
Ok(Self { encoder, decoder })
|
||||
}
|
||||
|
||||
pub fn encode(&mut self, pcm: &[i16], out: &mut [u8]) -> Result<usize> {
|
||||
self.encoder.encode(pcm, out).context("Opus encoding failed")
|
||||
self.encoder
|
||||
.encode(pcm, out)
|
||||
.context("Opus encoding failed")
|
||||
}
|
||||
|
||||
pub fn decode(&mut self, opus: &[u8], out: &mut [i16]) -> Result<usize> {
|
||||
self.decoder.decode(opus, out, false).context("Opus decoding failed")
|
||||
self.decoder
|
||||
.decode(opus, out, false)
|
||||
.context("Opus decoding failed")
|
||||
}
|
||||
|
||||
// Packet Loss Concealment
|
||||
pub fn decode_loss(&mut self, out: &mut [i16]) -> Result<usize> {
|
||||
// Fallback to silence as opus crate 0.3 doesn't expose safe PLC
|
||||
out.fill(0);
|
||||
Ok(out.len())
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+403
-395
@@ -3,88 +3,101 @@
|
||||
use anyhow::{Context, Result};
|
||||
use hello::audio::{AudioPlayer, OpusCodec};
|
||||
use hello::config::AudioConfig;
|
||||
use hello::net::{AudioPacket, ChannelRole, ControlPacket, DISCOVERY_PORT, SERVER_PORT};
|
||||
use hello::sync::{ClockSync, now_us};
|
||||
use std::collections::VecDeque;
|
||||
use hello::stereo_core::jitter_buffer::JitterBuffer;
|
||||
use hello::stereo_core::protocol::{AudioPacket, ChannelRole, ControlPacket};
|
||||
use hello::stereo_core::sync::{ClockSync, now_us};
|
||||
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::signal;
|
||||
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";
|
||||
|
||||
fn setup_alsa_config() -> Result<()> {
|
||||
cleanup_alsa_config();
|
||||
// Default ports
|
||||
const DISCOVERY_PORT: u16 = 53530;
|
||||
const SERVER_TCP_PORT: u16 = 53531;
|
||||
// UDP port will be dynamic or fixed
|
||||
|
||||
let original_conf = fs::read_to_string(REAL_ASOUND_CONF)?;
|
||||
struct AlsaRedirector;
|
||||
|
||||
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"
|
||||
}}
|
||||
impl AlsaRedirector {
|
||||
fn new() -> Result<Self> {
|
||||
Self::cleanup(); // Ensure clean state
|
||||
|
||||
pcm.stereo_interceptor {{
|
||||
type file
|
||||
slave.pcm "null"
|
||||
file "{}"
|
||||
format "raw"
|
||||
}}
|
||||
"#,
|
||||
FIFO_PATH
|
||||
));
|
||||
let original_conf = fs::read_to_string(REAL_ASOUND_CONF).unwrap_or_default();
|
||||
|
||||
fs::write(TEMP_ASOUND_CONF, new_conf)?;
|
||||
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
|
||||
));
|
||||
|
||||
// 挂载覆盖 /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")?;
|
||||
fs::write(TEMP_ASOUND_CONF, new_conf)?;
|
||||
|
||||
if !status.success() {
|
||||
return Err(anyhow::anyhow!("Failed to mount asound.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)
|
||||
}
|
||||
|
||||
// 创建 FIFO
|
||||
let _ = Command::new("mkfifo").arg(FIFO_PATH).status();
|
||||
let _ = Command::new("chmod").arg("666").arg(FIFO_PATH).status();
|
||||
|
||||
println!("Successfully redirected ALSA output to {}", FIFO_PATH);
|
||||
Ok(())
|
||||
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.");
|
||||
}
|
||||
}
|
||||
|
||||
fn cleanup_alsa_config() {
|
||||
println!("Cleaning up ALSA configurations...");
|
||||
let _ = Command::new("umount")
|
||||
.arg("-l")
|
||||
.arg(REAL_ASOUND_CONF)
|
||||
.status();
|
||||
let _ = fs::remove_file(TEMP_ASOUND_CONF);
|
||||
let _ = fs::remove_file(FIFO_PATH);
|
||||
impl Drop for AlsaRedirector {
|
||||
fn drop(&mut self) {
|
||||
Self::cleanup();
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> Result<()> {
|
||||
let args: Vec<String> = env::args().collect();
|
||||
if args.len() < 3 {
|
||||
eprintln!("用法: {} [master|slave] [left|right]", args[0]);
|
||||
eprintln!("示例:");
|
||||
eprintln!(" 主设备: {} master left", args[0]);
|
||||
eprintln!(" 从设备: {} slave right", args[0]);
|
||||
eprintln!("Usage: {} [master|slave] [left|right]", args[0]);
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
@@ -95,379 +108,374 @@ async fn main() -> Result<()> {
|
||||
ChannelRole::Right
|
||||
};
|
||||
|
||||
let result = tokio::select! {
|
||||
res = async {
|
||||
if mode == "master" {
|
||||
run_master(role).await
|
||||
} else {
|
||||
run_slave(role).await
|
||||
}
|
||||
} => res,
|
||||
_ = signal::ctrl_c() => {
|
||||
println!("\n收到退出信号");
|
||||
Ok(())
|
||||
}
|
||||
};
|
||||
// Explicit cleanup at start just in case
|
||||
AlsaRedirector::cleanup();
|
||||
|
||||
cleanup_alsa_config();
|
||||
return result;
|
||||
}
|
||||
|
||||
async fn run_master(role: ChannelRole) -> Result<()> {
|
||||
let config = AudioConfig {
|
||||
sample_rate: 48000,
|
||||
channels: 2, // 拦截的是立体声
|
||||
frame_size: 960,
|
||||
bitrate: 48000,
|
||||
..AudioConfig::default()
|
||||
};
|
||||
|
||||
println!("--- 主设备模式 ---");
|
||||
println!("本地声道: {:?}", role);
|
||||
|
||||
loop {
|
||||
println!("正在启动发现服务 (UDP Broadcast)...");
|
||||
// 1. 启动发现广播
|
||||
let socket = UdpSocket::bind("0.0.0.0:0").await?;
|
||||
socket.set_broadcast(true)?;
|
||||
let target_addr: SocketAddr = format!("255.255.255.255:{}", DISCOVERY_PORT).parse()?;
|
||||
let hello = ControlPacket::ServerHello { port: SERVER_PORT };
|
||||
let msg = postcard::to_allocvec(&hello)?;
|
||||
|
||||
let broadcast_handle = tokio::spawn(async move {
|
||||
loop {
|
||||
let _ = socket.send_to(&msg, target_addr).await;
|
||||
tokio::time::sleep(Duration::from_secs(1)).await;
|
||||
}
|
||||
});
|
||||
|
||||
// 2. 等待从设备连接
|
||||
let listener = TcpListener::bind(format!("0.0.0.0:{}", SERVER_PORT)).await?;
|
||||
println!("等待从设备连接于端口 {}...", SERVER_PORT);
|
||||
|
||||
let (mut slave_stream, addr) = tokio::select! {
|
||||
res = listener.accept() => res?,
|
||||
_ = signal::ctrl_c() => {
|
||||
broadcast_handle.abort();
|
||||
return Ok(());
|
||||
}
|
||||
};
|
||||
|
||||
println!("从设备已连接: {}", addr);
|
||||
broadcast_handle.abort();
|
||||
|
||||
// 3. 握手与时间同步
|
||||
let mut buf = [0u8; 1024];
|
||||
let n = slave_stream.read(&mut buf).await?;
|
||||
if let Ok(ControlPacket::ClientIdentify { role: slave_role }) =
|
||||
postcard::from_bytes::<ControlPacket>(&buf[..n])
|
||||
{
|
||||
println!("从设备识别为: {:?}", slave_role);
|
||||
}
|
||||
|
||||
let n = slave_stream.read(&mut buf).await?;
|
||||
if let Ok(ControlPacket::Ping { client_ts }) =
|
||||
postcard::from_bytes::<ControlPacket>(&buf[..n])
|
||||
{
|
||||
let pong = ControlPacket::Pong {
|
||||
client_ts,
|
||||
server_ts: now_us(),
|
||||
};
|
||||
slave_stream
|
||||
.write_all(&postcard::to_allocvec(&pong)?)
|
||||
.await?;
|
||||
println!("时钟同步完成");
|
||||
}
|
||||
|
||||
// 4. 初始化音频并开始拦截
|
||||
setup_alsa_config()?;
|
||||
|
||||
let res = run_master_audio_loop(&mut slave_stream, role.clone(), &config).await;
|
||||
|
||||
cleanup_alsa_config();
|
||||
|
||||
if let Err(e) = res {
|
||||
eprintln!("主设备音频循环出错: {:?}, 准备重连...", e);
|
||||
} else {
|
||||
println!("从设备正常断开,准备下一次连接...");
|
||||
}
|
||||
|
||||
tokio::time::sleep(Duration::from_secs(1)).await;
|
||||
if mode == "master" {
|
||||
run_master(role).await
|
||||
} else {
|
||||
run_slave(role).await
|
||||
}
|
||||
}
|
||||
|
||||
async fn run_master_audio_loop(
|
||||
slave_stream: &mut TcpStream,
|
||||
role: ChannelRole,
|
||||
config: &AudioConfig,
|
||||
) -> Result<()> {
|
||||
// 打开 FIFO (使用 tokio::fs 以支持异步读取)
|
||||
let mut fifo = tokio::fs::File::open(FIFO_PATH)
|
||||
.await
|
||||
.context("Failed to open FIFO for reading")?;
|
||||
async fn run_master(role: ChannelRole) -> Result<()> {
|
||||
println!("--- Master Mode ({:?}) ---", role);
|
||||
|
||||
// 使用 plug:original_default 以支持系统主音量控制
|
||||
let playback_config = AudioConfig {
|
||||
channels: 1,
|
||||
playback_device: "plug:original_default".to_string(),
|
||||
..config.clone()
|
||||
// 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 player = AudioPlayer::new(&playback_config).context("无法打开本地播放设备")?;
|
||||
|
||||
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).context("无法初始化编码器")?;
|
||||
|
||||
let mut stereo_raw_buf = vec![0i16; config.frame_size * 2];
|
||||
let mut opus_buf = vec![0u8; 2048];
|
||||
|
||||
let delay_us = 100_000; // 100ms 缓冲
|
||||
let mut current_ts = 0;
|
||||
let frame_duration =
|
||||
(config.frame_size as f64 / config.sample_rate as f64 * 1_000_000.0) as u128;
|
||||
|
||||
println!("开始拦截并分发立体声音频...");
|
||||
|
||||
let mut byte_buf = vec![0u8; config.frame_size * 2 * 2];
|
||||
|
||||
loop {
|
||||
// 从 FIFO 读取 (这是同步源)
|
||||
if let Err(e) = fifo.read_exact(&mut byte_buf).await {
|
||||
if e.kind() == std::io::ErrorKind::UnexpectedEof {
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
continue;
|
||||
}
|
||||
return Err(e.into());
|
||||
}
|
||||
|
||||
let now = now_us();
|
||||
if current_ts == 0 {
|
||||
current_ts = now + delay_us;
|
||||
}
|
||||
|
||||
for i in 0..stereo_raw_buf.len() {
|
||||
stereo_raw_buf[i] = i16::from_le_bytes([byte_buf[i * 2], byte_buf[i * 2 + 1]]);
|
||||
}
|
||||
|
||||
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, r) = (stereo_raw_buf[i * 2], stereo_raw_buf[i * 2 + 1]);
|
||||
if role == ChannelRole::Left {
|
||||
local_pcm.push(l);
|
||||
remote_pcm.push(r);
|
||||
} else {
|
||||
local_pcm.push(r);
|
||||
remote_pcm.push(l);
|
||||
}
|
||||
}
|
||||
|
||||
// 先发送网络包
|
||||
let opus_len = codec.encode(&remote_pcm, &mut opus_buf)?;
|
||||
let packet = AudioPacket {
|
||||
timestamp: current_ts,
|
||||
data: opus_buf[..opus_len].to_vec(),
|
||||
};
|
||||
let packet_data = postcard::to_allocvec(&packet)?;
|
||||
|
||||
if slave_stream
|
||||
.write_all(&(packet_data.len() as u32).to_be_bytes())
|
||||
.await
|
||||
.is_err()
|
||||
|| slave_stream.write_all(&packet_data).await.is_err()
|
||||
{
|
||||
return Err(anyhow::anyhow!("从设备断开连接"));
|
||||
}
|
||||
|
||||
// 本地播放也需要同步
|
||||
let now = now_us();
|
||||
if now < current_ts {
|
||||
let wait = current_ts - now;
|
||||
if wait > 1000 {
|
||||
tokio::time::sleep(Duration::from_micros(wait as u64)).await;
|
||||
}
|
||||
}
|
||||
player.write(&local_pcm)?;
|
||||
|
||||
current_ts += frame_duration;
|
||||
|
||||
// 漂移校正
|
||||
let now = now_us();
|
||||
if now > current_ts + 200_000 {
|
||||
current_ts = now + 50_000;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn run_slave(role: ChannelRole) -> Result<()> {
|
||||
let config = AudioConfig {
|
||||
sample_rate: 48000,
|
||||
let mut codec = OpusCodec::new(&mono_config)?;
|
||||
let playback_config = AudioConfig {
|
||||
channels: 1,
|
||||
frame_size: 960,
|
||||
bitrate: 32000,
|
||||
..AudioConfig::default()
|
||||
playback_device: "plug:original_default".into(),
|
||||
..config.clone()
|
||||
};
|
||||
let player = AudioPlayer::new(&playback_config)?;
|
||||
|
||||
println!("--- 从设备模式 ---");
|
||||
println!("本地声道: {:?}", role);
|
||||
let mut raw_buf = vec![0u8; config.frame_size * 2 * 2];
|
||||
let mut opus_out = vec![0u8; 1500];
|
||||
let mut seq = 0u32;
|
||||
|
||||
loop {
|
||||
// 1. 发现主设备
|
||||
println!("正在搜索主设备...");
|
||||
let master_addr = match tokio::select! {
|
||||
addr = hello::net::Discovery::client_discover_server() => addr,
|
||||
_ = signal::ctrl_c() => return Ok(()),
|
||||
} {
|
||||
Ok(addr) => addr,
|
||||
Err(e) => {
|
||||
eprintln!("搜索主设备失败: {:?}, 1秒后重试...", e);
|
||||
tokio::time::sleep(Duration::from_secs(1)).await;
|
||||
continue;
|
||||
}
|
||||
};
|
||||
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!("发现主设备: {}", master_addr);
|
||||
let mut stream = match TcpStream::connect(master_addr).await {
|
||||
Ok(s) => s,
|
||||
Err(e) => {
|
||||
eprintln!("连接主设备失败: {:?}, 1秒后重试...", e);
|
||||
tokio::time::sleep(Duration::from_secs(1)).await;
|
||||
continue;
|
||||
}
|
||||
};
|
||||
println!("Starting session loop...");
|
||||
|
||||
// 2. 身份识别
|
||||
if let Err(e) = stream
|
||||
.write_all(&postcard::to_allocvec(&ControlPacket::ClientIdentify {
|
||||
role: role.clone(),
|
||||
})?)
|
||||
.await
|
||||
{
|
||||
eprintln!("发送身份识别失败: {:?}, 准备重连...", e);
|
||||
continue;
|
||||
}
|
||||
|
||||
// 3. 时间同步
|
||||
let mut clock = ClockSync::new();
|
||||
let t1 = now_us();
|
||||
if let Err(e) = stream
|
||||
.write_all(&postcard::to_allocvec(&ControlPacket::Ping {
|
||||
client_ts: t1,
|
||||
})?)
|
||||
.await
|
||||
{
|
||||
eprintln!("发送 Ping 失败: {:?}, 准备重连...", e);
|
||||
continue;
|
||||
}
|
||||
// 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];
|
||||
let n = match stream.read(&mut buf).await {
|
||||
Ok(n) => n,
|
||||
Err(e) => {
|
||||
eprintln!("读取 Pong 失败: {:?}, 准备重连...", e);
|
||||
continue;
|
||||
}
|
||||
};
|
||||
if let Ok(ControlPacket::Pong {
|
||||
client_ts,
|
||||
server_ts,
|
||||
}) = postcard::from_bytes::<ControlPacket>(&buf[..n])
|
||||
{
|
||||
let t4 = now_us();
|
||||
clock.update(client_ts, server_ts, t4);
|
||||
println!("时钟同步完成. 偏移: {}us, RTT: {}us", clock.offset, t4 - t1);
|
||||
} else {
|
||||
eprintln!("收到无效的 Pong 响应, 准备重连...");
|
||||
continue;
|
||||
}
|
||||
|
||||
let res = run_slave_audio_loop(stream, &config, clock).await;
|
||||
|
||||
if let Err(e) = res {
|
||||
eprintln!("从设备音频循环出错: {:?}, 准备重连...", e);
|
||||
} else {
|
||||
println!("主设备已断开,准备重连...");
|
||||
}
|
||||
|
||||
tokio::time::sleep(Duration::from_secs(1)).await;
|
||||
}
|
||||
}
|
||||
|
||||
async fn run_slave_audio_loop(
|
||||
stream: TcpStream,
|
||||
config: &AudioConfig,
|
||||
clock: ClockSync,
|
||||
) -> Result<()> {
|
||||
let player = AudioPlayer::new(config).context("无法打开播放设备")?;
|
||||
let mut codec = OpusCodec::new(config).context("无法初始化解码器")?;
|
||||
let mut jitter_buffer: VecDeque<AudioPacket> = VecDeque::new();
|
||||
let mut pcm_buf = vec![0i16; config.frame_size];
|
||||
|
||||
let (tx, mut rx) = tokio::sync::mpsc::channel(100);
|
||||
|
||||
// 网络接收线程
|
||||
let mut stream_read = stream;
|
||||
let receive_handle = tokio::spawn(async move {
|
||||
loop {
|
||||
let mut s_buf = [0u8; 4];
|
||||
if stream_read.read_exact(&mut s_buf).await.is_err() {
|
||||
break;
|
||||
}
|
||||
let size = u32::from_be_bytes(s_buf) as usize;
|
||||
let mut data = vec![0u8; size];
|
||||
if stream_read.read_exact(&mut data).await.is_err() {
|
||||
break;
|
||||
}
|
||||
if let Ok(p) = postcard::from_bytes::<AudioPacket>(&data) {
|
||||
if tx.send(p).await.is_err() {
|
||||
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;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
println!("开始接收并播放音频...");
|
||||
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?,
|
||||
};
|
||||
|
||||
let res = loop {
|
||||
while let Ok(p) = rx.try_recv() {
|
||||
jitter_buffer.push_back(p);
|
||||
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);
|
||||
}
|
||||
|
||||
if let Some(p) = jitter_buffer.front() {
|
||||
let target_client_time = clock.to_client_time(p.timestamp);
|
||||
let now = now_us();
|
||||
let now = now_us();
|
||||
let current_server_time = {
|
||||
let c = clock.lock().await;
|
||||
c.to_server_time(now)
|
||||
};
|
||||
|
||||
if now >= target_client_time {
|
||||
let packet = jitter_buffer.pop_front().unwrap();
|
||||
|
||||
// 如果包太旧了(延迟超过150ms),跳过以赶上进度,防止累积卡顿
|
||||
if now > target_client_time + 150_000 {
|
||||
continue;
|
||||
}
|
||||
|
||||
let len = codec.decode(&packet.data, &mut pcm_buf)?;
|
||||
player.write(&pcm_buf[..len])?;
|
||||
} else {
|
||||
let wait = (target_client_time - now) as u64;
|
||||
if wait > 1000 {
|
||||
// 如果时间差太大,可能是时钟跳变,清空缓冲重新同步
|
||||
if wait > 1_000_000 {
|
||||
jitter_buffer.clear();
|
||||
} else {
|
||||
tokio::time::sleep(Duration::from_micros(wait)).await;
|
||||
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]);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
} else {
|
||||
if receive_handle.is_finished() {
|
||||
break Ok(());
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(5)).await;
|
||||
}
|
||||
};
|
||||
last_seq = Some(seq);
|
||||
|
||||
receive_handle.abort();
|
||||
res
|
||||
let len = codec.decode(&data, &mut pcm_buf)?;
|
||||
player.write(&pcm_buf[..len])?;
|
||||
} else {
|
||||
tokio::time::sleep(Duration::from_millis(1)).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
pub mod audio;
|
||||
pub mod config;
|
||||
pub mod net;
|
||||
pub mod stereo_core;
|
||||
pub mod sync;
|
||||
|
||||
@@ -0,0 +1,64 @@
|
||||
use super::protocol::AudioPacket;
|
||||
use std::collections::BinaryHeap;
|
||||
use std::cmp::Ordering;
|
||||
|
||||
#[derive(Debug)]
|
||||
struct OrderedPacket(AudioPacket);
|
||||
|
||||
impl PartialEq for OrderedPacket {
|
||||
fn eq(&self, other: &Self) -> bool {
|
||||
self.0.seq == other.0.seq
|
||||
}
|
||||
}
|
||||
impl Eq for OrderedPacket {}
|
||||
impl PartialOrd for OrderedPacket {
|
||||
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
|
||||
Some(other.0.seq.cmp(&self.0.seq))
|
||||
}
|
||||
}
|
||||
impl Ord for OrderedPacket {
|
||||
fn cmp(&self, other: &Self) -> Ordering {
|
||||
other.0.seq.cmp(&self.0.seq)
|
||||
}
|
||||
}
|
||||
|
||||
pub struct JitterBuffer {
|
||||
buffer: BinaryHeap<OrderedPacket>,
|
||||
last_played_seq: Option<u32>,
|
||||
pub target_delay_us: u128,
|
||||
}
|
||||
|
||||
impl JitterBuffer {
|
||||
pub fn new(target_delay_us: u128) -> Self {
|
||||
Self {
|
||||
buffer: BinaryHeap::new(),
|
||||
last_played_seq: None,
|
||||
target_delay_us,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn push(&mut self, packet: AudioPacket) {
|
||||
if let Some(last) = self.last_played_seq {
|
||||
if packet.seq <= last {
|
||||
return;
|
||||
}
|
||||
}
|
||||
self.buffer.push(OrderedPacket(packet));
|
||||
}
|
||||
|
||||
pub fn pop_frame(&mut self, current_time: u128) -> Option<(u32, Vec<u8>)> {
|
||||
if let Some(OrderedPacket(pkt)) = self.buffer.peek() {
|
||||
if current_time >= pkt.timestamp {
|
||||
let pkt = self.buffer.pop().unwrap().0;
|
||||
self.last_played_seq = Some(pkt.seq);
|
||||
return Some((pkt.seq, pkt.data));
|
||||
}
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
pub fn clear(&mut self) {
|
||||
self.buffer.clear();
|
||||
self.last_played_seq = None;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,3 @@
|
||||
pub mod jitter_buffer;
|
||||
pub mod protocol;
|
||||
pub mod sync;
|
||||
@@ -0,0 +1,39 @@
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq, Copy)]
|
||||
pub enum ChannelRole {
|
||||
Left,
|
||||
Right,
|
||||
}
|
||||
|
||||
#[derive(Serialize, Deserialize, Debug, Clone)]
|
||||
pub enum ControlPacket {
|
||||
// Discovery
|
||||
ServerHello {
|
||||
udp_port: u16, // Port for UDP audio stream
|
||||
},
|
||||
// Handshake
|
||||
ClientIdentify {
|
||||
role: ChannelRole,
|
||||
},
|
||||
// Time Sync (Continuous)
|
||||
Ping {
|
||||
client_ts: u128,
|
||||
seq: u32,
|
||||
},
|
||||
Pong {
|
||||
client_ts: u128,
|
||||
server_ts: u128,
|
||||
seq: u32,
|
||||
},
|
||||
// Control
|
||||
Volume(u8), // 0-100
|
||||
}
|
||||
|
||||
#[derive(Serialize, Deserialize, Debug, Clone)]
|
||||
pub struct AudioPacket {
|
||||
pub seq: u32, // Sequence number for packet loss detection
|
||||
pub timestamp: u128, // Target playback time (server time)
|
||||
pub data: Vec<u8>, // Opus encoded data
|
||||
}
|
||||
|
||||
@@ -0,0 +1,60 @@
|
||||
use std::collections::VecDeque;
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
pub fn now_us() -> u128 {
|
||||
SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.expect("Time went backwards")
|
||||
.as_micros()
|
||||
}
|
||||
|
||||
pub struct ClockSync {
|
||||
offsets: VecDeque<i128>,
|
||||
pub current_offset: i128,
|
||||
window_size: usize,
|
||||
}
|
||||
|
||||
impl ClockSync {
|
||||
pub fn new(window_size: usize) -> Self {
|
||||
Self {
|
||||
offsets: VecDeque::with_capacity(window_size),
|
||||
current_offset: 0,
|
||||
window_size,
|
||||
}
|
||||
}
|
||||
|
||||
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;
|
||||
// Basic filter: ignore if RTT is absurdly large (e.g. > 100ms on LAN)
|
||||
if rtt > 100_000 {
|
||||
return;
|
||||
}
|
||||
|
||||
// Clock offset = server_time - client_time
|
||||
// server_ts is at time (send + recv)/2
|
||||
let estimated_server_time = server_ts as i128 + rtt / 2;
|
||||
let offset = estimated_server_time - client_recv_ts as i128;
|
||||
|
||||
self.offsets.push_back(offset);
|
||||
if self.offsets.len() > self.window_size {
|
||||
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();
|
||||
sorted.sort_unstable();
|
||||
if !sorted.is_empty() {
|
||||
self.current_offset = sorted[sorted.len() / 2];
|
||||
}
|
||||
}
|
||||
|
||||
pub fn to_server_time(&self, client_time: u128) -> u128 {
|
||||
(client_time as i128 + self.current_offset) as u128
|
||||
}
|
||||
|
||||
pub fn to_client_time(&self, server_time: u128) -> u128 {
|
||||
(server_time as i128 - self.current_offset) as u128
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user