From 959a2060fb276d601a147c19f45a1b51fe6dfac1 Mon Sep 17 00:00:00 2001
From: cai <cai@nbcai.cc>
Date: Sat, 08 Aug 2026 18:50:05 +0800
Subject: [PATCH] test: cover observer session metadata refresh
---
src/audio.rs | 351 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
1 files changed, 350 insertions(+), 1 deletions(-)
diff --git a/src/audio.rs b/src/audio.rs
index 1aadcb4..cec22f0 100644
--- a/src/audio.rs
+++ b/src/audio.rs
@@ -142,6 +142,262 @@
)
}
+pub fn pcm_s16le_bytes_to_frames(
+ pcm_bytes: &[u8],
+ source_sample_rate_hz: u32,
+ source_num_channels: u32,
+ target_sample_rate_hz: u32,
+ target_num_channels: u16,
+) -> Result<Vec<PcmFrame>> {
+ if source_sample_rate_hz == 0 {
+ return Err(anyhow!("pcm_s16le source sample rate is zero"));
+ }
+ if source_num_channels == 0 {
+ return Err(anyhow!("pcm_s16le source channel count is zero"));
+ }
+ if pcm_bytes.len() % 2 != 0 {
+ return Err(anyhow!("pcm_s16le payload has odd byte length"));
+ }
+ let samples: Vec<i16> = pcm_bytes
+ .chunks_exact(2)
+ .map(|chunk| i16::from_le_bytes([chunk[0], chunk[1]]))
+ .collect();
+ if samples.is_empty() {
+ return Ok(Vec::new());
+ }
+ let source_channels = u16::try_from(source_num_channels)
+ .map_err(|_| anyhow!("unsupported pcm_s16le channel count {source_num_channels}"))?;
+ let mut target_samples = remap_channels(samples, source_channels, target_num_channels);
+ target_samples = resample_linear(
+ target_samples,
+ source_sample_rate_hz,
+ target_sample_rate_hz,
+ target_num_channels,
+ );
+ Ok(chunk_pcm_samples(
+ target_samples,
+ target_sample_rate_hz,
+ target_num_channels,
+ ))
+}
+
+#[derive(Debug, Clone, Default)]
+pub struct PcmS16leStreamChunkResult {
+ pub frames: Vec<PcmFrame>,
+ pub dropped_tail_bytes: usize,
+ pub buffered_source_bytes: usize,
+ pub buffered_source_samples: usize,
+}
+
+#[derive(Debug, Clone, Default)]
+pub struct PcmS16leStreamDebugDump {
+ pub debug_source_path: Option<String>,
+ pub debug_pcm_wav_path: Option<String>,
+ pub debug_pcm_wav_size_bytes: Option<u64>,
+}
+
+#[derive(Debug, Clone)]
+pub struct PcmS16leStreamDecoder {
+ source_sample_rate_hz: u32,
+ source_num_channels: u16,
+ target_sample_rate_hz: u32,
+ target_num_channels: u16,
+ source_byte_buffer: Vec<u8>,
+ source_sample_buffer: Vec<i16>,
+ debug_source_samples: Vec<i16>,
+ debug_target_samples: Vec<i16>,
+}
+
+impl PcmS16leStreamDecoder {
+ pub fn new(
+ source_sample_rate_hz: u32,
+ source_num_channels: u32,
+ target_sample_rate_hz: u32,
+ target_num_channels: u16,
+ ) -> Result<Self> {
+ if source_sample_rate_hz == 0 {
+ return Err(anyhow!("pcm_s16le source sample rate is zero"));
+ }
+ if source_num_channels == 0 {
+ return Err(anyhow!("pcm_s16le source channel count is zero"));
+ }
+ let source_num_channels = u16::try_from(source_num_channels)
+ .map_err(|_| anyhow!("unsupported pcm_s16le channel count {source_num_channels}"))?;
+ Ok(Self {
+ source_sample_rate_hz,
+ source_num_channels,
+ target_sample_rate_hz,
+ target_num_channels,
+ source_byte_buffer: Vec::new(),
+ source_sample_buffer: Vec::new(),
+ debug_source_samples: Vec::new(),
+ debug_target_samples: Vec::new(),
+ })
+ }
+
+ pub fn push_bytes(
+ &mut self,
+ pcm_bytes: &[u8],
+ source_sample_rate_hz: u32,
+ source_num_channels: u32,
+ last: bool,
+ ) -> Result<PcmS16leStreamChunkResult> {
+ let source_num_channels = u16::try_from(source_num_channels)
+ .map_err(|_| anyhow!("unsupported pcm_s16le channel count {source_num_channels}"))?;
+ if source_sample_rate_hz != self.source_sample_rate_hz
+ || source_num_channels != self.source_num_channels
+ {
+ return Err(anyhow!(
+ "pcm_s16le stream format changed from {}Hz/{}ch to {}Hz/{}ch",
+ self.source_sample_rate_hz,
+ self.source_num_channels,
+ source_sample_rate_hz,
+ source_num_channels
+ ));
+ }
+
+ self.source_byte_buffer.extend_from_slice(pcm_bytes);
+ let source_frame_alignment_bytes = usize::from(self.source_num_channels) * 2;
+ let aligned_len = self.source_byte_buffer.len()
+ - self.source_byte_buffer.len() % source_frame_alignment_bytes;
+ if aligned_len > 0 {
+ let aligned_bytes: Vec<u8> = self.source_byte_buffer.drain(..aligned_len).collect();
+ let source_samples: Vec<i16> = aligned_bytes
+ .chunks_exact(2)
+ .map(|chunk| i16::from_le_bytes([chunk[0], chunk[1]]))
+ .collect();
+ self.debug_source_samples.extend_from_slice(&source_samples);
+ let remapped = remap_channels(
+ source_samples,
+ self.source_num_channels,
+ self.target_num_channels,
+ );
+ self.source_sample_buffer.extend(remapped);
+ }
+
+ let dropped_tail_bytes = if last && !self.source_byte_buffer.is_empty() {
+ let dropped = self.source_byte_buffer.len();
+ self.source_byte_buffer.clear();
+ dropped
+ } else {
+ 0
+ };
+
+ let mut frames = self.drain_full_frames();
+ if last {
+ if let Some(frame) = self.drain_final_partial_frame() {
+ frames.push(frame);
+ }
+ }
+
+ Ok(PcmS16leStreamChunkResult {
+ frames,
+ dropped_tail_bytes,
+ buffered_source_bytes: self.source_byte_buffer.len(),
+ buffered_source_samples: self.source_sample_buffer.len(),
+ })
+ }
+
+ pub fn write_debug_dump(
+ &self,
+ dir: &str,
+ call_id: &str,
+ debug_label: &str,
+ ) -> Result<PcmS16leStreamDebugDump> {
+ let base_dir = PathBuf::from(dir);
+ fs::create_dir_all(&base_dir)
+ .with_context(|| format!("failed to create audio debug dump dir {dir}"))?;
+ let safe_call_id = sanitize_file_segment(call_id);
+ let safe_label = sanitize_file_segment(debug_label);
+
+ let source_path = base_dir.join(format!("{safe_call_id}-{safe_label}-source.wav"));
+ write_debug_wav(
+ &source_path,
+ &self.debug_source_samples,
+ self.source_sample_rate_hz,
+ self.source_num_channels,
+ )?;
+
+ let pcm_path = base_dir.join(format!("{safe_call_id}-{safe_label}-target.wav"));
+ write_debug_wav(
+ &pcm_path,
+ &self.debug_target_samples,
+ self.target_sample_rate_hz,
+ self.target_num_channels,
+ )?;
+ let pcm_size = fs::metadata(&pcm_path)
+ .with_context(|| format!("failed to stat audio debug wav {}", pcm_path.display()))?
+ .len();
+
+ Ok(PcmS16leStreamDebugDump {
+ debug_source_path: Some(source_path.to_string_lossy().to_string()),
+ debug_pcm_wav_path: Some(pcm_path.to_string_lossy().to_string()),
+ debug_pcm_wav_size_bytes: Some(pcm_size),
+ })
+ }
+
+ fn drain_full_frames(&mut self) -> Vec<PcmFrame> {
+ let source_samples_per_frame = self.source_samples_per_20ms_frame();
+ let mut frames = Vec::new();
+ while self.source_sample_buffer.len() >= source_samples_per_frame {
+ let source_samples: Vec<i16> = self
+ .source_sample_buffer
+ .drain(..source_samples_per_frame)
+ .collect();
+ frames.push(self.convert_source_samples_to_frame(source_samples, true));
+ }
+ frames
+ }
+
+ fn drain_final_partial_frame(&mut self) -> Option<PcmFrame> {
+ if self.source_sample_buffer.is_empty() {
+ return None;
+ }
+ let source_samples: Vec<i16> = self.source_sample_buffer.drain(..).collect();
+ Some(self.convert_source_samples_to_frame(source_samples, true))
+ }
+
+ fn convert_source_samples_to_frame(
+ &mut self,
+ source_samples: Vec<i16>,
+ pad_to_frame: bool,
+ ) -> PcmFrame {
+ let mut target_samples = resample_linear(
+ source_samples,
+ self.source_sample_rate_hz,
+ self.target_sample_rate_hz,
+ self.target_num_channels,
+ );
+ let target_frame_samples = self.target_samples_per_20ms_frame();
+ if target_samples.len() > target_frame_samples {
+ target_samples.truncate(target_frame_samples);
+ } else if pad_to_frame && target_samples.len() < target_frame_samples {
+ target_samples.resize(target_frame_samples, 0);
+ }
+ self.debug_target_samples.extend_from_slice(&target_samples);
+ PcmFrame::new(
+ target_samples,
+ self.target_sample_rate_hz,
+ u32::from(self.target_num_channels),
+ self.target_samples_per_channel_per_20ms_frame() as u32,
+ )
+ }
+
+ fn source_samples_per_20ms_frame(&self) -> usize {
+ let source_frames =
+ ((f64::from(self.source_sample_rate_hz) / 50.0).round() as usize).max(1);
+ source_frames * usize::from(self.target_num_channels)
+ }
+
+ fn target_samples_per_20ms_frame(&self) -> usize {
+ self.target_samples_per_channel_per_20ms_frame() * usize::from(self.target_num_channels)
+ }
+
+ fn target_samples_per_channel_per_20ms_frame(&self) -> usize {
+ ((self.target_sample_rate_hz / 1000) * 20).max(1) as usize
+ }
+}
+
fn decode_audio_frames(
audio_bytes: &[u8],
source_kind: &'static str,
@@ -674,7 +930,7 @@
#[cfg(test)]
mod tests {
- use super::load_pre_recorded_frames;
+ use super::{PcmS16leStreamDecoder, load_pre_recorded_frames, pcm_s16le_bytes_to_frames};
#[tokio::test]
async fn load_pre_recorded_frames_reads_local_wav() {
@@ -717,4 +973,97 @@
assert!(loaded.diagnostics.debug_source_path.is_some());
assert!(loaded.diagnostics.debug_pcm_wav_path.is_some());
}
+
+ #[test]
+ fn pcm_s16le_bytes_to_frames_resamples_elevenlabs_pcm_to_livekit_frames() {
+ let mut pcm = Vec::new();
+ for index in 0..4410 {
+ let sample = if index % 2 == 0 { 1024_i16 } else { -1024_i16 };
+ pcm.extend_from_slice(&sample.to_le_bytes());
+ }
+
+ let frames = pcm_s16le_bytes_to_frames(&pcm, 44_100, 1, 48_000, 1).expect("pcm frames");
+
+ assert!(!frames.is_empty());
+ assert_eq!(frames[0].sample_rate, 48_000);
+ assert_eq!(frames[0].num_channels, 1);
+ assert_eq!(frames[0].samples_per_channel, 960);
+ }
+
+ #[test]
+ fn pcm_s16le_stream_decoder_does_not_pad_each_network_chunk() {
+ let mut decoder = PcmS16leStreamDecoder::new(44_100, 1, 48_000, 1).expect("stream decoder");
+ let mut pcm = Vec::new();
+ for index in 0..4410 {
+ let sample = if index % 2 == 0 { 1024_i16 } else { -1024_i16 };
+ pcm.extend_from_slice(&sample.to_le_bytes());
+ }
+
+ let mut frames = Vec::new();
+ let chunk_size = 698 * 2;
+ for (index, chunk) in pcm.chunks(chunk_size).enumerate() {
+ let last = (index + 1) * chunk_size >= pcm.len();
+ let result = decoder
+ .push_bytes(chunk, 44_100, 1, last)
+ .expect("push pcm chunk");
+ frames.extend(result.frames);
+ }
+
+ assert_eq!(5, frames.len());
+ assert_eq!(
+ 4_800,
+ frames.iter().map(|frame| frame.data.len()).sum::<usize>()
+ );
+ assert!(frames.iter().all(|frame| frame.sample_rate == 48_000));
+ assert!(frames.iter().all(|frame| frame.num_channels == 1));
+ assert!(frames.iter().all(|frame| frame.samples_per_channel == 960));
+ }
+
+ #[test]
+ fn pcm_s16le_stream_decoder_waits_until_full_20ms_frame() {
+ let mut decoder = PcmS16leStreamDecoder::new(44_100, 1, 48_000, 1).expect("stream decoder");
+ let ten_ms = vec![0_u8; 441 * 2];
+
+ let first = decoder
+ .push_bytes(&ten_ms, 44_100, 1, false)
+ .expect("first chunk");
+ assert!(first.frames.is_empty());
+ assert_eq!(441, first.buffered_source_samples);
+
+ let second = decoder
+ .push_bytes(&ten_ms, 44_100, 1, false)
+ .expect("second chunk");
+ assert_eq!(1, second.frames.len());
+ assert_eq!(960, second.frames[0].samples_per_channel);
+ assert_eq!(0, second.buffered_source_samples);
+ }
+
+ #[test]
+ fn pcm_s16le_stream_decoder_can_keep_16k_native_audio_profile() {
+ let mut decoder = PcmS16leStreamDecoder::new(16_000, 1, 16_000, 1).expect("stream decoder");
+ let mut pcm = Vec::new();
+ for index in 0..1_600 {
+ let sample = if index % 2 == 0 { 768_i16 } else { -768_i16 };
+ pcm.extend_from_slice(&sample.to_le_bytes());
+ }
+
+ let mut frames = Vec::new();
+ let chunk_size = 250 * 2;
+ for (index, chunk) in pcm.chunks(chunk_size).enumerate() {
+ let last = (index + 1) * chunk_size >= pcm.len();
+ let result = decoder
+ .push_bytes(chunk, 16_000, 1, last)
+ .expect("push pcm chunk");
+ frames.extend(result.frames);
+ }
+
+ assert_eq!(5, frames.len());
+ assert!(frames.iter().all(|frame| frame.sample_rate == 16_000));
+ assert!(frames.iter().all(|frame| frame.num_channels == 1));
+ assert!(frames.iter().all(|frame| frame.samples_per_channel == 320));
+ assert_eq!(
+ 1_600,
+ frames.iter().map(|frame| frame.data.len()).sum::<usize>()
+ );
+ }
}
--
Gitblit v1.9.3