| | |
| | | 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( |
| | | start_realtime_session_for_new_speech( |
| | | http.clone(), |
| | | turn_bridge_config.realtime_asr_config(), |
| | | &call_id, |
| | | &trace_id, |
| | | &turn_id, |
| | | &vad.speech_samples, |
| | | vad, |
| | | || 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(_) => {} |
| | | } |
| | | &mut realtime_asr_upload, |
| | | turn_bridge_config.asr_realtime_enabled, |
| | | ); |
| | | } else if was_in_speech { |
| | | let push_failed = realtime_asr_upload |
| | | .as_mut() |
| | |
| | | "runtime helper user_audio_stream_ended" |
| | | ); |
| | | }) |
| | | } |
| | | |
| | | 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>, |
| | | realtime_enabled: bool, |
| | | ) { |
| | | let turn_id = format!("turn-{:04}", vad.turn_index); |
| | | match RealtimeAsrUpload::start_with_participant_attributes( |
| | | http, |
| | | config, |
| | | call_id, |
| | | trace_id, |
| | | &turn_id, |
| | | &vad.speech_samples, |
| | | read_attributes, |
| | | ) { |
| | | Ok(upload) => { |
| | | 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(_) => {} |
| | | } |
| | | } |
| | | |
| | | struct DrainedUserAudioFrame { |
| | |
| | | assert!(is_bound_user_participant("participant-any", None)); |
| | | } |
| | | |
| | | #[tokio::test] |
| | | async fn production_observer_vad_to_session_entry_reads_each_updated_attribute() { |
| | | 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 was = vad.in_speech; |
| | | vad.observe_frame( |
| | | "call-001", |
| | | "trace-001", |
| | | "participant", |
| | | "track", |
| | | 1, |
| | | 1_000, |
| | | &frame, |
| | | ); |
| | | assert!(!was && vad.in_speech); |
| | | start_realtime_session_for_new_speech( |
| | | 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, |
| | | }, |
| | | "call-001", |
| | | "trace-001", |
| | | &vad, |
| | | || attrs.clone(), |
| | | &mut upload, |
| | | true, |
| | | ); |
| | | 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 = vad.in_speech; |
| | | vad.observe_frame( |
| | | "call-001", |
| | | "trace-001", |
| | | "participant", |
| | | "track", |
| | | 2, |
| | | 2_000, |
| | | &frame, |
| | | ); |
| | | assert!(!was && vad.in_speech); |
| | | assert_eq!( |
| | | "fixture-02", |
| | | AudioIngressMetadata::from_participant(&attrs) |
| | | .expect("valid attributes") |
| | | .expect("bound") |
| | | .client_fixture_sequence |
| | | ); |
| | | } |
| | | |
| | | #[test] |
| | | fn reply_chunk_marker_state_emits_turn_first_once_and_later_segment_first_once() { |
| | | let mut state = ReplyChunkMarkerState::default(); |