| | |
| | | borrow::Cow, |
| | | collections::HashSet, |
| | | env, fs, |
| | | future::Future, |
| | | path::{Path, PathBuf}, |
| | | sync::{ |
| | | Arc, |
| | |
| | | const CONTROLLED_FIXTURE_ACK_TOPIC: &str = "controlled_fixture_attribute_ack"; |
| | | const CONTROLLED_FIXTURE_PROTOCOL_VERSION: u64 = 1; |
| | | const CONTROLLED_FIXTURE_GENERATION: u64 = 1; |
| | | const CONTROLLED_FIXTURE_PROBE_RECHECKS: usize = 3; |
| | | const CONTROLLED_FIXTURE_PROBE_RECHECK_DELAY: Duration = Duration::from_millis(25); |
| | | const CONTROLLED_FIXTURE_PROBE_TTL: Duration = Duration::from_millis(250); |
| | | |
| | | const CONTROLLED_FIXTURE_ACK_RESULTS: [&str; 4] = |
| | | ["observed", "rejected", "timeout", "publish_failed"]; |
| | | const CONTROLLED_FIXTURE_REJECT_REASONS: [&str; 10] = [ |
| | | "missing_attributes", |
| | | "wrong_source", |
| | | "missing_sequence", |
| | | "wrong_sequence", |
| | | "wrong_participant", |
| | | "expired", |
| | | "duplicate_or_old_sequence", |
| | | "no_current_participant", |
| | | "ack_publish_failed", |
| | | "unknown", |
| | | ]; |
| | | |
| | | #[derive(Debug, Deserialize)] |
| | | #[serde(rename_all = "camelCase")] |
| | |
| | | .iter() |
| | | .map(|byte| format!("{byte:02x}")) |
| | | .collect() |
| | | } |
| | | |
| | | fn controlled_fixture_ack_classification( |
| | | decision: Result<(), &'static str>, |
| | | ) -> (&'static str, Option<&'static str>, bool) { |
| | | match decision { |
| | | Ok(()) => ("observed", None, true), |
| | | Err("timeout") => ("timeout", Some("expired"), false), |
| | | Err(reason) if CONTROLLED_FIXTURE_REJECT_REASONS.contains(&reason) => { |
| | | ("rejected", Some(reason), false) |
| | | } |
| | | Err(_) => ("rejected", Some("unknown"), false), |
| | | } |
| | | } |
| | | |
| | | fn record_controlled_fixture_probe_event( |
| | | stage: &'static str, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | sequence: &str, |
| | | observed: bool, |
| | | ack_result: Option<&'static str>, |
| | | reject_reason: Option<&'static str>, |
| | | ) { |
| | | debug_assert!(ack_result.is_none_or(|value| CONTROLLED_FIXTURE_ACK_RESULTS.contains(&value))); |
| | | 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", |
| | | stage, |
| | | 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 record_controlled_fixture_attribute_decision( |
| | | decision: Result<(), &'static str>, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | 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, |
| | | sequence, |
| | | classification.2, |
| | | Some(classification.0), |
| | | classification.1, |
| | | ); |
| | | classification |
| | | } |
| | | |
| | | async fn complete_controlled_fixture_ack_publish<F, E>( |
| | | publish: F, |
| | | observed: bool, |
| | | ack_result: &'static str, |
| | | reject_reason: Option<&'static str>, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | sequence: &str, |
| | | acknowledged_probe_sequences: &mut HashSet<String>, |
| | | ) -> bool |
| | | where |
| | | F: Future<Output = Result<(), E>>, |
| | | { |
| | | if publish.await.is_ok() { |
| | | record_controlled_fixture_probe_event( |
| | | "ack_publish_completed", |
| | | call_id, |
| | | trace_id, |
| | | sequence, |
| | | observed, |
| | | Some(ack_result), |
| | | reject_reason, |
| | | ); |
| | | acknowledged_probe_sequences.insert(sequence.to_string()); |
| | | observed |
| | | } else { |
| | | record_controlled_fixture_probe_event( |
| | | "ack_publish_completed", |
| | | call_id, |
| | | trace_id, |
| | | sequence, |
| | | false, |
| | | Some("publish_failed"), |
| | | Some("ack_publish_failed"), |
| | | ); |
| | | false |
| | | } |
| | | } |
| | | |
| | | fn controlled_fixture_probe( |
| | |
| | | return Err("wrong_sequence"); |
| | | } |
| | | Ok(()) |
| | | } |
| | | |
| | | async fn observe_controlled_fixture_attributes<F>( |
| | | expires_at: Instant, |
| | | actual_participant: &str, |
| | | expected_participant: Option<&str>, |
| | | requested_sequence: &str, |
| | | mut read_attributes: F, |
| | | ) -> Result<(), &'static str> |
| | | where |
| | | F: FnMut() -> std::collections::HashMap<String, String>, |
| | | { |
| | | loop { |
| | | if Instant::now() > expires_at { |
| | | return Err("timeout"); |
| | | } |
| | | let decision = classify_controlled_fixture_attributes( |
| | | actual_participant, |
| | | expected_participant, |
| | | &read_attributes(), |
| | | requested_sequence, |
| | | ); |
| | | if decision.is_ok() || !matches!(decision, Err("missing_attributes")) { |
| | | return decision; |
| | | } |
| | | let Some(next_check) = Instant::now().checked_add(CONTROLLED_FIXTURE_PROBE_RECHECK_DELAY) |
| | | else { |
| | | return Err("timeout"); |
| | | }; |
| | | if next_check > expires_at { |
| | | return Err("timeout"); |
| | | } |
| | | sleep(CONTROLLED_FIXTURE_PROBE_RECHECK_DELAY).await; |
| | | } |
| | | } |
| | | |
| | | fn controlled_fixture_observer_gate(pending_probe: bool, probe_result: Option<bool>) -> bool { |
| | | !pending_probe || probe_result == Some(true) |
| | | } |
| | | |
| | | fn start_observer_after_controlled_fixture_probe<F>( |
| | | pending_probe: bool, |
| | | probe_result: Option<bool>, |
| | | spawn: F, |
| | | ) -> bool |
| | | where |
| | | F: FnOnce(), |
| | | { |
| | | if !controlled_fixture_observer_gate(pending_probe, probe_result) { |
| | | return false; |
| | | } |
| | | spawn(); |
| | | true |
| | | } |
| | | |
| | | #[tokio::main(flavor = "multi_thread")] |
| | |
| | | track_source = %track_source, |
| | | "runtime helper user_track_subscribed" |
| | | ); |
| | | let pending_sequence = pending_probe.as_ref().map(|probe| probe.sequence.clone()); |
| | | let participant_for_probe = participant.clone(); |
| | | spawn_user_audio_frame_observer( |
| | | track, |
| | | call_id.clone(), |
| | | trace_id.clone(), |
| | | participant_alias, |
| | | track_sid_alias, |
| | | simple_vad_enabled, |
| | | simple_vad_config.clone(), |
| | | vad_enabled_gate.clone(), |
| | | turn_bridge_config.clone(), |
| | | http.clone(), |
| | | sink.clone(), |
| | | user_participant_identity.clone(), |
| | | participant, |
| | | ); |
| | | current_user_participant = Some(participant_for_probe); |
| | | process_controlled_fixture_probe( |
| | | let probe_result = process_controlled_fixture_probe( |
| | | &mut pending_probe, |
| | | &mut acknowledged_probe_sequences, |
| | | current_user_participant.as_ref(), |
| | | 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(), |
| | | probe_result, |
| | | || { |
| | | spawn_user_audio_frame_observer( |
| | | track, |
| | | call_id.clone(), |
| | | trace_id.clone(), |
| | | participant_alias, |
| | | track_sid_alias, |
| | | simple_vad_enabled, |
| | | simple_vad_config.clone(), |
| | | vad_enabled_gate.clone(), |
| | | turn_bridge_config.clone(), |
| | | http.clone(), |
| | | sink.clone(), |
| | | user_participant_identity.clone(), |
| | | participant, |
| | | ); |
| | | }, |
| | | ); |
| | | if let Some(sequence) = pending_sequence.as_deref() { |
| | | 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( |
| | | if observer_started { |
| | | "audio_observer_allowed" |
| | | } else { |
| | | "audio_observer_blocked" |
| | | }, |
| | | &call_id, |
| | | &trace_id, |
| | | sequence, |
| | | observed, |
| | | ack_result, |
| | | reject_reason, |
| | | ); |
| | | } |
| | | if !observer_started { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | "runtime helper withheld audio observer until controlled fixture ACK" |
| | | ); |
| | | continue; |
| | | } |
| | | current_user_participant = Some(participant_for_probe); |
| | | } |
| | | RoomEvent::DataReceived { |
| | | payload, |
| | |
| | | serde_json::from_slice::<ControlledFixtureAttributeProbe>(&payload) |
| | | { |
| | | if acknowledged_probe_sequences.contains(&probe.client_fixture_sequence) { |
| | | record_controlled_fixture_probe_event( |
| | | "data_received", |
| | | &call_id, |
| | | &trace_id, |
| | | &probe.client_fixture_sequence, |
| | | false, |
| | | Some("rejected"), |
| | | Some("duplicate_or_old_sequence"), |
| | | ); |
| | | continue; |
| | | } |
| | | } |
| | |
| | | &sender.identity(), |
| | | user_participant_identity.as_deref(), |
| | | ); |
| | | if let Some(probe) = pending_probe.as_ref() { |
| | | record_controlled_fixture_probe_event( |
| | | "data_received", |
| | | &call_id, |
| | | &trace_id, |
| | | &probe.sequence, |
| | | false, |
| | | None, |
| | | None, |
| | | ); |
| | | } |
| | | process_controlled_fixture_probe( |
| | | &mut pending_probe, |
| | | &mut acknowledged_probe_sequences, |
| | |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | expected_participant: Option<&str>, |
| | | ) { |
| | | ) -> Option<bool> { |
| | | let Some(probe) = pending_probe.take() else { |
| | | return; |
| | | return None; |
| | | }; |
| | | if acknowledged_probe_sequences.contains(&probe.sequence) || Instant::now() > probe.expires_at { |
| | | return; |
| | | if acknowledged_probe_sequences.contains(&probe.sequence) { |
| | | record_controlled_fixture_probe_event( |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | &probe.sequence, |
| | | false, |
| | | Some("rejected"), |
| | | Some("duplicate_or_old_sequence"), |
| | | ); |
| | | return Some(false); |
| | | } |
| | | if Instant::now() > probe.expires_at { |
| | | record_controlled_fixture_probe_event( |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | &probe.sequence, |
| | | false, |
| | | Some("timeout"), |
| | | Some("expired"), |
| | | ); |
| | | return Some(false); |
| | | } |
| | | let Some(participant) = participant else { |
| | | record_controlled_fixture_probe_event( |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | &probe.sequence, |
| | | false, |
| | | None, |
| | | Some("no_current_participant"), |
| | | ); |
| | | *pending_probe = Some(probe); |
| | | return; |
| | | return None; |
| | | }; |
| | | if participant.identity() != probe.sender { |
| | | return; |
| | | } |
| | | let mut decision = Err("timeout"); |
| | | for attempt in 0..CONTROLLED_FIXTURE_PROBE_RECHECKS { |
| | | if Instant::now() > probe.expires_at { |
| | | break; |
| | | } |
| | | let attributes = participant.attributes(); |
| | | decision = classify_controlled_fixture_attributes( |
| | | &participant.identity().to_string(), |
| | | expected_participant, |
| | | &attributes, |
| | | record_controlled_fixture_probe_event( |
| | | "attributes_classified", |
| | | call_id, |
| | | trace_id, |
| | | &probe.sequence, |
| | | false, |
| | | Some("rejected"), |
| | | Some("wrong_participant"), |
| | | ); |
| | | if decision.is_ok() || !matches!(decision, Err("missing_attributes")) { |
| | | break; |
| | | } |
| | | if attempt + 1 < CONTROLLED_FIXTURE_PROBE_RECHECKS { |
| | | sleep(CONTROLLED_FIXTURE_PROBE_RECHECK_DELAY).await; |
| | | } |
| | | return Some(false); |
| | | } |
| | | let (result, input_source_category, reject_reason) = match decision { |
| | | Ok(()) => ("observed", Some("controlled_fixture"), None), |
| | | Err(reason) => ("rejected", None, Some(reason)), |
| | | }; |
| | | let decision = observe_controlled_fixture_attributes( |
| | | probe.expires_at, |
| | | &participant.identity().to_string(), |
| | | expected_participant, |
| | | &probe.sequence, |
| | | || 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, |
| | |
| | | client_fixture_sequence: probe.sequence.clone(), |
| | | result, |
| | | input_source_category, |
| | | reject_reason, |
| | | reject_reason: protocol_reject_reason, |
| | | }; |
| | | let payload = match serde_json::to_vec(&ack) { |
| | | Ok(payload) => payload, |
| | | Err(_) => return, |
| | | Err(_) => return Some(false), |
| | | }; |
| | | let local_participant = sink.room.local_participant(); |
| | | record_controlled_fixture_probe_event( |
| | | "ack_publish_started", |
| | | call_id, |
| | | trace_id, |
| | | &probe.sequence, |
| | | observed, |
| | | Some(ack_result), |
| | | reject_reason, |
| | | ); |
| | | let publish = local_participant.publish_data(DataPacket { |
| | | payload, |
| | | topic: Some(CONTROLLED_FIXTURE_ACK_TOPIC.to_string()), |
| | | reliable: true, |
| | | destination_identities: vec![probe.sender], |
| | | }); |
| | | if publish.await.is_ok() { |
| | | acknowledged_probe_sequences.insert(probe.sequence); |
| | | } |
| | | Some( |
| | | complete_controlled_fixture_ack_publish( |
| | | publish, |
| | | observed, |
| | | ack_result, |
| | | reject_reason, |
| | | call_id, |
| | | trace_id, |
| | | &probe.sequence, |
| | | acknowledged_probe_sequences, |
| | | ) |
| | | .await, |
| | | ) |
| | | } |
| | | |
| | | async fn handle_finished_turn( |
| | |
| | | mod tests { |
| | | use super::*; |
| | | use std::{ |
| | | collections::HashSet, |
| | | collections::{HashMap, HashSet}, |
| | | io::{Read, Write}, |
| | | net::TcpListener, |
| | | sync::{ |
| | | Arc, |
| | | Arc, Mutex, |
| | | atomic::{AtomicUsize, Ordering}, |
| | | mpsc, |
| | | }, |
| | | thread, |
| | | time::Duration, |
| | | }; |
| | | |
| | | #[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()) |
| | | } |
| | | } |
| | | |
| | | fn drive_pre_audio_order_test_seam(events: &[PreAudioOrderEvent]) -> Vec<&'static str> { |
| | | let mut pending_sequence = 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::TrackSubscribed { |
| | | participant, |
| | | attributes, |
| | | } => { |
| | | let pending = pending_sequence.is_some(); |
| | | let probe_result = pending_sequence.map(|sequence| { |
| | | classify_controlled_fixture_attributes( |
| | | participant, |
| | | Some("user-1"), |
| | | attributes, |
| | | sequence, |
| | | ) |
| | | .is_ok() |
| | | }); |
| | | let mut ack_observed = false; |
| | | let mut observer_started = false; |
| | | start_observer_after_controlled_fixture_probe(pending, probe_result, || { |
| | | if pending && probe_result == Some(true) { |
| | | ack_observed = true; |
| | | } |
| | | observer_started = true; |
| | | }); |
| | | if ack_observed { |
| | | effects.push("ack_observed"); |
| | | } |
| | | if observer_started { |
| | | effects.push("observer_started"); |
| | | } |
| | | pending_sequence = None; |
| | | } |
| | | PreAudioOrderEvent::DataReceived { .. } => {} |
| | | } |
| | | } |
| | | effects |
| | | } |
| | | |
| | | #[test] |
| | | fn production_vad_session_boundary_reads_updated_attributes() { |
| | |
| | | } |
| | | |
| | | #[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(); |
| | | let call_id = "private-call-value"; |
| | | let trace_id = "private-trace-value"; |
| | | let sequence = "private-sequence-value"; |
| | | 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, |
| | | sequence, |
| | | observed, |
| | | Some(ack_result), |
| | | reject_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!(observer_starts, 0); |
| | | } |
| | | |
| | | #[test] |
| | | fn controlled_fixture_probe_observed_and_failure_enums_are_stable() { |
| | | assert_eq!( |
| | | controlled_fixture_ack_classification(Ok(())), |
| | | ("observed", None, true) |
| | | ); |
| | | assert_eq!( |
| | | controlled_fixture_ack_classification(Err("timeout")), |
| | | ("timeout", Some("expired"), false) |
| | | ); |
| | | for reason in [ |
| | | "missing_attributes", |
| | | "wrong_source", |
| | | "missing_sequence", |
| | | "wrong_sequence", |
| | | "wrong_participant", |
| | | ] { |
| | | assert_eq!( |
| | | controlled_fixture_ack_classification(Err(reason)), |
| | | ("rejected", Some(reason), false) |
| | | ); |
| | | } |
| | | assert_eq!( |
| | | controlled_fixture_ack_classification(Err("unclassified")), |
| | | ("rejected", Some("unknown"), false) |
| | | ); |
| | | } |
| | | |
| | | #[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::<(), ()>(()) }, |
| | | true, |
| | | "observed", |
| | | None, |
| | | "call-publish-failure", |
| | | "trace-publish-failure", |
| | | "fixture-01", |
| | | &mut acknowledged, |
| | | ) |
| | | .await; |
| | | let mut observer_starts = 0; |
| | | let mut session_starts = 0; |
| | | let mut audio_starts = 0; |
| | | let mut speaking_starts = 0; |
| | | assert!(!start_observer_after_controlled_fixture_probe( |
| | | true, |
| | | Some(probe_result), |
| | | || { |
| | | observer_starts += 1; |
| | | session_starts += 1; |
| | | audio_starts += 1; |
| | | speaking_starts += 1; |
| | | }, |
| | | )); |
| | | assert!(!probe_result); |
| | | assert!(acknowledged.is_empty()); |
| | | assert_eq!(observer_starts, 0); |
| | | assert_eq!(session_starts, 0); |
| | | assert_eq!(audio_starts, 0); |
| | | assert_eq!(speaking_starts, 0); |
| | | |
| | | let observed_result = complete_controlled_fixture_ack_publish( |
| | | async { Ok::<(), ()>(()) }, |
| | | true, |
| | | "observed", |
| | | None, |
| | | "call-publish-success", |
| | | "trace-publish-success", |
| | | "fixture-02", |
| | | &mut acknowledged, |
| | | ) |
| | | .await; |
| | | let mut successful_observer_starts = 0; |
| | | assert!(start_observer_after_controlled_fixture_probe( |
| | | true, |
| | | Some(observed_result), |
| | | || successful_observer_starts += 1, |
| | | )); |
| | | assert!(observed_result); |
| | | assert!(acknowledged.contains("fixture-02")); |
| | | assert_eq!(successful_observer_starts, 1); |
| | | } |
| | | |
| | | #[tokio::test] |
| | | async fn production_attribute_observation_accepts_server_visibility_within_probe_ttl() { |
| | | let expected = HashMap::from([ |
| | | ( |
| | | "inputSourceCategory".to_string(), |
| | | "controlled_fixture".to_string(), |
| | | ), |
| | | ( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-01".to_string(), |
| | | ), |
| | | ]); |
| | | let mut reads = 0; |
| | | let result = observe_controlled_fixture_attributes( |
| | | Instant::now() + CONTROLLED_FIXTURE_PROBE_TTL, |
| | | "user-1", |
| | | Some("user-1"), |
| | | "fixture-01", |
| | | || { |
| | | reads += 1; |
| | | if reads <= 3 { |
| | | HashMap::new() |
| | | } else { |
| | | expected.clone() |
| | | } |
| | | }, |
| | | ) |
| | | .await; |
| | | assert_eq!(result, Ok(())); |
| | | assert_eq!(reads, 4); |
| | | } |
| | | |
| | | #[tokio::test] |
| | | async fn production_attribute_observation_rejects_wrong_sequence_without_audio_effect() { |
| | | let attributes = HashMap::from([ |
| | | ( |
| | | "inputSourceCategory".to_string(), |
| | | "controlled_fixture".to_string(), |
| | | ), |
| | | ( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-02".to_string(), |
| | | ), |
| | | ]); |
| | | let result = observe_controlled_fixture_attributes( |
| | | Instant::now() + CONTROLLED_FIXTURE_PROBE_TTL, |
| | | "user-1", |
| | | Some("user-1"), |
| | | "fixture-01", |
| | | || attributes.clone(), |
| | | ) |
| | | .await; |
| | | let mut observer_starts = 0; |
| | | assert_eq!(result, Err("wrong_sequence")); |
| | | assert!(!start_observer_after_controlled_fixture_probe( |
| | | true, |
| | | Some(result.is_ok()), |
| | | || observer_starts += 1, |
| | | )); |
| | | assert_eq!(observer_starts, 0); |
| | | } |
| | | |
| | | #[test] |
| | | fn controlled_fixture_ack_payload_is_reliable_and_redacted() { |
| | | let ack = ControlledFixtureAttributeAck { |
| | | message_type: CONTROLLED_FIXTURE_ACK_TOPIC, |
| | |
| | | "controlled_fixture_attribute_ack" |
| | | ); |
| | | } |
| | | |
| | | #[test] |
| | | fn controlled_fixture_probe_must_be_observed_before_audio_observer() { |
| | | assert!(controlled_fixture_observer_gate(false, None)); |
| | | assert!(controlled_fixture_observer_gate(true, Some(true))); |
| | | assert!(!controlled_fixture_observer_gate(true, Some(false))); |
| | | assert!(!controlled_fixture_observer_gate(true, None)); |
| | | } |
| | | |
| | | #[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(), |
| | | }, |
| | | 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(), |
| | | }, |
| | | PreAudioOrderEvent::TrackSubscribed { |
| | | participant: "user-1".to_string(), |
| | | attributes: invalid, |
| | | }, |
| | | ]); |
| | | assert!(effects.is_empty()); |
| | | } |
| | | |
| | | #[test] |
| | | fn production_audio_branch_orders_probe_before_spawn_callsite() { |
| | | let source = include_str!("main.rs"); |
| | | let branch = source |
| | | .find("RoomEvent::TrackSubscribed {\n track: RemoteTrack::Audio") |
| | | .expect("audio TrackSubscribed production branch"); |
| | | let branch_source = &source[branch..]; |
| | | let probe = branch_source |
| | | .find("let probe_result = process_controlled_fixture_probe") |
| | | .expect("probe must be processed in audio branch"); |
| | | let spawn = branch_source |
| | | .find("start_observer_after_controlled_fixture_probe") |
| | | .expect("spawn must use shared order entry"); |
| | | assert!( |
| | | probe < spawn, |
| | | "probe must precede shared observer spawn entry" |
| | | ); |
| | | } |
| | | } |