From eeb442d5b8710e849119423b64f9043eda50324b Mon Sep 17 00:00:00 2001 From: Del Wang Date: Sat, 3 Jan 2026 19:26:35 +0800 Subject: [PATCH] =?UTF-8?q?refactor:=20=E5=88=86=E7=A6=BB=E9=9F=B3?= =?UTF-8?q?=E9=A2=91=E7=AE=A1=E7=90=86=E6=A8=A1=E5=9D=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../client-v2/src/app/client/audio_manager.rs | 151 +++++++++++++ packages/client-v2/src/app/client/mod.rs | 154 ++------------ .../client-v2/src/app/server/audio_manager.rs | 201 ++++++++++++++++++ packages/client-v2/src/app/server/mod.rs | 188 +++------------- 4 files changed, 394 insertions(+), 300 deletions(-) create mode 100644 packages/client-v2/src/app/client/audio_manager.rs create mode 100644 packages/client-v2/src/app/server/audio_manager.rs diff --git a/packages/client-v2/src/app/client/audio_manager.rs b/packages/client-v2/src/app/client/audio_manager.rs new file mode 100644 index 0000000..b71a82d --- /dev/null +++ b/packages/client-v2/src/app/client/audio_manager.rs @@ -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, + server_audio_addr: SocketAddr, + session_cancel: CancellationToken, + record_cancel: RwLock>, + play_cancel: RwLock>, +} + +impl ClientAudioManager { + pub fn new( + audio_socket: Arc, + 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::>(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::>(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."); + }); + } +} diff --git a/packages/client-v2/src/app/client/mod.rs b/packages/client-v2/src/app/client/mod.rs index 8c21e78..8970f56 100644 --- a/packages/client-v2/src/app/client/mod.rs +++ b/packages/client-v2/src/app/client/mod.rs @@ -1,14 +1,13 @@ #![cfg(target_os = "linux")] -use crate::audio::codec::OpusCodec; -use crate::audio::config::AudioConfig; -use crate::audio::player::AudioPlayer; -use crate::audio::recorder::AudioRecorder; +mod audio_manager; + use crate::net::discovery::Discovery; 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 anyhow::{Result, anyhow}; +use audio_manager::ClientAudioManager; use std::net::SocketAddr; use std::sync::Arc; use tokio::sync::{RwLock, mpsc}; @@ -17,11 +16,8 @@ use tokio_util::sync::CancellationToken; /// 内部连接上下文,包含了音频流所需的全部信息 struct ActiveSession { conn: Arc, - audio_socket: Arc, - server_audio_addr: SocketAddr, + audio_manager: ClientAudioManager, session_cancel: CancellationToken, // 控制整个 Session 的生命周期 - record_cancel: RwLock>, - play_cancel: RwLock>, } pub struct Client { @@ -126,13 +122,15 @@ impl Client { ); // --- 初始化 Session --- + let session_cancel = CancellationToken::new(); let session = Arc::new(ActiveSession { conn: conn.clone(), - audio_socket, - server_audio_addr, - session_cancel: CancellationToken::new(), - record_cancel: RwLock::new(None), - play_cancel: RwLock::new(None), + audio_manager: ClientAudioManager::new( + audio_socket, + server_audio_addr, + session_cancel.clone(), + ), + session_cancel, }); *self.session.write().await = Some(session.clone()); @@ -206,22 +204,16 @@ impl Client { .await?; } ControlPacket::StartRecording { config } => { - self.stop_recorder(session).await; // 开启前先停止旧的,防止资源冲突 - let token = session.session_cancel.child_token(); - *session.record_cancel.write().await = Some(token.clone()); - self.spawn_recorder(session.clone(), config, token); + session.audio_manager.start_recording(config).await; } ControlPacket::StartPlayback { config } => { - self.stop_player(session).await; - let token = session.session_cancel.child_token(); - *session.play_cancel.write().await = Some(token.clone()); - self.spawn_player(session.clone(), config, token); + session.audio_manager.start_playback(config).await; } ControlPacket::StopRecording => { - self.stop_recorder(session).await; + session.audio_manager.stop_recorder().await; } 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, - config: AudioConfig, - token: CancellationToken, - ) { - tokio::spawn(async move { - let (pcm_tx, mut pcm_rx) = mpsc::channel::>(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, - config: AudioConfig, - token: CancellationToken, - ) { - tokio::spawn(async move { - let (pcm_tx, mut pcm_rx) = mpsc::channel::>(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) -> Result { let (id, rx) = self.rpc.register(); let session_guard = self.session.read().await; diff --git a/packages/client-v2/src/app/server/audio_manager.rs b/packages/client-v2/src/app/server/audio_manager.rs new file mode 100644 index 0000000..d0b3b5d --- /dev/null +++ b/packages/client-v2/src/app/server/audio_manager.rs @@ -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>, + play_cancel: Mutex>, + audio_tx: mpsc::Sender, + recorder_tx: mpsc::Sender, + tracker: TaskTracker, +} + +impl ServerAudioManager { + pub fn new(session_cancel: CancellationToken, tracker: TaskTracker) -> (Self, mpsc::Receiver, mpsc::Receiver) { + 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 { + 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, + 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, mut recorder_rx: mpsc::Receiver) { + 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(); + } + }); + } +} + diff --git a/packages/client-v2/src/app/server/mod.rs b/packages/client-v2/src/app/server/mod.rs index 0e61a8e..cebfc68 100644 --- a/packages/client-v2/src/app/server/mod.rs +++ b/packages/client-v2/src/app/server/mod.rs @@ -1,27 +1,19 @@ -use crate::audio::codec::OpusCodec; +mod audio_manager; + use crate::audio::config::AudioConfig; -use crate::audio::wav::{WavReader, WavWriter}; +use crate::audio::wav::WavReader; use crate::net::discovery::Discovery; 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 audio_manager::ServerAudioManager; use anyhow::{Context, Result, anyhow}; use dashmap::DashMap; -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 Session { pub info: ClientInfo, pub conn: Arc, @@ -29,10 +21,7 @@ pub struct Session { pub tcp_addr: SocketAddr, pub audio_addr: SocketAddr, pub session_cancel: CancellationToken, - pub record_cancel: Mutex>, - pub play_cancel: Mutex>, - pub audio_tx: mpsc::Sender, - pub recorder_tx: mpsc::Sender, + pub audio_manager: ServerAudioManager, pub tracker: TaskTracker, } @@ -67,7 +56,7 @@ impl Server { if let Some(tcp_addr) = this.udp_to_tcp.get(&src_addr) { if let Some(session) = this.sessions.get(tcp_addr.value()) { // 使用 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 ); - let (audio_tx, audio_rx) = mpsc::channel(1024); - let (recorder_tx, recorder_rx) = mpsc::channel(64); 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 { info, conn: conn.clone(), rpc: Arc::new(RpcManager::new()), tcp_addr: addr, audio_addr, - session_cancel: CancellationToken::new(), - record_cancel: Mutex::new(None), - play_cancel: Mutex::new(None), - audio_tx, - recorder_tx, + session_cancel, + audio_manager, tracker: tracker.clone(), }); + session.audio_manager.spawn_audio_processor(audio_rx, recorder_rx); + self.sessions.insert(addr, session.clone()); 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 --- // 统一心跳与超时处理,节省资源 let mut heartbeat = tokio::time::interval(std::time::Duration::from_secs(30)); @@ -328,27 +264,15 @@ impl Server { .map(|r| r.value().clone()) .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!( "temp/recorded_{}.wav", session.info.serial_number.replace(":", "") ); - // 通知 Audio Processor 开始录音 + // 通知 Audio Manager 开始录音 session - .recorder_tx - .send(RecorderCommand::Start { - config: config.clone(), - filename, - }) + .audio_manager + .start_recording(config.clone(), filename) .await?; session @@ -367,11 +291,7 @@ impl Server { .context("Session not found")?; // 停止录音任务 - if let Some(token) = session.record_cancel.lock().take() { - token.cancel(); - } - - let _ = session.recorder_tx.send(RecorderCommand::Stop).await; + session.audio_manager.stop_recording().await?; session.conn.send(&ControlPacket::StopRecording).await?; Ok(()) } @@ -397,17 +317,6 @@ impl Server { ..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 .conn .send(&ControlPacket::StartPlayback { @@ -415,54 +324,12 @@ impl Server { }) .await?; - let audio_socket = self.audio.clone(); - let target_addr = session.audio_addr; - let playback_session = session.clone(); - - session.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)); - - 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); - }); + session.audio_manager.start_playback( + config, + reader, + self.audio.clone(), + session.audio_addr, + ).await?; Ok(()) } @@ -475,10 +342,7 @@ impl Server { .context("Session not found")?; // 停止播放任务 - if let Some(token) = session.play_cancel.lock().take() { - token.cancel(); - } - + session.audio_manager.stop_playback(); session.conn.send(&ControlPacket::StopPlayback).await?; Ok(()) }