From 9c2d8f0dba462b347a706d0a539ccbe163e51b85 Mon Sep 17 00:00:00 2001
From: cai <cai@nbcai.cc>
Date: Mon, 06 Jul 2026 16:33:47 +0800
Subject: [PATCH] fix: tolerate split pcm samples in stream chunks

---
 src/main.rs |  105 ++++++++++++++++++++++++++++++++++++++++++++++++++++
 1 files changed, 104 insertions(+), 1 deletions(-)

diff --git a/src/main.rs b/src/main.rs
index e05f602..3b18943 100644
--- a/src/main.rs
+++ b/src/main.rs
@@ -1343,8 +1343,40 @@
         .trim()
         .to_ascii_lowercase();
     let frames = if format == "pcm_s16le" {
-        pcm_s16le_payload_to_frames(
+        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 {
+            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,
+                "runtime helper stream_audio_pcm_unaligned_tail_dropped"
+            );
+        }
+        if aligned_payload.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"
+            );
+            return Ok(0);
+        }
+        pcm_s16le_payload_to_frames(
+            &aligned_payload,
             audio_chunk.sample_rate.unwrap_or(TARGET_SAMPLE_RATE_HZ),
             audio_chunk
                 .channels
@@ -1472,6 +1504,75 @@
         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] {
@@ -1849,6 +1950,7 @@
     audio_chunk_count: u64,
     device_output_count: u64,
     encoded_audio_buffer: Vec<u8>,
+    pcm_audio_buffer: Vec<u8>,
 }
 
 impl Default for RuntimeTurnStreamState {
@@ -1862,6 +1964,7 @@
             audio_chunk_count: 0,
             device_output_count: 0,
             encoded_audio_buffer: Vec::new(),
+            pcm_audio_buffer: Vec::new(),
         }
     }
 }

--
Gitblit v1.9.3