| | |
| | | }; |
| | | |
| | | 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 { |
| | |
| | | .await |
| | | { |
| | | Ok(outcome) => { |
| | | let mut published_device_outputs = HashSet::new(); |
| | | for output in &outcome.device_outputs { |
| | | if !should_publish_device_output(&mut published_device_outputs, output) { |
| | | continue; |
| | | } |
| | | if let Err(error) = sink |
| | | .publish_device_output(call_id, trace_id, &turn.turn_id, output) |
| | | .await |
| | |
| | | } |
| | | Some("device_output") => { |
| | | if let Some(output) = event.device_output.as_ref() { |
| | | if !should_publish_device_output(&mut state.published_device_output_ids, output) { |
| | | return Ok(()); |
| | | } |
| | | sink.publish_device_output(call_id, trace_id, &turn.turn_id, output) |
| | | .await?; |
| | | state.device_output_count = state.device_output_count.saturating_add(1); |
| | |
| | | completed: bool, |
| | | audio_chunk_count: u64, |
| | | device_output_count: u64, |
| | | published_device_output_ids: HashSet<String>, |
| | | encoded_audio_buffer: Vec<u8>, |
| | | pcm_stream_decoder: Option<audio::PcmS16leStreamDecoder>, |
| | | pcm_stream_network_chunk_count: u64, |
| | |
| | | completed: false, |
| | | audio_chunk_count: 0, |
| | | device_output_count: 0, |
| | | published_device_output_ids: HashSet::new(), |
| | | encoded_audio_buffer: Vec::new(), |
| | | pcm_stream_decoder: None, |
| | | pcm_stream_network_chunk_count: 0, |
| | |
| | | params: Option<serde_json::Value>, |
| | | } |
| | | |
| | | fn should_publish_device_output( |
| | | published_ids: &mut HashSet<String>, |
| | | output: &RuntimeTurnDeviceOutput, |
| | | ) -> bool { |
| | | let Some(command_id) = output |
| | | .command_id |
| | | .as_deref() |
| | | .map(str::trim) |
| | | .filter(|value| !value.is_empty()) |
| | | else { |
| | | return false; |
| | | }; |
| | | if output |
| | | .command_code |
| | | .as_deref() |
| | | .map(str::trim) |
| | | .filter(|value| !value.is_empty()) |
| | | .is_none() |
| | | { |
| | | return false; |
| | | } |
| | | published_ids.insert(command_id.to_string()) |
| | | } |
| | | |
| | | fn require_safe_segment(value: &str) -> Result<()> { |
| | | if value.is_empty() |
| | | || value.contains('/') |
| | |
| | | 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); |
| | | let ingress_metadata = match AudioIngressMetadata::from_participant( |
| | | &participant.attributes(), |
| | | ) { |
| | | Ok(value) => value, |
| | | Err(reason) => { |
| | | warn!(call_id = %call_id, trace_id = %trace_id, |
| | | turn_id = %turn_id, metadata_status = "invalid", reason = reason, |
| | | "runtime helper ignored invalid audio ingress metadata"); |
| | | None |
| | | } |
| | | }; |
| | | match RealtimeAsrUpload::start( |
| | | http.clone(), |
| | | turn_bridge_config.realtime_asr_config(), |
| | |
| | | &trace_id, |
| | | &turn_id, |
| | | &vad.speech_samples, |
| | | ingress_metadata.as_ref(), |
| | | ) { |
| | | Ok(upload) => { |
| | | info!( |
| | |
| | | #[cfg(test)] |
| | | mod tests { |
| | | use super::{ |
| | | ReplyChunkMarker, ReplyChunkMarkerState, RuntimeTurnStreamEvent, RuntimeTurnStreamState, |
| | | RuntimeTurnStreamTimingPhase, runtime_session_nonce_hash, |
| | | ReplyChunkMarker, ReplyChunkMarkerState, RuntimeTurnDeviceOutput, RuntimeTurnStreamEvent, |
| | | RuntimeTurnStreamState, RuntimeTurnStreamTimingPhase, runtime_session_nonce_hash, |
| | | should_publish_device_output, |
| | | }; |
| | | use std::collections::HashSet; |
| | | |
| | | #[test] |
| | | fn reply_chunk_marker_state_emits_turn_first_once_and_later_segment_first_once() { |
| | |
| | | ); |
| | | assert!(state.record_m7().is_none()); |
| | | } |
| | | |
| | | #[test] |
| | | fn device_output_contract_is_reliable_and_deduplicated() { |
| | | let output = RuntimeTurnDeviceOutput { |
| | | command_id: Some("cmd-1".to_string()), |
| | | command_code: Some("custom.app.DeviceLevelChange".to_string()), |
| | | params: Some(serde_json::json!({"level": 1})), |
| | | }; |
| | | let mut published = HashSet::new(); |
| | | assert!(should_publish_device_output(&mut published, &output)); |
| | | assert!(!should_publish_device_output(&mut published, &output)); |
| | | assert_eq!(published.len(), 1); |
| | | } |
| | | |
| | | #[test] |
| | | fn device_output_invalid_or_missing_command_is_fail_closed() { |
| | | for output in [ |
| | | RuntimeTurnDeviceOutput { |
| | | command_id: None, |
| | | command_code: Some("custom.app.DeviceLevelChange".to_string()), |
| | | params: None, |
| | | }, |
| | | RuntimeTurnDeviceOutput { |
| | | command_id: Some("cmd-1".to_string()), |
| | | command_code: None, |
| | | params: None, |
| | | }, |
| | | ] { |
| | | let mut published = HashSet::new(); |
| | | assert!(!should_publish_device_output(&mut published, &output)); |
| | | assert!(published.is_empty()); |
| | | } |
| | | } |
| | | } |