fix(helper): wrap fixture probe logs as activities
| | |
| | | } |
| | | |
| | | fn record_controlled_fixture_probe_event( |
| | | runtime_call_id: &str, |
| | | runtime_trace_id: &str, |
| | | stage: &'static str, |
| | | call_id_hash: &str, |
| | | trace_id_hash: &str, |
| | |
| | | println!( |
| | | "{}", |
| | | controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | stage, |
| | | call_id_hash, |
| | | trace_id_hash, |
| | |
| | | } |
| | | |
| | | fn controlled_fixture_probe_event( |
| | | runtime_call_id: &str, |
| | | runtime_trace_id: &str, |
| | | stage: &'static str, |
| | | call_id_hash: &str, |
| | | trace_id_hash: &str, |
| | |
| | | reject_reason: Option<&'static str>, |
| | | ) -> serde_json::Value { |
| | | json!({ |
| | | "event": "controlled_fixture_attribute_probe", |
| | | "type": "cv_activity", |
| | | "callId": runtime_call_id, |
| | | "traceId": runtime_trace_id, |
| | | "turnId": null, |
| | | "eventName": "controlled_fixture_attribute_probe", |
| | | "eventWallTimeMs": current_time_millis(), |
| | | "result": "ok", |
| | | "reasonCode": null, |
| | | "retryable": null, |
| | | "extension": { |
| | | "stage": stage, |
| | | "observed": observed, |
| | | "ack_result": ack_result, |
| | |
| | | "trace_id_hash": trace_id_hash, |
| | | "generation": generation, |
| | | "sequence_hash": sha256_hex(sequence), |
| | | }, |
| | | }) |
| | | } |
| | | |
| | | fn record_controlled_fixture_attribute_decision( |
| | | decision: Result<(), &'static str>, |
| | | runtime_call_id: &str, |
| | | runtime_trace_id: &str, |
| | | call_id_hash: &str, |
| | | trace_id_hash: &str, |
| | | generation: u64, |
| | |
| | | ) -> (&'static str, Option<&'static str>, bool) { |
| | | let classification = controlled_fixture_ack_classification(decision); |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "attributes_classified", |
| | | call_id_hash, |
| | | trace_id_hash, |
| | |
| | | |
| | | async fn complete_controlled_fixture_ack_publish<F, E>( |
| | | publish: F, |
| | | runtime_call_id: &str, |
| | | runtime_trace_id: &str, |
| | | observed: bool, |
| | | ack_result: &'static str, |
| | | reject_reason: Option<&'static str>, |
| | |
| | | { |
| | | if publish.await.is_ok() { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "ack_publish_completed", |
| | | call_id_hash, |
| | | trace_id_hash, |
| | |
| | | observed |
| | | } else { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "ack_publish_completed", |
| | | call_id_hash, |
| | | trace_id_hash, |
| | |
| | | Some(&participant_for_probe), |
| | | &sink, |
| | | user_participant_identity.as_deref(), |
| | | &call_id, |
| | | &trace_id, |
| | | ) |
| | | .await; |
| | | let observer_started = start_observer_after_controlled_fixture_probe( |
| | |
| | | None => (None, Some("no_current_participant"), false), |
| | | }; |
| | | record_controlled_fixture_probe_event( |
| | | &call_id, |
| | | &trace_id, |
| | | if observer_started { |
| | | "audio_observer_allowed" |
| | | } else { |
| | |
| | | { |
| | | if acknowledged_probe_sequences.contains(&probe.client_fixture_sequence) { |
| | | record_controlled_fixture_probe_event( |
| | | &call_id, |
| | | &trace_id, |
| | | "data_received", |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | |
| | | ); |
| | | if let Some(probe) = pending_probe.as_ref() { |
| | | record_controlled_fixture_probe_event( |
| | | &call_id, |
| | | &trace_id, |
| | | "data_received", |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | |
| | | current_user_participant.as_ref(), |
| | | &sink, |
| | | user_participant_identity.as_deref(), |
| | | &call_id, |
| | | &trace_id, |
| | | ) |
| | | .await; |
| | | } |
| | |
| | | participant: Option<&RemoteParticipant>, |
| | | sink: &BotAudioOutputSink, |
| | | expected_participant: Option<&str>, |
| | | runtime_call_id: &str, |
| | | runtime_trace_id: &str, |
| | | ) -> Option<bool> { |
| | | let Some(probe) = pending_probe.take() else { |
| | | return None; |
| | | }; |
| | | if acknowledged_probe_sequences.contains(&probe.sequence) { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "attributes_classified", |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | |
| | | } |
| | | if Instant::now() > probe.expires_at { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "attributes_classified", |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | |
| | | } |
| | | let Some(participant) = participant else { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "attributes_classified", |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | |
| | | }; |
| | | if participant.identity() != probe.sender { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "attributes_classified", |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | |
| | | .await; |
| | | let (ack_result, reject_reason, observed) = record_controlled_fixture_attribute_decision( |
| | | decision, |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | |
| | | }; |
| | | let local_participant = sink.room.local_participant(); |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "ack_publish_started", |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | |
| | | Some( |
| | | complete_controlled_fixture_ack_publish( |
| | | publish, |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | observed, |
| | | ack_result, |
| | | reject_reason, |
| | |
| | | ]; |
| | | for (stage, result, reason) in stages { |
| | | let event = controlled_fixture_probe_event( |
| | | call_id, |
| | | trace_id, |
| | | stage, |
| | | &call_id_hash, |
| | | &trace_id_hash, |
| | |
| | | result, |
| | | reason, |
| | | ); |
| | | assert_eq!(event["call_id_hash"], call_id_hash); |
| | | assert_eq!(event["trace_id_hash"], trace_id_hash); |
| | | assert_eq!(event["stage"], stage); |
| | | assert_eq!(event["type"], "cv_activity"); |
| | | assert_eq!(event["eventName"], "controlled_fixture_attribute_probe"); |
| | | assert_eq!(event["extension"]["call_id_hash"], call_id_hash); |
| | | assert_eq!(event["extension"]["trace_id_hash"], trace_id_hash); |
| | | assert_eq!(event["extension"]["stage"], stage); |
| | | let output = event.to_string(); |
| | | assert!(!output.contains(call_id)); |
| | | assert!(!output.contains(trace_id)); |
| | | assert!(!output.contains(sequence)); |
| | | } |
| | | assert!(!start_observer_after_controlled_fixture_probe( |
| | |
| | | |
| | | let (ack_result, reject_reason, observed) = controlled_fixture_ack_classification(Ok(())); |
| | | let allowed = controlled_fixture_probe_event( |
| | | call_id, |
| | | trace_id, |
| | | "audio_observer_allowed", |
| | | &call_id_hash, |
| | | &trace_id_hash, |
| | |
| | | Some(ack_result), |
| | | reject_reason, |
| | | ); |
| | | assert_eq!(allowed["trace_id_hash"], trace_id_hash); |
| | | assert_eq!(allowed["extension"]["trace_id_hash"], trace_id_hash); |
| | | assert!(start_observer_after_controlled_fixture_probe( |
| | | true, |
| | | Some(observed), |
| | |
| | | let mut acknowledged = HashSet::new(); |
| | | let probe_result = complete_controlled_fixture_ack_publish( |
| | | async { Err::<(), ()>(()) }, |
| | | "runtime-call-publish-failure", |
| | | "runtime-trace-publish-failure", |
| | | true, |
| | | "observed", |
| | | None, |
| | |
| | | |
| | | let observed_result = complete_controlled_fixture_ack_publish( |
| | | async { Ok::<(), ()>(()) }, |
| | | "runtime-call-publish-success", |
| | | "runtime-trace-publish-success", |
| | | true, |
| | | "observed", |
| | | None, |
| | |
| | | continue; |
| | | }; |
| | | println!("{line}"); |
| | | if let Ok(value) = serde_json::from_str::<Value>(&line) { |
| | | if let Some(value) = project_worker_stdout_event(&call_id, &line) { |
| | | apply_worker_event(&call_id, &state, &value); |
| | | } |
| | | } |
| | | }); |
| | | } |
| | | |
| | | fn project_worker_stdout_event(call_id: &str, line: &str) -> Option<Value> { |
| | | let value = serde_json::from_str::<Value>(line).ok()?; |
| | | (value.get("type").and_then(Value::as_str) == Some("cv_activity") |
| | | && value.get("callId").and_then(Value::as_str) == Some(call_id) |
| | | && value.get("eventName").and_then(Value::as_str).is_some()) |
| | | .then_some(value) |
| | | } |
| | | |
| | | fn apply_worker_event(call_id: &str, state: &Arc<ServiceState>, value: &Value) { |
| | |
| | | |
| | | #[cfg(test)] |
| | | mod tests { |
| | | use super::{EventCallbackDispatch, deliver_event_callback_with, run_event_callback_worker}; |
| | | use super::{ |
| | | EventCallbackDispatch, deliver_event_callback_with, project_worker_stdout_event, |
| | | run_event_callback_worker, |
| | | }; |
| | | use crate::controlled_fixture_probe_event; |
| | | use anyhow::Result; |
| | | use serde_json::{Value, json}; |
| | | use std::{cell::RefCell, rc::Rc, sync::mpsc}; |
| | | |
| | | #[test] |
| | | fn controlled_fixture_activity_requires_production_stdout_envelope() { |
| | | let call_id = "runtime-call"; |
| | | let runtime_trace_id = "runtime-trace"; |
| | | let trace_hash = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; |
| | | let valid = controlled_fixture_probe_event( |
| | | call_id, |
| | | runtime_trace_id, |
| | | "ack_publish_completed", |
| | | "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", |
| | | trace_hash, |
| | | 1, |
| | | "fixture-01", |
| | | false, |
| | | Some("rejected"), |
| | | Some("wrong_source"), |
| | | ); |
| | | let projected = project_worker_stdout_event(call_id, &valid.to_string()) |
| | | .expect("valid activity must enter the production stdout projection"); |
| | | assert_eq!(projected["eventName"], "controlled_fixture_attribute_probe"); |
| | | assert_eq!(projected["extension"]["trace_id_hash"], trace_hash); |
| | | assert_eq!(projected["extension"]["reject_reason"], "wrong_source"); |
| | | |
| | | for invalid in [ |
| | | json!({ |
| | | "event": "controlled_fixture_attribute_probe", |
| | | "stage": "ack_publish_completed", |
| | | "observed": false, |
| | | "ack_result": "rejected", |
| | | "reject_reason": "wrong_source", |
| | | "call_id_hash": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", |
| | | "trace_id_hash": trace_hash, |
| | | "generation": 1, |
| | | "sequence_hash": "cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc" |
| | | }), |
| | | json!({"type": "not_activity", "callId": call_id, "eventName": "controlled_fixture_attribute_probe"}), |
| | | json!({"type": "cv_activity", "callId": call_id}), |
| | | json!({"type": "cv_activity", "callId": "wrong-call", "eventName": "controlled_fixture_attribute_probe"}), |
| | | ] { |
| | | assert!(project_worker_stdout_event(call_id, &invalid.to_string()).is_none()); |
| | | } |
| | | } |
| | | |
| | | #[test] |
| | | fn same_session_callback_retries_m7_before_terminal() { |
| | | let (sender, receiver) = mpsc::channel(); |
| | | let dispatch = |payload| EventCallbackDispatch { |