| | |
| | | struct ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: &'static str, |
| | | visibility_source: &'static str, |
| | | visibility_result: &'static str, |
| | | binding_matched: bool, |
| | | } |
| | | |
| | |
| | | } |
| | | } |
| | | |
| | | async fn observe_controlled_fixture_post_expiry<F>( |
| | | async fn observe_controlled_fixture_post_expiry_views<F, G>( |
| | | received_at: Instant, |
| | | observation_deadline: Instant, |
| | | actual_participant: &str, |
| | | held_participant: &str, |
| | | expected_participant: Option<&str>, |
| | | requested_sequence: &str, |
| | | lifecycle_active: Arc<AtomicBool>, |
| | | mut read_attributes: F, |
| | | mut read_held_attributes: F, |
| | | mut read_current_participant: G, |
| | | ) -> Option<ControlledFixtureVisibilityEvidence> |
| | | where |
| | | F: FnMut() -> std::collections::HashMap<String, String>, |
| | | G: FnMut() -> Option<(String, 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, |
| | | let held_decision = classify_controlled_fixture_attributes( |
| | | held_participant, |
| | | expected_participant, |
| | | &read_attributes(), |
| | | &read_held_attributes(), |
| | | requested_sequence, |
| | | ); |
| | | match decision { |
| | | match held_decision { |
| | | Ok(()) => { |
| | | return Some(ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: controlled_fixture_visibility_bucket( |
| | | now.saturating_duration_since(received_at), |
| | | ), |
| | | visibility_source: "participant_attributes_poll", |
| | | visibility_result: "held_visible", |
| | | binding_matched: true, |
| | | }); |
| | | } |
| | | Err("missing_attributes") if now < observation_deadline => {} |
| | | Err("missing_attributes") => { |
| | | Err("missing_attributes") => {} |
| | | Err(_) => return None, |
| | | } |
| | | let Some((current_identity, current_attributes)) = read_current_participant() else { |
| | | return None; |
| | | }; |
| | | match classify_controlled_fixture_attributes( |
| | | ¤t_identity, |
| | | expected_participant, |
| | | ¤t_attributes, |
| | | requested_sequence, |
| | | ) { |
| | | Ok(()) => { |
| | | return Some(ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: controlled_fixture_visibility_bucket( |
| | | now.saturating_duration_since(received_at), |
| | | ), |
| | | visibility_source: "current_room_lookup", |
| | | visibility_result: "held_stale_current_visible", |
| | | binding_matched: true, |
| | | }); |
| | | } |
| | | Err("missing_attributes") => {} |
| | | Err(_) => return None, |
| | | } |
| | | if now >= observation_deadline { |
| | | return Some(ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: "never_visible_within_observation_window", |
| | | visibility_source: "participant_attributes_poll", |
| | | visibility_source: "held_and_current_room_lookup", |
| | | visibility_result: "unavailable_both", |
| | | binding_matched: true, |
| | | }); |
| | | } |
| | | Err(_) => return None, |
| | | } |
| | | sleep(CONTROLLED_FIXTURE_PROBE_RECHECK_DELAY).await; |
| | | } |
| | |
| | | "stage": "post_expiry_visibility", |
| | | "first_visible_bucket": evidence.first_visible_bucket, |
| | | "visibility_source": evidence.visibility_source, |
| | | "visibility_result": evidence.visibility_result, |
| | | "binding_matched": evidence.binding_matched, |
| | | "call_id_hash": probe.call_id_hash, |
| | | "trace_id_hash": probe.call_trace_id_hash, |
| | |
| | | fn spawn_controlled_fixture_post_expiry_observation( |
| | | probe: PendingControlledFixtureProbe, |
| | | participant: RemoteParticipant, |
| | | room: Arc<Room>, |
| | | expected_participant: Option<String>, |
| | | lifecycle_active: Arc<AtomicBool>, |
| | | runtime_call_id: String, |
| | |
| | | ) { |
| | | tokio::spawn(async move { |
| | | let participant_identity = participant.identity().to_string(); |
| | | let evidence = observe_controlled_fixture_post_expiry( |
| | | let evidence = observe_controlled_fixture_post_expiry_views( |
| | | probe.received_at, |
| | | probe.received_at + CONTROLLED_FIXTURE_POST_EXPIRY_WINDOW, |
| | | &participant_identity, |
| | |
| | | &probe.sequence, |
| | | lifecycle_active.clone(), |
| | | || participant.attributes(), |
| | | || { |
| | | room.remote_participants() |
| | | .get(&probe.sender) |
| | | .map(|current| (current.identity().to_string(), current.attributes())) |
| | | }, |
| | | ) |
| | | .await; |
| | | if lifecycle_active.load(Ordering::Acquire) { |
| | |
| | | spawn_controlled_fixture_post_expiry_observation( |
| | | probe, |
| | | participant.clone(), |
| | | sink.room.clone(), |
| | | expected_participant.map(str::to_string), |
| | | lifecycle_active, |
| | | runtime_call_id.to_string(), |
| | |
| | | published: true, |
| | | } |
| | | ); |
| | | let evidence = observe_controlled_fixture_post_expiry( |
| | | let evidence = observe_controlled_fixture_post_expiry_views( |
| | | started_at, |
| | | started_at + CONTROLLED_FIXTURE_POST_EXPIRY_WINDOW, |
| | | "user-1", |
| | |
| | | HashMap::new() |
| | | } |
| | | }, |
| | | || Some(("user-1".to_string(), HashMap::new())), |
| | | ) |
| | | .await; |
| | | assert_eq!(acknowledged.len(), 1); |
| | |
| | | Some(ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: "250_500ms", |
| | | visibility_source: "participant_attributes_poll", |
| | | visibility_result: "held_visible", |
| | | binding_matched: true, |
| | | }) |
| | | ); |
| | | |
| | | let never_started_at = Instant::now(); |
| | | assert_eq!( |
| | | observe_controlled_fixture_post_expiry( |
| | | observe_controlled_fixture_post_expiry_views( |
| | | never_started_at, |
| | | never_started_at + Duration::from_millis(40), |
| | | "user-1", |
| | |
| | | "fixture-01", |
| | | Arc::new(AtomicBool::new(true)), |
| | | HashMap::new, |
| | | || Some(("user-1".to_string(), HashMap::new())), |
| | | ) |
| | | .await, |
| | | Some(ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: "never_visible_within_observation_window", |
| | | visibility_source: "participant_attributes_poll", |
| | | visibility_source: "held_and_current_room_lookup", |
| | | visibility_result: "unavailable_both", |
| | | binding_matched: true, |
| | | }) |
| | | ); |
| | | |
| | | assert_eq!( |
| | | observe_controlled_fixture_post_expiry( |
| | | observe_controlled_fixture_post_expiry_views( |
| | | Instant::now(), |
| | | Instant::now() + Duration::from_millis(50), |
| | | "cross-call-user", |
| | |
| | | "fixture-01", |
| | | Arc::new(AtomicBool::new(true)), |
| | | || expected.clone(), |
| | | || Some(("user-1".to_string(), expected.clone())), |
| | | ) |
| | | .await, |
| | | None |
| | |
| | | |
| | | let inactive = Arc::new(AtomicBool::new(false)); |
| | | assert_eq!( |
| | | observe_controlled_fixture_post_expiry( |
| | | observe_controlled_fixture_post_expiry_views( |
| | | Instant::now(), |
| | | Instant::now() + Duration::from_millis(50), |
| | | "user-1", |
| | |
| | | "fixture-01", |
| | | inactive, |
| | | HashMap::new, |
| | | || Some(("user-1".to_string(), HashMap::new())), |
| | | ) |
| | | .await, |
| | | None |
| | |
| | | "sequence_hash", |
| | | "stage", |
| | | "trace_id_hash", |
| | | "visibility_result", |
| | | "visibility_source", |
| | | ] |
| | | ); |
| | |
| | | } |
| | | } |
| | | |
| | | #[tokio::test] |
| | | async fn production_post_expiry_observation_distinguishes_held_stale_from_current_room_view() { |
| | | let started_at = Instant::now(); |
| | | let current_attributes = 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-current-view", |
| | | "runtime-trace-current-view", |
| | | false, |
| | | "timeout", |
| | | Some("expired"), |
| | | "call-current-view", |
| | | "trace-current-view", |
| | | CONTROLLED_FIXTURE_GENERATION, |
| | | "fixture-01", |
| | | &mut acknowledged, |
| | | ) |
| | | .await; |
| | | let evidence = observe_controlled_fixture_post_expiry_views( |
| | | started_at, |
| | | started_at + Duration::from_millis(100), |
| | | "user-1", |
| | | Some("user-1"), |
| | | "fixture-01", |
| | | Arc::new(AtomicBool::new(true)), |
| | | HashMap::new, |
| | | || Some(("user-1".to_string(), current_attributes.clone())), |
| | | ) |
| | | .await; |
| | | |
| | | assert_eq!( |
| | | expired_ack, |
| | | ControlledFixtureAckPublishOutcome { |
| | | observed: false, |
| | | published: true, |
| | | } |
| | | ); |
| | | assert_eq!(acknowledged.len(), 1); |
| | | assert_eq!( |
| | | evidence, |
| | | Some(ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: "lte_250ms", |
| | | visibility_source: "current_room_lookup", |
| | | visibility_result: "held_stale_current_visible", |
| | | binding_matched: true, |
| | | }) |
| | | ); |
| | | let mut observer_starts = 0; |
| | | assert!(!start_observer_after_controlled_fixture_probe( |
| | | true, |
| | | Some(expired_ack.observed), |
| | | || observer_starts += 1, |
| | | )); |
| | | assert_eq!(observer_starts, 0); |
| | | |
| | | for current_view in [ |
| | | None, |
| | | Some(("cross-call-user".to_string(), current_attributes.clone())), |
| | | Some(( |
| | | "user-1".to_string(), |
| | | HashMap::from([ |
| | | ( |
| | | "inputSourceCategory".to_string(), |
| | | "controlled_fixture".to_string(), |
| | | ), |
| | | ( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-old".to_string(), |
| | | ), |
| | | ]), |
| | | )), |
| | | ] { |
| | | assert_eq!( |
| | | observe_controlled_fixture_post_expiry_views( |
| | | Instant::now(), |
| | | Instant::now() + Duration::from_millis(20), |
| | | "user-1", |
| | | Some("user-1"), |
| | | "fixture-01", |
| | | Arc::new(AtomicBool::new(true)), |
| | | HashMap::new, |
| | | || current_view.clone(), |
| | | ) |
| | | .await, |
| | | None |
| | | ); |
| | | } |
| | | } |
| | | |
| | | #[test] |
| | | fn controlled_fixture_ack_payload_is_reliable_and_redacted() { |
| | | let ack = ControlledFixtureAttributeAck { |