cai
2026-08-12 f3237616a0bc95cf50fd790f07b72128c624b947
src/main.rs
@@ -111,9 +111,12 @@
    reject_reason: Option<&'static str>,
}
#[derive(Debug)]
#[derive(Clone, Debug)]
struct PendingControlledFixtureProbe {
    sender: ParticipantIdentity,
    call_id_hash: String,
    call_trace_id_hash: String,
    generation: u64,
    sequence: String,
    expires_at: Instant,
}
@@ -143,8 +146,9 @@
fn record_controlled_fixture_probe_event(
    stage: &'static str,
    call_id: &str,
    trace_id: &str,
    call_id_hash: &str,
    trace_id_hash: &str,
    generation: u64,
    sequence: &str,
    observed: bool,
    ack_result: Option<&'static str>,
@@ -161,9 +165,9 @@
            observed,
            ack_result,
            reject_reason = reject_reason.unwrap_or("unknown"),
            call_id_hash = %sha256_hex(call_id),
            trace_id_hash = %sha256_hex(trace_id),
            generation = CONTROLLED_FIXTURE_GENERATION,
            call_id_hash,
            trace_id_hash,
            generation,
            sequence_hash = %sha256_hex(sequence),
            "runtime helper controlled fixture probe state"
        );
@@ -173,9 +177,9 @@
            stage,
            observed,
            reject_reason,
            call_id_hash = %sha256_hex(call_id),
            trace_id_hash = %sha256_hex(trace_id),
            generation = CONTROLLED_FIXTURE_GENERATION,
            call_id_hash,
            trace_id_hash,
            generation,
            sequence_hash = %sha256_hex(sequence),
            "runtime helper controlled fixture probe state"
        );
@@ -184,9 +188,9 @@
            event = "controlled_fixture_attribute_probe",
            stage,
            observed,
            call_id_hash = %sha256_hex(call_id),
            trace_id_hash = %sha256_hex(trace_id),
            generation = CONTROLLED_FIXTURE_GENERATION,
            call_id_hash,
            trace_id_hash,
            generation,
            sequence_hash = %sha256_hex(sequence),
            "runtime helper controlled fixture probe state"
        );
@@ -195,15 +199,17 @@
fn record_controlled_fixture_attribute_decision(
    decision: Result<(), &'static str>,
    call_id: &str,
    trace_id: &str,
    call_id_hash: &str,
    trace_id_hash: &str,
    generation: u64,
    sequence: &str,
) -> (&'static str, Option<&'static str>, bool) {
    let classification = controlled_fixture_ack_classification(decision);
    record_controlled_fixture_probe_event(
        "attributes_classified",
        call_id,
        trace_id,
        call_id_hash,
        trace_id_hash,
        generation,
        sequence,
        classification.2,
        Some(classification.0),
@@ -217,8 +223,9 @@
    observed: bool,
    ack_result: &'static str,
    reject_reason: Option<&'static str>,
    call_id: &str,
    trace_id: &str,
    call_id_hash: &str,
    trace_id_hash: &str,
    generation: u64,
    sequence: &str,
    acknowledged_probe_sequences: &mut HashSet<String>,
) -> bool
@@ -228,8 +235,9 @@
    if publish.await.is_ok() {
        record_controlled_fixture_probe_event(
            "ack_publish_completed",
            call_id,
            trace_id,
            call_id_hash,
            trace_id_hash,
            generation,
            sequence,
            observed,
            Some(ack_result),
@@ -240,14 +248,33 @@
    } else {
        record_controlled_fixture_probe_event(
            "ack_publish_completed",
            call_id,
            trace_id,
            call_id_hash,
            trace_id_hash,
            generation,
            sequence,
            false,
            Some("publish_failed"),
            Some("ack_publish_failed"),
        );
        false
    }
}
fn controlled_fixture_ack_from_probe(
    probe: &PendingControlledFixtureProbe,
    observed: bool,
    reject_reason: Option<&'static str>,
) -> ControlledFixtureAttributeAck {
    ControlledFixtureAttributeAck {
        message_type: CONTROLLED_FIXTURE_ACK_TOPIC,
        protocol_version: CONTROLLED_FIXTURE_PROTOCOL_VERSION,
        call_id_hash: probe.call_id_hash.clone(),
        call_trace_id_hash: probe.call_trace_id_hash.clone(),
        generation: probe.generation,
        client_fixture_sequence: probe.sequence.clone(),
        result: if observed { "observed" } else { "rejected" },
        input_source_category: observed.then_some("controlled_fixture"),
        reject_reason,
    }
}
@@ -273,6 +300,9 @@
    }
    Some(PendingControlledFixtureProbe {
        sender: sender.clone(),
        call_id_hash: probe.call_id_hash,
        call_trace_id_hash: probe.call_trace_id_hash,
        generation: probe.generation,
        sequence: probe.client_fixture_sequence,
        expires_at: Instant::now() + CONTROLLED_FIXTURE_PROBE_TTL,
    })
