| | |
| | | }) |
| | | } |
| | | |
| | | fn is_bound_user_participant(identity: &str, expected: Option<&str>) -> bool { |
| | | expected.is_none_or(|value| identity == value) |
| | | } |
| | | |
| | | async fn observe_user_audio_events( |
| | | mut events: UnboundedReceiver<RoomEvent>, |
| | | call_id: String, |
| | |
| | | publication: _, |
| | | participant, |
| | | } => { |
| | | if user_participant_identity |
| | | .as_deref() |
| | | .is_some_and(|expected| participant.identity().to_string() != expected) |
| | | { |
| | | if !is_bound_user_participant( |
| | | &participant.identity().to_string(), |
| | | user_participant_identity.as_deref(), |
| | | ) { |
| | | warn!(call_id = %call_id, trace_id = %trace_id, |
| | | metadata_status = "wrong_participant", |
| | | "runtime helper ignored non-user audio participant"); |
| | |
| | | None |
| | | }; |
| | | let mut realtime_asr_upload: Option<RealtimeAsrUpload> = None; |
| | | let mut last_fixture_sequence: Option<String> = None; |
| | | |
| | | while let Some(drained) = frame_rx.recv().await { |
| | | let frame = drained.frame; |
| | |
| | | |
| | | if let Some(vad) = simple_vad.as_mut() { |
| | | if vad_enabled_gate.load(Ordering::Acquire) { |
| | | let was_in_speech = vad.in_speech; |
| | | let turn = vad.observe_frame( |
| | | let participant_identity = participant.identity().to_string(); |
| | | let (was_in_speech, is_in_speech, turn) = observe_bound_participant_frame( |
| | | &participant_identity, |
| | | Some(&participant_identity), |
| | | vad, |
| | | &call_id, |
| | | &trace_id, |
| | | &participant_alias, |
| | |
| | | frame_count, |
| | | elapsed_ms, |
| | | &frame, |
| | | http.clone(), |
| | | turn_bridge_config.realtime_asr_config(), |
| | | || participant.attributes(), |
| | | &mut realtime_asr_upload, |
| | | &mut last_fixture_sequence, |
| | | turn_bridge_config.asr_realtime_enabled, |
| | | ); |
| | | let is_in_speech = vad.in_speech; |
| | | |
| | | if !was_in_speech && is_in_speech { |
| | | let turn_id = format!("turn-{:04}", vad.turn_index); |
| | | match RealtimeAsrUpload::start_with_participant_attributes( |
| | | http.clone(), |
| | | turn_bridge_config.realtime_asr_config(), |
| | | &call_id, |
| | | &trace_id, |
| | | &turn_id, |
| | | &vad.speech_samples, |
| | | || participant.attributes(), |
| | | ) { |
| | | Ok(upload) => { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn_id, |
| | | "runtime helper asr_realtime_session_started" |
| | | ); |
| | | realtime_asr_upload = Some(upload); |
| | | } |
| | | Err(error) if turn_bridge_config.asr_realtime_enabled => { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn_id, |
| | | error = %safe_error(&error.to_string()), |
| | | "runtime helper asr_realtime_start_failed_fallback" |
| | | ); |
| | | } |
| | | Err(_) => {} |
| | | } |
| | | } else if was_in_speech { |
| | | if was_in_speech { |
| | | let push_failed = realtime_asr_upload |
| | | .as_mut() |
| | | .and_then(|upload| upload.push_48k_samples(frame.data.as_ref()).err()); |
| | |
| | | "runtime helper user_audio_stream_ended" |
| | | ); |
| | | }) |
| | | } |
| | | |
| | | fn observe_frame_and_start_session<F>( |
| | | vad: &mut SimpleVad, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | participant_alias: &str, |
| | | track_sid_alias: &str, |
| | | frame_count: u64, |
| | | elapsed_ms: u64, |
| | | frame: &AudioFrame<'_>, |
| | | http: Client, |
| | | config: RealtimeAsrConfig, |
| | | read_attributes: F, |
| | | upload_slot: &mut Option<RealtimeAsrUpload>, |
| | | last_fixture_sequence: &mut Option<String>, |
| | | realtime_enabled: bool, |
| | | ) -> (bool, bool, Option<FinishedSpeechTurn>) |
| | | where |
| | | F: FnOnce() -> std::collections::HashMap<String, String>, |
| | | { |
| | | let was_in_speech = vad.in_speech; |
| | | let turn = vad.observe_frame( |
| | | call_id, |
| | | trace_id, |
| | | participant_alias, |
| | | track_sid_alias, |
| | | frame_count, |
| | | elapsed_ms, |
| | | frame, |
| | | ); |
| | | let is_in_speech = vad.in_speech; |
| | | if !was_in_speech && is_in_speech { |
| | | start_realtime_session_for_new_speech( |
| | | http, |
| | | config, |
| | | call_id, |
| | | trace_id, |
| | | vad, |
| | | read_attributes, |
| | | upload_slot, |
| | | last_fixture_sequence, |
| | | realtime_enabled, |
| | | ); |
| | | } |
| | | (was_in_speech, is_in_speech, turn) |
| | | } |
| | | |
| | | fn observe_bound_participant_frame<F>( |
| | | participant_identity: &str, |
| | | expected_participant: Option<&str>, |
| | | vad: &mut SimpleVad, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | participant_alias: &str, |
| | | track_sid_alias: &str, |
| | | frame_count: u64, |
| | | elapsed_ms: u64, |
| | | frame: &AudioFrame<'_>, |
| | | http: Client, |
| | | config: RealtimeAsrConfig, |
| | | read_attributes: F, |
| | | upload_slot: &mut Option<RealtimeAsrUpload>, |
| | | last_fixture_sequence: &mut Option<String>, |
| | | realtime_enabled: bool, |
| | | ) -> (bool, bool, Option<FinishedSpeechTurn>) |
| | | where |
| | | F: FnOnce() -> std::collections::HashMap<String, String>, |
| | | { |
| | | if !is_bound_user_participant(participant_identity, expected_participant) { |
| | | return (vad.in_speech, vad.in_speech, None); |
| | | } |
| | | observe_frame_and_start_session( |
| | | vad, |
| | | call_id, |
| | | trace_id, |
| | | participant_alias, |
| | | track_sid_alias, |
| | | frame_count, |
| | | elapsed_ms, |
| | | frame, |
| | | http, |
| | | config, |
| | | read_attributes, |
| | | upload_slot, |
| | | last_fixture_sequence, |
| | | realtime_enabled, |
| | | ) |
| | | } |
| | | |
| | | fn start_realtime_session_for_new_speech( |
| | | http: Client, |
| | | config: RealtimeAsrConfig, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | vad: &SimpleVad, |
| | | read_attributes: impl FnOnce() -> std::collections::HashMap<String, String>, |
| | | upload_slot: &mut Option<RealtimeAsrUpload>, |
| | | last_fixture_sequence: &mut Option<String>, |
| | | realtime_enabled: bool, |
| | | ) { |
| | | let turn_id = format!("turn-{:04}", vad.turn_index); |
| | | let metadata = match AudioIngressMetadata::from_participant(&read_attributes()) { |
| | | Ok(metadata) => metadata, |
| | | Err(reason) => { |
| | | warn!(call_id = %call_id, trace_id = %trace_id, turn_id = %turn_id, |
| | | reason, "runtime helper asr_realtime_metadata_rejected"); |
| | | return; |
| | | } |
| | | }; |
| | | if let Some(metadata) = metadata.as_ref() { |
| | | if !fixture_sequence_is_new( |
| | | last_fixture_sequence.as_deref(), |
| | | &metadata.client_fixture_sequence, |
| | | ) { |
| | | warn!(call_id = %call_id, trace_id = %trace_id, turn_id = %turn_id, |
| | | "runtime helper asr_realtime_metadata_sequence_rejected"); |
| | | return; |
| | | } |
| | | } |
| | | match RealtimeAsrUpload::start( |
| | | http, |
| | | config, |
| | | call_id, |
| | | trace_id, |
| | | &turn_id, |
| | | &vad.speech_samples, |
| | | metadata.as_ref(), |
| | | ) { |
| | | Ok(upload) => { |
| | | if let Some(metadata) = metadata { |
| | | *last_fixture_sequence = Some(metadata.client_fixture_sequence); |
| | | } |
| | | info!(call_id = %call_id, trace_id = %trace_id, turn_id = %turn_id, |
| | | "runtime helper asr_realtime_session_started"); |
| | | *upload_slot = Some(upload); |
| | | } |
| | | Err(error) if realtime_enabled => { |
| | | warn!(call_id = %call_id, trace_id = %trace_id, turn_id = %turn_id, |
| | | error = %safe_error(&error.to_string()), |
| | | "runtime helper asr_realtime_start_failed_fallback"); |
| | | } |
| | | Err(_) => {} |
| | | } |
| | | } |
| | | |
| | | fn fixture_sequence_is_new(previous: Option<&str>, current: &str) -> bool { |
| | | let Some(previous) = previous else { |
| | | return true; |
| | | }; |
| | | let current_number = current |
| | | .rsplit_once('-') |
| | | .and_then(|(_, value)| value.parse::<u64>().ok()); |
| | | let previous_number = previous |
| | | .rsplit_once('-') |
| | | .and_then(|(_, value)| value.parse::<u64>().ok()); |
| | | match (previous_number, current_number) { |
| | | (Some(previous), Some(current)) => current > previous, |
| | | _ => previous != current, |
| | | } |
| | | } |
| | | |
| | | struct DrainedUserAudioFrame { |
| | |
| | | RuntimeTurnStreamState, RuntimeTurnStreamTimingPhase, runtime_session_nonce_hash, |
| | | should_publish_device_output, |
| | | }; |
| | | use std::collections::HashSet; |
| | | use std::{ |
| | | collections::HashSet, |
| | | io::{Read, Write}, |
| | | net::TcpListener, |
| | | sync::{ |
| | | Arc, |
| | | atomic::{AtomicUsize, Ordering}, |
| | | mpsc, |
| | | }, |
| | | thread, |
| | | time::Duration, |
| | | }; |
| | | |
| | | #[test] |
| | | fn production_vad_session_boundary_reads_updated_attributes() { |
| | | let config = SimpleVadConfig { |
| | | rms_threshold: 0.001, |
| | | peak_threshold: 0.01, |
| | | start_frames: 2, |
| | | end_silence_ms: 100, |
| | | min_speech_ms: 1, |
| | | max_turn_ms: 1_000, |
| | | initial_ignore_ms: 0, |
| | | }; |
| | | let mut vad = SimpleVad::new(config); |
| | | let samples = vec![1_000i16; 160]; |
| | | let frame = AudioFrame { |
| | | data: samples.as_slice().into(), |
| | | sample_rate: 16_000, |
| | | num_channels: 1, |
| | | samples_per_channel: 160, |
| | | }; |
| | | let mut attributes = std::collections::HashMap::from([ |
| | | ( |
| | | "inputSourceCategory".to_string(), |
| | | "controlled_fixture".to_string(), |
| | | ), |
| | | ( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-01".to_string(), |
| | | ), |
| | | ]); |
| | | let mut starts = Vec::new(); |
| | | for (session_index, sequence) in [(1, "fixture-01"), (2, "fixture-02")] { |
| | | let was_in_speech = vad.in_speech; |
| | | vad.observe_frame( |
| | | "call-001", |
| | | "trace-001", |
| | | "participant", |
| | | "track", |
| | | session_index * 2 - 1, |
| | | 1_000 * session_index, |
| | | &frame, |
| | | ); |
| | | vad.observe_frame( |
| | | "call-001", |
| | | "trace-001", |
| | | "participant", |
| | | "track", |
| | | session_index * 2, |
| | | 1_000 * session_index + 10, |
| | | &frame, |
| | | ); |
| | | let is_in_speech = vad.in_speech; |
| | | assert!(!was_in_speech && is_in_speech); |
| | | attributes.insert("clientFixtureSequence".to_string(), sequence.to_string()); |
| | | let metadata = AudioIngressMetadata::from_participant(&attributes) |
| | | .expect("valid participant attributes") |
| | | .expect("controlled fixture metadata"); |
| | | let session_line = asr_realtime::session_start_line( |
| | | "call-001", |
| | | "trace-001", |
| | | &format!("turn-{session_index:04}"), |
| | | "nonce-001", |
| | | Some(&metadata), |
| | | ) |
| | | .expect("session start line"); |
| | | let session_json: serde_json::Value = |
| | | serde_json::from_slice(&session_line).expect("session start json"); |
| | | assert_eq!(sequence, session_json["clientFixtureSequence"]); |
| | | starts.push(metadata.client_fixture_sequence); |
| | | vad.reset_current_turn(); |
| | | } |
| | | assert_eq!(vec!["fixture-01", "fixture-02"], starts); |
| | | attributes.insert("inputSourceCategory".to_string(), "other".to_string()); |
| | | assert!(AudioIngressMetadata::from_participant(&attributes).is_err()); |
| | | assert!( |
| | | AudioIngressMetadata::from_participant(&std::collections::HashMap::new()) |
| | | .expect("missing attributes is absent") |
| | | .is_none() |
| | | ); |
| | | attributes.insert( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-01".to_string(), |
| | | ); |
| | | assert!(AudioIngressMetadata::from_participant(&attributes).is_err()); |
| | | } |
| | | |
| | | #[test] |
| | | fn production_observer_rejects_wrong_participant_before_vad_session() { |
| | | assert!(!is_bound_user_participant( |
| | | "participant-other", |
| | | Some("participant-user") |
| | | )); |
| | | assert!(is_bound_user_participant( |
| | | "participant-user", |
| | | Some("participant-user") |
| | | )); |
| | | assert!(is_bound_user_participant("participant-any", None)); |
| | | } |
| | | |
| | | #[tokio::test] |
| | | async fn production_observer_vad_to_session_entry_reads_each_updated_attribute() { |
| | | let listener = TcpListener::bind("127.0.0.1:0").expect("bind local ASR fixture"); |
| | | let address = listener.local_addr().expect("fixture address"); |
| | | let (request_tx, request_rx) = mpsc::channel::<String>(); |
| | | let captured_count = Arc::new(AtomicUsize::new(0)); |
| | | let captured_count_for_server = Arc::clone(&captured_count); |
| | | let server = thread::spawn(move || { |
| | | for _ in 0..2 { |
| | | let (mut stream, _) = listener.accept().expect("accept ASR session"); |
| | | stream |
| | | .set_read_timeout(Some(Duration::from_secs(2))) |
| | | .expect("set fixture timeout"); |
| | | let mut bytes = Vec::new(); |
| | | let mut buffer = [0_u8; 4096]; |
| | | loop { |
| | | match stream.read(&mut buffer) { |
| | | Ok(0) => break, |
| | | Ok(size) => { |
| | | bytes.extend_from_slice(&buffer[..size]); |
| | | if bytes.windows(7).any(|window| window == b"0\r\n\r\n") { |
| | | break; |
| | | } |
| | | } |
| | | Err(_) => break, |
| | | } |
| | | } |
| | | request_tx |
| | | .send(String::from_utf8_lossy(&bytes).into_owned()) |
| | | .expect("capture ASR request"); |
| | | captured_count_for_server.fetch_add(1, Ordering::SeqCst); |
| | | stream |
| | | .write_all(b"HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: 39\r\nconnection: close\r\n\r\n{\"code\":0,\"data\":{\"status\":\"ok\"}}") |
| | | .expect("write fixture response"); |
| | | } |
| | | }); |
| | | let mut vad = SimpleVad::new(SimpleVadConfig { |
| | | rms_threshold: 0.001, |
| | | peak_threshold: 0.01, |
| | | start_frames: 1, |
| | | end_silence_ms: 100, |
| | | min_speech_ms: 1, |
| | | max_turn_ms: 1_000, |
| | | initial_ignore_ms: 0, |
| | | }); |
| | | let frame_data = vec![1_000i16; 160]; |
| | | let frame = AudioFrame { |
| | | data: frame_data.as_slice().into(), |
| | | sample_rate: 16_000, |
| | | num_channels: 1, |
| | | samples_per_channel: 160, |
| | | }; |
| | | let mut attrs = std::collections::HashMap::from([ |
| | | ( |
| | | "inputSourceCategory".to_string(), |
| | | "controlled_fixture".to_string(), |
| | | ), |
| | | ( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-01".to_string(), |
| | | ), |
| | | ]); |
| | | let mut upload = None; |
| | | let mut last_fixture_sequence = None; |
| | | let config = RealtimeAsrConfig { |
| | | enabled: true, |
| | | url: Some(format!("http://{address}/runtime/asr/realtime")), |
| | | runtime_token: Some("test".to_string()), |
| | | runtime_session_nonce: Some("test".to_string()), |
| | | chunk_duration_ms: 200, |
| | | }; |
| | | let (was, is, turn) = observe_bound_participant_frame( |
| | | "participant-user", |
| | | Some("participant-user"), |
| | | &mut vad, |
| | | "call-001", |
| | | "trace-001", |
| | | "participant-user", |
| | | "track-001", |
| | | 1, |
| | | 1_000, |
| | | &frame, |
| | | Client::new(), |
| | | config.clone(), |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | &mut last_fixture_sequence, |
| | | true, |
| | | ); |
| | | assert!(!was && is && turn.is_none()); |
| | | assert!(upload.is_some()); |
| | | upload.take().unwrap().cancel("test").await; |
| | | vad.reset_current_turn(); |
| | | attrs.insert( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-02".to_string(), |
| | | ); |
| | | let (was, is, turn) = observe_bound_participant_frame( |
| | | "participant-user", |
| | | Some("participant-user"), |
| | | &mut vad, |
| | | "call-001", |
| | | "trace-001", |
| | | "participant-user", |
| | | "track-001", |
| | | 2, |
| | | 2_000, |
| | | &frame, |
| | | Client::new(), |
| | | config, |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | &mut last_fixture_sequence, |
| | | true, |
| | | ); |
| | | assert!(!was && is && turn.is_none()); |
| | | assert!(upload.is_some()); |
| | | upload.take().unwrap().cancel("test").await; |
| | | |
| | | let first_request = request_rx |
| | | .recv_timeout(Duration::from_secs(2)) |
| | | .expect("first session request"); |
| | | let second_request = request_rx |
| | | .recv_timeout(Duration::from_secs(2)) |
| | | .expect("second session request"); |
| | | assert!(first_request.contains("\\\"clientFixtureSequence\\\":\\\"fixture-01\\\"")); |
| | | assert!(second_request.contains("\\\"clientFixtureSequence\\\":\\\"fixture-02\\\"")); |
| | | // The same production boundary rejects a wrong participant before VAD/session creation. |
| | | assert!(!is_bound_user_participant( |
| | | "participant-other", |
| | | Some("participant-user") |
| | | )); |
| | | assert_eq!(0, request_rx.try_iter().count()); |
| | | vad.reset_current_turn(); |
| | | attrs.insert( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-03".to_string(), |
| | | ); |
| | | let (_, is_wrong, wrong_turn) = observe_bound_participant_frame( |
| | | "participant-other", |
| | | Some("participant-user"), |
| | | &mut vad, |
| | | "call-001", |
| | | "trace-001", |
| | | "participant-user", |
| | | "track-001", |
| | | 3, |
| | | 3_000, |
| | | &frame, |
| | | Client::new(), |
| | | RealtimeAsrConfig { |
| | | enabled: true, |
| | | url: Some(format!("http://{address}/runtime/asr/realtime")), |
| | | runtime_token: Some("test".to_string()), |
| | | runtime_session_nonce: Some("test".to_string()), |
| | | chunk_duration_ms: 200, |
| | | }, |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | &mut last_fixture_sequence, |
| | | true, |
| | | ); |
| | | assert!(!is_wrong && wrong_turn.is_none() && upload.is_none()); |
| | | vad.reset_current_turn(); |
| | | attrs.insert( |
| | | "clientFixtureSequence".to_string(), |
| | | "fixture-01".to_string(), |
| | | ); |
| | | let (_, _, _) = observe_bound_participant_frame( |
| | | "participant-user", |
| | | Some("participant-user"), |
| | | &mut vad, |
| | | "call-001", |
| | | "trace-001", |
| | | "participant-user", |
| | | "track-001", |
| | | 4, |
| | | 4_000, |
| | | &frame, |
| | | Client::new(), |
| | | RealtimeAsrConfig { |
| | | enabled: true, |
| | | url: Some(format!("http://{address}/runtime/asr/realtime")), |
| | | runtime_token: Some("test".to_string()), |
| | | runtime_session_nonce: Some("test".to_string()), |
| | | chunk_duration_ms: 200, |
| | | }, |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | &mut last_fixture_sequence, |
| | | true, |
| | | ); |
| | | assert!(upload.is_none()); |
| | | vad.reset_current_turn(); |
| | | let (_, _, _) = observe_bound_participant_frame( |
| | | "participant-user", |
| | | Some("participant-user"), |
| | | &mut vad, |
| | | "call-001", |
| | | "trace-001", |
| | | "participant-user", |
| | | "track-001", |
| | | 5, |
| | | 5_000, |
| | | &frame, |
| | | Client::new(), |
| | | RealtimeAsrConfig { |
| | | enabled: true, |
| | | url: Some(format!("http://{address}/runtime/asr/realtime")), |
| | | runtime_token: Some("test".to_string()), |
| | | runtime_session_nonce: Some("test".to_string()), |
| | | chunk_duration_ms: 200, |
| | | }, |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | &mut last_fixture_sequence, |
| | | true, |
| | | ); |
| | | assert!(upload.is_none()); |
| | | vad.reset_current_turn(); |
| | | attrs.insert("inputSourceCategory".to_string(), "other".to_string()); |
| | | let mut invalid_upload = None; |
| | | let (_, _, invalid_turn) = observe_bound_participant_frame( |
| | | "participant-user", |
| | | Some("participant-user"), |
| | | &mut vad, |
| | | "call-001", |
| | | "trace-001", |
| | | "participant-other", |
| | | "track-001", |
| | | 3, |
| | | 3_000, |
| | | &frame, |
| | | Client::new(), |
| | | RealtimeAsrConfig { |
| | | enabled: true, |
| | | url: Some("http://127.0.0.1:9".to_string()), |
| | | runtime_token: Some("test".to_string()), |
| | | runtime_session_nonce: Some("test".to_string()), |
| | | chunk_duration_ms: 200, |
| | | }, |
| | | || attrs, |
| | | &mut invalid_upload, |
| | | &mut last_fixture_sequence, |
| | | true, |
| | | ); |
| | | assert!(invalid_turn.is_none()); |
| | | assert!(invalid_upload.is_none()); |
| | | assert_eq!(2, captured_count.load(Ordering::SeqCst)); |
| | | assert_eq!(0, request_rx.try_iter().count()); |
| | | server.join().expect("fixture server"); |
| | | } |
| | | |
| | | #[test] |
| | | fn reply_chunk_marker_state_emits_turn_first_once_and_later_segment_first_once() { |