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