#[cfg(feature = "ble")]
use crate::kernel::msbc::MsbcDecoder;
use std::sync::atomic::{AtomicU64, Ordering};
#[cfg(feature = "ble")]
use std::sync::{Arc, Mutex};
pub trait AudioFrameSink: Send + Sync {
fn on_msbc_frame(&self, frame: &[u8]);
}
pub trait PcmSink: Send + Sync {
fn on_pcm(&self, samples: &[f32]);
}
pub struct CountingSink {
pub frames: AtomicU64,
pub bytes: AtomicU64,
}
impl CountingSink {
pub fn new() -> Self {
Self {
frames: AtomicU64::new(0),
bytes: AtomicU64::new(0),
}
}
pub fn frame_count(&self) -> u64 {
self.frames.load(Ordering::Relaxed)
}
pub fn byte_count(&self) -> u64 {
self.bytes.load(Ordering::Relaxed)
}
}
impl Default for CountingSink {
fn default() -> Self {
Self::new()
}
}
impl AudioFrameSink for CountingSink {
fn on_msbc_frame(&self, frame: &[u8]) {
self.frames.fetch_add(1, Ordering::Relaxed);
self.bytes.fetch_add(frame.len() as u64, Ordering::Relaxed);
}
}
#[cfg(feature = "ble")]
pub struct MsbcDecoderSink {
inner: Mutex<MsbcDecoderSinkInner>,
pcm_sink: Arc<dyn PcmSink>,
}
#[cfg(feature = "ble")]
struct MsbcDecoderSinkInner {
decoder: MsbcDecoder,
pcm_buf: Vec<f32>,
}
#[cfg(feature = "ble")]
impl MsbcDecoderSink {
pub fn new(pcm_sink: Arc<dyn PcmSink>) -> Self {
Self {
inner: Mutex::new(MsbcDecoderSinkInner {
decoder: MsbcDecoder::new(),
pcm_buf: Vec::with_capacity(320),
}),
pcm_sink,
}
}
}
#[cfg(feature = "ble")]
impl AudioFrameSink for MsbcDecoderSink {
fn on_msbc_frame(&self, frame: &[u8]) {
let mut inner = match self.inner.lock() {
Ok(g) => g,
Err(e) => e.into_inner(), };
match inner.decoder.decode_frame(frame) {
Ok(pcm_i16) => {
inner.pcm_buf.clear();
inner
.pcm_buf
.extend(pcm_i16.iter().map(|&s| s as f32 / 32768.0));
self.pcm_sink.on_pcm(&inner.pcm_buf);
}
Err(e) => log::warn!(target: "audio", "mSBC 帧解码失败: {}", e),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[cfg(feature = "ble")]
use std::sync::Mutex as StdMutex;
#[cfg(feature = "ble")]
struct VecSink(StdMutex<Vec<f32>>);
#[cfg(feature = "ble")]
impl PcmSink for VecSink {
fn on_pcm(&self, samples: &[f32]) {
self.0.lock().unwrap().extend_from_slice(samples);
}
}
#[test]
fn counting_sink_counts() {
let sink = CountingSink::new();
sink.on_msbc_frame(&[0u8; 57]);
sink.on_msbc_frame(&[0u8; 57]);
assert_eq!(sink.frame_count(), 2);
assert_eq!(sink.byte_count(), 114);
}
#[cfg(feature = "ble")]
#[test]
fn decoder_sink_rejects_bad_frame() {
let collected = Arc::new(VecSink(StdMutex::new(Vec::new())));
let sink = MsbcDecoderSink::new(collected.clone());
sink.on_msbc_frame(&[0u8; 57]);
assert!(collected.0.lock().unwrap().is_empty());
}
}