@@ -1105,20 +1135,18 @@
                    track_source = %track_source,
                    "runtime helper user_track_subscribed"
                );
                let pending_sequence = pending_probe.as_ref().map(|probe| probe.sequence.clone());
                let pending_binding = pending_probe.clone();
                let participant_for_probe = participant.clone();
                let probe_result = process_controlled_fixture_probe(
                    &mut pending_probe,
                    &mut acknowledged_probe_sequences,
                    Some(&participant_for_probe),
                    &sink,
                    &call_id,
                    &trace_id,
                    user_participant_identity.as_deref(),
                )
                .await;
                let observer_started = start_observer_after_controlled_fixture_probe(
                    pending_sequence.is_some(),
                    pending_binding.is_some(),
                    probe_result,
                    || {
                        spawn_user_audio_frame_observer(
@@ -1138,7 +1166,7 @@
                        );
                    },
                );
                if let Some(sequence) = pending_sequence.as_deref() {
                if let Some(probe) = pending_binding.as_ref() {
                    let (ack_result, reject_reason, observed) = match probe_result {
                        Some(true) => (Some("observed"), None, true),
                        Some(false) => (Some("rejected"), Some("unknown"), false),
@@ -1150,9 +1178,10 @@
                        } else {
                            "audio_observer_blocked"
                        },
                        &call_id,
                        &trace_id,
                        sequence,
                        &probe.call_id_hash,
                        &probe.call_trace_id_hash,
                        probe.generation,
                        &probe.sequence,
                        observed,
                        ack_result,
                        reject_reason,
@@ -1180,8 +1209,9 @@
                    if acknowledged_probe_sequences.contains(&probe.client_fixture_sequence) {
                        record_controlled_fixture_probe_event(
                            "data_received",
                            &call_id,
                            &trace_id,
                            &probe.call_id_hash,
                            &probe.call_trace_id_hash,
                            probe.generation,
                            &probe.client_fixture_sequence,
                            false,
                            Some("rejected"),
@@ -1200,8 +1230,9 @@
                if let Some(probe) = pending_probe.as_ref() {
                    record_controlled_fixture_probe_event(
                        "data_received",
                        &call_id,
                        &trace_id,
                        &probe.call_id_hash,
                        &probe.call_trace_id_hash,
                        probe.generation,
                        &probe.sequence,
                        false,
                        None,
@@ -1213,8 +1244,6 @@
                    &mut acknowledged_probe_sequences,
                    current_user_participant.as_ref(),
                    &sink,
                    &call_id,
                    &trace_id,
                    user_participant_identity.as_deref(),
                )
                .await;
@@ -1265,8 +1294,6 @@
    acknowledged_probe_sequences: &mut HashSet<String>,
    participant: Option<&RemoteParticipant>,
    sink: &BotAudioOutputSink,
    call_id: &str,
    trace_id: &str,
    expected_participant: Option<&str>,
) -> Option<bool> {
    let Some(probe) = pending_probe.take() else {
@@ -1275,8 +1302,9 @@
    if acknowledged_probe_sequences.contains(&probe.sequence) {
        record_controlled_fixture_probe_event(
            "attributes_classified",
            call_id,
            trace_id,
            &probe.call_id_hash,
            &probe.call_trace_id_hash,
            probe.generation,
            &probe.sequence,
            false,
            Some("rejected"),
@@ -1287,8 +1315,9 @@
    if Instant::now() > probe.expires_at {
        record_controlled_fixture_probe_event(
            "attributes_classified",
            call_id,
            trace_id,
            &probe.call_id_hash,
            &probe.call_trace_id_hash,
            probe.generation,
            &probe.sequence,
            false,
            Some("timeout"),
@@ -1299,8 +1328,9 @@
    let Some(participant) = participant else {
        record_controlled_fixture_probe_event(
            "attributes_classified",
            call_id,
            trace_id,
            &probe.call_id_hash,
            &probe.call_trace_id_hash,
            probe.generation,
            &probe.sequence,
            false,
            None,
@@ -1312,8 +1342,9 @@
    if participant.identity() != probe.sender {
        record_controlled_fixture_probe_event(
            "attributes_classified",
            call_id,
            trace_id,
            &probe.call_id_hash,
            &probe.call_trace_id_hash,
            probe.generation,
            &probe.sequence,
            false,
            Some("rejected"),
@@ -1329,22 +1360,14 @@
        || participant.attributes(),
    )
    .await;
    let protocol_reject_reason = decision.err();
    let (ack_result, reject_reason, observed) =
        record_controlled_fixture_attribute_decision(decision, call_id, trace_id, &probe.sequence);
    let result = if observed { "observed" } else { "rejected" };
    let input_source_category = observed.then_some("controlled_fixture");
    let ack = ControlledFixtureAttributeAck {
        message_type: CONTROLLED_FIXTURE_ACK_TOPIC,
        protocol_version: CONTROLLED_FIXTURE_PROTOCOL_VERSION,
        call_id_hash: sha256_hex(call_id),
        call_trace_id_hash: sha256_hex(trace_id),
        generation: CONTROLLED_FIXTURE_GENERATION,
        client_fixture_sequence: probe.sequence.clone(),
        result,
        input_source_category,
        reject_reason: protocol_reject_reason,
    };
    let (ack_result, reject_reason, observed) = record_controlled_fixture_attribute_decision(
        decision,
        &probe.call_id_hash,
        &probe.call_trace_id_hash,
        probe.generation,
        &probe.sequence,
    );
    let ack = controlled_fixture_ack_from_probe(&probe, observed, reject_reason);
    let payload = match serde_json::to_vec(&ack) {
        Ok(payload) => payload,
        Err(_) => return Some(false),
@@ -1352,8 +1375,9 @@
    let local_participant = sink.room.local_participant();
    record_controlled_fixture_probe_event(
        "ack_publish_started",
        call_id,
        trace_id,
        &probe.call_id_hash,
        &probe.call_trace_id_hash,
        probe.generation,
        &probe.sequence,
        observed,
        Some(ack_result),
@@ -1371,8 +1395,9 @@
            observed,
            ack_result,
            reject_reason,
            call_id,
            trace_id,
            &probe.call_id_hash,
            &probe.call_trace_id_hash,
            probe.generation,
            &probe.sequence,
            acknowledged_probe_sequences,
        )
@@ -5253,6 +5278,9 @@
            controlled_fixture_probe(&payload, call_id, trace_id, &sender, Some("user-1"))
                .expect("valid probe");
        assert_eq!(pending.sequence, "fixture-01");
        assert_eq!(pending.call_id_hash, sha256_hex(call_id));
        assert_eq!(pending.call_trace_id_hash, sha256_hex(trace_id));
        assert_eq!(pending.generation, CONTROLLED_FIXTURE_GENERATION);
        assert!(
            controlled_fixture_probe(&payload, call_id, "other-trace", &sender, Some("user-1"),)
                .is_none()
@@ -5335,15 +5363,24 @@
        let call_id = "private-call-value";
        let trace_id = "private-trace-value";
        let sequence = "private-sequence-value";
        let call_id_hash = sha256_hex(call_id);
        let trace_id_hash = sha256_hex(trace_id);
        let mut observer_starts = 0;
        tracing::subscriber::with_default(subscriber, || {
            let decision = Err("wrong_source");
            let (ack_result, reject_reason, observed) =
                record_controlled_fixture_attribute_decision(decision, call_id, trace_id, sequence);
                record_controlled_fixture_attribute_decision(
                    decision,
                    &call_id_hash,
                    &trace_id_hash,
                    CONTROLLED_FIXTURE_GENERATION,
                    sequence,
                );
            record_controlled_fixture_probe_event(
                "audio_observer_blocked",
                call_id,
                trace_id,
                &call_id_hash,
                &trace_id_hash,
                CONTROLLED_FIXTURE_GENERATION,
                sequence,
                observed,
                Some(ack_result),
@@ -5396,6 +5433,51 @@
        );
    }
    #[test]
    fn production_probe_binding_drives_rejected_and_observed_ack_without_local_rehash() {
        let request_call_hash = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa";
        let request_trace_hash = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb";
        let probe = PendingControlledFixtureProbe {
            sender: ParticipantIdentity("user-1".to_string()),
            call_id_hash: request_call_hash.to_string(),
            call_trace_id_hash: request_trace_hash.to_string(),
            generation: CONTROLLED_FIXTURE_GENERATION,
            sequence: "fixture-01".to_string(),
            expires_at: Instant::now() + CONTROLLED_FIXTURE_PROBE_TTL,
        };
        let (ack_result, reject_reason, observed) =
            controlled_fixture_ack_classification(Err("wrong_sequence"));
        let rejected_ack = controlled_fixture_ack_from_probe(&probe, observed, reject_reason);
        let rejected_json = serde_json::to_value(&rejected_ack).unwrap();
        let mut rejected_observer_starts = 0;
        assert_eq!(ack_result, "rejected");
        assert_eq!(rejected_json["callIdHash"], request_call_hash);
        assert_eq!(rejected_json["callTraceIdHash"], request_trace_hash);
        assert_eq!(rejected_json["rejectReason"], "wrong_sequence");
        assert!(!start_observer_after_controlled_fixture_probe(
            true,
            Some(observed),
            || rejected_observer_starts += 1,
        ));
        assert_eq!(rejected_observer_starts, 0);
        let (ack_result, reject_reason, observed) = controlled_fixture_ack_classification(Ok(()));
        let observed_ack = controlled_fixture_ack_from_probe(&probe, observed, reject_reason);
        let observed_json = serde_json::to_value(&observed_ack).unwrap();
        let mut observed_observer_starts = 0;
        assert_eq!(ack_result, "observed");
        assert_eq!(observed_json["callIdHash"], request_call_hash);
        assert_eq!(observed_json["callTraceIdHash"], request_trace_hash);
        assert!(observed_json.get("rejectReason").is_none());
        assert!(start_observer_after_controlled_fixture_probe(
            true,
            Some(observed),
            || observed_observer_starts += 1,
        ));
        assert_eq!(observed_observer_starts, 1);
    }
    #[tokio::test]
    async fn production_ack_publish_failure_keeps_observer_session_and_audio_closed() {
        let mut acknowledged = HashSet::new();
@@ -5406,6 +5488,7 @@
            None,
            "call-publish-failure",
            "trace-publish-failure",
            CONTROLLED_FIXTURE_GENERATION,
            "fixture-01",
            &mut acknowledged,
        )
@@ -5438,6 +5521,7 @@
            None,
            "call-publish-success",
            "trace-publish-success",
            CONTROLLED_FIXTURE_GENERATION,
            "fixture-02",
            &mut acknowledged,
        )