| | |
| | | reject_reason: Option<&'static str>, |
| | | } |
| | | |
| | | #[derive(Debug)] |
| | | #[derive(Clone, Debug)] |
| | | struct PendingControlledFixtureProbe { |
| | | sender: ParticipantIdentity, |
| | | call_id_hash: String, |
| | | call_trace_id_hash: String, |
| | | generation: u64, |
| | | sequence: String, |
| | | expires_at: Instant, |
| | | } |
| | |
| | | |
| | | fn record_controlled_fixture_probe_event( |
| | | stage: &'static str, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | call_id_hash: &str, |
| | | trace_id_hash: &str, |
| | | generation: u64, |
| | | sequence: &str, |
| | | observed: bool, |
| | | ack_result: Option<&'static str>, |
| | |
| | | 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( |
| | | 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( |
| | | 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!({ |
| | | "event": "controlled_fixture_attribute_probe", |
| | | "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, |
| | | call_id_hash: &str, |
| | | trace_id_hash: &str, |
| | | generation: u64, |
| | | sequence: &str, |
| | | ) -> (&'static str, Option<&'static str>, bool) { |
| | | let classification = controlled_fixture_ack_classification(decision); |
| | | record_controlled_fixture_probe_event( |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | call_id_hash, |
| | | trace_id_hash, |
| | | generation, |
| | | sequence, |
| | | classification.2, |
| | | Some(classification.0), |
| | |
| | | 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 |
| | |
| | | if publish.await.is_ok() { |
| | | record_controlled_fixture_probe_event( |
| | | "ack_publish_completed", |
| | | call_id, |
| | | trace_id, |
| | | call_id_hash, |
| | | trace_id_hash, |
| | | generation, |
| | | sequence, |
| | | observed, |
| | | Some(ack_result), |
| | |
| | | } else { |
| | | record_controlled_fixture_probe_event( |
| | | "ack_publish_completed", |
| | | call_id, |
| | | trace_id, |
| | | call_id_hash, |
| | | trace_id_hash, |
| | | generation, |
| | | sequence, |
| | | false, |
| | | Some("publish_failed"), |
| | | Some("ack_publish_failed"), |
| | | ); |
| | | false |
| | | } |
| | | } |
| | | |
| | | fn controlled_fixture_ack_from_probe( |
| | | probe: &PendingControlledFixtureProbe, |
| | | observed: bool, |
| | | reject_reason: Option<&'static str>, |
| | | ) -> ControlledFixtureAttributeAck { |
| | | ControlledFixtureAttributeAck { |
| | | message_type: CONTROLLED_FIXTURE_ACK_TOPIC, |
| | | protocol_version: CONTROLLED_FIXTURE_PROTOCOL_VERSION, |
| | | call_id_hash: probe.call_id_hash.clone(), |
| | | call_trace_id_hash: probe.call_trace_id_hash.clone(), |
| | | generation: probe.generation, |
| | | client_fixture_sequence: probe.sequence.clone(), |
| | | result: if observed { "observed" } else { "rejected" }, |
| | | input_source_category: observed.then_some("controlled_fixture"), |
| | | reject_reason, |
| | | } |
| | | } |
| | | |
| | |
| | | } |
| | | 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, |
| | | }) |
| | |
| | | track_source = %track_source, |
| | | "runtime helper user_track_subscribed" |
| | | ); |
| | | let pending_sequence = pending_probe.as_ref().map(|probe| probe.sequence.clone()); |
| | | let pending_binding = pending_probe.clone(); |
| | | let participant_for_probe = participant.clone(); |
| | | let probe_result = process_controlled_fixture_probe( |
| | | &mut pending_probe, |
| | | &mut acknowledged_probe_sequences, |
| | | Some(&participant_for_probe), |
| | | &sink, |
| | | &call_id, |
| | | &trace_id, |
| | | user_participant_identity.as_deref(), |
| | | ) |
| | | .await; |
| | | let observer_started = start_observer_after_controlled_fixture_probe( |
| | | pending_sequence.is_some(), |
| | | pending_binding.is_some(), |
| | | probe_result, |
| | | || { |
| | | spawn_user_audio_frame_observer( |
| | |
| | | ); |
| | | }, |
| | | ); |
| | | 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), |
| | |
| | | } 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, |
| | |
| | | if acknowledged_probe_sequences.contains(&probe.client_fixture_sequence) { |
| | | record_controlled_fixture_probe_event( |
| | | "data_received", |
| | | &call_id, |
| | | &trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.client_fixture_sequence, |
| | | false, |
| | | Some("rejected"), |
| | |
| | | if let Some(probe) = pending_probe.as_ref() { |
| | | record_controlled_fixture_probe_event( |
| | | "data_received", |
| | | &call_id, |
| | | &trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | false, |
| | | None, |
| | |
| | | &mut acknowledged_probe_sequences, |
| | | current_user_participant.as_ref(), |
| | | &sink, |
| | | &call_id, |
| | | &trace_id, |
| | | user_participant_identity.as_deref(), |
| | | ) |
| | | .await; |
| | |
| | | acknowledged_probe_sequences: &mut HashSet<String>, |
| | | participant: Option<&RemoteParticipant>, |
| | | sink: &BotAudioOutputSink, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | expected_participant: Option<&str>, |
| | | ) -> Option<bool> { |
| | | let Some(probe) = pending_probe.take() else { |
| | |
| | | if acknowledged_probe_sequences.contains(&probe.sequence) { |
| | | record_controlled_fixture_probe_event( |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | false, |
| | | Some("rejected"), |
| | |
| | | if Instant::now() > probe.expires_at { |
| | | record_controlled_fixture_probe_event( |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | false, |
| | | Some("timeout"), |
| | |
| | | let Some(participant) = participant else { |
| | | record_controlled_fixture_probe_event( |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | false, |
| | | None, |
| | |
| | | if participant.identity() != probe.sender { |
| | | record_controlled_fixture_probe_event( |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | false, |
| | | Some("rejected"), |
| | |
| | | || participant.attributes(), |
| | | ) |
| | | .await; |
| | | let protocol_reject_reason = decision.err(); |
| | | let (ack_result, reject_reason, observed) = |
| | | record_controlled_fixture_attribute_decision(decision, call_id, trace_id, &probe.sequence); |
| | | let result = if observed { "observed" } else { "rejected" }; |
| | | let input_source_category = observed.then_some("controlled_fixture"); |
| | | let ack = ControlledFixtureAttributeAck { |
| | | message_type: CONTROLLED_FIXTURE_ACK_TOPIC, |
| | | protocol_version: CONTROLLED_FIXTURE_PROTOCOL_VERSION, |
| | | call_id_hash: sha256_hex(call_id), |
| | | call_trace_id_hash: sha256_hex(trace_id), |
| | | generation: CONTROLLED_FIXTURE_GENERATION, |
| | | client_fixture_sequence: probe.sequence.clone(), |
| | | result, |
| | | input_source_category, |
| | | reject_reason: protocol_reject_reason, |
| | | }; |
| | | let (ack_result, reject_reason, observed) = record_controlled_fixture_attribute_decision( |
| | | decision, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | ); |
| | | let ack = controlled_fixture_ack_from_probe(&probe, observed, reject_reason); |
| | | let payload = match serde_json::to_vec(&ack) { |
| | | Ok(payload) => payload, |
| | | Err(_) => return Some(false), |
| | |
| | | let local_participant = sink.room.local_participant(); |
| | | record_controlled_fixture_probe_event( |
| | | "ack_publish_started", |
| | | call_id, |
| | | trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | observed, |
| | | Some(ack_result), |
| | |
| | | observed, |
| | | ack_result, |
| | | reject_reason, |
| | | call_id, |
| | | trace_id, |
| | | &probe.call_id_hash, |
| | | &probe.call_trace_id_hash, |
| | | probe.generation, |
| | | &probe.sequence, |
| | | acknowledged_probe_sequences, |
| | | ) |
| | |
| | | io::{Read, Write}, |
| | | net::TcpListener, |
| | | sync::{ |
| | | Arc, Mutex, |
| | | Arc, |
| | | atomic::{AtomicUsize, Ordering}, |
| | | mpsc, |
| | | }, |
| | |
| | | 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()) |
| | | } |
| | | } |
| | | |
| | | fn drive_pre_audio_order_test_seam(events: &[PreAudioOrderEvent]) -> Vec<&'static str> { |
| | |
| | | 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", |
| | | call_id, |
| | | trace_id, |
| | | 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( |
| | | 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["call_id_hash"], call_id_hash); |
| | | assert_eq!(event["trace_id_hash"], trace_id_hash); |
| | | assert_eq!(event["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( |
| | | 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( |
| | | "audio_observer_allowed", |
| | | &call_id_hash, |
| | | &trace_id_hash, |
| | | CONTROLLED_FIXTURE_GENERATION, |
| | | sequence, |
| | | observed, |
| | | Some(ack_result), |
| | | reject_reason, |
| | | ); |
| | | assert_eq!(allowed["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(), |
| | | 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(); |
| | |
| | | None, |
| | | "call-publish-failure", |
| | | "trace-publish-failure", |
| | | CONTROLLED_FIXTURE_GENERATION, |
| | | "fixture-01", |
| | | &mut acknowledged, |
| | | ) |
| | |
| | | None, |
| | | "call-publish-success", |
| | | "trace-publish-success", |
| | | CONTROLLED_FIXTURE_GENERATION, |
| | | "fixture-02", |
| | | &mut acknowledged, |
| | | ) |