| | |
| | | }; |
| | | |
| | | use anyhow::{Context, Result, anyhow}; |
| | | use asr_realtime::{RealtimeAsrConfig, RealtimeAsrOutcome, RealtimeAsrUpload}; |
| | | use asr_realtime::{ |
| | | AudioIngressMetadata, RealtimeAsrConfig, RealtimeAsrOutcome, RealtimeAsrUpload, |
| | | }; |
| | | use audio::{AudioDiagnostics, load_pre_recorded_frames}; |
| | | use base64::{Engine as _, engine::general_purpose}; |
| | | use futures_util::StreamExt; |
| | |
| | | options::TrackPublishOptions, |
| | | prelude::{ |
| | | DataPacket, LocalAudioTrack, LocalTrack, ParticipantIdentity, RemoteAudioTrack, |
| | | RemoteTrack, Room, RoomEvent, RoomOptions, |
| | | RemoteParticipant, RemoteTrack, Room, RoomEvent, RoomOptions, |
| | | }, |
| | | }; |
| | | use reqwest::Client; |
| | |
| | | let enabled = config.user_audio_observer_enabled; |
| | | let simple_vad_enabled = config.simple_vad_enabled; |
| | | let simple_vad_config = config.simple_vad_config.clone(); |
| | | let user_participant_identity = config.user_participant_identity.clone(); |
| | | let turn_bridge_config = TurnBridgeConfig::from_config(config); |
| | | |
| | | tokio::spawn(async move { |
| | |
| | | turn_bridge_config, |
| | | http, |
| | | sink, |
| | | user_participant_identity, |
| | | ) |
| | | .await; |
| | | }) |
| | |
| | | turn_bridge_config: TurnBridgeConfig, |
| | | http: Client, |
| | | sink: Arc<BotAudioOutputSink>, |
| | | user_participant_identity: Option<String>, |
| | | ) { |
| | | info!( |
| | | call_id = %call_id, |
| | |
| | | publication: _, |
| | | participant, |
| | | } => { |
| | | if user_participant_identity |
| | | .as_deref() |
| | | .is_some_and(|expected| participant.identity().to_string() != expected) |
| | | { |
| | | warn!(call_id = %call_id, trace_id = %trace_id, |
| | | metadata_status = "wrong_participant", |
| | | "runtime helper ignored non-user audio participant"); |
| | | continue; |
| | | } |
| | | let participant_alias = redact(&participant.identity().to_string()); |
| | | let track_sid_alias = redact(&track.sid().to_string()); |
| | | let track_name = track.name(); |
| | |
| | | turn_bridge_config.clone(), |
| | | http.clone(), |
| | | sink.clone(), |
| | | participant, |
| | | ); |
| | | } |
| | | RoomEvent::TrackSubscribed { |
| | |
| | | turn_bridge_config: TurnBridgeConfig, |
| | | http: Client, |
| | | sink: Arc<BotAudioOutputSink>, |
| | | participant: RemoteParticipant, |
| | | ) -> JoinHandle<()> { |
| | | tokio::spawn(async move { |
| | | let mut stream = NativeAudioStream::new( |
| | |
| | | |
| | | if !was_in_speech && is_in_speech { |
| | | let turn_id = format!("turn-{:04}", vad.turn_index); |
| | | match RealtimeAsrUpload::start( |
| | | 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!( |
| | |
| | | use std::collections::HashSet; |
| | | |
| | | #[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"); |
| | | 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 reply_chunk_marker_state_emits_turn_first_once_and_later_segment_first_once() { |
| | | let mut state = ReplyChunkMarkerState::default(); |
| | | |