cai
2026-07-11 76ab325311aa14c9d3bf55293047afbba1166403
feat: mark reply audio segment boundaries
1 files modified
109 ■■■■■ changed files
src/main.rs 109 ●●●●● patch | view | raw | blame | history
src/main.rs
@@ -4,6 +4,7 @@
use std::{
    borrow::Cow,
    collections::HashSet,
    env, fs,
    path::{Path, PathBuf},
    sync::{
@@ -1719,6 +1720,42 @@
        .unwrap_or("pcm_s16le")
        .trim()
        .to_ascii_lowercase();
    if !matches!(format.as_str(), "pcm_s16le" | "mp3" | "mpeg" | "wav") {
        return Err(anyhow!("unsupported reply_audio_chunk format {format}"));
    }
    match state.reply_chunk_markers.observe(audio_chunk.segment_seq) {
        ReplyChunkMarker::FirstReply => emit_activity(
            call_id,
            trace_id,
            Some(&turn.turn_id),
            "helper_first_reply_audio_chunk_received",
            "ok",
            None,
            None,
            json!({
                "segmentSeq": audio_chunk.segment_seq,
                "chunkSeq": audio_chunk.chunk_seq,
                "format": format.as_str(),
                "bytes": payload.len(),
            }),
        ),
        ReplyChunkMarker::SegmentFirst => emit_activity(
            call_id,
            trace_id,
            Some(&turn.turn_id),
            "helper_segment_first_audio_chunk_received",
            "ok",
            None,
            None,
            json!({
                "segmentSeq": audio_chunk.segment_seq,
                "chunkSeq": audio_chunk.chunk_seq,
                "format": format.as_str(),
                "bytes": payload.len(),
            }),
        ),
        ReplyChunkMarker::None => {}
    }
    let frames = if format == "pcm_s16le" {
        let sample_rate = audio_chunk.sample_rate.unwrap_or(sink.sample_rate_hz);
        let channels = audio_chunk.channels.unwrap_or(u32::from(sink.num_channels));
@@ -1810,7 +1847,7 @@
            Err(error) => return Err(error).context("failed to decode final stream audio chunk"),
        }
    } else {
        return Err(anyhow!("unsupported reply_audio_chunk format {format}"));
        unreachable!("supported encoded format checked above")
    };
    if frames.is_empty() {
        return Ok(0);
@@ -2309,6 +2346,7 @@
    encoded_audio_buffer: Vec<u8>,
    pcm_stream_decoder: Option<audio::PcmS16leStreamDecoder>,
    pcm_stream_network_chunk_count: u64,
    reply_chunk_markers: ReplyChunkMarkerState,
}
impl Default for RuntimeTurnStreamState {
@@ -2324,7 +2362,37 @@
            encoded_audio_buffer: Vec::new(),
            pcm_stream_decoder: None,
            pcm_stream_network_chunk_count: 0,
            reply_chunk_markers: ReplyChunkMarkerState::default(),
        }
    }
}
#[derive(Debug, PartialEq, Eq)]
enum ReplyChunkMarker {
    FirstReply,
    SegmentFirst,
    None,
}
#[derive(Default)]
struct ReplyChunkMarkerState {
    first_reply_seen: bool,
    seen_segments: HashSet<u64>,
}
impl ReplyChunkMarkerState {
    fn observe(&mut self, segment_seq: Option<u64>) -> ReplyChunkMarker {
        let first_for_segment = segment_seq
            .map(|value| self.seen_segments.insert(value))
            .unwrap_or(false);
        if !self.first_reply_seen {
            self.first_reply_seen = true;
            return ReplyChunkMarker::FirstReply;
        }
        if first_for_segment {
            return ReplyChunkMarker::SegmentFirst;
        }
        ReplyChunkMarker::None
    }
}
@@ -2379,6 +2447,8 @@
struct RuntimeTurnStreamAudioChunk {
    #[serde(rename = "chunkSeq", alias = "seq")]
    chunk_seq: Option<u64>,
    #[serde(rename = "segmentSeq")]
    segment_seq: Option<u64>,
    format: Option<String>,
    #[serde(rename = "sampleRate")]
    sample_rate: Option<u32>,
@@ -3489,3 +3559,40 @@
        Ok(())
    }
}
#[cfg(test)]
mod tests {
    use super::{ReplyChunkMarker, ReplyChunkMarkerState, RuntimeTurnStreamEvent};
    #[test]
    fn reply_chunk_marker_state_emits_turn_first_once_and_later_segment_first_once() {
        let mut state = ReplyChunkMarkerState::default();
        assert_eq!(ReplyChunkMarker::FirstReply, state.observe(Some(1)));
        assert_eq!(ReplyChunkMarker::None, state.observe(Some(1)));
        assert_eq!(ReplyChunkMarker::SegmentFirst, state.observe(Some(2)));
        assert_eq!(ReplyChunkMarker::None, state.observe(Some(2)));
        assert_eq!(ReplyChunkMarker::SegmentFirst, state.observe(Some(3)));
    }
    #[test]
    fn reply_chunk_marker_state_without_segment_only_emits_turn_first() {
        let mut state = ReplyChunkMarkerState::default();
        assert_eq!(ReplyChunkMarker::FirstReply, state.observe(None));
        assert_eq!(ReplyChunkMarker::None, state.observe(None));
    }
    #[test]
    fn runtime_turn_stream_audio_chunk_reads_segment_seq() {
        let event: RuntimeTurnStreamEvent = serde_json::from_str(
            r#"{"type":"reply_audio_chunk","audioChunk":{"chunkSeq":4,"segmentSeq":2,"format":"pcm_s16le","payloadBase64":"AA==","last":false}}"#,
        )
        .expect("turn stream event");
        assert_eq!(
            Some(2),
            event.audio_chunk.and_then(|chunk| chunk.segment_seq)
        );
    }
}