refactor: 分离音频管理模块
This commit is contained in:
@@ -0,0 +1,151 @@
|
|||||||
|
use crate::audio::codec::OpusCodec;
|
||||||
|
use crate::audio::config::AudioConfig;
|
||||||
|
use crate::audio::player::AudioPlayer;
|
||||||
|
use crate::audio::recorder::AudioRecorder;
|
||||||
|
use crate::net::network::AudioSocket;
|
||||||
|
use crate::net::protocol::AudioPacket;
|
||||||
|
use std::net::SocketAddr;
|
||||||
|
use std::sync::Arc;
|
||||||
|
use tokio::sync::{RwLock, mpsc};
|
||||||
|
use tokio_util::sync::CancellationToken;
|
||||||
|
|
||||||
|
pub struct ClientAudioManager {
|
||||||
|
audio_socket: Arc<AudioSocket>,
|
||||||
|
server_audio_addr: SocketAddr,
|
||||||
|
session_cancel: CancellationToken,
|
||||||
|
record_cancel: RwLock<Option<CancellationToken>>,
|
||||||
|
play_cancel: RwLock<Option<CancellationToken>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ClientAudioManager {
|
||||||
|
pub fn new(
|
||||||
|
audio_socket: Arc<AudioSocket>,
|
||||||
|
server_audio_addr: SocketAddr,
|
||||||
|
session_cancel: CancellationToken,
|
||||||
|
) -> Self {
|
||||||
|
Self {
|
||||||
|
audio_socket,
|
||||||
|
server_audio_addr,
|
||||||
|
session_cancel,
|
||||||
|
record_cancel: RwLock::new(None),
|
||||||
|
play_cancel: RwLock::new(None),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn stop_recorder(&self) {
|
||||||
|
let mut cancel_guard = self.record_cancel.write().await;
|
||||||
|
if let Some(token) = cancel_guard.take() {
|
||||||
|
token.cancel();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn stop_player(&self) {
|
||||||
|
let mut cancel_guard = self.play_cancel.write().await;
|
||||||
|
if let Some(token) = cancel_guard.take() {
|
||||||
|
token.cancel();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn start_recording(&self, config: AudioConfig) {
|
||||||
|
self.stop_recorder().await;
|
||||||
|
let token = self.session_cancel.child_token();
|
||||||
|
*self.record_cancel.write().await = Some(token.clone());
|
||||||
|
|
||||||
|
let audio_socket = self.audio_socket.clone();
|
||||||
|
let server_audio_addr = self.server_audio_addr;
|
||||||
|
|
||||||
|
tokio::spawn(async move {
|
||||||
|
let (pcm_tx, mut pcm_rx) = mpsc::channel::<Vec<i16>>(32);
|
||||||
|
let conf = config.clone();
|
||||||
|
|
||||||
|
// 录音线程 (ALSA 阻塞)
|
||||||
|
std::thread::spawn(move || {
|
||||||
|
let recorder = match AudioRecorder::new(&conf) {
|
||||||
|
Ok(r) => r,
|
||||||
|
Err(e) => {
|
||||||
|
eprintln!("Failed to start recorder: {}", e);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let mut buf = vec![0i16; conf.frame_size];
|
||||||
|
while let Ok(n) = recorder.read(&mut buf) {
|
||||||
|
if pcm_tx.blocking_send(buf[..n].to_vec()).is_err() {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
let mut codec = match OpusCodec::new(&config) {
|
||||||
|
Ok(c) => c,
|
||||||
|
Err(e) => {
|
||||||
|
eprintln!("Failed to init opus codec: {}", e);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
println!("Recording started...");
|
||||||
|
loop {
|
||||||
|
tokio::select! {
|
||||||
|
_ = token.cancelled() => break,
|
||||||
|
Some(pcm) = pcm_rx.recv() => {
|
||||||
|
let mut out = vec![0u8; 4096];
|
||||||
|
if let Ok(len) = codec.encode(&pcm, &mut out) {
|
||||||
|
let _ = audio_socket.send(&AudioPacket { data: out[..len].to_vec() }, server_audio_addr).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
println!("Recording stopped.");
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn start_playback(&self, config: AudioConfig) {
|
||||||
|
self.stop_player().await;
|
||||||
|
let token = self.session_cancel.child_token();
|
||||||
|
*self.play_cancel.write().await = Some(token.clone());
|
||||||
|
|
||||||
|
let audio_socket = self.audio_socket.clone();
|
||||||
|
|
||||||
|
tokio::spawn(async move {
|
||||||
|
let (pcm_tx, mut pcm_rx) = mpsc::channel::<Vec<i16>>(32);
|
||||||
|
let conf = config.clone();
|
||||||
|
|
||||||
|
// 播放线程 (ALSA 阻塞)
|
||||||
|
std::thread::spawn(move || {
|
||||||
|
let player = match AudioPlayer::new(&conf) {
|
||||||
|
Ok(p) => p,
|
||||||
|
Err(e) => {
|
||||||
|
eprintln!("Failed to start player: {}", e);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
while let Some(pcm) = pcm_rx.blocking_recv() {
|
||||||
|
let _ = player.write(&pcm);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
let mut codec = match OpusCodec::new(&config) {
|
||||||
|
Ok(c) => c,
|
||||||
|
Err(e) => {
|
||||||
|
eprintln!("Failed to init opus codec: {}", e);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let mut udp_buf = vec![0u8; 4096];
|
||||||
|
println!("Playback started...");
|
||||||
|
loop {
|
||||||
|
tokio::select! {
|
||||||
|
_ = token.cancelled() => break,
|
||||||
|
res = audio_socket.recv(&mut udp_buf) => {
|
||||||
|
if let Ok((packet, _)) = res {
|
||||||
|
let mut pcm = vec![0i16; config.frame_size];
|
||||||
|
if let Ok(n) = codec.decode(&packet.data, &mut pcm) {
|
||||||
|
let _ = pcm_tx.send(pcm[..n].to_vec()).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
println!("Playback stopped.");
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,14 +1,13 @@
|
|||||||
#![cfg(target_os = "linux")]
|
#![cfg(target_os = "linux")]
|
||||||
|
|
||||||
use crate::audio::codec::OpusCodec;
|
mod audio_manager;
|
||||||
use crate::audio::config::AudioConfig;
|
|
||||||
use crate::audio::player::AudioPlayer;
|
|
||||||
use crate::audio::recorder::AudioRecorder;
|
|
||||||
use crate::net::discovery::Discovery;
|
use crate::net::discovery::Discovery;
|
||||||
use crate::net::network::{AudioSocket, Connection};
|
use crate::net::network::{AudioSocket, Connection};
|
||||||
use crate::net::protocol::{AudioPacket, ClientInfo, ControlPacket, RpcResult};
|
use crate::net::protocol::{ClientInfo, ControlPacket, RpcResult};
|
||||||
use crate::net::rpc::RpcManager;
|
use crate::net::rpc::RpcManager;
|
||||||
use anyhow::{Result, anyhow};
|
use anyhow::{Result, anyhow};
|
||||||
|
use audio_manager::ClientAudioManager;
|
||||||
use std::net::SocketAddr;
|
use std::net::SocketAddr;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use tokio::sync::{RwLock, mpsc};
|
use tokio::sync::{RwLock, mpsc};
|
||||||
@@ -17,11 +16,8 @@ use tokio_util::sync::CancellationToken;
|
|||||||
/// 内部连接上下文,包含了音频流所需的全部信息
|
/// 内部连接上下文,包含了音频流所需的全部信息
|
||||||
struct ActiveSession {
|
struct ActiveSession {
|
||||||
conn: Arc<Connection>,
|
conn: Arc<Connection>,
|
||||||
audio_socket: Arc<AudioSocket>,
|
audio_manager: ClientAudioManager,
|
||||||
server_audio_addr: SocketAddr,
|
|
||||||
session_cancel: CancellationToken, // 控制整个 Session 的生命周期
|
session_cancel: CancellationToken, // 控制整个 Session 的生命周期
|
||||||
record_cancel: RwLock<Option<CancellationToken>>,
|
|
||||||
play_cancel: RwLock<Option<CancellationToken>>,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
pub struct Client {
|
pub struct Client {
|
||||||
@@ -126,13 +122,15 @@ impl Client {
|
|||||||
);
|
);
|
||||||
|
|
||||||
// --- 初始化 Session ---
|
// --- 初始化 Session ---
|
||||||
|
let session_cancel = CancellationToken::new();
|
||||||
let session = Arc::new(ActiveSession {
|
let session = Arc::new(ActiveSession {
|
||||||
conn: conn.clone(),
|
conn: conn.clone(),
|
||||||
audio_socket,
|
audio_manager: ClientAudioManager::new(
|
||||||
server_audio_addr,
|
audio_socket,
|
||||||
session_cancel: CancellationToken::new(),
|
server_audio_addr,
|
||||||
record_cancel: RwLock::new(None),
|
session_cancel.clone(),
|
||||||
play_cancel: RwLock::new(None),
|
),
|
||||||
|
session_cancel,
|
||||||
});
|
});
|
||||||
*self.session.write().await = Some(session.clone());
|
*self.session.write().await = Some(session.clone());
|
||||||
|
|
||||||
@@ -206,22 +204,16 @@ impl Client {
|
|||||||
.await?;
|
.await?;
|
||||||
}
|
}
|
||||||
ControlPacket::StartRecording { config } => {
|
ControlPacket::StartRecording { config } => {
|
||||||
self.stop_recorder(session).await; // 开启前先停止旧的,防止资源冲突
|
session.audio_manager.start_recording(config).await;
|
||||||
let token = session.session_cancel.child_token();
|
|
||||||
*session.record_cancel.write().await = Some(token.clone());
|
|
||||||
self.spawn_recorder(session.clone(), config, token);
|
|
||||||
}
|
}
|
||||||
ControlPacket::StartPlayback { config } => {
|
ControlPacket::StartPlayback { config } => {
|
||||||
self.stop_player(session).await;
|
session.audio_manager.start_playback(config).await;
|
||||||
let token = session.session_cancel.child_token();
|
|
||||||
*session.play_cancel.write().await = Some(token.clone());
|
|
||||||
self.spawn_player(session.clone(), config, token);
|
|
||||||
}
|
}
|
||||||
ControlPacket::StopRecording => {
|
ControlPacket::StopRecording => {
|
||||||
self.stop_recorder(session).await;
|
session.audio_manager.stop_recorder().await;
|
||||||
}
|
}
|
||||||
ControlPacket::StopPlayback => {
|
ControlPacket::StopPlayback => {
|
||||||
self.stop_player(session).await;
|
session.audio_manager.stop_player().await;
|
||||||
}
|
}
|
||||||
_ => {}
|
_ => {}
|
||||||
}
|
}
|
||||||
@@ -256,120 +248,6 @@ impl Client {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn stop_recorder(&self, session: &ActiveSession) {
|
|
||||||
let mut cancel_guard = session.record_cancel.write().await;
|
|
||||||
if let Some(token) = cancel_guard.take() {
|
|
||||||
token.cancel();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn stop_player(&self, session: &ActiveSession) {
|
|
||||||
let mut cancel_guard = session.play_cancel.write().await;
|
|
||||||
if let Some(token) = cancel_guard.take() {
|
|
||||||
token.cancel();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
fn spawn_recorder(
|
|
||||||
&self,
|
|
||||||
session: Arc<ActiveSession>,
|
|
||||||
config: AudioConfig,
|
|
||||||
token: CancellationToken,
|
|
||||||
) {
|
|
||||||
tokio::spawn(async move {
|
|
||||||
let (pcm_tx, mut pcm_rx) = mpsc::channel::<Vec<i16>>(32);
|
|
||||||
let conf = config.clone();
|
|
||||||
|
|
||||||
// 录音线程 (ALSA 阻塞)
|
|
||||||
std::thread::spawn(move || {
|
|
||||||
let recorder = match AudioRecorder::new(&conf) {
|
|
||||||
Ok(r) => r,
|
|
||||||
Err(e) => {
|
|
||||||
eprintln!("Failed to start recorder: {}", e);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
};
|
|
||||||
let mut buf = vec![0i16; conf.frame_size];
|
|
||||||
while let Ok(n) = recorder.read(&mut buf) {
|
|
||||||
if pcm_tx.blocking_send(buf[..n].to_vec()).is_err() {
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
let mut codec = match OpusCodec::new(&config) {
|
|
||||||
Ok(c) => c,
|
|
||||||
Err(e) => {
|
|
||||||
eprintln!("Failed to init opus codec: {}", e);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
};
|
|
||||||
println!("Recording started...");
|
|
||||||
loop {
|
|
||||||
tokio::select! {
|
|
||||||
_ = token.cancelled() => break,
|
|
||||||
Some(pcm) = pcm_rx.recv() => {
|
|
||||||
let mut out = vec![0u8; 4096];
|
|
||||||
if let Ok(len) = codec.encode(&pcm, &mut out) {
|
|
||||||
let _ = session.audio_socket.send(&AudioPacket { data: out[..len].to_vec() }, session.server_audio_addr).await;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
println!("Recording stopped.");
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
fn spawn_player(
|
|
||||||
&self,
|
|
||||||
session: Arc<ActiveSession>,
|
|
||||||
config: AudioConfig,
|
|
||||||
token: CancellationToken,
|
|
||||||
) {
|
|
||||||
tokio::spawn(async move {
|
|
||||||
let (pcm_tx, mut pcm_rx) = mpsc::channel::<Vec<i16>>(32);
|
|
||||||
let conf = config.clone();
|
|
||||||
|
|
||||||
// 播放线程 (ALSA 阻塞)
|
|
||||||
std::thread::spawn(move || {
|
|
||||||
let player = match AudioPlayer::new(&conf) {
|
|
||||||
Ok(p) => p,
|
|
||||||
Err(e) => {
|
|
||||||
eprintln!("Failed to start player: {}", e);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
};
|
|
||||||
while let Some(pcm) = pcm_rx.blocking_recv() {
|
|
||||||
let _ = player.write(&pcm);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
let mut codec = match OpusCodec::new(&config) {
|
|
||||||
Ok(c) => c,
|
|
||||||
Err(e) => {
|
|
||||||
eprintln!("Failed to init opus codec: {}", e);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
};
|
|
||||||
let mut udp_buf = vec![0u8; 4096];
|
|
||||||
println!("Playback started...");
|
|
||||||
loop {
|
|
||||||
tokio::select! {
|
|
||||||
_ = token.cancelled() => break,
|
|
||||||
res = session.audio_socket.recv(&mut udp_buf) => {
|
|
||||||
if let Ok((packet, _)) = res {
|
|
||||||
let mut pcm = vec![0i16; config.frame_size];
|
|
||||||
if let Ok(n) = codec.decode(&packet.data, &mut pcm) {
|
|
||||||
let _ = pcm_tx.send(pcm[..n].to_vec()).await;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
println!("Playback stopped.");
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
pub async fn call(&self, method: &str, args: Vec<String>) -> Result<RpcResult> {
|
pub async fn call(&self, method: &str, args: Vec<String>) -> Result<RpcResult> {
|
||||||
let (id, rx) = self.rpc.register();
|
let (id, rx) = self.rpc.register();
|
||||||
let session_guard = self.session.read().await;
|
let session_guard = self.session.read().await;
|
||||||
|
|||||||
@@ -0,0 +1,201 @@
|
|||||||
|
use crate::audio::codec::OpusCodec;
|
||||||
|
use crate::audio::config::AudioConfig;
|
||||||
|
use crate::audio::wav::{WavReader, WavWriter};
|
||||||
|
use crate::net::network::AudioSocket;
|
||||||
|
use crate::net::protocol::AudioPacket;
|
||||||
|
use anyhow::Result;
|
||||||
|
use parking_lot::Mutex;
|
||||||
|
use std::net::SocketAddr;
|
||||||
|
use std::sync::Arc;
|
||||||
|
use tokio::sync::mpsc;
|
||||||
|
use tokio_util::sync::CancellationToken;
|
||||||
|
use tokio_util::task::TaskTracker;
|
||||||
|
|
||||||
|
pub enum RecorderCommand {
|
||||||
|
Start {
|
||||||
|
config: AudioConfig,
|
||||||
|
filename: String,
|
||||||
|
},
|
||||||
|
Stop,
|
||||||
|
}
|
||||||
|
|
||||||
|
pub struct ServerAudioManager {
|
||||||
|
session_cancel: CancellationToken,
|
||||||
|
record_cancel: Mutex<Option<CancellationToken>>,
|
||||||
|
play_cancel: Mutex<Option<CancellationToken>>,
|
||||||
|
audio_tx: mpsc::Sender<AudioPacket>,
|
||||||
|
recorder_tx: mpsc::Sender<RecorderCommand>,
|
||||||
|
tracker: TaskTracker,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ServerAudioManager {
|
||||||
|
pub fn new(session_cancel: CancellationToken, tracker: TaskTracker) -> (Self, mpsc::Receiver<AudioPacket>, mpsc::Receiver<RecorderCommand>) {
|
||||||
|
let (audio_tx, audio_rx) = mpsc::channel(1024);
|
||||||
|
let (recorder_tx, recorder_rx) = mpsc::channel(64);
|
||||||
|
|
||||||
|
let manager = Self {
|
||||||
|
session_cancel,
|
||||||
|
record_cancel: Mutex::new(None),
|
||||||
|
play_cancel: Mutex::new(None),
|
||||||
|
audio_tx,
|
||||||
|
recorder_tx,
|
||||||
|
tracker,
|
||||||
|
};
|
||||||
|
|
||||||
|
(manager, audio_rx, recorder_rx)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn audio_tx(&self) -> mpsc::Sender<AudioPacket> {
|
||||||
|
self.audio_tx.clone()
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn start_recording(&self, config: AudioConfig, filename: String) -> Result<()> {
|
||||||
|
// 仅停止之前的录音任务
|
||||||
|
{
|
||||||
|
let mut guard = self.record_cancel.lock();
|
||||||
|
if let Some(token) = guard.take() {
|
||||||
|
token.cancel();
|
||||||
|
}
|
||||||
|
*guard = Some(self.session_cancel.child_token());
|
||||||
|
}
|
||||||
|
|
||||||
|
self.recorder_tx
|
||||||
|
.send(RecorderCommand::Start {
|
||||||
|
config,
|
||||||
|
filename,
|
||||||
|
})
|
||||||
|
.await?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn stop_recording(&self) -> Result<()> {
|
||||||
|
if let Some(token) = self.record_cancel.lock().take() {
|
||||||
|
token.cancel();
|
||||||
|
}
|
||||||
|
let _ = self.recorder_tx.send(RecorderCommand::Stop).await;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn start_playback(
|
||||||
|
&self,
|
||||||
|
config: AudioConfig,
|
||||||
|
reader: WavReader,
|
||||||
|
audio_socket: Arc<AudioSocket>,
|
||||||
|
target_addr: SocketAddr,
|
||||||
|
) -> Result<()> {
|
||||||
|
let token = {
|
||||||
|
let mut guard = self.play_cancel.lock();
|
||||||
|
if let Some(token) = guard.take() {
|
||||||
|
token.cancel();
|
||||||
|
}
|
||||||
|
let token = self.session_cancel.child_token();
|
||||||
|
*guard = Some(token.clone());
|
||||||
|
token
|
||||||
|
};
|
||||||
|
|
||||||
|
let session_cancel = self.session_cancel.clone();
|
||||||
|
|
||||||
|
self.tracker.spawn(async move {
|
||||||
|
let mut reader = reader;
|
||||||
|
let mut codec = match OpusCodec::new(&config) {
|
||||||
|
Ok(c) => c,
|
||||||
|
Err(e) => {
|
||||||
|
eprintln!("Failed to create opus codec: {}", e);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let mut pcm = vec![0i16; config.frame_size];
|
||||||
|
let mut opus = vec![0u8; 4096];
|
||||||
|
let mut interval = tokio::time::interval(std::time::Duration::from_millis(20));
|
||||||
|
|
||||||
|
loop {
|
||||||
|
tokio::select! {
|
||||||
|
_ = token.cancelled() => break,
|
||||||
|
_ = session_cancel.cancelled() => break,
|
||||||
|
_ = interval.tick() => {
|
||||||
|
match reader.read_samples(&mut pcm) {
|
||||||
|
Ok(0) => break,
|
||||||
|
Ok(n) => {
|
||||||
|
if let Ok(len) = codec.encode(&pcm[..n], &mut opus) {
|
||||||
|
let _ = audio_socket
|
||||||
|
.send(
|
||||||
|
&AudioPacket {
|
||||||
|
data: opus[..len].to_vec(),
|
||||||
|
},
|
||||||
|
target_addr,
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
eprintln!("Failed to read samples: {}", e);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn stop_playback(&self) {
|
||||||
|
if let Some(token) = self.play_cancel.lock().take() {
|
||||||
|
token.cancel();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn spawn_audio_processor(&self, mut audio_rx: mpsc::Receiver<AudioPacket>, mut recorder_rx: mpsc::Receiver<RecorderCommand>) {
|
||||||
|
let session_cancel = self.session_cancel.clone();
|
||||||
|
self.tracker.spawn(async move {
|
||||||
|
let mut active_recorder: Option<(WavWriter, OpusCodec, usize)> = None;
|
||||||
|
loop {
|
||||||
|
tokio::select! {
|
||||||
|
_ = session_cancel.cancelled() => break,
|
||||||
|
cmd = recorder_rx.recv() => {
|
||||||
|
match cmd {
|
||||||
|
Some(RecorderCommand::Start { config, filename }) => {
|
||||||
|
if let Some((writer, _, _)) = active_recorder.take() {
|
||||||
|
let _ = writer.finalize();
|
||||||
|
}
|
||||||
|
match WavWriter::create(&filename, config.sample_rate, config.channels) {
|
||||||
|
Ok(writer) => {
|
||||||
|
match OpusCodec::new(&config) {
|
||||||
|
Ok(codec) => active_recorder = Some((writer, codec, config.frame_size)),
|
||||||
|
Err(e) => eprintln!("Failed to create opus codec: {}", e),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Err(e) => eprintln!("Failed to create wav writer: {}", e),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Some(RecorderCommand::Stop) => {
|
||||||
|
if let Some((writer, _, _)) = active_recorder.take() {
|
||||||
|
let _ = writer.finalize();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
None => break,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
packet = audio_rx.recv() => {
|
||||||
|
match packet {
|
||||||
|
Some(packet) => {
|
||||||
|
if let Some((writer, codec, frame_size)) = &mut active_recorder {
|
||||||
|
let mut pcm = vec![0i16; *frame_size];
|
||||||
|
if let Ok(n) = codec.decode(&packet.data, &mut pcm) {
|
||||||
|
let _ = writer.write_samples(&pcm[..n]);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
None => break,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if let Some((writer, _, _)) = active_recorder.take() {
|
||||||
|
let _ = writer.finalize();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@@ -1,27 +1,19 @@
|
|||||||
use crate::audio::codec::OpusCodec;
|
mod audio_manager;
|
||||||
|
|
||||||
use crate::audio::config::AudioConfig;
|
use crate::audio::config::AudioConfig;
|
||||||
use crate::audio::wav::{WavReader, WavWriter};
|
use crate::audio::wav::WavReader;
|
||||||
use crate::net::discovery::Discovery;
|
use crate::net::discovery::Discovery;
|
||||||
use crate::net::network::{AudioSocket, Connection};
|
use crate::net::network::{AudioSocket, Connection};
|
||||||
use crate::net::protocol::{AudioPacket, ClientInfo, ControlPacket, RpcResult};
|
use crate::net::protocol::{ClientInfo, ControlPacket, RpcResult};
|
||||||
use crate::net::rpc::RpcManager;
|
use crate::net::rpc::RpcManager;
|
||||||
|
use audio_manager::ServerAudioManager;
|
||||||
use anyhow::{Context, Result, anyhow};
|
use anyhow::{Context, Result, anyhow};
|
||||||
use dashmap::DashMap;
|
use dashmap::DashMap;
|
||||||
use parking_lot::Mutex;
|
|
||||||
use std::net::SocketAddr;
|
use std::net::SocketAddr;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use tokio::sync::mpsc;
|
|
||||||
use tokio_util::sync::CancellationToken;
|
use tokio_util::sync::CancellationToken;
|
||||||
use tokio_util::task::TaskTracker;
|
use tokio_util::task::TaskTracker;
|
||||||
|
|
||||||
pub enum RecorderCommand {
|
|
||||||
Start {
|
|
||||||
config: AudioConfig,
|
|
||||||
filename: String,
|
|
||||||
},
|
|
||||||
Stop,
|
|
||||||
}
|
|
||||||
|
|
||||||
pub struct Session {
|
pub struct Session {
|
||||||
pub info: ClientInfo,
|
pub info: ClientInfo,
|
||||||
pub conn: Arc<Connection>,
|
pub conn: Arc<Connection>,
|
||||||
@@ -29,10 +21,7 @@ pub struct Session {
|
|||||||
pub tcp_addr: SocketAddr,
|
pub tcp_addr: SocketAddr,
|
||||||
pub audio_addr: SocketAddr,
|
pub audio_addr: SocketAddr,
|
||||||
pub session_cancel: CancellationToken,
|
pub session_cancel: CancellationToken,
|
||||||
pub record_cancel: Mutex<Option<CancellationToken>>,
|
pub audio_manager: ServerAudioManager,
|
||||||
pub play_cancel: Mutex<Option<CancellationToken>>,
|
|
||||||
pub audio_tx: mpsc::Sender<AudioPacket>,
|
|
||||||
pub recorder_tx: mpsc::Sender<RecorderCommand>,
|
|
||||||
pub tracker: TaskTracker,
|
pub tracker: TaskTracker,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -67,7 +56,7 @@ impl Server {
|
|||||||
if let Some(tcp_addr) = this.udp_to_tcp.get(&src_addr) {
|
if let Some(tcp_addr) = this.udp_to_tcp.get(&src_addr) {
|
||||||
if let Some(session) = this.sessions.get(tcp_addr.value()) {
|
if let Some(session) = this.sessions.get(tcp_addr.value()) {
|
||||||
// 使用 try_send 避免某一个客户端阻塞导致全局音频延迟
|
// 使用 try_send 避免某一个客户端阻塞导致全局音频延迟
|
||||||
let _ = session.audio_tx.try_send(packet);
|
let _ = session.audio_manager.audio_tx().try_send(packet);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -148,80 +137,27 @@ impl Server {
|
|||||||
info.model, info.serial_number, audio_addr
|
info.model, info.serial_number, audio_addr
|
||||||
);
|
);
|
||||||
|
|
||||||
let (audio_tx, audio_rx) = mpsc::channel(1024);
|
|
||||||
let (recorder_tx, recorder_rx) = mpsc::channel(64);
|
|
||||||
let tracker = TaskTracker::new();
|
let tracker = TaskTracker::new();
|
||||||
|
let session_cancel = CancellationToken::new();
|
||||||
|
let (audio_manager, audio_rx, recorder_rx) =
|
||||||
|
ServerAudioManager::new(session_cancel.clone(), tracker.clone());
|
||||||
|
|
||||||
let session = Arc::new(Session {
|
let session = Arc::new(Session {
|
||||||
info,
|
info,
|
||||||
conn: conn.clone(),
|
conn: conn.clone(),
|
||||||
rpc: Arc::new(RpcManager::new()),
|
rpc: Arc::new(RpcManager::new()),
|
||||||
tcp_addr: addr,
|
tcp_addr: addr,
|
||||||
audio_addr,
|
audio_addr,
|
||||||
session_cancel: CancellationToken::new(),
|
session_cancel,
|
||||||
record_cancel: Mutex::new(None),
|
audio_manager,
|
||||||
play_cancel: Mutex::new(None),
|
|
||||||
audio_tx,
|
|
||||||
recorder_tx,
|
|
||||||
tracker: tracker.clone(),
|
tracker: tracker.clone(),
|
||||||
});
|
});
|
||||||
|
|
||||||
|
session.audio_manager.spawn_audio_processor(audio_rx, recorder_rx);
|
||||||
|
|
||||||
self.sessions.insert(addr, session.clone());
|
self.sessions.insert(addr, session.clone());
|
||||||
self.udp_to_tcp.insert(audio_addr, addr);
|
self.udp_to_tcp.insert(audio_addr, addr);
|
||||||
|
|
||||||
// --- Audio Processor Task ---
|
|
||||||
// 优化 Audio Receiver 的 OwnerShip,并使用 Tracker 跟踪
|
|
||||||
let processor_session = session.clone();
|
|
||||||
tracker.spawn(async move {
|
|
||||||
let mut active_recorder: Option<(WavWriter, OpusCodec, usize)> = None;
|
|
||||||
let mut audio_rx = audio_rx;
|
|
||||||
let mut recorder_rx = recorder_rx;
|
|
||||||
loop {
|
|
||||||
tokio::select! {
|
|
||||||
_ = processor_session.session_cancel.cancelled() => break,
|
|
||||||
cmd = recorder_rx.recv() => {
|
|
||||||
match cmd {
|
|
||||||
Some(RecorderCommand::Start { config, filename }) => {
|
|
||||||
if let Some((writer, _, _)) = active_recorder.take() {
|
|
||||||
let _ = writer.finalize();
|
|
||||||
}
|
|
||||||
match WavWriter::create(&filename, config.sample_rate, config.channels) {
|
|
||||||
Ok(writer) => {
|
|
||||||
match OpusCodec::new(&config) {
|
|
||||||
Ok(codec) => active_recorder = Some((writer, codec, config.frame_size)),
|
|
||||||
Err(e) => eprintln!("Failed to create opus codec: {}", e),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Err(e) => eprintln!("Failed to create wav writer: {}", e),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Some(RecorderCommand::Stop) => {
|
|
||||||
if let Some((writer, _, _)) = active_recorder.take() {
|
|
||||||
let _ = writer.finalize();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
None => break,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
packet = audio_rx.recv() => {
|
|
||||||
match packet {
|
|
||||||
Some(packet) => {
|
|
||||||
if let Some((writer, codec, frame_size)) = &mut active_recorder {
|
|
||||||
let mut pcm = vec![0i16; *frame_size];
|
|
||||||
if let Ok(n) = codec.decode(&packet.data, &mut pcm) {
|
|
||||||
let _ = writer.write_samples(&pcm[..n]);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
None => break,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if let Some((writer, _, _)) = active_recorder.take() {
|
|
||||||
let _ = writer.finalize();
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
// --- Main Connection Loop ---
|
// --- Main Connection Loop ---
|
||||||
// 统一心跳与超时处理,节省资源
|
// 统一心跳与超时处理,节省资源
|
||||||
let mut heartbeat = tokio::time::interval(std::time::Duration::from_secs(30));
|
let mut heartbeat = tokio::time::interval(std::time::Duration::from_secs(30));
|
||||||
@@ -328,27 +264,15 @@ impl Server {
|
|||||||
.map(|r| r.value().clone())
|
.map(|r| r.value().clone())
|
||||||
.context("Session not found")?;
|
.context("Session not found")?;
|
||||||
|
|
||||||
// 仅停止之前的录音任务
|
|
||||||
{
|
|
||||||
let mut guard = session.record_cancel.lock();
|
|
||||||
if let Some(token) = guard.take() {
|
|
||||||
token.cancel();
|
|
||||||
}
|
|
||||||
*guard = Some(session.session_cancel.child_token());
|
|
||||||
}
|
|
||||||
|
|
||||||
let filename = format!(
|
let filename = format!(
|
||||||
"temp/recorded_{}.wav",
|
"temp/recorded_{}.wav",
|
||||||
session.info.serial_number.replace(":", "")
|
session.info.serial_number.replace(":", "")
|
||||||
);
|
);
|
||||||
|
|
||||||
// 通知 Audio Processor 开始录音
|
// 通知 Audio Manager 开始录音
|
||||||
session
|
session
|
||||||
.recorder_tx
|
.audio_manager
|
||||||
.send(RecorderCommand::Start {
|
.start_recording(config.clone(), filename)
|
||||||
config: config.clone(),
|
|
||||||
filename,
|
|
||||||
})
|
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
session
|
session
|
||||||
@@ -367,11 +291,7 @@ impl Server {
|
|||||||
.context("Session not found")?;
|
.context("Session not found")?;
|
||||||
|
|
||||||
// 停止录音任务
|
// 停止录音任务
|
||||||
if let Some(token) = session.record_cancel.lock().take() {
|
session.audio_manager.stop_recording().await?;
|
||||||
token.cancel();
|
|
||||||
}
|
|
||||||
|
|
||||||
let _ = session.recorder_tx.send(RecorderCommand::Stop).await;
|
|
||||||
session.conn.send(&ControlPacket::StopRecording).await?;
|
session.conn.send(&ControlPacket::StopRecording).await?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
@@ -397,17 +317,6 @@ impl Server {
|
|||||||
..AudioConfig::music_48k()
|
..AudioConfig::music_48k()
|
||||||
};
|
};
|
||||||
|
|
||||||
// 仅停止之前的播放任务
|
|
||||||
let token = {
|
|
||||||
let mut guard = session.play_cancel.lock();
|
|
||||||
if let Some(token) = guard.take() {
|
|
||||||
token.cancel();
|
|
||||||
}
|
|
||||||
let token = session.session_cancel.child_token();
|
|
||||||
*guard = Some(token.clone());
|
|
||||||
token
|
|
||||||
};
|
|
||||||
|
|
||||||
session
|
session
|
||||||
.conn
|
.conn
|
||||||
.send(&ControlPacket::StartPlayback {
|
.send(&ControlPacket::StartPlayback {
|
||||||
@@ -415,54 +324,12 @@ impl Server {
|
|||||||
})
|
})
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
let audio_socket = self.audio.clone();
|
session.audio_manager.start_playback(
|
||||||
let target_addr = session.audio_addr;
|
config,
|
||||||
let playback_session = session.clone();
|
reader,
|
||||||
|
self.audio.clone(),
|
||||||
session.tracker.spawn(async move {
|
session.audio_addr,
|
||||||
let mut reader = reader;
|
).await?;
|
||||||
let mut codec = match OpusCodec::new(&config) {
|
|
||||||
Ok(c) => c,
|
|
||||||
Err(e) => {
|
|
||||||
eprintln!("Failed to create opus codec: {}", e);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
};
|
|
||||||
let mut pcm = vec![0i16; config.frame_size];
|
|
||||||
let mut opus = vec![0u8; 4096];
|
|
||||||
let mut interval = tokio::time::interval(std::time::Duration::from_millis(20));
|
|
||||||
|
|
||||||
println!("Playback started for {}", playback_session.tcp_addr);
|
|
||||||
|
|
||||||
loop {
|
|
||||||
tokio::select! {
|
|
||||||
_ = token.cancelled() => break,
|
|
||||||
_ = playback_session.session_cancel.cancelled() => break,
|
|
||||||
_ = interval.tick() => {
|
|
||||||
match reader.read_samples(&mut pcm) {
|
|
||||||
Ok(0) => break,
|
|
||||||
Ok(n) => {
|
|
||||||
if let Ok(len) = codec.encode(&pcm[..n], &mut opus) {
|
|
||||||
let _ = audio_socket
|
|
||||||
.send(
|
|
||||||
&AudioPacket {
|
|
||||||
data: opus[..len].to_vec(),
|
|
||||||
},
|
|
||||||
target_addr,
|
|
||||||
)
|
|
||||||
.await;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Err(e) => {
|
|
||||||
eprintln!("Failed to read samples: {}", e);
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
println!("Playback finished for {}", playback_session.tcp_addr);
|
|
||||||
});
|
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
@@ -475,10 +342,7 @@ impl Server {
|
|||||||
.context("Session not found")?;
|
.context("Session not found")?;
|
||||||
|
|
||||||
// 停止播放任务
|
// 停止播放任务
|
||||||
if let Some(token) = session.play_cancel.lock().take() {
|
session.audio_manager.stop_playback();
|
||||||
token.cancel();
|
|
||||||
}
|
|
||||||
|
|
||||||
session.conn.send(&ControlPacket::StopPlayback).await?;
|
session.conn.send(&ControlPacket::StopPlayback).await?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user