cai
2026-08-11 dc98e7dbb2e909beab577f287da8dccac374c714
src/main.rs
@@ -809,6 +809,7 @@
                    turn_bridge_config.clone(),
                    http.clone(),
                    sink.clone(),
                    user_participant_identity.clone(),
                    participant,
                );
            }
@@ -2921,6 +2922,7 @@
    turn_bridge_config: TurnBridgeConfig,
    http: Client,
    sink: Arc<BotAudioOutputSink>,
    expected_participant_identity: Option<String>,
    participant: RemoteParticipant,
) -> JoinHandle<()> {
    tokio::spawn(async move {
@@ -3037,7 +3039,7 @@
                    let participant_identity = participant.identity().to_string();
                    let (was_in_speech, is_in_speech, turn) = observe_bound_participant_frame(
                        &participant_identity,
                        Some(&participant_identity),
                        expected_participant_identity.as_deref(),
                        vad,
                        &call_id,
                        &trace_id,
@@ -3241,6 +3243,10 @@
    F: FnOnce() -> std::collections::HashMap<String, String>,
{
    if !is_bound_user_participant(participant_identity, expected_participant) {
        warn!(
            "audioIngressOriginStatus" = "wrong_participant_or_track",
            "runtime helper rejected audio participant before VAD/session"
        );
        return (vad.in_speech, vad.in_speech, None);
    }
    observe_frame_and_start_session(
@@ -3273,11 +3279,14 @@
    realtime_enabled: bool,
) {
    let turn_id = format!("turn-{:04}", vad.turn_index);
    let metadata = match AudioIngressMetadata::from_participant(&read_attributes()) {
    let attributes = read_attributes();
    let origin_status = AudioIngressMetadata::origin_status(&attributes);
    let metadata = match AudioIngressMetadata::from_participant(&attributes) {
        Ok(metadata) => metadata,
        Err(reason) => {
            warn!(call_id = %call_id, trace_id = %trace_id, turn_id = %turn_id,
                reason, "runtime helper asr_realtime_metadata_rejected");
                reason, audioIngressOriginStatus = %AudioIngressMetadata::rejection_status(reason),
                "runtime helper asr_realtime_metadata_rejected");
            return;
        }
    };
@@ -3287,6 +3296,7 @@
            &metadata.client_fixture_sequence,
        ) {
            warn!(call_id = %call_id, trace_id = %trace_id, turn_id = %turn_id,
                audioIngressOriginStatus = "sequence_replayed_or_regressed",
                "runtime helper asr_realtime_metadata_sequence_rejected");
            return;
        }
@@ -3305,7 +3315,7 @@
                *last_fixture_sequence = Some(metadata.client_fixture_sequence);
            }
            info!(call_id = %call_id, trace_id = %trace_id, turn_id = %turn_id,
                "runtime helper asr_realtime_session_started");
                origin_status, "runtime helper asr_realtime_session_started");
            *upload_slot = Some(upload);
        }
        Err(error) if realtime_enabled => {
@@ -4085,11 +4095,7 @@
#[cfg(test)]
mod tests {
    use super::{
        ReplyChunkMarker, ReplyChunkMarkerState, RuntimeTurnDeviceOutput, RuntimeTurnStreamEvent,
        RuntimeTurnStreamState, RuntimeTurnStreamTimingPhase, runtime_session_nonce_hash,
        should_publish_device_output,
    };
    use super::*;
    use std::{
        collections::HashSet,
        io::{Read, Write},
@@ -4326,8 +4332,14 @@
        let second_request = request_rx
            .recv_timeout(Duration::from_secs(2))
            .expect("second session request");
        assert!(first_request.contains("\\\"clientFixtureSequence\\\":\\\"fixture-01\\\""));
        assert!(second_request.contains("\\\"clientFixtureSequence\\\":\\\"fixture-02\\\""));
        assert!(first_request.contains("\"clientFixtureSequence\":\"fixture-01\""));
        assert!(second_request.contains("\"clientFixtureSequence\":\"fixture-02\""));
        assert!(
            first_request.contains("\"audioIngressOriginStatus\":\"controlled_fixture_bound\"")
        );
        assert!(
            second_request.contains("\"audioIngressOriginStatus\":\"controlled_fixture_bound\"")
        );
        // The same production boundary rejects a wrong participant before VAD/session creation.
        assert!(!is_bound_user_participant(
            "participant-other",