| | |
| | | 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"]; |
| | | const CONTROLLED_FIXTURE_REJECT_REASONS: [&str; 10] = [ |
| | | const CONTROLLED_FIXTURE_REJECT_REASONS: [&str; 16] = [ |
| | | "missing_attributes", |
| | | "wrong_source", |
| | | "missing_sequence", |
| | |
| | | "duplicate_or_old_sequence", |
| | | "no_current_participant", |
| | | "ack_publish_failed", |
| | | "missing_generation", |
| | | "invalid_generation", |
| | | "wrong_generation", |
| | | "invalid_language", |
| | | "incomplete_metadata", |
| | | "invalid_source_or_sequence", |
| | | "unknown", |
| | | ]; |
| | | |
| | |
| | | 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, |
| | | received_at: Instant, |
| | | expires_at: Instant, |
| | | } |
| | | |
| | | #[derive(Debug, PartialEq, Eq)] |
| | | struct ControlledFixtureVisibilityEvidence { |
| | | first_visible_bucket: &'static str, |
| | | visibility_source: &'static str, |
| | | visibility_result: &'static str, |
| | | binding_matched: bool, |
| | | } |
| | | |
| | | #[derive(Debug, PartialEq, Eq)] |
| | | struct ControlledFixtureAckPublishOutcome { |
| | | observed: bool, |
| | | published: bool, |
| | | } |
| | | |
| | | fn sha256_hex(value: &str) -> String { |
| | |
| | | } |
| | | |
| | | fn record_controlled_fixture_probe_event( |
| | | runtime_call_id: &str, |
| | | runtime_trace_id: &str, |
| | | 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>, |
| | |
| | | debug_assert!( |
| | | reject_reason.is_none_or(|value| CONTROLLED_FIXTURE_REJECT_REASONS.contains(&value)) |
| | | ); |
| | | if let Some(ack_result) = ack_result { |
| | | info!( |
| | | event = "controlled_fixture_attribute_probe", |
| | | println!( |
| | | "{}", |
| | | controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | stage, |
| | | call_id_hash, |
| | | trace_id_hash, |
| | | generation, |
| | | sequence, |
| | | 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, |
| | | sequence_hash = %sha256_hex(sequence), |
| | | "runtime helper controlled fixture probe state" |
| | | ); |
| | | } else if let Some(reject_reason) = reject_reason { |
| | | info!( |
| | | event = "controlled_fixture_attribute_probe", |
| | | stage, |
| | | observed, |
| | | reject_reason, |
| | | call_id_hash = %sha256_hex(call_id), |
| | | trace_id_hash = %sha256_hex(trace_id), |
| | | generation = CONTROLLED_FIXTURE_GENERATION, |
| | | sequence_hash = %sha256_hex(sequence), |
| | | "runtime helper controlled fixture probe state" |
| | | ); |
| | | } else { |
| | | info!( |
| | | 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, |
| | | sequence_hash = %sha256_hex(sequence), |
| | | "runtime helper controlled fixture probe state" |
| | | ); |
| | | } |
| | | ) |
| | | ); |
| | | } |
| | | |
| | | fn controlled_fixture_probe_event( |
| | | runtime_call_id: &str, |
| | | runtime_trace_id: &str, |
| | | stage: &'static str, |
| | | call_id_hash: &str, |
| | | trace_id_hash: &str, |
| | | generation: u64, |
| | | sequence: &str, |
| | | observed: bool, |
| | | ack_result: Option<&'static str>, |
| | | reject_reason: Option<&'static str>, |
| | | ) -> 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": stage, |
| | | "observed": observed, |
| | | "ack_result": ack_result, |
| | | "reject_reason": reject_reason, |
| | | "call_id_hash": call_id_hash, |
| | | "trace_id_hash": trace_id_hash, |
| | | "generation": generation, |
| | | "sequence_hash": sha256_hex(sequence), |
| | | }, |
| | | }) |
| | | } |
| | | |
| | | fn record_controlled_fixture_attribute_decision( |
| | | decision: Result<(), &'static str>, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | runtime_call_id: &str, |
| | | runtime_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( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | call_id_hash, |
| | | trace_id_hash, |
| | | generation, |
| | | sequence, |
| | | classification.2, |
| | | Some(classification.0), |
| | |
| | | |
| | | 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>, |
| | | 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 |
| | | acknowledged_probe_sequences: &mut HashSet<(u64, String)>, |
| | | ) -> ControlledFixtureAckPublishOutcome |
| | | where |
| | | F: Future<Output = Result<(), E>>, |
| | | { |
| | | if publish.await.is_ok() { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "ack_publish_completed", |
| | | call_id, |
| | | trace_id, |
| | | call_id_hash, |
| | | trace_id_hash, |
| | | generation, |
| | | sequence, |
| | | observed, |
| | | Some(ack_result), |
| | | reject_reason, |
| | | ); |
| | | acknowledged_probe_sequences.insert(sequence.to_string()); |
| | | observed |
| | | acknowledged_probe_sequences.insert((generation, sequence.to_string())); |
| | | ControlledFixtureAckPublishOutcome { |
| | | observed, |
| | | published: true, |
| | | } |
| | | } else { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "ack_publish_completed", |
| | | call_id, |
| | | trace_id, |
| | | call_id_hash, |
| | | trace_id_hash, |
| | | generation, |
| | | sequence, |
| | | false, |
| | | Some("publish_failed"), |
| | | Some("ack_publish_failed"), |
| | | ); |
| | | false |
| | | ControlledFixtureAckPublishOutcome { |
| | | observed: false, |
| | | published: 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, |
| | | } |
| | | } |
| | | |
| | |
| | | let probe: ControlledFixtureAttributeProbe = serde_json::from_slice(payload).ok()?; |
| | | if probe.message_type != CONTROLLED_FIXTURE_PROBE_TOPIC |
| | | || probe.protocol_version != CONTROLLED_FIXTURE_PROTOCOL_VERSION |
| | | || probe.generation != CONTROLLED_FIXTURE_GENERATION |
| | | || !valid_input_generation(probe.generation) |
| | | || probe.call_id_hash != sha256_hex(call_id) |
| | | || probe.call_trace_id_hash != sha256_hex(trace_id) |
| | | || probe.client_fixture_sequence.trim().is_empty() |
| | | { |
| | | 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_probe_binding_decision( |
| | | probe_sender: &str, |
| | | current_audio_participant: Option<&str>, |
| | | expected_participant: Option<&str>, |
| | | ) -> Result<(), &'static str> { |
| | | if !is_bound_user_participant(probe_sender, expected_participant) { |
| | | return Err("wrong_participant"); |
| | | } |
| | | let Some(current_audio_participant) = current_audio_participant else { |
| | | return Err("no_current_participant"); |
| | | }; |
| | | if current_audio_participant != probe_sender |
| | | || !is_bound_user_participant(current_audio_participant, expected_participant) |
| | | { |
| | | return Err("wrong_participant"); |
| | | } |
| | | Ok(()) |
| | | } |
| | | |
| | | 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_views<F, G>( |
| | | received_at: Instant, |
| | | observation_deadline: Instant, |
| | | held_participant: &str, |
| | | expected_participant: Option<&str>, |
| | | requested_sequence: &str, |
| | | lifecycle_active: Arc<AtomicBool>, |
| | | 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 held_decision = classify_controlled_fixture_attributes( |
| | | held_participant, |
| | | expected_participant, |
| | | &read_held_attributes(), |
| | | requested_sequence, |
| | | ); |
| | | 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") => {} |
| | | 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: "held_and_current_room_lookup", |
| | | visibility_result: "unavailable_both", |
| | | binding_matched: true, |
| | | }); |
| | | } |
| | | 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, |
| | | "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, |
| | | "generation": probe.generation, |
| | | "sequence_hash": sha256_hex(&probe.sequence), |
| | | "evidence_count": 1, |
| | | }, |
| | | }) |
| | | } |
| | | |
| | | 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, |
| | | runtime_trace_id: String, |
| | | ) { |
| | | tokio::spawn(async move { |
| | | let participant_identity = participant.identity().to_string(); |
| | | let evidence = observe_controlled_fixture_post_expiry_views( |
| | | probe.received_at, |
| | | probe.received_at + CONTROLLED_FIXTURE_POST_EXPIRY_WINDOW, |
| | | &participant_identity, |
| | | expected_participant.as_deref(), |
| | | &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) { |
| | | if let Some(evidence) = evidence { |
| | | println!( |
| | | "{}", |
| | | controlled_fixture_visibility_event( |
| | | &runtime_call_id, |
| | | &runtime_trace_id, |
| | | &probe, |
| | | &evidence, |
| | | ) |
| | | ); |
| | | } |
| | | } |
| | | }); |
| | | } |
| | | |
| | | fn classify_controlled_fixture_attributes( |
| | |
| | | return Err("wrong_sequence"); |
| | | } |
| | | Ok(()) |
| | | } |
| | | |
| | | fn controlled_fixture_probe_attribute_decision( |
| | | actual_participant: &str, |
| | | expected_participant: Option<&str>, |
| | | attributes: &std::collections::HashMap<String, String>, |
| | | probe: &PendingControlledFixtureProbe, |
| | | ) -> Result<(), &'static str> { |
| | | if !is_bound_user_participant(actual_participant, expected_participant) { |
| | | return Err("wrong_participant"); |
| | | } |
| | | let metadata = match AudioIngressMetadata::from_participant(attributes) { |
| | | Ok(Some(metadata)) => metadata, |
| | | Ok(None) => return Err("missing_attributes"), |
| | | Err(reason) => return Err(reason), |
| | | }; |
| | | if metadata.client_fixture_sequence != probe.sequence { |
| | | return Err("wrong_sequence"); |
| | | } |
| | | if metadata.input_generation != probe.generation { |
| | | return Err("wrong_generation"); |
| | | } |
| | | Ok(()) |
| | | } |
| | | |
| | | async fn observe_controlled_fixture_probe_attributes<F>( |
| | | expires_at: Instant, |
| | | actual_participant: &str, |
| | | expected_participant: Option<&str>, |
| | | probe: &PendingControlledFixtureProbe, |
| | | mut read_attributes: F, |
| | | ) -> Result<(), &'static str> |
| | | where |
| | | F: FnMut() -> std::collections::HashMap<String, String>, |
| | | { |
| | | loop { |
| | | if Instant::now() > expires_at { |
| | | return Err("timeout"); |
| | | } |
| | | match controlled_fixture_probe_attribute_decision( |
| | | actual_participant, |
| | | expected_participant, |
| | | &read_attributes(), |
| | | probe, |
| | | ) { |
| | | Ok(()) => return Ok(()), |
| | | Err("missing_attributes") | Err("missing_generation") => { |
| | | sleep(CONTROLLED_FIXTURE_PROBE_RECHECK_DELAY).await; |
| | | } |
| | | Err(reason) => return Err(reason), |
| | | } |
| | | } |
| | | } |
| | | |
| | | async fn observe_controlled_fixture_attributes<F>( |
| | |
| | | expected.is_none_or(|value| identity == value) |
| | | } |
| | | |
| | | fn valid_input_generation(generation: u64) -> bool { |
| | | generation > 0 |
| | | } |
| | | |
| | | async fn observe_user_audio_events( |
| | | mut events: UnboundedReceiver<RoomEvent>, |
| | | call_id: String, |
| | |
| | | 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 { |
| | |
| | | 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, |
| | | user_participant_identity.as_deref(), |
| | | &call_id, |
| | | &trace_id, |
| | | user_participant_identity.as_deref(), |
| | | controlled_fixture_lifecycle_active.clone(), |
| | | ) |
| | | .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( |
| | |
| | | ); |
| | | }, |
| | | ); |
| | | 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), |
| | | None => (None, Some("no_current_participant"), false), |
| | | }; |
| | | record_controlled_fixture_probe_event( |
| | | &call_id, |
| | | &trace_id, |
| | | if observer_started { |
| | | "audio_observer_allowed" |
| | | } 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, |
| | |
| | | } |
| | | current_user_participant = Some(participant_for_probe); |
| | | } |
| | | RoomEvent::TrackUnsubscribed { |
| | | track: RemoteTrack::Audio(_), |
| | | publication: _, |
| | | participant, |
| | | } => { |
| | | if current_user_participant |
| | | .as_ref() |
| | | .is_some_and(|current| current.identity() == participant.identity()) |
| | | { |
| | | current_user_participant = None; |
| | | } |
| | | } |
| | | RoomEvent::DataReceived { |
| | | payload, |
| | | topic: Some(topic), |
| | |
| | | if let Ok(probe) = |
| | | serde_json::from_slice::<ControlledFixtureAttributeProbe>(&payload) |
| | | { |
| | | if acknowledged_probe_sequences.contains(&probe.client_fixture_sequence) { |
| | | if acknowledged_probe_sequences |
| | | .contains(&(probe.generation, probe.client_fixture_sequence.clone())) |
| | | { |
| | | record_controlled_fixture_probe_event( |
| | | "data_received", |
| | | &call_id, |
| | | &trace_id, |
| | | "data_received", |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.client_fixture_sequence, |
| | | false, |
| | | Some("rejected"), |
| | |
| | | ); |
| | | if let Some(probe) = pending_probe.as_ref() { |
| | | record_controlled_fixture_probe_event( |
| | | "data_received", |
| | | &call_id, |
| | | &trace_id, |
| | | "data_received", |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | false, |
| | | None, |
| | |
| | | &mut acknowledged_probe_sequences, |
| | | current_user_participant.as_ref(), |
| | | &sink, |
| | | user_participant_identity.as_deref(), |
| | | &call_id, |
| | | &trace_id, |
| | | user_participant_identity.as_deref(), |
| | | controlled_fixture_lifecycle_active.clone(), |
| | | ) |
| | | .await; |
| | | } |
| | |
| | | _ => {} |
| | | } |
| | | } |
| | | controlled_fixture_lifecycle_active.store(false, Ordering::Release); |
| | | } |
| | | |
| | | async fn process_controlled_fixture_probe( |
| | | pending_probe: &mut Option<PendingControlledFixtureProbe>, |
| | | acknowledged_probe_sequences: &mut HashSet<String>, |
| | | acknowledged_probe_sequences: &mut HashSet<(u64, String)>, |
| | | participant: Option<&RemoteParticipant>, |
| | | sink: &BotAudioOutputSink, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | 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; |
| | | }; |
| | | if acknowledged_probe_sequences.contains(&probe.sequence) { |
| | | if acknowledged_probe_sequences.contains(&(probe.generation, probe.sequence.clone())) { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | false, |
| | | Some("rejected"), |
| | |
| | | } |
| | | if Instant::now() > probe.expires_at { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | false, |
| | | Some("timeout"), |
| | |
| | | } |
| | | let Some(participant) = participant else { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | false, |
| | | None, |
| | |
| | | }; |
| | | if participant.identity() != probe.sender { |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | false, |
| | | Some("rejected"), |
| | |
| | | ); |
| | | return Some(false); |
| | | } |
| | | let decision = observe_controlled_fixture_attributes( |
| | | let decision = observe_controlled_fixture_probe_attributes( |
| | | probe.expires_at, |
| | | &participant.identity().to_string(), |
| | | participant.identity().as_str(), |
| | | expected_participant, |
| | | &probe.sequence, |
| | | &probe, |
| | | || 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, |
| | | }; |
| | | if let Err(reason) = decision { |
| | | record_controlled_fixture_attribute_decision( |
| | | Err(reason), |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | ); |
| | | return Some(false); |
| | | } |
| | | let (ack_result, reject_reason, observed) = record_controlled_fixture_attribute_decision( |
| | | Ok(()), |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | &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), |
| | | }; |
| | | let local_participant = sink.room.local_participant(); |
| | | record_controlled_fixture_probe_event( |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | "ack_publish_started", |
| | | call_id, |
| | | trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | observed, |
| | | Some(ack_result), |
| | |
| | | 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, |
| | | observed, |
| | | ack_result, |
| | | reject_reason, |
| | | call_id, |
| | | trace_id, |
| | | &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(), |
| | | sink.room.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( |
| | |
| | | }; |
| | | let mut realtime_asr_upload: Option<RealtimeAsrUpload> = None; |
| | | let mut last_fixture_sequence: Option<String> = None; |
| | | let mut last_fixture_generation: Option<u64> = None; |
| | | |
| | | while let Some(drained) = frame_rx.recv().await { |
| | | let frame = drained.frame; |
| | |
| | | || participant.attributes(), |
| | | &mut realtime_asr_upload, |
| | | &mut last_fixture_sequence, |
| | | &mut last_fixture_generation, |
| | | turn_bridge_config.asr_realtime_enabled, |
| | | ); |
| | | |
| | |
| | | read_attributes: F, |
| | | upload_slot: &mut Option<RealtimeAsrUpload>, |
| | | last_fixture_sequence: &mut Option<String>, |
| | | last_fixture_generation: &mut Option<u64>, |
| | | realtime_enabled: bool, |
| | | ) -> (bool, bool, Option<FinishedSpeechTurn>) |
| | | where |
| | |
| | | read_attributes, |
| | | upload_slot, |
| | | last_fixture_sequence, |
| | | last_fixture_generation, |
| | | realtime_enabled, |
| | | ); |
| | | } |
| | |
| | | read_attributes: F, |
| | | upload_slot: &mut Option<RealtimeAsrUpload>, |
| | | last_fixture_sequence: &mut Option<String>, |
| | | last_fixture_generation: &mut Option<u64>, |
| | | realtime_enabled: bool, |
| | | ) -> (bool, bool, Option<FinishedSpeechTurn>) |
| | | where |
| | |
| | | read_attributes, |
| | | upload_slot, |
| | | last_fixture_sequence, |
| | | last_fixture_generation, |
| | | realtime_enabled, |
| | | ) |
| | | } |
| | |
| | | read_attributes: impl FnOnce() -> std::collections::HashMap<String, String>, |
| | | upload_slot: &mut Option<RealtimeAsrUpload>, |
| | | last_fixture_sequence: &mut Option<String>, |
| | | last_fixture_generation: &mut Option<u64>, |
| | | realtime_enabled: bool, |
| | | ) { |
| | | let turn_id = format!("turn-{:04}", vad.turn_index); |
| | |
| | | } |
| | | }; |
| | | if let Some(metadata) = metadata.as_ref() { |
| | | if !fixture_sequence_is_new( |
| | | if !fixture_binding_is_new( |
| | | last_fixture_sequence.as_deref(), |
| | | *last_fixture_generation, |
| | | &metadata.client_fixture_sequence, |
| | | metadata.input_generation, |
| | | ) { |
| | | warn!(call_id = %call_id, trace_id = %trace_id, turn_id = %turn_id, |
| | | audioIngressOriginStatus = "sequence_replayed_or_regressed", |
| | |
| | | Ok(upload) => { |
| | | if let Some(metadata) = metadata { |
| | | *last_fixture_sequence = Some(metadata.client_fixture_sequence); |
| | | *last_fixture_generation = Some(metadata.input_generation); |
| | | } |
| | | info!(call_id = %call_id, trace_id = %trace_id, turn_id = %turn_id, |
| | | origin_status, "runtime helper asr_realtime_session_started"); |
| | |
| | | match (previous_number, current_number) { |
| | | (Some(previous), Some(current)) => current > previous, |
| | | _ => previous != current, |
| | | } |
| | | } |
| | | |
| | | fn fixture_binding_is_new( |
| | | previous_sequence: Option<&str>, |
| | | previous_generation: Option<u64>, |
| | | current_sequence: &str, |
| | | current_generation: u64, |
| | | ) -> bool { |
| | | if !valid_input_generation(current_generation) { |
| | | return false; |
| | | } |
| | | match previous_generation { |
| | | Some(previous) if current_generation < previous => false, |
| | | Some(previous) if current_generation > previous => true, |
| | | Some(_) => fixture_sequence_is_new(previous_sequence, current_sequence), |
| | | None => true, |
| | | } |
| | | } |
| | | |
| | |
| | | io::{Read, Write}, |
| | | net::TcpListener, |
| | | sync::{ |
| | | Arc, Mutex, |
| | | Arc, |
| | | atomic::{AtomicUsize, Ordering}, |
| | | mpsc, |
| | | }, |
| | |
| | | |
| | | #[derive(Debug)] |
| | | enum PreAudioOrderEvent { |
| | | DataReceived { |
| | | sender: String, |
| | | sequence: String, |
| | | }, |
| | | TrackSubscribed { |
| | | participant: String, |
| | | attributes: HashMap<String, String>, |
| | | }, |
| | | } |
| | | |
| | | #[derive(Clone, Default)] |
| | | struct CapturedLogs(Arc<Mutex<Vec<u8>>>); |
| | | |
| | | struct CapturedLogWriter(Arc<Mutex<Vec<u8>>>); |
| | | |
| | | impl Write for CapturedLogWriter { |
| | | fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> { |
| | | self.0.lock().unwrap().extend_from_slice(bytes); |
| | | Ok(bytes.len()) |
| | | } |
| | | |
| | | fn flush(&mut self) -> std::io::Result<()> { |
| | | Ok(()) |
| | | } |
| | | } |
| | | |
| | | impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for CapturedLogs { |
| | | type Writer = CapturedLogWriter; |
| | | |
| | | fn make_writer(&'a self) -> Self::Writer { |
| | | CapturedLogWriter(self.0.clone()) |
| | | } |
| | | DataReceived { sender: String, sequence: String }, |
| | | TrackSubscribed { participant: String }, |
| | | } |
| | | |
| | | fn drive_pre_audio_order_test_seam(events: &[PreAudioOrderEvent]) -> Vec<&'static str> { |
| | | let mut pending_sequence = None; |
| | | let call_id = "production-order-call"; |
| | | let trace_id = "production-order-trace"; |
| | | let mut pending_probe = None; |
| | | let mut effects = Vec::new(); |
| | | for event in events { |
| | | match event { |
| | | PreAudioOrderEvent::DataReceived { sender, sequence } if sender == "user-1" => { |
| | | pending_sequence = Some(sequence.as_str()); |
| | | PreAudioOrderEvent::DataReceived { sender, sequence } => { |
| | | let payload = serde_json::to_vec(&json!({ |
| | | "type": CONTROLLED_FIXTURE_PROBE_TOPIC, |
| | | "protocolVersion": CONTROLLED_FIXTURE_PROTOCOL_VERSION, |
| | | "callIdHash": sha256_hex(call_id), |
| | | "callTraceIdHash": sha256_hex(trace_id), |
| | | "generation": CONTROLLED_FIXTURE_GENERATION, |
| | | "clientFixtureSequence": sequence, |
| | | })) |
| | | .expect("production probe payload"); |
| | | pending_probe = controlled_fixture_probe( |
| | | &payload, |
| | | call_id, |
| | | trace_id, |
| | | &ParticipantIdentity(sender.clone()), |
| | | Some("user-1"), |
| | | ); |
| | | } |
| | | PreAudioOrderEvent::TrackSubscribed { |
| | | participant, |
| | | attributes, |
| | | } => { |
| | | let pending = pending_sequence.is_some(); |
| | | let probe_result = pending_sequence.map(|sequence| { |
| | | classify_controlled_fixture_attributes( |
| | | participant, |
| | | PreAudioOrderEvent::TrackSubscribed { participant } => { |
| | | let pending = pending_probe.is_some(); |
| | | let probe_result = pending_probe.as_ref().map(|probe| { |
| | | controlled_fixture_probe_binding_decision( |
| | | probe.sender.as_str(), |
| | | Some(participant), |
| | | Some("user-1"), |
| | | attributes, |
| | | sequence, |
| | | ) |
| | | .is_ok() |
| | | }); |
| | |
| | | if observer_started { |
| | | effects.push("observer_started"); |
| | | } |
| | | pending_sequence = None; |
| | | pending_probe = None; |
| | | } |
| | | PreAudioOrderEvent::DataReceived { .. } => {} |
| | | } |
| | | } |
| | | effects |
| | |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-01".to_string(), |
| | | ), |
| | | ("inputGeneration".to_string(), "1".to_string()), |
| | | ("language".to_string(), "ja-JP".to_string()), |
| | | ]); |
| | | let mut starts = Vec::new(); |
| | | for (session_index, sequence) in [(1, "fixture-01"), (2, "fixture-02")] { |
| | |
| | | let is_in_speech = vad.in_speech; |
| | | assert!(!was_in_speech && is_in_speech); |
| | | attributes.insert("clientFixtureSequence".to_string(), sequence.to_string()); |
| | | attributes.insert("language".to_string(), "ja-JP".to_string()); |
| | | let metadata = AudioIngressMetadata::from_participant(&attributes) |
| | | .expect("valid participant attributes") |
| | | .expect("controlled fixture metadata"); |
| | |
| | | let session_json: serde_json::Value = |
| | | serde_json::from_slice(&session_line).expect("session start json"); |
| | | assert_eq!(sequence, session_json["clientFixtureSequence"]); |
| | | assert_eq!("ja-JP", session_json["language"]); |
| | | starts.push(metadata.client_fixture_sequence); |
| | | vad.reset_current_turn(); |
| | | } |
| | |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-01".to_string(), |
| | | ), |
| | | ("inputGeneration".to_string(), "1".to_string()), |
| | | ("language".to_string(), "ja-JP".to_string()), |
| | | ]); |
| | | let mut upload = None; |
| | | let mut last_fixture_sequence = None; |
| | | let mut last_fixture_generation = None; |
| | | let config = RealtimeAsrConfig { |
| | | enabled: true, |
| | | url: Some(format!("http://{address}/runtime/asr/realtime")), |
| | |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | &mut last_fixture_sequence, |
| | | &mut last_fixture_generation, |
| | | true, |
| | | ); |
| | | assert!(!was && is && turn.is_none()); |
| | |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-02".to_string(), |
| | | ); |
| | | attrs.insert("inputGeneration".to_string(), "2".to_string()); |
| | | attrs.insert("language".to_string(), "zh-CN".to_string()); |
| | | let (was, is, turn) = observe_bound_participant_frame( |
| | | "participant-user", |
| | | Some("participant-user"), |
| | |
| | | 2_000, |
| | | &frame, |
| | | Client::new(), |
| | | config, |
| | | config.clone(), |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | &mut last_fixture_sequence, |
| | | &mut last_fixture_generation, |
| | | true, |
| | | ); |
| | | assert!(!was && is && turn.is_none()); |
| | |
| | | .expect("second session request"); |
| | | assert!(first_request.contains("\"clientFixtureSequence\":\"fixture-01\"")); |
| | | assert!(second_request.contains("\"clientFixtureSequence\":\"fixture-02\"")); |
| | | assert!(first_request.contains("\"callId\":\"call-001\"")); |
| | | assert!(first_request.contains("\"traceId\":\"trace-001\"")); |
| | | assert!(second_request.contains("\"callId\":\"call-001\"")); |
| | | assert!(second_request.contains("\"traceId\":\"trace-001\"")); |
| | | assert!(first_request.contains("\"inputSourceCategory\":\"controlled_fixture\"")); |
| | | assert!(second_request.contains("\"inputSourceCategory\":\"controlled_fixture\"")); |
| | | assert!(first_request.contains("\"inputGeneration\":1")); |
| | | assert!(second_request.contains("\"inputGeneration\":2")); |
| | | assert!(first_request.contains("\"language\":\"ja-JP\"")); |
| | | assert!(second_request.contains("\"language\":\"zh-CN\"")); |
| | | assert!( |
| | | first_request.contains("\"audioIngressOriginStatus\":\"controlled_fixture_bound\"") |
| | | ); |
| | |
| | | attrs.insert( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-03".to_string(), |
| | | ); |
| | | attrs.insert("inputGeneration".to_string(), "1".to_string()); |
| | | let (_, is_old_generation, old_generation_turn) = observe_bound_participant_frame( |
| | | "participant-user", |
| | | Some("participant-user"), |
| | | &mut vad, |
| | | "call-001", |
| | | "trace-001", |
| | | "participant-user", |
| | | "track-001", |
| | | 3, |
| | | 3_000, |
| | | &frame, |
| | | Client::new(), |
| | | config.clone(), |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | &mut last_fixture_sequence, |
| | | &mut last_fixture_generation, |
| | | true, |
| | | ); |
| | | assert!(is_old_generation && old_generation_turn.is_none() && upload.is_none()); |
| | | vad.reset_current_turn(); |
| | | attrs.insert( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-04".to_string(), |
| | | ); |
| | | let (_, is_wrong, wrong_turn) = observe_bound_participant_frame( |
| | | "participant-other", |
| | |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | &mut last_fixture_sequence, |
| | | &mut last_fixture_generation, |
| | | true, |
| | | ); |
| | | assert!(!is_wrong && wrong_turn.is_none() && upload.is_none()); |
| | |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | &mut last_fixture_sequence, |
| | | &mut last_fixture_generation, |
| | | true, |
| | | ); |
| | | assert!(upload.is_none()); |
| | |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | &mut last_fixture_sequence, |
| | | &mut last_fixture_generation, |
| | | true, |
| | | ); |
| | | assert!(upload.is_none()); |
| | |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | &mut last_fixture_sequence, |
| | | &mut last_fixture_generation, |
| | | true, |
| | | ); |
| | | assert!(missing_sequence_turn.is_none() && upload.is_none()); |
| | |
| | | || attrs, |
| | | &mut invalid_upload, |
| | | &mut last_fixture_sequence, |
| | | &mut last_fixture_generation, |
| | | true, |
| | | ); |
| | | assert!(invalid_turn.is_none()); |
| | |
| | | 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() |
| | |
| | | } |
| | | |
| | | #[test] |
| | | fn controlled_fixture_probe_failure_log_is_fixed_redacted_and_has_no_audio_effect() { |
| | | let logs = CapturedLogs::default(); |
| | | let subscriber = tracing_subscriber::fmt() |
| | | .without_time() |
| | | .with_ansi(false) |
| | | .with_writer(logs.clone()) |
| | | .finish(); |
| | | fn controlled_fixture_probe_runtime_projection_binds_request_hashes_and_audio_gate() { |
| | | 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_probe_event( |
| | | "audio_observer_blocked", |
| | | let (ack_result, reject_reason, observed) = |
| | | controlled_fixture_ack_classification(Err("wrong_source")); |
| | | let stages = [ |
| | | ("data_received", None, None), |
| | | ("attributes_classified", Some(ack_result), reject_reason), |
| | | ("ack_publish_started", Some(ack_result), reject_reason), |
| | | ("ack_publish_completed", Some(ack_result), reject_reason), |
| | | ("audio_observer_blocked", Some(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, |
| | | CONTROLLED_FIXTURE_GENERATION, |
| | | sequence, |
| | | observed, |
| | | Some(ack_result), |
| | | reject_reason, |
| | | result, |
| | | reason, |
| | | ); |
| | | assert!(!start_observer_after_controlled_fixture_probe( |
| | | true, |
| | | Some(observed), |
| | | || observer_starts += 1, |
| | | )); |
| | | }); |
| | | let output = String::from_utf8(logs.0.lock().unwrap().clone()).unwrap(); |
| | | assert!(output.contains("ack_result=\"rejected\"")); |
| | | assert!(output.contains("reject_reason=\"wrong_source\"")); |
| | | assert!(output.contains("observed=false")); |
| | | assert!(output.contains(&sha256_hex(call_id))); |
| | | assert!(output.contains(&sha256_hex(trace_id))); |
| | | assert!(output.contains(&sha256_hex(sequence))); |
| | | assert!(!output.contains(call_id)); |
| | | assert!(!output.contains(trace_id)); |
| | | assert!(!output.contains(sequence)); |
| | | 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(sequence)); |
| | | } |
| | | assert!(!start_observer_after_controlled_fixture_probe( |
| | | true, |
| | | Some(observed), |
| | | || observer_starts += 1, |
| | | )); |
| | | assert_eq!(observer_starts, 0); |
| | | |
| | | 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, |
| | | CONTROLLED_FIXTURE_GENERATION, |
| | | sequence, |
| | | observed, |
| | | Some(ack_result), |
| | | reject_reason, |
| | | ); |
| | | assert_eq!(allowed["extension"]["trace_id_hash"], trace_id_hash); |
| | | assert!(start_observer_after_controlled_fixture_probe( |
| | | true, |
| | | Some(observed), |
| | | || observer_starts += 1, |
| | | )); |
| | | assert_eq!(observer_starts, 1); |
| | | } |
| | | |
| | | #[test] |
| | |
| | | ); |
| | | } |
| | | |
| | | #[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(), |
| | | received_at: Instant::now(), |
| | | 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(); |
| | | let probe_result = complete_controlled_fixture_ack_publish( |
| | | async { Err::<(), ()>(()) }, |
| | | "runtime-call-publish-failure", |
| | | "runtime-trace-publish-failure", |
| | | true, |
| | | "observed", |
| | | None, |
| | | "call-publish-failure", |
| | | "trace-publish-failure", |
| | | CONTROLLED_FIXTURE_GENERATION, |
| | | "fixture-01", |
| | | &mut acknowledged, |
| | | ) |
| | |
| | | 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 observed_result = complete_controlled_fixture_ack_publish( |
| | | async { Ok::<(), ()>(()) }, |
| | | "runtime-call-publish-success", |
| | | "runtime-trace-publish-success", |
| | | true, |
| | | "observed", |
| | | None, |
| | | "call-publish-success", |
| | | "trace-publish-success", |
| | | CONTROLLED_FIXTURE_GENERATION, |
| | | "fixture-02", |
| | | &mut acknowledged, |
| | | ) |
| | |
| | | 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!(acknowledged.contains("fixture-02")); |
| | | assert_eq!( |
| | | observed_result, |
| | | ControlledFixtureAckPublishOutcome { |
| | | observed: true, |
| | | published: true, |
| | | } |
| | | ); |
| | | assert!(acknowledged.contains(&(CONTROLLED_FIXTURE_GENERATION, "fixture-02".to_string()))); |
| | | assert_eq!(successful_observer_starts, 1); |
| | | } |
| | | |
| | |
| | | } |
| | | |
| | | #[tokio::test] |
| | | async fn production_probe_waits_for_generation_attribute_before_binding() { |
| | | let received_at = Instant::now(); |
| | | let probe = PendingControlledFixtureProbe { |
| | | sender: ParticipantIdentity("user-1".to_string()), |
| | | call_id_hash: "call-hash".to_string(), |
| | | call_trace_id_hash: "trace-hash".to_string(), |
| | | generation: 2, |
| | | sequence: "fixture-02".to_string(), |
| | | received_at, |
| | | expires_at: received_at + CONTROLLED_FIXTURE_PROBE_TTL, |
| | | }; |
| | | let expected = HashMap::from([ |
| | | ( |
| | | "inputSourceCategory".to_string(), |
| | | "controlled_fixture".to_string(), |
| | | ), |
| | | ( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-02".to_string(), |
| | | ), |
| | | ("inputGeneration".to_string(), "2".to_string()), |
| | | ]); |
| | | let mut reads = 0; |
| | | let result = observe_controlled_fixture_probe_attributes( |
| | | probe.expires_at, |
| | | "user-1", |
| | | Some("user-1"), |
| | | &probe, |
| | | || { |
| | | reads += 1; |
| | | if reads == 1 { |
| | | HashMap::new() |
| | | } else { |
| | | expected.clone() |
| | | } |
| | | }, |
| | | ) |
| | | .await; |
| | | assert_eq!(result, Ok(())); |
| | | assert_eq!(reads, 2); |
| | | } |
| | | |
| | | #[tokio::test] |
| | | async fn production_attribute_observation_rejects_wrong_sequence_without_audio_effect() { |
| | | let attributes = HashMap::from([ |
| | | ( |
| | |
| | | || 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_views( |
| | | 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() |
| | | } |
| | | }, |
| | | || Some(("user-1".to_string(), HashMap::new())), |
| | | ) |
| | | .await; |
| | | assert_eq!(acknowledged.len(), 1); |
| | | assert_eq!( |
| | | evidence, |
| | | 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_views( |
| | | never_started_at, |
| | | never_started_at + Duration::from_millis(40), |
| | | "user-1", |
| | | Some("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: "held_and_current_room_lookup", |
| | | visibility_result: "unavailable_both", |
| | | binding_matched: true, |
| | | }) |
| | | ); |
| | | |
| | | assert_eq!( |
| | | observe_controlled_fixture_post_expiry_views( |
| | | Instant::now(), |
| | | Instant::now() + Duration::from_millis(50), |
| | | "cross-call-user", |
| | | Some("user-1"), |
| | | "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_views( |
| | | Instant::now(), |
| | | Instant::now() + Duration::from_millis(50), |
| | | "user-1", |
| | | Some("user-1"), |
| | | "fixture-01", |
| | | inactive, |
| | | HashMap::new, |
| | | || Some(("user-1".to_string(), 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_result", |
| | | "visibility_source", |
| | | ] |
| | | ); |
| | | let encoded = event.to_string(); |
| | | for forbidden in [ |
| | | "\"participant\":", |
| | | "\"room\":", |
| | | "\"track\":", |
| | | "\"payload\":", |
| | | "\"audio\":", |
| | | ] { |
| | | assert!(!encoded.contains(forbidden)); |
| | | } |
| | | } |
| | | |
| | | #[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] |
| | |
| | | |
| | | #[test] |
| | | fn production_event_order_probe_then_track_publishes_ack_before_observer() { |
| | | let attributes = HashMap::from([ |
| | | ( |
| | | "inputSourceCategory".to_string(), |
| | | "controlled_fixture".to_string(), |
| | | ), |
| | | ( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-01".to_string(), |
| | | ), |
| | | ]); |
| | | let effects = drive_pre_audio_order_test_seam(&[ |
| | | PreAudioOrderEvent::DataReceived { |
| | | sender: "user-1".to_string(), |
| | | sequence: "fixture-01".to_string(), |
| | | sequence: "1".to_string(), |
| | | }, |
| | | PreAudioOrderEvent::TrackSubscribed { |
| | | participant: "user-1".to_string(), |
| | | attributes, |
| | | }, |
| | | ]); |
| | | assert_eq!(effects, ["ack_observed", "observer_started"]); |
| | |
| | | |
| | | #[test] |
| | | fn production_event_order_negative_probe_has_no_observer_or_session_effect() { |
| | | let mut invalid = HashMap::new(); |
| | | invalid.insert( |
| | | "inputSourceCategory".to_string(), |
| | | "ordinary_mic".to_string(), |
| | | ); |
| | | let effects = drive_pre_audio_order_test_seam(&[ |
| | | PreAudioOrderEvent::DataReceived { |
| | | sender: "user-1".to_string(), |
| | | sequence: "fixture-01".to_string(), |
| | | sequence: "1".to_string(), |
| | | }, |
| | | PreAudioOrderEvent::TrackSubscribed { |
| | | participant: "user-1".to_string(), |
| | | attributes: invalid, |
| | | participant: "cross-call-user".to_string(), |
| | | }, |
| | | ]); |
| | | assert!(effects.is_empty()); |