From a62af4b377643d3b76318ee782763b3833b27456 Mon Sep 17 00:00:00 2001
From: cai <cai@nbcai.cc>
Date: Tue, 07 Jul 2026 15:38:44 +0800
Subject: [PATCH] fix: keep pcm stream continuous
---
src/main.rs | 180 +++++++++++++++++++++++------------------------------------
1 files changed, 71 insertions(+), 109 deletions(-)
diff --git a/src/main.rs b/src/main.rs
index 3b18943..64eedd4 100644
--- a/src/main.rs
+++ b/src/main.rs
@@ -12,7 +12,7 @@
};
use anyhow::{Context, Result, anyhow};
-use audio::{AudioDiagnostics, PcmFrame, load_pre_recorded_frames};
+use audio::{AudioDiagnostics, load_pre_recorded_frames};
use base64::{Engine as _, engine::general_purpose};
use futures_util::StreamExt;
use libwebrtc::{
@@ -1343,45 +1343,54 @@
.trim()
.to_ascii_lowercase();
let frames = if format == "pcm_s16le" {
- let frame_alignment_bytes = pcm_s16le_frame_alignment_bytes(
- audio_chunk
- .channels
- .unwrap_or(u32::from(TARGET_NUM_CHANNELS)),
- )?;
- let (aligned_payload, dropped_tail_bytes) = take_aligned_pcm_payload(
- &mut state.pcm_audio_buffer,
- &payload,
- frame_alignment_bytes,
- audio_chunk.last.unwrap_or(false),
- );
- if dropped_tail_bytes > 0 {
+ let sample_rate = audio_chunk.sample_rate.unwrap_or(TARGET_SAMPLE_RATE_HZ);
+ let channels = audio_chunk
+ .channels
+ .unwrap_or(u32::from(TARGET_NUM_CHANNELS));
+ if state.pcm_stream_decoder.is_none() {
+ state.pcm_stream_decoder = Some(audio::PcmS16leStreamDecoder::new(
+ sample_rate,
+ channels,
+ TARGET_SAMPLE_RATE_HZ,
+ TARGET_NUM_CHANNELS,
+ )?);
+ }
+ state.pcm_stream_network_chunk_count =
+ state.pcm_stream_network_chunk_count.saturating_add(1);
+ let stream_result = state
+ .pcm_stream_decoder
+ .as_mut()
+ .expect("pcm stream decoder initialized")
+ .push_bytes(
+ &payload,
+ sample_rate,
+ channels,
+ audio_chunk.last.unwrap_or(false),
+ )?;
+ if stream_result.dropped_tail_bytes > 0 {
warn!(
call_id = %call_id,
trace_id = %trace_id,
turn_id = %turn.turn_id,
chunk_seq = audio_chunk.chunk_seq.unwrap_or_default(),
- dropped_tail_bytes,
+ dropped_tail_bytes = stream_result.dropped_tail_bytes,
"runtime helper stream_audio_pcm_unaligned_tail_dropped"
);
}
- if aligned_payload.is_empty() {
+ if stream_result.frames.is_empty() {
warn!(
call_id = %call_id,
trace_id = %trace_id,
turn_id = %turn.turn_id,
chunk_seq = audio_chunk.chunk_seq.unwrap_or_default(),
- buffered_bytes = state.pcm_audio_buffer.len(),
- "runtime helper stream_audio_pcm_waiting_for_sample_boundary"
+ buffered_source_bytes = stream_result.buffered_source_bytes,
+ buffered_source_samples = stream_result.buffered_source_samples,
+ network_chunk_count = state.pcm_stream_network_chunk_count,
+ "runtime helper stream_audio_pcm_waiting_for_20ms_frame"
);
return Ok(0);
}
- pcm_s16le_payload_to_frames(
- &aligned_payload,
- audio_chunk.sample_rate.unwrap_or(TARGET_SAMPLE_RATE_HZ),
- audio_chunk
- .channels
- .unwrap_or(u32::from(TARGET_NUM_CHANNELS)),
- )?
+ stream_result.frames
} else if matches!(format.as_str(), "mp3" | "mpeg" | "wav") {
state.encoded_audio_buffer.extend_from_slice(&payload);
match audio::decode_audio_bytes_to_frames(
@@ -1473,6 +1482,34 @@
sleep_until(pacing_started_at + Duration::from_millis(((index + 1) as u64) * 20)).await;
}
if audio_chunk.last.unwrap_or(false) {
+ let mut debug_source_path = None;
+ let mut debug_pcm_wav_path = None;
+ let mut debug_pcm_wav_size_bytes = None;
+ if format == "pcm_s16le" {
+ if let (Some(debug_dump_dir), Some(decoder)) = (
+ bridge_config.audio_debug_dump_dir.as_deref(),
+ state.pcm_stream_decoder.as_ref(),
+ ) {
+ match decoder.write_debug_dump(
+ debug_dump_dir,
+ call_id,
+ &format!("stream-reply-{}", turn.turn_id),
+ ) {
+ Ok(debug_dump) => {
+ debug_source_path = debug_dump.debug_source_path;
+ debug_pcm_wav_path = debug_dump.debug_pcm_wav_path;
+ debug_pcm_wav_size_bytes = debug_dump.debug_pcm_wav_size_bytes;
+ }
+ Err(error) => warn!(
+ call_id = %call_id,
+ trace_id = %trace_id,
+ turn_id = %turn.turn_id,
+ error = %safe_error(&error.to_string()),
+ "runtime helper stream_audio_pcm_debug_dump_failed"
+ ),
+ }
+ }
+ }
emit_activity(
call_id,
trace_id,
@@ -1485,94 +1522,17 @@
"replyPlaybackMode": state.reply_playback_mode.as_str(),
"format": format.as_str(),
"chunkSeq": audio_chunk.chunk_seq,
+ "networkChunkCount": state.pcm_stream_network_chunk_count,
+ "debugSourcePath": debug_source_path,
+ "debugPcmWavPath": debug_pcm_wav_path,
+ "debugPcmWavSizeBytes": debug_pcm_wav_size_bytes,
+ "sampleRate": audio_chunk.sample_rate,
+ "channels": audio_chunk.channels,
"replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64,
}),
);
}
Ok(frames.len())
-}
-
-fn pcm_s16le_payload_to_frames(
- payload: &[u8],
- sample_rate: u32,
- channels: u32,
-) -> Result<Vec<PcmFrame>> {
- audio::pcm_s16le_bytes_to_frames(
- payload,
- sample_rate,
- channels,
- TARGET_SAMPLE_RATE_HZ,
- TARGET_NUM_CHANNELS,
- )
-}
-
-fn pcm_s16le_frame_alignment_bytes(channels: u32) -> Result<usize> {
- if channels == 0 {
- return Err(anyhow!("pcm_s16le channel count is zero"));
- }
- let channel_count = usize::try_from(channels)
- .map_err(|_| anyhow!("unsupported pcm_s16le channel count {channels}"))?;
- Ok(2 * channel_count)
-}
-
-fn take_aligned_pcm_payload(
- buffer: &mut Vec<u8>,
- payload: &[u8],
- frame_alignment_bytes: usize,
- last: bool,
-) -> (Vec<u8>, usize) {
- let alignment = frame_alignment_bytes.max(2);
- buffer.extend_from_slice(payload);
- let aligned_len = buffer.len() - buffer.len() % alignment;
- let aligned_payload = if aligned_len == 0 {
- Vec::new()
- } else {
- buffer.drain(..aligned_len).collect()
- };
- let dropped_tail_bytes = if last && !buffer.is_empty() {
- let dropped = buffer.len();
- buffer.clear();
- dropped
- } else {
- 0
- };
- (aligned_payload, dropped_tail_bytes)
-}
-
-#[cfg(test)]
-mod pcm_stream_tests {
- use super::take_aligned_pcm_payload;
-
- #[test]
- fn take_aligned_pcm_payload_buffers_split_sample_bytes() {
- let mut buffer = Vec::new();
-
- let (first, dropped) = take_aligned_pcm_payload(&mut buffer, &[0x01], 2, false);
- assert!(first.is_empty());
- assert_eq!(0, dropped);
- assert_eq!(vec![0x01], buffer);
-
- let (second, dropped) = take_aligned_pcm_payload(&mut buffer, &[0x02, 0x03], 2, false);
- assert_eq!(vec![0x01, 0x02], second);
- assert_eq!(0, dropped);
- assert_eq!(vec![0x03], buffer);
-
- let (third, dropped) = take_aligned_pcm_payload(&mut buffer, &[0x04], 2, true);
- assert_eq!(vec![0x03, 0x04], third);
- assert_eq!(0, dropped);
- assert!(buffer.is_empty());
- }
-
- #[test]
- fn take_aligned_pcm_payload_drops_final_half_sample() {
- let mut buffer = Vec::new();
-
- let (payload, dropped) =
- take_aligned_pcm_payload(&mut buffer, &[0x01, 0x02, 0x03], 2, true);
- assert_eq!(vec![0x01, 0x02], payload);
- assert_eq!(1, dropped);
- assert!(buffer.is_empty());
- }
}
fn trim_ascii_whitespace(value: &[u8]) -> &[u8] {
@@ -1950,7 +1910,8 @@
audio_chunk_count: u64,
device_output_count: u64,
encoded_audio_buffer: Vec<u8>,
- pcm_audio_buffer: Vec<u8>,
+ pcm_stream_decoder: Option<audio::PcmS16leStreamDecoder>,
+ pcm_stream_network_chunk_count: u64,
}
impl Default for RuntimeTurnStreamState {
@@ -1964,7 +1925,8 @@
audio_chunk_count: 0,
device_output_count: 0,
encoded_audio_buffer: Vec::new(),
- pcm_audio_buffer: Vec::new(),
+ pcm_stream_decoder: None,
+ pcm_stream_network_chunk_count: 0,
}
}
}
--
Gitblit v1.9.3