| | |
| | | }; |
| | | |
| | | 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; |
| | |
| | | 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 ingress_metadata = |
| | | match AudioIngressMetadata::from_participant(&participant.attributes()) { |
| | | Ok(value) => value, |
| | | Err(reason) => { |
| | | warn!(call_id = %call_id, trace_id = %trace_id, |
| | | metadata_status = "invalid", reason = reason, |
| | | "runtime helper ignored invalid audio ingress metadata"); |
| | | None |
| | | } |
| | | }; |
| | | let participant_alias = redact(&participant.identity().to_string()); |
| | | let track_sid_alias = redact(&track.sid().to_string()); |
| | | let track_name = track.name(); |
| | |
| | | track_sid_alias = %track_sid_alias, |
| | | track_name = %track_name, |
| | | track_source = %track_source, |
| | | metadata_status = if ingress_metadata.is_some() { "bound" } else { "absent" }, |
| | | metadata_source = ingress_metadata.as_ref().map(|_| "controlled_fixture"), |
| | | metadata_sequence_present = ingress_metadata.is_some(), |
| | | "runtime helper user_track_subscribed" |
| | | ); |
| | | spawn_user_audio_frame_observer( |
| | |
| | | &trace_id, |
| | | &turn_id, |
| | | &vad.speech_samples, |
| | | ingress_metadata.as_ref(), |
| | | ) { |
| | | Ok(upload) => { |
| | | info!( |