| | |
| | | const CONTROLLED_FIXTURE_GENERATION: u64 = 1; |
| | | const CONTROLLED_FIXTURE_PROBE_RECHECK_DELAY: Duration = Duration::from_millis(25); |
| | | const CONTROLLED_FIXTURE_PROBE_TTL: Duration = Duration::from_millis(250); |
| | | const CONTROLLED_FIXTURE_POST_EXPIRY_WINDOW: Duration = Duration::from_millis(2_000); |
| | | |
| | | const CONTROLLED_FIXTURE_ACK_RESULTS: [&str; 4] = |
| | | ["observed", "rejected", "timeout", "publish_failed"]; |
| | |
| | | call_trace_id_hash: String, |
| | | generation: u64, |
| | | sequence: String, |
| | | received_at: Instant, |
| | | expires_at: Instant, |
| | | } |
| | | |
| | | #[derive(Debug, PartialEq, Eq)] |
| | | struct ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: &'static str, |
| | | visibility_source: &'static str, |
| | | binding_matched: bool, |
| | | } |
| | | |
| | | #[derive(Debug, PartialEq, Eq)] |
| | | struct ControlledFixtureAckPublishOutcome { |
| | | observed: bool, |
| | | published: bool, |
| | | } |
| | | |
| | | fn sha256_hex(value: &str) -> String { |
| | |
| | | generation: u64, |
| | | sequence: &str, |
| | | acknowledged_probe_sequences: &mut HashSet<String>, |
| | | ) -> bool |
| | | ) -> ControlledFixtureAckPublishOutcome |
| | | where |
| | | F: Future<Output = Result<(), E>>, |
| | | { |
| | |
| | | reject_reason, |
| | | ); |
| | | acknowledged_probe_sequences.insert(sequence.to_string()); |
| | | observed |
| | | ControlledFixtureAckPublishOutcome { |
| | | observed, |
| | | published: true, |
| | | } |
| | | } else { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | |
| | | Some("publish_failed"), |
| | | Some("ack_publish_failed"), |
| | | ); |
| | | false |
| | | ControlledFixtureAckPublishOutcome { |
| | | observed: false, |
| | | published: false, |
| | | } |
| | | } |
| | | } |
| | | |
| | |
| | | { |
| | | return None; |
| | | } |
| | | let received_at = Instant::now(); |
| | | 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, |
| | | received_at, |
| | | expires_at: received_at + CONTROLLED_FIXTURE_PROBE_TTL, |
| | | }) |
| | | } |
| | | |
| | | fn controlled_fixture_visibility_bucket(elapsed: Duration) -> &'static str { |
| | | if elapsed <= Duration::from_millis(250) { |
| | | "lte_250ms" |
| | | } else if elapsed <= Duration::from_millis(500) { |
| | | "250_500ms" |
| | | } else if elapsed <= Duration::from_millis(1_000) { |
| | | "500_1000ms" |
| | | } else { |
| | | "1000_2000ms" |
| | | } |
| | | } |
| | | |
| | | async fn observe_controlled_fixture_post_expiry<F>( |
| | | received_at: Instant, |
| | | observation_deadline: Instant, |
| | | actual_participant: &str, |
| | | expected_participant: Option<&str>, |
| | | requested_sequence: &str, |
| | | lifecycle_active: Arc<AtomicBool>, |
| | | mut read_attributes: F, |
| | | ) -> Option<ControlledFixtureVisibilityEvidence> |
| | | where |
| | | F: FnMut() -> std::collections::HashMap<String, String>, |
| | | { |
| | | loop { |
| | | if !lifecycle_active.load(Ordering::Acquire) { |
| | | return None; |
| | | } |
| | | let now = Instant::now(); |
| | | let decision = classify_controlled_fixture_attributes( |
| | | actual_participant, |
| | | expected_participant, |
| | | &read_attributes(), |
| | | requested_sequence, |
| | | ); |
| | | match decision { |
| | | Ok(()) => { |
| | | return Some(ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: controlled_fixture_visibility_bucket( |
| | | now.saturating_duration_since(received_at), |
| | | ), |
| | | visibility_source: "participant_attributes_poll", |
| | | binding_matched: true, |
| | | }); |
| | | } |
| | | Err("missing_attributes") if now < observation_deadline => {} |
| | | Err("missing_attributes") => { |
| | | return Some(ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: "never_visible_within_observation_window", |
| | | visibility_source: "participant_attributes_poll", |
| | | binding_matched: true, |
| | | }); |
| | | } |
| | | Err(_) => return None, |
| | | } |
| | | sleep(CONTROLLED_FIXTURE_PROBE_RECHECK_DELAY).await; |
| | | } |
| | | } |
| | | |
| | | fn controlled_fixture_visibility_event( |
| | | runtime_call_id: &str, |
| | | runtime_trace_id: &str, |
| | | probe: &PendingControlledFixtureProbe, |
| | | evidence: &ControlledFixtureVisibilityEvidence, |
| | | ) -> serde_json::Value { |
| | | json!({ |
| | | "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": "post_expiry_visibility", |
| | | "first_visible_bucket": evidence.first_visible_bucket, |
| | | "visibility_source": evidence.visibility_source, |
| | | "binding_matched": evidence.binding_matched, |
| | | "call_id_hash": probe.call_id_hash, |
| | | "trace_id_hash": probe.call_trace_id_hash, |
| | | "generation": probe.generation, |
| | | "sequence_hash": sha256_hex(&probe.sequence), |
| | | "evidence_count": 1, |
| | | }, |
| | | }) |
| | | } |
| | | |
| | | fn spawn_controlled_fixture_post_expiry_observation( |
| | | probe: PendingControlledFixtureProbe, |
| | | participant: RemoteParticipant, |
| | | expected_participant: Option<String>, |
| | | lifecycle_active: Arc<AtomicBool>, |
| | | runtime_call_id: String, |
| | | runtime_trace_id: String, |
| | | ) { |
| | | tokio::spawn(async move { |
| | | let participant_identity = participant.identity().to_string(); |
| | | let evidence = observe_controlled_fixture_post_expiry( |
| | | probe.received_at, |
| | | probe.received_at + CONTROLLED_FIXTURE_POST_EXPIRY_WINDOW, |
| | | &participant_identity, |
| | | expected_participant.as_deref(), |
| | | &probe.sequence, |
| | | lifecycle_active.clone(), |
| | | || participant.attributes(), |
| | | ) |
| | | .await; |
| | | if lifecycle_active.load(Ordering::Acquire) { |
| | | if let Some(evidence) = evidence { |
| | | println!( |
| | | "{}", |
| | | controlled_fixture_visibility_event( |
| | | &runtime_call_id, |
| | | &runtime_trace_id, |
| | | &probe, |
| | | &evidence, |
| | | ) |
| | | ); |
| | | } |
| | | } |
| | | }); |
| | | } |
| | | |
| | | fn classify_controlled_fixture_attributes( |
| | |
| | | let mut current_user_participant: Option<RemoteParticipant> = None; |
| | | let mut pending_probe: Option<PendingControlledFixtureProbe> = None; |
| | | let mut acknowledged_probe_sequences = HashSet::new(); |
| | | let controlled_fixture_lifecycle_active = Arc::new(AtomicBool::new(true)); |
| | | |
| | | while let Some(event) = events.recv().await { |
| | | match event { |
| | |
| | | user_participant_identity.as_deref(), |
| | | &call_id, |
| | | &trace_id, |
| | | controlled_fixture_lifecycle_active.clone(), |
| | | ) |
| | | .await; |
| | | let observer_started = start_observer_after_controlled_fixture_probe( |
| | |
| | | user_participant_identity.as_deref(), |
| | | &call_id, |
| | | &trace_id, |
| | | controlled_fixture_lifecycle_active.clone(), |
| | | ) |
| | | .await; |
| | | } |
| | |
| | | _ => {} |
| | | } |
| | | } |
| | | controlled_fixture_lifecycle_active.store(false, Ordering::Release); |
| | | } |
| | | |
| | | async fn process_controlled_fixture_probe( |
| | |
| | | expected_participant: Option<&str>, |
| | | runtime_call_id: &str, |
| | | runtime_trace_id: &str, |
| | | lifecycle_active: Arc<AtomicBool>, |
| | | ) -> Option<bool> { |
| | | let Some(probe) = pending_probe.take() else { |
| | | return None; |
| | |
| | | payload, |
| | | topic: Some(CONTROLLED_FIXTURE_ACK_TOPIC.to_string()), |
| | | reliable: true, |
| | | destination_identities: vec![probe.sender], |
| | | destination_identities: vec![probe.sender.clone()], |
| | | }); |
| | | Some( |
| | | complete_controlled_fixture_ack_publish( |
| | | publish, |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | observed, |
| | | ack_result, |
| | | reject_reason, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | acknowledged_probe_sequences, |
| | | ) |
| | | .await, |
| | | let outcome = complete_controlled_fixture_ack_publish( |
| | | publish, |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | observed, |
| | | ack_result, |
| | | reject_reason, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | acknowledged_probe_sequences, |
| | | ) |
| | | .await; |
| | | if outcome.published && ack_result == "timeout" && reject_reason == Some("expired") { |
| | | spawn_controlled_fixture_post_expiry_observation( |
| | | probe, |
| | | participant.clone(), |
| | | expected_participant.map(str::to_string), |
| | | lifecycle_active, |
| | | runtime_call_id.to_string(), |
| | | runtime_trace_id.to_string(), |
| | | ); |
| | | } |
| | | Some(outcome.observed) |
| | | } |
| | | |
| | | async fn handle_finished_turn( |
| | |
| | | call_trace_id_hash: request_trace_hash.to_string(), |
| | | generation: CONTROLLED_FIXTURE_GENERATION, |
| | | sequence: "fixture-01".to_string(), |
| | | received_at: Instant::now(), |
| | | expires_at: Instant::now() + CONTROLLED_FIXTURE_PROBE_TTL, |
| | | }; |
| | | |
| | |
| | | let mut speaking_starts = 0; |
| | | assert!(!start_observer_after_controlled_fixture_probe( |
| | | true, |
| | | Some(probe_result), |
| | | Some(probe_result.observed), |
| | | || { |
| | | observer_starts += 1; |
| | | session_starts += 1; |
| | |
| | | speaking_starts += 1; |
| | | }, |
| | | )); |
| | | assert!(!probe_result); |
| | | assert_eq!( |
| | | probe_result, |
| | | ControlledFixtureAckPublishOutcome { |
| | | observed: false, |
| | | published: false, |
| | | } |
| | | ); |
| | | assert!(acknowledged.is_empty()); |
| | | assert_eq!(observer_starts, 0); |
| | | assert_eq!(session_starts, 0); |
| | |
| | | let mut successful_observer_starts = 0; |
| | | assert!(start_observer_after_controlled_fixture_probe( |
| | | true, |
| | | Some(observed_result), |
| | | Some(observed_result.observed), |
| | | || successful_observer_starts += 1, |
| | | )); |
| | | assert!(observed_result); |
| | | assert_eq!( |
| | | observed_result, |
| | | ControlledFixtureAckPublishOutcome { |
| | | observed: true, |
| | | published: true, |
| | | } |
| | | ); |
| | | assert!(acknowledged.contains("fixture-02")); |
| | | assert_eq!(successful_observer_starts, 1); |
| | | } |
| | |
| | | assert_eq!(observer_starts, 0); |
| | | } |
| | | |
| | | #[tokio::test] |
| | | async fn production_post_expiry_observation_records_bounded_visibility_without_second_ack() { |
| | | let started_at = Instant::now(); |
| | | let active = Arc::new(AtomicBool::new(true)); |
| | | let expected = HashMap::from([ |
| | | ( |
| | | "inputSourceCategory".to_string(), |
| | | "controlled_fixture".to_string(), |
| | | ), |
| | | ( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-01".to_string(), |
| | | ), |
| | | ]); |
| | | let mut acknowledged = HashSet::new(); |
| | | let expired_ack = complete_controlled_fixture_ack_publish( |
| | | async { Ok::<(), ()>(()) }, |
| | | "runtime-call-post-expiry", |
| | | "runtime-trace-post-expiry", |
| | | false, |
| | | "timeout", |
| | | Some("expired"), |
| | | "call-post-expiry", |
| | | "trace-post-expiry", |
| | | CONTROLLED_FIXTURE_GENERATION, |
| | | "fixture-01", |
| | | &mut acknowledged, |
| | | ) |
| | | .await; |
| | | assert_eq!( |
| | | expired_ack, |
| | | ControlledFixtureAckPublishOutcome { |
| | | observed: false, |
| | | published: true, |
| | | } |
| | | ); |
| | | let evidence = observe_controlled_fixture_post_expiry( |
| | | started_at, |
| | | started_at + CONTROLLED_FIXTURE_POST_EXPIRY_WINDOW, |
| | | "user-1", |
| | | Some("user-1"), |
| | | "fixture-01", |
| | | active, |
| | | || { |
| | | if started_at.elapsed() >= Duration::from_millis(300) { |
| | | expected.clone() |
| | | } else { |
| | | HashMap::new() |
| | | } |
| | | }, |
| | | ) |
| | | .await; |
| | | assert_eq!(acknowledged.len(), 1); |
| | | assert_eq!( |
| | | evidence, |
| | | Some(ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: "250_500ms", |
| | | visibility_source: "participant_attributes_poll", |
| | | binding_matched: true, |
| | | }) |
| | | ); |
| | | |
| | | let never_started_at = Instant::now(); |
| | | assert_eq!( |
| | | observe_controlled_fixture_post_expiry( |
| | | never_started_at, |
| | | never_started_at + Duration::from_millis(40), |
| | | "user-1", |
| | | Some("user-1"), |
| | | "fixture-01", |
| | | Arc::new(AtomicBool::new(true)), |
| | | HashMap::new, |
| | | ) |
| | | .await, |
| | | Some(ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: "never_visible_within_observation_window", |
| | | visibility_source: "participant_attributes_poll", |
| | | binding_matched: true, |
| | | }) |
| | | ); |
| | | |
| | | assert_eq!( |
| | | observe_controlled_fixture_post_expiry( |
| | | Instant::now(), |
| | | Instant::now() + Duration::from_millis(50), |
| | | "cross-call-user", |
| | | Some("user-1"), |
| | | "fixture-01", |
| | | Arc::new(AtomicBool::new(true)), |
| | | || expected.clone(), |
| | | ) |
| | | .await, |
| | | None |
| | | ); |
| | | |
| | | let inactive = Arc::new(AtomicBool::new(false)); |
| | | assert_eq!( |
| | | observe_controlled_fixture_post_expiry( |
| | | Instant::now(), |
| | | Instant::now() + Duration::from_millis(50), |
| | | "user-1", |
| | | Some("user-1"), |
| | | "fixture-01", |
| | | inactive, |
| | | HashMap::new, |
| | | ) |
| | | .await, |
| | | None |
| | | ); |
| | | assert_eq!(acknowledged.len(), 1); |
| | | |
| | | let probe = PendingControlledFixtureProbe { |
| | | sender: ParticipantIdentity("user-1".to_string()), |
| | | call_id_hash: "call-post-expiry".to_string(), |
| | | call_trace_id_hash: "trace-post-expiry".to_string(), |
| | | generation: CONTROLLED_FIXTURE_GENERATION, |
| | | sequence: "fixture-01".to_string(), |
| | | received_at: started_at, |
| | | expires_at: started_at + CONTROLLED_FIXTURE_PROBE_TTL, |
| | | }; |
| | | let event = controlled_fixture_visibility_event( |
| | | "runtime-call-post-expiry", |
| | | "runtime-trace-post-expiry", |
| | | &probe, |
| | | &evidence.unwrap(), |
| | | ); |
| | | let extension = event["extension"].as_object().unwrap(); |
| | | let mut keys = extension.keys().map(String::as_str).collect::<Vec<_>>(); |
| | | keys.sort_unstable(); |
| | | assert_eq!( |
| | | keys, |
| | | vec![ |
| | | "binding_matched", |
| | | "call_id_hash", |
| | | "evidence_count", |
| | | "first_visible_bucket", |
| | | "generation", |
| | | "sequence_hash", |
| | | "stage", |
| | | "trace_id_hash", |
| | | "visibility_source", |
| | | ] |
| | | ); |
| | | let encoded = event.to_string(); |
| | | for forbidden in [ |
| | | "\"participant\":", |
| | | "\"room\":", |
| | | "\"track\":", |
| | | "\"payload\":", |
| | | "\"audio\":", |
| | | ] { |
| | | assert!(!encoded.contains(forbidden)); |
| | | } |
| | | } |
| | | |
| | | #[test] |
| | | fn controlled_fixture_ack_payload_is_reliable_and_redacted() { |
| | | let ack = ControlledFixtureAttributeAck { |