| | |
| | | |
| | | 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", |
| | | ]; |
| | | |
| | |
| | | trace_id_hash: &str, |
| | | generation: u64, |
| | | sequence: &str, |
| | | acknowledged_probe_sequences: &mut HashSet<String>, |
| | | acknowledged_probe_sequences: &mut HashSet<(u64, String)>, |
| | | ) -> ControlledFixtureAckPublishOutcome |
| | | where |
| | | F: Future<Output = Result<(), E>>, |
| | |
| | | Some(ack_result), |
| | | reject_reason, |
| | | ); |
| | | acknowledged_probe_sequences.insert(sequence.to_string()); |
| | | acknowledged_probe_sequences.insert((generation, sequence.to_string())); |
| | | ControlledFixtureAckPublishOutcome { |
| | | observed, |
| | | published: true, |
| | |
| | | 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() |
| | |
| | | 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 { |
| | |
| | | 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, |
| | |
| | | } |
| | | 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( |
| | | &call_id, |
| | | &trace_id, |
| | |
| | | |
| | | 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, |
| | | expected_participant: Option<&str>, |
| | |
| | | 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, |
| | |
| | | ); |
| | | 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; |
| | | 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( |
| | | decision, |
| | | Ok(()), |
| | | runtime_call_id, |
| | | runtime_trace_id, |
| | | &probe.call_id_hash, |
| | |
| | | }; |
| | | 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, |
| | | } |
| | | } |
| | | |
| | |
| | | |
| | | #[derive(Debug)] |
| | | enum PreAudioOrderEvent { |
| | | DataReceived { |
| | | sender: String, |
| | | sequence: String, |
| | | }, |
| | | TrackSubscribed { |
| | | participant: String, |
| | | attributes: HashMap<String, String>, |
| | | }, |
| | | 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()); |
| | |
| | | published: true, |
| | | } |
| | | ); |
| | | assert!(acknowledged.contains("fixture-02")); |
| | | assert!(acknowledged.contains(&(CONTROLLED_FIXTURE_GENERATION, "fixture-02".to_string()))); |
| | | assert_eq!(successful_observer_starts, 1); |
| | | } |
| | | |
| | |
| | | .await; |
| | | assert_eq!(result, Ok(())); |
| | | assert_eq!(reads, 4); |
| | | } |
| | | |
| | | #[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] |
| | |
| | | |
| | | #[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()); |