| | |
| | | mod audio; |
| | | mod service; |
| | | |
| | | use std::{env, sync::Arc, time::Duration}; |
| | | use std::{ |
| | | env, fs, |
| | | path::{Path, PathBuf}, |
| | | sync::{ |
| | | Arc, |
| | | atomic::{AtomicBool, Ordering}, |
| | | }, |
| | | time::{Duration, Instant, SystemTime, UNIX_EPOCH}, |
| | | }; |
| | | |
| | | use anyhow::{Context, Result, anyhow}; |
| | | use audio::load_pre_recorded_frames; |
| | | use audio::{AudioDiagnostics, PcmFrame, load_pre_recorded_frames}; |
| | | use base64::{Engine as _, engine::general_purpose}; |
| | | use futures_util::StreamExt; |
| | | use libwebrtc::{ |
| | | audio_source::native::NativeAudioSource, |
| | | audio_stream::native::NativeAudioStream, |
| | | prelude::{AudioFrame, AudioSourceOptions, RtcAudioSource}, |
| | | }; |
| | | use livekit::{ |
| | | options::TrackPublishOptions, |
| | | prelude::{DataPacket, LocalAudioTrack, LocalTrack, ParticipantIdentity, Room, RoomOptions}, |
| | | prelude::{ |
| | | DataPacket, LocalAudioTrack, LocalTrack, ParticipantIdentity, RemoteAudioTrack, |
| | | RemoteTrack, Room, RoomEvent, RoomOptions, |
| | | }, |
| | | }; |
| | | use reqwest::Client; |
| | | use serde::Serialize; |
| | | use tokio::time::sleep; |
| | | use serde::{Deserialize, Serialize}; |
| | | use serde_json::json; |
| | | use tokio::time::{sleep, sleep_until}; |
| | | use tokio::{sync::mpsc::UnboundedReceiver, task::JoinHandle}; |
| | | use tracing::{info, warn}; |
| | | |
| | | const TARGET_SAMPLE_RATE_HZ: u32 = 48_000; |
| | | const TARGET_NUM_CHANNELS: u16 = 1; |
| | | const TRACK_NAME: &str = "bot-main-audio"; |
| | | const DEVICE_OUTPUT_TOPIC: &str = "device_output"; |
| | | |
| | | #[tokio::main(flavor = "multi_thread")] |
| | | async fn main() -> Result<()> { |
| | | init_tracing(); |
| | | if service::service_mode_enabled() { |
| | | return service::run_service().await; |
| | | } |
| | | run_worker().await |
| | | } |
| | | |
| | | async fn run_worker() -> Result<()> { |
| | | let config = Config::from_env()?; |
| | | let http = Client::builder() |
| | | .use_rustls_tls() |
| | |
| | | config.greeting_audio_url.as_deref(), |
| | | TARGET_SAMPLE_RATE_HZ, |
| | | TARGET_NUM_CHANNELS, |
| | | config.audio_debug_dump_dir.as_deref(), |
| | | &config.call_id, |
| | | "greeting", |
| | | ) |
| | | .await?; |
| | | |
| | | let (room, _events) = Room::connect( |
| | | let (room, events) = Room::connect( |
| | | config.livekit_url.as_str(), |
| | | config.bot_token.as_str(), |
| | | RoomOptions::default(), |
| | |
| | | info!( |
| | | call_id = %config.call_id, |
| | | trace_id = %config.trace_id, |
| | | room_id = %config.room_id, |
| | | room_alias = %redact(&config.room_id), |
| | | participant_alias = %redact(&config.bot_participant_identity), |
| | | greeting_source = %config.greeting_source, |
| | | "combrabo voice runtime helper connected" |
| | | ); |
| | | emit_activity( |
| | | &config.call_id, |
| | | &config.trace_id, |
| | | None, |
| | | "bot_participant_joined", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({"participantAlias": redact(&config.bot_participant_identity)}), |
| | | ); |
| | | |
| | | let vad_enabled_gate = Arc::new(AtomicBool::new( |
| | | !config.simple_vad_enabled || !config.simple_vad_gate_until_greeting_done, |
| | | )); |
| | | if config.simple_vad_enabled && config.simple_vad_gate_until_greeting_done { |
| | | info!( |
| | | call_id = %config.call_id, |
| | | trace_id = %config.trace_id, |
| | | post_greeting_delay_ms = config.simple_vad_post_greeting_delay_ms, |
| | | "runtime helper vad_disabled_greeting" |
| | | ); |
| | | } else if config.simple_vad_enabled { |
| | | info!( |
| | | call_id = %config.call_id, |
| | | trace_id = %config.trace_id, |
| | | "runtime helper vad_enabled" |
| | | ); |
| | | emit_activity( |
| | | &config.call_id, |
| | | &config.trace_id, |
| | | None, |
| | | "vad_enabled", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({"reason": "gate_disabled"}), |
| | | ); |
| | | } |
| | | |
| | | let sink = BotAudioOutputSink::publish( |
| | | room.clone(), |
| | | &config.room_id, |
| | | &config.bot_participant_identity, |
| | | &config.call_id, |
| | | &config.trace_id, |
| | | TRACK_NAME, |
| | | TARGET_SAMPLE_RATE_HZ, |
| | | u32::from(TARGET_NUM_CHANNELS), |
| | | config.user_participant_identity.clone(), |
| | | ) |
| | | .await?; |
| | | let sink = Arc::new(sink); |
| | | |
| | | if config.device_output_smoke_enabled { |
| | | publish_device_output_smoke(room.as_ref(), &config).await?; |
| | | } |
| | | let user_audio_observer = spawn_user_audio_observer( |
| | | events, |
| | | &config, |
| | | vad_enabled_gate.clone(), |
| | | http.clone(), |
| | | sink.clone(), |
| | | ); |
| | | |
| | | if let Some(frames) = greeting_frames { |
| | | if let Some(loaded_audio) = greeting_frames { |
| | | let diagnostics = loaded_audio.diagnostics; |
| | | let frames = loaded_audio.frames; |
| | | log_audio_diagnostics(&config, &diagnostics); |
| | | info!( |
| | | call_id = %config.call_id, |
| | | frame_count = frames.len(), |
| | | theoretical_duration_ms = diagnostics.target_duration_ms, |
| | | greeting_source = %config.greeting_source, |
| | | "runtime helper starting greeting playback" |
| | | ); |
| | | for frame in &frames { |
| | | emit_activity( |
| | | &config.call_id, |
| | | &config.trace_id, |
| | | None, |
| | | "greeting_write_started", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "frameCount": frames.len(), |
| | | "greetingTheoreticalDurationMs": diagnostics.target_duration_ms, |
| | | }), |
| | | ); |
| | | let playback_started_at = Instant::now(); |
| | | let pacing_started_at = tokio::time::Instant::now(); |
| | | for (index, frame) in frames.iter().enumerate() { |
| | | sink.write_pcm_frame(frame).await?; |
| | | sleep(Duration::from_millis(20)).await; |
| | | sleep_until(pacing_started_at + Duration::from_millis(((index + 1) as u64) * 20)).await; |
| | | } |
| | | let playback_wall_ms = playback_started_at.elapsed().as_millis() as i64; |
| | | let drift_ms = playback_wall_ms - diagnostics.target_duration_ms as i64; |
| | | sink.clear_buffer(); |
| | | info!( |
| | | call_id = %config.call_id, |
| | | greeting_source = %config.greeting_source, |
| | | frame_count = frames.len(), |
| | | theoretical_duration_ms = diagnostics.target_duration_ms, |
| | | push_wall_duration_ms = playback_wall_ms, |
| | | push_drift_ms = drift_ms, |
| | | "runtime helper finished greeting playback" |
| | | ); |
| | | emit_activity( |
| | | &config.call_id, |
| | | &config.trace_id, |
| | | None, |
| | | "greeting_write_finished", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "frameCount": frames.len(), |
| | | "greetingTheoreticalDurationMs": diagnostics.target_duration_ms, |
| | | "greetingPushWallDurationMs": playback_wall_ms, |
| | | "greetingPushDriftMs": drift_ms, |
| | | }), |
| | | ); |
| | | schedule_vad_gate_enable( |
| | | vad_enabled_gate, |
| | | &config, |
| | | config.simple_vad_post_greeting_delay_ms, |
| | | "greeting_finished", |
| | | ); |
| | | } else { |
| | | warn!( |
| | |
| | | greeting_source = %config.greeting_source, |
| | | "runtime helper started without greeting audio; keeping published track alive" |
| | | ); |
| | | schedule_vad_gate_enable(vad_enabled_gate, &config, 0, "no_greeting_audio"); |
| | | } |
| | | |
| | | wait_for_shutdown_signal().await?; |
| | | user_audio_observer.abort(); |
| | | let _ = user_audio_observer.await; |
| | | |
| | | if let Err(error) = sink.close().await { |
| | | warn!( |
| | |
| | | room_id: String, |
| | | bot_token: String, |
| | | bot_participant_identity: String, |
| | | user_participant_identity: Option<String>, |
| | | greeting_source: String, |
| | | greeting_audio_file: Option<String>, |
| | | greeting_audio_url: Option<String>, |
| | | role_id: String, |
| | | device_output_smoke_enabled: bool, |
| | | device_output_destination_identities: Vec<ParticipantIdentity>, |
| | | audio_debug_dump_dir: Option<String>, |
| | | runtime_turn_bridge_url: Option<String>, |
| | | runtime_turn_bridge_token: Option<String>, |
| | | runtime_turn_bridge_mode: String, |
| | | runtime_turn_artifact_dir: Option<String>, |
| | | runtime_session_nonce: Option<String>, |
| | | user_audio_observer_enabled: bool, |
| | | simple_vad_enabled: bool, |
| | | simple_vad_gate_until_greeting_done: bool, |
| | | simple_vad_post_greeting_delay_ms: u64, |
| | | simple_vad_config: SimpleVadConfig, |
| | | } |
| | | |
| | | #[derive(Clone)] |
| | | struct TurnBridgeConfig { |
| | | bridge_url: Option<String>, |
| | | bridge_token: Option<String>, |
| | | bridge_mode: String, |
| | | artifact_dir: Option<String>, |
| | | runtime_session_nonce: Option<String>, |
| | | audio_debug_dump_dir: Option<String>, |
| | | } |
| | | |
| | | impl TurnBridgeConfig { |
| | | fn from_config(config: &Config) -> Self { |
| | | Self { |
| | | bridge_url: config.runtime_turn_bridge_url.clone(), |
| | | bridge_token: config.runtime_turn_bridge_token.clone(), |
| | | bridge_mode: config.runtime_turn_bridge_mode.clone(), |
| | | artifact_dir: config.runtime_turn_artifact_dir.clone(), |
| | | runtime_session_nonce: config.runtime_session_nonce.clone(), |
| | | audio_debug_dump_dir: config.audio_debug_dump_dir.clone(), |
| | | } |
| | | } |
| | | |
| | | fn is_ready(&self) -> bool { |
| | | self.bridge_url |
| | | .as_ref() |
| | | .is_some_and(|value| !value.is_empty()) |
| | | && self |
| | | .bridge_token |
| | | .as_ref() |
| | | .is_some_and(|value| !value.is_empty()) |
| | | && self |
| | | .artifact_dir |
| | | .as_ref() |
| | | .is_some_and(|value| !value.is_empty()) |
| | | && self |
| | | .runtime_session_nonce |
| | | .as_ref() |
| | | .is_some_and(|value| !value.is_empty()) |
| | | } |
| | | |
| | | fn is_stream_mode(&self) -> bool { |
| | | self.bridge_mode.eq_ignore_ascii_case("stream") |
| | | || self |
| | | .bridge_url |
| | | .as_deref() |
| | | .is_some_and(|value| value.trim_end_matches('/').ends_with("/stream")) |
| | | } |
| | | } |
| | | |
| | | impl Config { |
| | |
| | | room_id: required_env("CV_LIVEKIT_ROOM_ID")?, |
| | | bot_token: required_env("CV_LIVEKIT_BOT_TOKEN")?, |
| | | bot_participant_identity: required_env("CV_LIVEKIT_BOT_PARTICIPANT_IDENTITY")?, |
| | | user_participant_identity: optional_env("CV_LIVEKIT_USER_PARTICIPANT_IDENTITY"), |
| | | greeting_source: env::var("CV_GREETING_SOURCE") |
| | | .unwrap_or_else(|_| "no_audio".to_string()), |
| | | greeting_audio_file: optional_env("CV_GREETING_AUDIO_FILE"), |
| | | greeting_audio_url: optional_env("CV_GREETING_AUDIO_URL"), |
| | | role_id: env::var("CV_ROLE_ID").unwrap_or_else(|_| "90".to_string()), |
| | | device_output_smoke_enabled: bool_env("CV_DEVICE_OUTPUT_SMOKE_ENABLED"), |
| | | device_output_destination_identities: device_output_destinations( |
| | | optional_env("CV_DEVICE_OUTPUT_DESTINATION_IDENTITIES"), |
| | | optional_env("CV_LIVEKIT_USER_PARTICIPANT_IDENTITY"), |
| | | ), |
| | | audio_debug_dump_dir: optional_env("CV_AUDIO_DEBUG_DUMP_DIR"), |
| | | runtime_turn_bridge_url: optional_env("CV_RUNTIME_TURN_BRIDGE_URL"), |
| | | runtime_turn_bridge_token: optional_env("CV_RUNTIME_TURN_BRIDGE_TOKEN"), |
| | | runtime_turn_bridge_mode: env::var("CV_RUNTIME_TURN_BRIDGE_MODE") |
| | | .unwrap_or_else(|_| "json".to_string()), |
| | | runtime_turn_artifact_dir: optional_env("CV_RUNTIME_TURN_ARTIFACT_DIR"), |
| | | runtime_session_nonce: optional_env("CV_RUNTIME_SESSION_NONCE"), |
| | | user_audio_observer_enabled: bool_env("CV_ENABLE_USER_AUDIO_OBSERVER", true), |
| | | simple_vad_enabled: bool_env("CV_ENABLE_SIMPLE_VAD", true), |
| | | simple_vad_gate_until_greeting_done: bool_env("CV_VAD_GATE_UNTIL_GREETING_DONE", true), |
| | | simple_vad_post_greeting_delay_ms: u64_env("CV_VAD_POST_GREETING_DELAY_MS", 800), |
| | | simple_vad_config: SimpleVadConfig::from_env(), |
| | | }) |
| | | } |
| | | } |
| | | |
| | | fn schedule_vad_gate_enable( |
| | | gate: Arc<AtomicBool>, |
| | | config: &Config, |
| | | delay_ms: u64, |
| | | reason: &'static str, |
| | | ) { |
| | | if !config.simple_vad_enabled || !config.simple_vad_gate_until_greeting_done { |
| | | return; |
| | | } |
| | | |
| | | let call_id = config.call_id.clone(); |
| | | let trace_id = config.trace_id.clone(); |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | delay_ms, |
| | | reason, |
| | | "runtime helper vad_enable_scheduled" |
| | | ); |
| | | emit_activity( |
| | | &call_id, |
| | | &trace_id, |
| | | None, |
| | | "vad_enable_scheduled", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "vadEnableDelayMs": delay_ms, |
| | | "reason": reason, |
| | | }), |
| | | ); |
| | | |
| | | tokio::spawn(async move { |
| | | if delay_ms > 0 { |
| | | sleep(Duration::from_millis(delay_ms)).await; |
| | | } |
| | | gate.store(true, Ordering::Release); |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | delay_ms, |
| | | reason, |
| | | "runtime helper vad_enabled" |
| | | ); |
| | | emit_activity( |
| | | &call_id, |
| | | &trace_id, |
| | | None, |
| | | "vad_enabled", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "vadEnableDelayMs": delay_ms, |
| | | "reason": reason, |
| | | }), |
| | | ); |
| | | }); |
| | | } |
| | | |
| | | fn log_audio_diagnostics(config: &Config, diagnostics: &AudioDiagnostics) { |
| | | info!( |
| | | call_id = %config.call_id, |
| | | trace_id = %config.trace_id, |
| | | greeting_source = %config.greeting_source, |
| | | source_kind = diagnostics.source_kind, |
| | | source_format = diagnostics.source_format, |
| | | source_bytes = diagnostics.source_bytes, |
| | | source_sample_rate_hz = diagnostics.source_sample_rate_hz, |
| | | source_num_channels = diagnostics.source_num_channels, |
| | | decoded_sample_count = diagnostics.decoded_sample_count, |
| | | decoded_duration_ms = diagnostics.decoded_duration_ms, |
| | | target_sample_rate_hz = diagnostics.target_sample_rate_hz, |
| | | target_num_channels = diagnostics.target_num_channels, |
| | | target_sample_count = diagnostics.target_sample_count, |
| | | target_duration_ms = diagnostics.target_duration_ms, |
| | | frame_count = diagnostics.frame_count, |
| | | rms = diagnostics.rms, |
| | | peak = diagnostics.peak, |
| | | clipped_sample_count = diagnostics.clipped_sample_count, |
| | | silence_ratio = diagnostics.silence_ratio, |
| | | mp3_skipped_data_count = diagnostics.mp3_skipped_data_count, |
| | | mp3_insufficient_data_count = diagnostics.mp3_insufficient_data_count, |
| | | debug_source_path = diagnostics.debug_source_path.as_deref().unwrap_or(""), |
| | | debug_pcm_wav_path = diagnostics.debug_pcm_wav_path.as_deref().unwrap_or(""), |
| | | "runtime helper greeting audio quality diagnostics" |
| | | ); |
| | | } |
| | | |
| | | fn required_env(key: &str) -> Result<String> { |
| | |
| | | }) |
| | | } |
| | | |
| | | fn bool_env(key: &str) -> bool { |
| | | matches!( |
| | | env::var(key) |
| | | .unwrap_or_default() |
| | | .trim() |
| | | .to_ascii_lowercase() |
| | | .as_str(), |
| | | "1" | "true" | "yes" | "y" | "on" |
| | | ) |
| | | fn bool_env(key: &str, default_value: bool) -> bool { |
| | | match env::var(key) { |
| | | Ok(value) => match value.trim().to_ascii_lowercase().as_str() { |
| | | "1" | "true" | "yes" | "on" => true, |
| | | "0" | "false" | "no" | "off" => false, |
| | | _ => default_value, |
| | | }, |
| | | Err(_) => default_value, |
| | | } |
| | | } |
| | | |
| | | fn device_output_destinations( |
| | | configured: Option<String>, |
| | | user_identity: Option<String>, |
| | | ) -> Vec<ParticipantIdentity> { |
| | | let identities = configured |
| | | .filter(|value| !value.trim().is_empty()) |
| | | .map(|value| { |
| | | value |
| | | .split(',') |
| | | .map(str::trim) |
| | | .filter(|identity| !identity.is_empty()) |
| | | .map(ToOwned::to_owned) |
| | | .collect::<Vec<_>>() |
| | | }) |
| | | .or_else(|| user_identity.map(|identity| vec![identity])) |
| | | .unwrap_or_default(); |
| | | identities.into_iter().map(Into::into).collect() |
| | | fn f64_env(key: &str, default_value: f64) -> f64 { |
| | | env::var(key) |
| | | .ok() |
| | | .and_then(|value| value.trim().parse::<f64>().ok()) |
| | | .filter(|value| value.is_finite() && *value >= 0.0) |
| | | .unwrap_or(default_value) |
| | | } |
| | | |
| | | fn u64_env(key: &str, default_value: u64) -> u64 { |
| | | env::var(key) |
| | | .ok() |
| | | .and_then(|value| value.trim().parse::<u64>().ok()) |
| | | .unwrap_or(default_value) |
| | | } |
| | | |
| | | fn u32_env(key: &str, default_value: u32) -> u32 { |
| | | env::var(key) |
| | | .ok() |
| | | .and_then(|value| value.trim().parse::<u32>().ok()) |
| | | .unwrap_or(default_value) |
| | | } |
| | | |
| | | fn redact(value: &str) -> String { |
| | |
| | | format!("{}***{}", &value[..4], &value[value.len() - 4..]) |
| | | } |
| | | |
| | | #[derive(Serialize)] |
| | | struct DeviceOutputMessage { |
| | | #[serde(rename = "type")] |
| | | message_type: &'static str, |
| | | #[serde(rename = "schemaVersion")] |
| | | schema_version: &'static str, |
| | | #[serde(rename = "callId")] |
| | | fn safe_error(value: &str) -> String { |
| | | let sanitized = value.replace(['\r', '\n'], " "); |
| | | let trimmed = sanitized.trim(); |
| | | let mut output: String = trimmed.chars().take(180).collect(); |
| | | if trimmed.chars().count() > 180 { |
| | | output.push_str("..."); |
| | | } |
| | | output |
| | | } |
| | | |
| | | fn current_time_millis() -> u64 { |
| | | SystemTime::now() |
| | | .duration_since(UNIX_EPOCH) |
| | | .map(|value| value.as_millis() as u64) |
| | | .unwrap_or_default() |
| | | } |
| | | |
| | | fn emit_activity( |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | turn_id: Option<&str>, |
| | | event_name: &str, |
| | | result: &str, |
| | | reason_code: Option<&str>, |
| | | retryable: Option<bool>, |
| | | extension: serde_json::Value, |
| | | ) { |
| | | let payload = json!({ |
| | | "type": "cv_activity", |
| | | "callId": call_id, |
| | | "traceId": trace_id, |
| | | "turnId": turn_id, |
| | | "eventName": event_name, |
| | | "result": result, |
| | | "reasonCode": reason_code, |
| | | "retryable": retryable, |
| | | "extension": extension, |
| | | }); |
| | | println!("{payload}"); |
| | | } |
| | | |
| | | fn spawn_user_audio_observer( |
| | | events: UnboundedReceiver<RoomEvent>, |
| | | config: &Config, |
| | | vad_enabled_gate: Arc<AtomicBool>, |
| | | http: Client, |
| | | sink: Arc<BotAudioOutputSink>, |
| | | ) -> JoinHandle<()> { |
| | | let call_id = config.call_id.clone(); |
| | | let trace_id = config.trace_id.clone(); |
| | | 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 turn_bridge_config = TurnBridgeConfig::from_config(config); |
| | | |
| | | tokio::spawn(async move { |
| | | if !enabled { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | "runtime helper user audio observer disabled" |
| | | ); |
| | | return; |
| | | } |
| | | observe_user_audio_events( |
| | | events, |
| | | call_id, |
| | | trace_id, |
| | | simple_vad_enabled, |
| | | simple_vad_config, |
| | | vad_enabled_gate, |
| | | turn_bridge_config, |
| | | http, |
| | | sink, |
| | | ) |
| | | .await; |
| | | }) |
| | | } |
| | | |
| | | async fn observe_user_audio_events( |
| | | mut events: UnboundedReceiver<RoomEvent>, |
| | | call_id: String, |
| | | #[serde(rename = "roleId")] |
| | | role_id: String, |
| | | #[serde(rename = "traceId")] |
| | | trace_id: String, |
| | | #[serde(rename = "ackMode")] |
| | | ack_mode: &'static str, |
| | | #[serde(rename = "sensorInstructions")] |
| | | sensor_instructions: Vec<SensorInstruction>, |
| | | simple_vad_enabled: bool, |
| | | simple_vad_config: SimpleVadConfig, |
| | | vad_enabled_gate: Arc<AtomicBool>, |
| | | turn_bridge_config: TurnBridgeConfig, |
| | | http: Client, |
| | | sink: Arc<BotAudioOutputSink>, |
| | | ) { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | simple_vad_enabled, |
| | | vad_rms_threshold = simple_vad_config.rms_threshold, |
| | | vad_peak_threshold = simple_vad_config.peak_threshold, |
| | | vad_start_frames = simple_vad_config.start_frames, |
| | | vad_end_silence_ms = simple_vad_config.end_silence_ms, |
| | | vad_min_speech_ms = simple_vad_config.min_speech_ms, |
| | | vad_max_turn_ms = simple_vad_config.max_turn_ms, |
| | | vad_initial_ignore_ms = simple_vad_config.initial_ignore_ms, |
| | | vad_gate_enabled = vad_enabled_gate.load(Ordering::Acquire), |
| | | "runtime helper user_track_subscribe_requested" |
| | | ); |
| | | |
| | | while let Some(event) = events.recv().await { |
| | | match event { |
| | | RoomEvent::TrackSubscribed { |
| | | track: RemoteTrack::Audio(track), |
| | | publication: _, |
| | | participant, |
| | | } => { |
| | | let participant_alias = redact(&participant.identity().to_string()); |
| | | let track_sid_alias = redact(&track.sid().to_string()); |
| | | let track_name = track.name(); |
| | | let track_source = format!("{:?}", track.source()); |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | participant_alias = %participant_alias, |
| | | track_sid_alias = %track_sid_alias, |
| | | track_name = %track_name, |
| | | track_source = %track_source, |
| | | "runtime helper user_track_subscribed" |
| | | ); |
| | | spawn_user_audio_frame_observer( |
| | | track, |
| | | call_id.clone(), |
| | | trace_id.clone(), |
| | | participant_alias, |
| | | track_sid_alias, |
| | | simple_vad_enabled, |
| | | simple_vad_config.clone(), |
| | | vad_enabled_gate.clone(), |
| | | turn_bridge_config.clone(), |
| | | http.clone(), |
| | | sink.clone(), |
| | | ); |
| | | } |
| | | RoomEvent::TrackSubscribed { |
| | | track: RemoteTrack::Video(track), |
| | | publication: _, |
| | | participant, |
| | | } => { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | participant_alias = %redact(&participant.identity().to_string()), |
| | | track_sid_alias = %redact(&track.sid().to_string()), |
| | | "runtime helper ignored non-audio subscribed track" |
| | | ); |
| | | } |
| | | RoomEvent::TrackSubscriptionFailed { |
| | | participant, |
| | | error, |
| | | track_sid, |
| | | } => { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | participant_alias = %redact(&participant.identity().to_string()), |
| | | track_sid_alias = %redact(&track_sid.to_string()), |
| | | error = %error, |
| | | "runtime helper user_track_subscription_failed" |
| | | ); |
| | | } |
| | | RoomEvent::Disconnected { reason } => { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | reason = ?reason, |
| | | "runtime helper room event stream disconnected" |
| | | ); |
| | | break; |
| | | } |
| | | _ => {} |
| | | } |
| | | } |
| | | } |
| | | |
| | | #[derive(Serialize)] |
| | | struct SensorInstruction { |
| | | #[serde(rename = "commandId")] |
| | | command_id: String, |
| | | #[serde(rename = "sensorType")] |
| | | sensor_type: &'static str, |
| | | #[serde(rename = "operationType")] |
| | | operation_type: &'static str, |
| | | step: i32, |
| | | #[serde(rename = "durationSec", skip_serializing_if = "Option::is_none")] |
| | | duration_sec: Option<i32>, |
| | | extension: String, |
| | | async fn handle_finished_turn( |
| | | http: &Client, |
| | | bridge_config: &TurnBridgeConfig, |
| | | sink: &BotAudioOutputSink, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | turn: FinishedSpeechTurn, |
| | | ) { |
| | | let turn_pipeline_started_at = Instant::now(); |
| | | if !bridge_config.is_ready() { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | turn_index = turn.turn_index, |
| | | duration_ms = turn.duration_ms, |
| | | "runtime helper turn_bridge_skipped" |
| | | ); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "turn_bridge_skipped", |
| | | "skipped", |
| | | Some("TURN_BRIDGE_NOT_CONFIGURED"), |
| | | Some(true), |
| | | json!({ |
| | | "durationMs": turn.duration_ms, |
| | | "frameCount": turn.frame_count, |
| | | "sampleCount": turn.sample_count, |
| | | }), |
| | | ); |
| | | return; |
| | | } |
| | | |
| | | let artifact_root = PathBuf::from(bridge_config.artifact_dir.as_deref().unwrap_or_default()); |
| | | match write_user_turn_artifact(&artifact_root, call_id, &turn) { |
| | | Ok((path_ref, byte_size)) => { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | turn_index = turn.turn_index, |
| | | duration_ms = turn.duration_ms, |
| | | frame_count = turn.frame_count, |
| | | sample_count = turn.sample_count, |
| | | byte_size, |
| | | end_reason = %turn.end_reason, |
| | | "runtime helper turn_artifact_written" |
| | | ); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "turn_bridge_requested", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "turnDurationMs": turn.duration_ms, |
| | | "turnArtifactBytes": byte_size, |
| | | "frameCount": turn.frame_count, |
| | | "sampleCount": turn.sample_count, |
| | | "endReason": turn.end_reason, |
| | | }), |
| | | ); |
| | | if bridge_config.is_stream_mode() { |
| | | match request_turn_bridge_stream( |
| | | http, |
| | | bridge_config, |
| | | sink, |
| | | call_id, |
| | | trace_id, |
| | | &turn, |
| | | &path_ref, |
| | | byte_size, |
| | | turn_pipeline_started_at, |
| | | ) |
| | | .await |
| | | { |
| | | Ok(outcome) => { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | audio_chunk_count = outcome.audio_chunk_count, |
| | | device_output_count = outcome.device_output_count, |
| | | "runtime helper turn_stream_completed" |
| | | ); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "turn_completed", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64, |
| | | "audioChunkCount": outcome.audio_chunk_count, |
| | | "deviceOutputCount": outcome.device_output_count, |
| | | "replyPlaybackMode": outcome.reply_playback_mode, |
| | | }), |
| | | ); |
| | | } |
| | | Err(error) => { |
| | | let safe = safe_error(&error.to_string()); |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | error = %safe, |
| | | "runtime helper turn_stream_failed" |
| | | ); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "turn_failed", |
| | | "failed", |
| | | Some("TURN_STREAM_FAILED"), |
| | | Some(true), |
| | | json!({ |
| | | "stage": "turn_bridge_stream", |
| | | "error": safe, |
| | | }), |
| | | ); |
| | | } |
| | | } |
| | | return; |
| | | } |
| | | match request_turn_bridge( |
| | | http, |
| | | bridge_config, |
| | | call_id, |
| | | trace_id, |
| | | &turn, |
| | | &path_ref, |
| | | byte_size, |
| | | turn_pipeline_started_at, |
| | | ) |
| | | .await |
| | | { |
| | | Ok(outcome) => { |
| | | for output in &outcome.device_outputs { |
| | | if let Err(error) = sink |
| | | .publish_device_output(call_id, trace_id, &turn.turn_id, output) |
| | | .await |
| | | { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | reason_code = "BOT_DATA_WRITE_FAILED", |
| | | error = %safe_error(&error.to_string()), |
| | | "runtime helper device_output_failed" |
| | | ); |
| | | } |
| | | } |
| | | let Some(reply_audio_artifact) = outcome.reply_audio_artifact else { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | reason_code = "REPLY_AUDIO_MISSING", |
| | | "runtime helper turn_failed" |
| | | ); |
| | | return; |
| | | }; |
| | | if let Err(error) = write_reply_audio_artifact( |
| | | http, |
| | | bridge_config, |
| | | sink, |
| | | call_id, |
| | | trace_id, |
| | | &turn, |
| | | &reply_audio_artifact, |
| | | turn_pipeline_started_at, |
| | | ) |
| | | .await |
| | | { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | reason_code = "BOT_AUDIO_WRITE_FAILED", |
| | | error = %safe_error(&error.to_string()), |
| | | "runtime helper turn_failed" |
| | | ); |
| | | return; |
| | | } |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | "runtime helper turn_completed" |
| | | ); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "turn_completed", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64, |
| | | }), |
| | | ); |
| | | } |
| | | Err(error) => warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | error = %safe_error(&error.to_string()), |
| | | "runtime helper turn_bridge_failed" |
| | | ), |
| | | } |
| | | } |
| | | Err(error) => warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | error = %safe_error(&error.to_string()), |
| | | "runtime helper turn_artifact_write_failed" |
| | | ), |
| | | } |
| | | } |
| | | |
| | | async fn publish_device_output_smoke(room: &Room, config: &Config) -> Result<()> { |
| | | let payload = build_device_output_smoke_payload(config); |
| | | let payload_json = |
| | | serde_json::to_vec(&payload).context("failed to serialize device output smoke payload")?; |
| | | let destination_count = config.device_output_destination_identities.len(); |
| | | room.local_participant() |
| | | .publish_data(DataPacket { |
| | | reliable: true, |
| | | payload: payload_json, |
| | | topic: Some(DEVICE_OUTPUT_TOPIC.to_string()), |
| | | destination_identities: config.device_output_destination_identities.clone(), |
| | | }) |
| | | fn write_user_turn_artifact( |
| | | artifact_root: &Path, |
| | | call_id: &str, |
| | | turn: &FinishedSpeechTurn, |
| | | ) -> Result<(String, u64)> { |
| | | require_safe_segment(call_id)?; |
| | | require_safe_segment(&turn.turn_id)?; |
| | | if turn.samples.is_empty() { |
| | | return Err(anyhow!("empty turn samples")); |
| | | } |
| | | let path_ref = format!("{}/{}/user.wav", call_id, turn.turn_id); |
| | | let output_path = normalize_path_lexically(&artifact_root.join(&path_ref)); |
| | | let root = normalize_path_lexically(artifact_root); |
| | | if !output_path.starts_with(&root) { |
| | | return Err(anyhow!("turn artifact path escapes root")); |
| | | } |
| | | if let Some(parent) = output_path.parent() { |
| | | fs::create_dir_all(parent).context("failed to create turn artifact dir")?; |
| | | } |
| | | audio::write_pcm_wav( |
| | | &output_path, |
| | | &turn.samples, |
| | | TARGET_SAMPLE_RATE_HZ, |
| | | TARGET_NUM_CHANNELS, |
| | | )?; |
| | | let byte_size = fs::metadata(&output_path) |
| | | .context("failed to stat turn artifact")? |
| | | .len(); |
| | | Ok((path_ref, byte_size)) |
| | | } |
| | | |
| | | async fn request_turn_bridge_stream( |
| | | http: &Client, |
| | | bridge_config: &TurnBridgeConfig, |
| | | sink: &BotAudioOutputSink, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | turn: &FinishedSpeechTurn, |
| | | path_ref: &str, |
| | | byte_size: u64, |
| | | turn_pipeline_started_at: Instant, |
| | | ) -> Result<RuntimeTurnStreamOutcome> { |
| | | let bridge_started_at = Instant::now(); |
| | | let request = RuntimeTurnRequest { |
| | | call_id: call_id.to_string(), |
| | | trace_id: trace_id.to_string(), |
| | | turn_id: turn.turn_id.clone(), |
| | | audio_artifact: RuntimeTurnAudioArtifact { |
| | | artifact_type: "local_file".to_string(), |
| | | path_ref: path_ref.to_string(), |
| | | format: "wav".to_string(), |
| | | sample_rate: TARGET_SAMPLE_RATE_HZ, |
| | | channels: u32::from(TARGET_NUM_CHANNELS), |
| | | duration_ms: turn.duration_ms, |
| | | byte_size, |
| | | }, |
| | | }; |
| | | let response = http |
| | | .post(bridge_config.bridge_url.as_deref().unwrap_or_default()) |
| | | .header( |
| | | "X-CV-Runtime-Token", |
| | | bridge_config.bridge_token.as_deref().unwrap_or_default(), |
| | | ) |
| | | .header("X-CV-Call-Id", call_id) |
| | | .header("X-CV-Trace-Id", trace_id) |
| | | .header( |
| | | "X-CV-Runtime-Session-Nonce", |
| | | bridge_config |
| | | .runtime_session_nonce |
| | | .as_deref() |
| | | .unwrap_or_default(), |
| | | ) |
| | | .json(&request) |
| | | .send() |
| | | .await |
| | | .map_err(|error| anyhow!("failed to publish livekit device output data: {error}"))?; |
| | | .context("failed to post turn stream bridge")?; |
| | | let status = response.status(); |
| | | if !status.is_success() { |
| | | let body_len = response.text().await.map(|body| body.len()).unwrap_or(0); |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | http_status = status.as_u16(), |
| | | body_len, |
| | | "runtime helper turn_stream_http_failed" |
| | | ); |
| | | return Err(anyhow!("turn stream bridge http failed")); |
| | | } |
| | | |
| | | info!( |
| | | call_id = %config.call_id, |
| | | trace_id = %config.trace_id, |
| | | topic = DEVICE_OUTPUT_TOPIC, |
| | | reliable = true, |
| | | command_count = payload.sensor_instructions.len(), |
| | | destination_count, |
| | | "runtime helper published device output smoke data" |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | "runtime helper turn_stream_connected" |
| | | ); |
| | | let mut byte_stream = response.bytes_stream(); |
| | | let mut line_buffer: Vec<u8> = Vec::new(); |
| | | let mut state = RuntimeTurnStreamState::default(); |
| | | while let Some(chunk) = byte_stream.next().await { |
| | | let chunk = chunk.context("failed to read turn stream chunk")?; |
| | | line_buffer.extend_from_slice(&chunk); |
| | | while let Some(newline_index) = line_buffer.iter().position(|value| *value == b'\n') { |
| | | let line: Vec<u8> = line_buffer.drain(..=newline_index).collect(); |
| | | if let Some(event) = parse_turn_stream_event_line(&line)? { |
| | | handle_turn_stream_event( |
| | | bridge_config, |
| | | sink, |
| | | call_id, |
| | | trace_id, |
| | | turn, |
| | | event, |
| | | &mut state, |
| | | turn_pipeline_started_at, |
| | | ) |
| | | .await?; |
| | | } |
| | | } |
| | | } |
| | | if !line_buffer.is_empty() { |
| | | if let Some(event) = parse_turn_stream_event_line(&line_buffer)? { |
| | | handle_turn_stream_event( |
| | | bridge_config, |
| | | sink, |
| | | call_id, |
| | | trace_id, |
| | | turn, |
| | | event, |
| | | &mut state, |
| | | turn_pipeline_started_at, |
| | | ) |
| | | .await?; |
| | | } |
| | | } |
| | | if !state.completed { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | audio_chunk_count = state.audio_chunk_count, |
| | | "runtime helper turn_stream_completed_without_final_event" |
| | | ); |
| | | return Err(anyhow!("turn stream ended without turn_completed")); |
| | | } |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "turn_bridge_completed", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "bridgeWallDurationMs": bridge_started_at.elapsed().as_millis() as u64, |
| | | "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64, |
| | | "audioChunkCount": state.audio_chunk_count, |
| | | "deviceOutputCount": state.device_output_count, |
| | | "replyPlaybackMode": state.reply_playback_mode.as_str(), |
| | | }), |
| | | ); |
| | | Ok(RuntimeTurnStreamOutcome { |
| | | reply_playback_mode: state.reply_playback_mode, |
| | | audio_chunk_count: state.audio_chunk_count, |
| | | device_output_count: state.device_output_count, |
| | | }) |
| | | } |
| | | |
| | | fn parse_turn_stream_event_line(line: &[u8]) -> Result<Option<RuntimeTurnStreamEvent>> { |
| | | let line = trim_ascii_whitespace(line); |
| | | if line.is_empty() { |
| | | return Ok(None); |
| | | } |
| | | serde_json::from_slice(line) |
| | | .context("failed to parse turn stream event") |
| | | .map(Some) |
| | | } |
| | | |
| | | fn diagnostic_str<'a>(diagnostics: Option<&'a serde_json::Value>, key: &str) -> Option<&'a str> { |
| | | diagnostics?.get(key)?.as_str() |
| | | } |
| | | |
| | | fn diagnostic_bool(diagnostics: Option<&serde_json::Value>, key: &str) -> Option<bool> { |
| | | diagnostics?.get(key)?.as_bool() |
| | | } |
| | | |
| | | fn diagnostic_u64(diagnostics: Option<&serde_json::Value>, key: &str) -> Option<u64> { |
| | | diagnostics?.get(key)?.as_u64() |
| | | } |
| | | |
| | | async fn handle_turn_stream_event( |
| | | bridge_config: &TurnBridgeConfig, |
| | | sink: &BotAudioOutputSink, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | turn: &FinishedSpeechTurn, |
| | | event: RuntimeTurnStreamEvent, |
| | | state: &mut RuntimeTurnStreamState, |
| | | turn_pipeline_started_at: Instant, |
| | | ) -> Result<()> { |
| | | let event_type = event.event_type(); |
| | | match event_type.as_deref() { |
| | | Some("reply_playback_mode_selected") => { |
| | | if let Some(reply_playback_mode) = event.reply_playback_mode.as_deref() { |
| | | state.reply_playback_mode = reply_playback_mode.to_string(); |
| | | } |
| | | let diagnostics = event.diagnostics.as_ref(); |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | reply_playback_mode = %state.reply_playback_mode, |
| | | streaming_enabled = ?diagnostic_bool(diagnostics, "streamingEnabled"), |
| | | provider_streaming_supported = ?diagnostic_bool(diagnostics, "providerStreamingSupported"), |
| | | first_chunk_received = ?diagnostic_bool(diagnostics, "firstChunkReceived"), |
| | | stream_chunk_count = ?diagnostic_u64(diagnostics, "streamChunkCount"), |
| | | fallback_reason = %diagnostic_str(diagnostics, "fallbackReason").unwrap_or("none"), |
| | | fallback_stage = %diagnostic_str(diagnostics, "fallbackStage").unwrap_or("none"), |
| | | stream_bridge_mode = %diagnostic_str(diagnostics, "streamBridgeMode").unwrap_or("unknown"), |
| | | tts_provider = %diagnostic_str(diagnostics, "ttsProvider").unwrap_or("unknown"), |
| | | "runtime helper reply_playback_mode_selected" |
| | | ); |
| | | } |
| | | Some("reply_state") => { |
| | | if let Some(reply_state) = event.state.as_deref() { |
| | | publish_reply_state_from_stream_event( |
| | | sink, |
| | | call_id, |
| | | trace_id, |
| | | &turn.turn_id, |
| | | &state.reply_playback_mode, |
| | | reply_state, |
| | | event.seq, |
| | | ) |
| | | .await?; |
| | | if reply_state == "reply_playback_started" { |
| | | state.playback_started_sent = true; |
| | | } |
| | | } |
| | | } |
| | | Some("reply_audio_chunk") => { |
| | | let audio_chunk = event |
| | | .audio_chunk |
| | | .as_ref() |
| | | .ok_or_else(|| anyhow!("reply_audio_chunk event missing audioChunk"))?; |
| | | if !state.playback_started_sent { |
| | | state.reply_state_seq = state.reply_state_seq.saturating_add(1); |
| | | sink.publish_reply_state( |
| | | call_id, |
| | | trace_id, |
| | | &turn.turn_id, |
| | | &state.reply_playback_mode, |
| | | "reply_playback_started", |
| | | state.reply_state_seq, |
| | | ) |
| | | .await?; |
| | | state.playback_started_sent = true; |
| | | } |
| | | let written_frames = write_stream_audio_chunk( |
| | | bridge_config, |
| | | sink, |
| | | call_id, |
| | | trace_id, |
| | | turn, |
| | | audio_chunk, |
| | | state, |
| | | turn_pipeline_started_at, |
| | | ) |
| | | .await?; |
| | | if written_frames > 0 { |
| | | state.audio_chunk_count = state.audio_chunk_count.saturating_add(1); |
| | | } |
| | | } |
| | | Some("device_output") => { |
| | | if let Some(output) = event.device_output.as_ref() { |
| | | sink.publish_device_output(call_id, trace_id, &turn.turn_id, output) |
| | | .await?; |
| | | state.device_output_count = state.device_output_count.saturating_add(1); |
| | | } |
| | | } |
| | | Some("turn_completed") => { |
| | | state.completed = true; |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | message_id_alias = %event.completion.as_ref().and_then(|value| value.message_id.as_deref()).map(redact).unwrap_or_else(|| "none".to_string()), |
| | | audio_chunk_count = state.audio_chunk_count, |
| | | "runtime helper turn_stream_final_received" |
| | | ); |
| | | } |
| | | Some("turn_failed") => { |
| | | let error = event.error.as_ref(); |
| | | let reason_code = error |
| | | .and_then(|value| value.reason_code.as_deref()) |
| | | .unwrap_or("TURN_STREAM_FAILED"); |
| | | let stage = error |
| | | .and_then(|value| value.stage.as_deref()) |
| | | .unwrap_or("turn_bridge_stream"); |
| | | let retryable = error.and_then(|value| value.retryable).unwrap_or(false); |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | reason_code = %reason_code, |
| | | stage = %stage, |
| | | retryable = retryable, |
| | | "runtime helper turn_stream_failed_event" |
| | | ); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "turn_failed", |
| | | "failed", |
| | | Some(reason_code), |
| | | Some(retryable), |
| | | json!({ |
| | | "stage": stage, |
| | | "replyPlaybackMode": state.reply_playback_mode.as_str(), |
| | | }), |
| | | ); |
| | | return Err(anyhow!("turn stream failed event")); |
| | | } |
| | | Some("turn_cancelled") => { |
| | | state.completed = true; |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | "runtime helper turn_stream_cancelled_event" |
| | | ); |
| | | publish_reply_state_from_stream_event( |
| | | sink, |
| | | call_id, |
| | | trace_id, |
| | | &turn.turn_id, |
| | | &state.reply_playback_mode, |
| | | "reply_playback_cancelled", |
| | | event.seq, |
| | | ) |
| | | .await?; |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "turn_cancelled", |
| | | "ok", |
| | | None, |
| | | Some(false), |
| | | json!({ |
| | | "replyPlaybackMode": state.reply_playback_mode.as_str(), |
| | | "audioChunkCount": state.audio_chunk_count, |
| | | }), |
| | | ); |
| | | } |
| | | Some("activity") => { |
| | | if let Some(activity) = event.activity.as_ref() { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | activity_event = %activity.event_type.as_deref().unwrap_or("unknown"), |
| | | stage = %activity.stage.as_deref().unwrap_or("unknown"), |
| | | reason_code = %activity.reason_code.as_deref().unwrap_or("none"), |
| | | "runtime helper turn_stream_activity" |
| | | ); |
| | | } |
| | | } |
| | | Some(other) => { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | event_type = %other, |
| | | "runtime helper ignored turn stream event" |
| | | ); |
| | | } |
| | | None => { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | "runtime helper ignored turn stream event without type" |
| | | ); |
| | | } |
| | | } |
| | | Ok(()) |
| | | } |
| | | |
| | | async fn publish_reply_state_from_stream_event( |
| | | sink: &BotAudioOutputSink, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | turn_id: &str, |
| | | reply_playback_mode: &str, |
| | | state: &str, |
| | | seq: Option<u64>, |
| | | ) -> Result<()> { |
| | | sink.publish_reply_state( |
| | | call_id, |
| | | trace_id, |
| | | turn_id, |
| | | reply_playback_mode, |
| | | state, |
| | | seq.unwrap_or(0), |
| | | ) |
| | | .await |
| | | } |
| | | |
| | | async fn write_stream_audio_chunk( |
| | | bridge_config: &TurnBridgeConfig, |
| | | sink: &BotAudioOutputSink, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | turn: &FinishedSpeechTurn, |
| | | audio_chunk: &RuntimeTurnStreamAudioChunk, |
| | | state: &mut RuntimeTurnStreamState, |
| | | turn_pipeline_started_at: Instant, |
| | | ) -> Result<usize> { |
| | | let payload_base64 = audio_chunk |
| | | .payload_base64 |
| | | .as_deref() |
| | | .map(str::trim) |
| | | .filter(|value| !value.is_empty()) |
| | | .ok_or_else(|| anyhow!("reply_audio_chunk payloadBase64 missing"))?; |
| | | let payload = general_purpose::STANDARD |
| | | .decode(payload_base64) |
| | | .context("failed to decode reply_audio_chunk payloadBase64")?; |
| | | let format = audio_chunk |
| | | .format |
| | | .as_deref() |
| | | .unwrap_or("pcm_s16le") |
| | | .trim() |
| | | .to_ascii_lowercase(); |
| | | let frames = if format == "pcm_s16le" { |
| | | let frame_alignment_bytes = pcm_s16le_frame_alignment_bytes( |
| | | audio_chunk |
| | | .channels |
| | | .unwrap_or(u32::from(TARGET_NUM_CHANNELS)), |
| | | )?; |
| | | let (aligned_payload, dropped_tail_bytes) = take_aligned_pcm_payload( |
| | | &mut state.pcm_audio_buffer, |
| | | &payload, |
| | | frame_alignment_bytes, |
| | | audio_chunk.last.unwrap_or(false), |
| | | ); |
| | | if dropped_tail_bytes > 0 { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | chunk_seq = audio_chunk.chunk_seq.unwrap_or_default(), |
| | | dropped_tail_bytes, |
| | | "runtime helper stream_audio_pcm_unaligned_tail_dropped" |
| | | ); |
| | | } |
| | | if aligned_payload.is_empty() { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | chunk_seq = audio_chunk.chunk_seq.unwrap_or_default(), |
| | | buffered_bytes = state.pcm_audio_buffer.len(), |
| | | "runtime helper stream_audio_pcm_waiting_for_sample_boundary" |
| | | ); |
| | | return Ok(0); |
| | | } |
| | | pcm_s16le_payload_to_frames( |
| | | &aligned_payload, |
| | | audio_chunk.sample_rate.unwrap_or(TARGET_SAMPLE_RATE_HZ), |
| | | audio_chunk |
| | | .channels |
| | | .unwrap_or(u32::from(TARGET_NUM_CHANNELS)), |
| | | )? |
| | | } else if matches!(format.as_str(), "mp3" | "mpeg" | "wav") { |
| | | state.encoded_audio_buffer.extend_from_slice(&payload); |
| | | match audio::decode_audio_bytes_to_frames( |
| | | &state.encoded_audio_buffer, |
| | | "stream_chunk", |
| | | TARGET_SAMPLE_RATE_HZ, |
| | | TARGET_NUM_CHANNELS, |
| | | bridge_config.audio_debug_dump_dir.as_deref(), |
| | | call_id, |
| | | &format!("stream-reply-{}", turn.turn_id), |
| | | ) { |
| | | Ok(loaded) => { |
| | | state.encoded_audio_buffer.clear(); |
| | | loaded.frames |
| | | } |
| | | Err(error) if !audio_chunk.last.unwrap_or(false) => { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | chunk_seq = audio_chunk.chunk_seq.unwrap_or_default(), |
| | | format = %format, |
| | | error = %safe_error(&error.to_string()), |
| | | "runtime helper stream_audio_chunk_decode_waiting_for_more_data" |
| | | ); |
| | | return Ok(0); |
| | | } |
| | | Err(error) if state.first_audio_frame_written => { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | chunk_seq = audio_chunk.chunk_seq.unwrap_or_default(), |
| | | format = %format, |
| | | buffered_bytes = state.encoded_audio_buffer.len(), |
| | | error = %safe_error(&error.to_string()), |
| | | "runtime helper stream_audio_final_chunk_decode_ignored" |
| | | ); |
| | | state.encoded_audio_buffer.clear(); |
| | | return Ok(0); |
| | | } |
| | | Err(error) => return Err(error).context("failed to decode final stream audio chunk"), |
| | | } |
| | | } else { |
| | | return Err(anyhow!("unsupported reply_audio_chunk format {format}")); |
| | | }; |
| | | if frames.is_empty() { |
| | | return Ok(0); |
| | | } |
| | | |
| | | if !state.first_audio_frame_written { |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "bot_reply_audio_write_started", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "replyPlaybackMode": state.reply_playback_mode.as_str(), |
| | | "format": format.as_str(), |
| | | "chunkSeq": audio_chunk.chunk_seq, |
| | | "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64, |
| | | }), |
| | | ); |
| | | } |
| | | let pacing_started_at = tokio::time::Instant::now(); |
| | | for (index, frame) in frames.iter().enumerate() { |
| | | sink.write_pcm_frame(frame).await?; |
| | | if !state.first_audio_frame_written { |
| | | state.first_audio_frame_written = true; |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "bot_reply_first_audio_frame_written", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "replyPlaybackMode": state.reply_playback_mode.as_str(), |
| | | "format": format.as_str(), |
| | | "chunkSeq": audio_chunk.chunk_seq, |
| | | "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64, |
| | | }), |
| | | ); |
| | | } |
| | | sleep_until(pacing_started_at + Duration::from_millis(((index + 1) as u64) * 20)).await; |
| | | } |
| | | if audio_chunk.last.unwrap_or(false) { |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "bot_reply_audio_write_finished", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "replyPlaybackMode": state.reply_playback_mode.as_str(), |
| | | "format": format.as_str(), |
| | | "chunkSeq": audio_chunk.chunk_seq, |
| | | "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64, |
| | | }), |
| | | ); |
| | | } |
| | | Ok(frames.len()) |
| | | } |
| | | |
| | | fn pcm_s16le_payload_to_frames( |
| | | payload: &[u8], |
| | | sample_rate: u32, |
| | | channels: u32, |
| | | ) -> Result<Vec<PcmFrame>> { |
| | | audio::pcm_s16le_bytes_to_frames( |
| | | payload, |
| | | sample_rate, |
| | | channels, |
| | | TARGET_SAMPLE_RATE_HZ, |
| | | TARGET_NUM_CHANNELS, |
| | | ) |
| | | } |
| | | |
| | | fn pcm_s16le_frame_alignment_bytes(channels: u32) -> Result<usize> { |
| | | if channels == 0 { |
| | | return Err(anyhow!("pcm_s16le channel count is zero")); |
| | | } |
| | | let channel_count = usize::try_from(channels) |
| | | .map_err(|_| anyhow!("unsupported pcm_s16le channel count {channels}"))?; |
| | | Ok(2 * channel_count) |
| | | } |
| | | |
| | | fn take_aligned_pcm_payload( |
| | | buffer: &mut Vec<u8>, |
| | | payload: &[u8], |
| | | frame_alignment_bytes: usize, |
| | | last: bool, |
| | | ) -> (Vec<u8>, usize) { |
| | | let alignment = frame_alignment_bytes.max(2); |
| | | buffer.extend_from_slice(payload); |
| | | let aligned_len = buffer.len() - buffer.len() % alignment; |
| | | let aligned_payload = if aligned_len == 0 { |
| | | Vec::new() |
| | | } else { |
| | | buffer.drain(..aligned_len).collect() |
| | | }; |
| | | let dropped_tail_bytes = if last && !buffer.is_empty() { |
| | | let dropped = buffer.len(); |
| | | buffer.clear(); |
| | | dropped |
| | | } else { |
| | | 0 |
| | | }; |
| | | (aligned_payload, dropped_tail_bytes) |
| | | } |
| | | |
| | | #[cfg(test)] |
| | | mod pcm_stream_tests { |
| | | use super::take_aligned_pcm_payload; |
| | | |
| | | #[test] |
| | | fn take_aligned_pcm_payload_buffers_split_sample_bytes() { |
| | | let mut buffer = Vec::new(); |
| | | |
| | | let (first, dropped) = take_aligned_pcm_payload(&mut buffer, &[0x01], 2, false); |
| | | assert!(first.is_empty()); |
| | | assert_eq!(0, dropped); |
| | | assert_eq!(vec![0x01], buffer); |
| | | |
| | | let (second, dropped) = take_aligned_pcm_payload(&mut buffer, &[0x02, 0x03], 2, false); |
| | | assert_eq!(vec![0x01, 0x02], second); |
| | | assert_eq!(0, dropped); |
| | | assert_eq!(vec![0x03], buffer); |
| | | |
| | | let (third, dropped) = take_aligned_pcm_payload(&mut buffer, &[0x04], 2, true); |
| | | assert_eq!(vec![0x03, 0x04], third); |
| | | assert_eq!(0, dropped); |
| | | assert!(buffer.is_empty()); |
| | | } |
| | | |
| | | #[test] |
| | | fn take_aligned_pcm_payload_drops_final_half_sample() { |
| | | let mut buffer = Vec::new(); |
| | | |
| | | let (payload, dropped) = |
| | | take_aligned_pcm_payload(&mut buffer, &[0x01, 0x02, 0x03], 2, true); |
| | | assert_eq!(vec![0x01, 0x02], payload); |
| | | assert_eq!(1, dropped); |
| | | assert!(buffer.is_empty()); |
| | | } |
| | | } |
| | | |
| | | fn trim_ascii_whitespace(value: &[u8]) -> &[u8] { |
| | | let mut start = 0; |
| | | let mut end = value.len(); |
| | | while start < end && value[start].is_ascii_whitespace() { |
| | | start += 1; |
| | | } |
| | | while end > start && value[end - 1].is_ascii_whitespace() { |
| | | end -= 1; |
| | | } |
| | | &value[start..end] |
| | | } |
| | | |
| | | async fn request_turn_bridge( |
| | | http: &Client, |
| | | bridge_config: &TurnBridgeConfig, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | turn: &FinishedSpeechTurn, |
| | | path_ref: &str, |
| | | byte_size: u64, |
| | | turn_pipeline_started_at: Instant, |
| | | ) -> Result<RuntimeTurnBridgeOutcome> { |
| | | let bridge_started_at = Instant::now(); |
| | | let request = RuntimeTurnRequest { |
| | | call_id: call_id.to_string(), |
| | | trace_id: trace_id.to_string(), |
| | | turn_id: turn.turn_id.clone(), |
| | | audio_artifact: RuntimeTurnAudioArtifact { |
| | | artifact_type: "local_file".to_string(), |
| | | path_ref: path_ref.to_string(), |
| | | format: "wav".to_string(), |
| | | sample_rate: TARGET_SAMPLE_RATE_HZ, |
| | | channels: u32::from(TARGET_NUM_CHANNELS), |
| | | duration_ms: turn.duration_ms, |
| | | byte_size, |
| | | }, |
| | | }; |
| | | let response = http |
| | | .post(bridge_config.bridge_url.as_deref().unwrap_or_default()) |
| | | .header( |
| | | "X-CV-Runtime-Token", |
| | | bridge_config.bridge_token.as_deref().unwrap_or_default(), |
| | | ) |
| | | .header("X-CV-Call-Id", call_id) |
| | | .header("X-CV-Trace-Id", trace_id) |
| | | .header( |
| | | "X-CV-Runtime-Session-Nonce", |
| | | bridge_config |
| | | .runtime_session_nonce |
| | | .as_deref() |
| | | .unwrap_or_default(), |
| | | ) |
| | | .json(&request) |
| | | .send() |
| | | .await |
| | | .context("failed to post turn bridge")?; |
| | | let status = response.status(); |
| | | let body = response |
| | | .text() |
| | | .await |
| | | .context("failed to read turn bridge response")?; |
| | | let body_len = body.len(); |
| | | let parsed: Option<RuntimeTurnCommonResult<RuntimeTurnResponseData>> = |
| | | serde_json::from_str(&body).ok(); |
| | | if !status.is_success() { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | http_status = status.as_u16(), |
| | | body_len, |
| | | code = parsed.as_ref().map(|value| value.code).unwrap_or_default(), |
| | | reason_code = %parsed.as_ref().and_then(|value| value.msg.as_deref()).unwrap_or("unknown"), |
| | | "runtime helper turn_bridge_http_failed" |
| | | ); |
| | | return Err(anyhow!("turn bridge http failed")); |
| | | } |
| | | let Some(result) = parsed else { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | body_len, |
| | | "runtime helper turn_bridge_response_invalid" |
| | | ); |
| | | return Err(anyhow!("turn bridge response invalid")); |
| | | }; |
| | | if result.code != 0 { |
| | | warn!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | code = result.code, |
| | | reason_code = %result.msg.as_deref().unwrap_or("unknown"), |
| | | retryable = result.retryable.unwrap_or(false), |
| | | stage = %result.stage.as_deref().unwrap_or("unknown"), |
| | | "runtime helper turn_bridge_business_failed" |
| | | ); |
| | | return Err(anyhow!("turn bridge business failed")); |
| | | } |
| | | let data = result.data; |
| | | let reply_audio_artifact = data |
| | | .as_ref() |
| | | .and_then(|value| value.reply_audio_artifact.as_ref()); |
| | | let device_outputs = data |
| | | .as_ref() |
| | | .and_then(|value| value.device_outputs.clone()) |
| | | .unwrap_or_default(); |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | response_turn_id = %data.as_ref().and_then(|value| value.turn_id.as_deref()).unwrap_or("unknown"), |
| | | message_id_alias = %data.as_ref().and_then(|value| value.message_id.as_deref()).map(redact).unwrap_or_else(|| "none".to_string()), |
| | | reply_audio_present = reply_audio_artifact.is_some(), |
| | | reply_audio_type = %reply_audio_artifact.and_then(|value| value.artifact_type.as_deref()).unwrap_or("none"), |
| | | reply_audio_path_alias = %reply_audio_artifact.and_then(|value| value.path_ref.as_deref()).map(redact).unwrap_or_else(|| "none".to_string()), |
| | | reply_audio_format = %reply_audio_artifact.and_then(|value| value.format.as_deref()).unwrap_or("none"), |
| | | device_output_count = device_outputs.len(), |
| | | "runtime helper turn_bridge_completed" |
| | | ); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "turn_bridge_completed", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "bridgeWallDurationMs": bridge_started_at.elapsed().as_millis() as u64, |
| | | "replyAudioPresent": reply_audio_artifact.is_some(), |
| | | "deviceOutputCount": device_outputs.len(), |
| | | "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64, |
| | | }), |
| | | ); |
| | | Ok(RuntimeTurnBridgeOutcome { |
| | | reply_audio_artifact: reply_audio_artifact.cloned(), |
| | | device_outputs, |
| | | }) |
| | | } |
| | | |
| | | async fn write_reply_audio_artifact( |
| | | http: &Client, |
| | | bridge_config: &TurnBridgeConfig, |
| | | sink: &BotAudioOutputSink, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | turn: &FinishedSpeechTurn, |
| | | artifact: &RuntimeTurnReplyAudioArtifact, |
| | | turn_pipeline_started_at: Instant, |
| | | ) -> Result<()> { |
| | | let artifact_type = artifact.artifact_type.as_deref().unwrap_or_default().trim(); |
| | | if artifact_type != "local_file" { |
| | | return Err(anyhow!("unsupported reply audio artifact type")); |
| | | } |
| | | let path_ref = artifact |
| | | .path_ref |
| | | .as_deref() |
| | | .map(str::trim) |
| | | .filter(|value| !value.is_empty()) |
| | | .ok_or_else(|| anyhow!("reply audio pathRef missing"))?; |
| | | let artifact_root = PathBuf::from(bridge_config.artifact_dir.as_deref().unwrap_or_default()); |
| | | let audio_path = resolve_artifact_path(&artifact_root, path_ref)?; |
| | | let audio_path_string = audio_path.to_string_lossy().to_string(); |
| | | let reply_audio_format = artifact.format.as_deref().unwrap_or("unknown"); |
| | | |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | reply_audio_type = %artifact_type, |
| | | reply_audio_path_alias = %redact(path_ref), |
| | | reply_audio_format = %reply_audio_format, |
| | | "runtime helper bot_reply_audio_write_started" |
| | | ); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "bot_reply_audio_write_started", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "replyAudioFormat": reply_audio_format, |
| | | "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64, |
| | | }), |
| | | ); |
| | | |
| | | let reply_debug_label = format!("reply-{}", turn.turn_id); |
| | | let loaded_audio = load_pre_recorded_frames( |
| | | http, |
| | | Some(&audio_path_string), |
| | | None, |
| | | TARGET_SAMPLE_RATE_HZ, |
| | | TARGET_NUM_CHANNELS, |
| | | bridge_config.audio_debug_dump_dir.as_deref(), |
| | | call_id, |
| | | &reply_debug_label, |
| | | ) |
| | | .await? |
| | | .ok_or_else(|| anyhow!("reply audio artifact decode returned empty"))?; |
| | | let AudioDiagnostics { |
| | | source_format, |
| | | source_bytes, |
| | | source_sample_rate_hz, |
| | | source_num_channels, |
| | | target_duration_ms, |
| | | frame_count, |
| | | rms, |
| | | peak, |
| | | clipped_sample_count, |
| | | silence_ratio, |
| | | mp3_skipped_data_count, |
| | | mp3_insufficient_data_count, |
| | | debug_source_path, |
| | | debug_pcm_wav_path, |
| | | .. |
| | | } = loaded_audio.diagnostics; |
| | | let frames = loaded_audio.frames; |
| | | |
| | | let playback_started_at = Instant::now(); |
| | | let pacing_started_at = tokio::time::Instant::now(); |
| | | for (index, frame) in frames.iter().enumerate() { |
| | | sink.write_pcm_frame(frame).await?; |
| | | sleep_until(pacing_started_at + Duration::from_millis(((index + 1) as u64) * 20)).await; |
| | | } |
| | | let playback_wall_ms = playback_started_at.elapsed().as_millis() as i64; |
| | | let drift_ms = playback_wall_ms - target_duration_ms as i64; |
| | | sink.clear_buffer(); |
| | | |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn.turn_id, |
| | | source_format, |
| | | source_bytes, |
| | | source_sample_rate_hz, |
| | | source_num_channels, |
| | | frame_count, |
| | | theoretical_duration_ms = target_duration_ms, |
| | | push_wall_duration_ms = playback_wall_ms, |
| | | push_drift_ms = drift_ms, |
| | | rms = round4(rms), |
| | | peak = round4(peak), |
| | | clipped_sample_count, |
| | | silence_ratio = round4(silence_ratio), |
| | | mp3_skipped_data_count, |
| | | mp3_insufficient_data_count, |
| | | debug_source_path = debug_source_path.as_deref().unwrap_or(""), |
| | | debug_pcm_wav_path = debug_pcm_wav_path.as_deref().unwrap_or(""), |
| | | "runtime helper bot_reply_audio_write_finished" |
| | | ); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn.turn_id), |
| | | "bot_reply_audio_write_finished", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "sourceFormat": source_format, |
| | | "sourceBytes": source_bytes, |
| | | "sourceSampleRateHz": source_sample_rate_hz, |
| | | "sourceNumChannels": source_num_channels, |
| | | "frameCount": frame_count, |
| | | "theoreticalDurationMs": target_duration_ms, |
| | | "pushWallDurationMs": playback_wall_ms, |
| | | "pushDriftMs": drift_ms, |
| | | "rms": round4(rms), |
| | | "peak": round4(peak), |
| | | "clippedSampleCount": clipped_sample_count, |
| | | "silenceRatio": round4(silence_ratio), |
| | | "mp3SkippedDataCount": mp3_skipped_data_count, |
| | | "mp3InsufficientDataCount": mp3_insufficient_data_count, |
| | | "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64, |
| | | }), |
| | | ); |
| | | Ok(()) |
| | | } |
| | | |
| | | fn build_device_output_smoke_payload(config: &Config) -> DeviceOutputMessage { |
| | | DeviceOutputMessage { |
| | | message_type: DEVICE_OUTPUT_TOPIC, |
| | | schema_version: "1.0", |
| | | call_id: config.call_id.clone(), |
| | | role_id: config.role_id.clone(), |
| | | trace_id: config.trace_id.clone(), |
| | | ack_mode: "http", |
| | | sensor_instructions: vec![ |
| | | SensorInstruction { |
| | | command_id: command_id(&config.call_id, 1), |
| | | sensor_type: "Vibrator", |
| | | operation_type: "VibratorStart", |
| | | step: 1, |
| | | duration_sec: None, |
| | | extension: r#"{"levelList":[{"level":1,"percent":"0.6"}]}"#.to_string(), |
| | | }, |
| | | SensorInstruction { |
| | | command_id: command_id(&config.call_id, 2), |
| | | sensor_type: "Vibrator", |
| | | operation_type: "VibratorUp", |
| | | step: 1, |
| | | duration_sec: None, |
| | | extension: r#"{"delta":1}"#.to_string(), |
| | | }, |
| | | SensorInstruction { |
| | | command_id: command_id(&config.call_id, 3), |
| | | sensor_type: "Pump", |
| | | operation_type: "JiaStart", |
| | | step: 1, |
| | | duration_sec: None, |
| | | extension: r#"{"levelList":[{"level":1,"percent":"0.5"}]}"#.to_string(), |
| | | }, |
| | | SensorInstruction { |
| | | command_id: command_id(&config.call_id, 4), |
| | | sensor_type: "Heating", |
| | | operation_type: "HeatingStart", |
| | | step: 1, |
| | | duration_sec: Some(3), |
| | | extension: r#"{"target":"warm"}"#.to_string(), |
| | | }, |
| | | ], |
| | | fn resolve_artifact_path(artifact_root: &Path, path_ref: &str) -> Result<PathBuf> { |
| | | let path_ref = path_ref.trim(); |
| | | if path_ref.is_empty() { |
| | | return Err(anyhow!("empty artifact pathRef")); |
| | | } |
| | | let relative = Path::new(path_ref); |
| | | if relative.is_absolute() { |
| | | return Err(anyhow!("absolute artifact pathRef is not allowed")); |
| | | } |
| | | let mut safe_relative = PathBuf::new(); |
| | | for component in relative.components() { |
| | | match component { |
| | | std::path::Component::Normal(segment) => { |
| | | let segment = segment |
| | | .to_str() |
| | | .ok_or_else(|| anyhow!("non-utf8 artifact pathRef segment"))?; |
| | | require_safe_segment(segment)?; |
| | | safe_relative.push(segment); |
| | | } |
| | | std::path::Component::CurDir => {} |
| | | _ => return Err(anyhow!("unsafe artifact pathRef component")), |
| | | } |
| | | } |
| | | if safe_relative.as_os_str().is_empty() { |
| | | return Err(anyhow!("artifact pathRef has no safe components")); |
| | | } |
| | | let root = normalize_path_lexically(artifact_root); |
| | | let output_path = normalize_path_lexically(&root.join(safe_relative)); |
| | | if !output_path.starts_with(&root) { |
| | | return Err(anyhow!("reply artifact path escapes root")); |
| | | } |
| | | Ok(output_path) |
| | | } |
| | | |
| | | #[derive(Serialize)] |
| | | struct RuntimeTurnRequest { |
| | | #[serde(rename = "callId")] |
| | | call_id: String, |
| | | #[serde(rename = "traceId")] |
| | | trace_id: String, |
| | | #[serde(rename = "turnId")] |
| | | turn_id: String, |
| | | #[serde(rename = "audioArtifact")] |
| | | audio_artifact: RuntimeTurnAudioArtifact, |
| | | } |
| | | |
| | | #[derive(Serialize)] |
| | | struct RuntimeTurnAudioArtifact { |
| | | #[serde(rename = "type")] |
| | | artifact_type: String, |
| | | #[serde(rename = "pathRef")] |
| | | path_ref: String, |
| | | format: String, |
| | | #[serde(rename = "sampleRate")] |
| | | sample_rate: u32, |
| | | channels: u32, |
| | | #[serde(rename = "durationMs")] |
| | | duration_ms: u64, |
| | | #[serde(rename = "byteSize")] |
| | | byte_size: u64, |
| | | } |
| | | |
| | | #[derive(Deserialize)] |
| | | struct RuntimeTurnCommonResult<T> { |
| | | code: i64, |
| | | msg: Option<String>, |
| | | data: Option<T>, |
| | | stage: Option<String>, |
| | | retryable: Option<bool>, |
| | | } |
| | | |
| | | #[derive(Default)] |
| | | struct RuntimeTurnBridgeOutcome { |
| | | reply_audio_artifact: Option<RuntimeTurnReplyAudioArtifact>, |
| | | device_outputs: Vec<RuntimeTurnDeviceOutput>, |
| | | } |
| | | |
| | | struct RuntimeTurnStreamOutcome { |
| | | reply_playback_mode: String, |
| | | audio_chunk_count: u64, |
| | | device_output_count: u64, |
| | | } |
| | | |
| | | struct RuntimeTurnStreamState { |
| | | reply_playback_mode: String, |
| | | reply_state_seq: u64, |
| | | playback_started_sent: bool, |
| | | first_audio_frame_written: bool, |
| | | completed: bool, |
| | | audio_chunk_count: u64, |
| | | device_output_count: u64, |
| | | encoded_audio_buffer: Vec<u8>, |
| | | pcm_audio_buffer: Vec<u8>, |
| | | } |
| | | |
| | | impl Default for RuntimeTurnStreamState { |
| | | fn default() -> Self { |
| | | Self { |
| | | reply_playback_mode: "full_tts_fallback".to_string(), |
| | | reply_state_seq: 0, |
| | | playback_started_sent: false, |
| | | first_audio_frame_written: false, |
| | | completed: false, |
| | | audio_chunk_count: 0, |
| | | device_output_count: 0, |
| | | encoded_audio_buffer: Vec::new(), |
| | | pcm_audio_buffer: Vec::new(), |
| | | } |
| | | } |
| | | } |
| | | |
| | | fn command_id(call_id: &str, sequence: u8) -> String { |
| | | format!("{call_id}-device-smoke-{sequence:03}") |
| | | #[derive(Deserialize)] |
| | | struct RuntimeTurnResponseData { |
| | | #[serde(rename = "turnId")] |
| | | turn_id: Option<String>, |
| | | #[serde(rename = "messageId")] |
| | | message_id: Option<String>, |
| | | #[serde(rename = "replyAudioArtifact")] |
| | | reply_audio_artifact: Option<RuntimeTurnReplyAudioArtifact>, |
| | | #[serde(rename = "deviceOutputs")] |
| | | device_outputs: Option<Vec<RuntimeTurnDeviceOutput>>, |
| | | } |
| | | |
| | | #[derive(Deserialize)] |
| | | struct RuntimeTurnStreamEvent { |
| | | #[serde(rename = "type", alias = "event")] |
| | | event_type: Option<String>, |
| | | #[serde(rename = "seq")] |
| | | seq: Option<u64>, |
| | | #[serde(rename = "replyPlaybackMode")] |
| | | reply_playback_mode: Option<String>, |
| | | state: Option<String>, |
| | | #[serde(rename = "audioChunk")] |
| | | audio_chunk: Option<RuntimeTurnStreamAudioChunk>, |
| | | activity: Option<RuntimeTurnStreamActivity>, |
| | | #[serde(rename = "deviceOutput")] |
| | | device_output: Option<RuntimeTurnDeviceOutput>, |
| | | error: Option<RuntimeTurnStreamError>, |
| | | completion: Option<RuntimeTurnStreamCompletion>, |
| | | diagnostics: Option<serde_json::Value>, |
| | | } |
| | | |
| | | impl RuntimeTurnStreamEvent { |
| | | fn event_type(&self) -> Option<String> { |
| | | self.event_type.as_deref().map(|value| match value { |
| | | "reply_playback_mode_selected" => "reply_playback_mode_selected".to_string(), |
| | | "reply_state" => "reply_state".to_string(), |
| | | "reply_audio_chunk" => "reply_audio_chunk".to_string(), |
| | | "device_output" => "device_output".to_string(), |
| | | "turn_completed" => "turn_completed".to_string(), |
| | | "turn_failed" => "turn_failed".to_string(), |
| | | "turn_cancelled" => "turn_cancelled".to_string(), |
| | | "activity" => "activity".to_string(), |
| | | other => other.to_string(), |
| | | }) |
| | | } |
| | | } |
| | | |
| | | #[derive(Deserialize)] |
| | | struct RuntimeTurnStreamAudioChunk { |
| | | #[serde(rename = "chunkSeq", alias = "seq")] |
| | | chunk_seq: Option<u64>, |
| | | format: Option<String>, |
| | | #[serde(rename = "sampleRate")] |
| | | sample_rate: Option<u32>, |
| | | channels: Option<u32>, |
| | | #[serde(rename = "payloadBase64", alias = "audioBase64")] |
| | | payload_base64: Option<String>, |
| | | last: Option<bool>, |
| | | } |
| | | |
| | | #[derive(Deserialize)] |
| | | struct RuntimeTurnStreamActivity { |
| | | #[serde(rename = "eventType", alias = "event")] |
| | | event_type: Option<String>, |
| | | stage: Option<String>, |
| | | #[serde(rename = "reasonCode")] |
| | | reason_code: Option<String>, |
| | | } |
| | | |
| | | #[derive(Deserialize)] |
| | | struct RuntimeTurnStreamError { |
| | | #[serde(rename = "reasonCode")] |
| | | reason_code: Option<String>, |
| | | stage: Option<String>, |
| | | retryable: Option<bool>, |
| | | } |
| | | |
| | | #[derive(Deserialize)] |
| | | struct RuntimeTurnStreamCompletion { |
| | | #[serde(rename = "messageId")] |
| | | message_id: Option<String>, |
| | | } |
| | | |
| | | #[derive(Clone, Deserialize)] |
| | | struct RuntimeTurnReplyAudioArtifact { |
| | | #[serde(rename = "type")] |
| | | artifact_type: Option<String>, |
| | | #[serde(rename = "pathRef")] |
| | | path_ref: Option<String>, |
| | | format: Option<String>, |
| | | } |
| | | |
| | | #[derive(Clone, Deserialize)] |
| | | struct RuntimeTurnDeviceOutput { |
| | | #[serde(rename = "commandId")] |
| | | command_id: Option<String>, |
| | | #[serde(rename = "commandCode")] |
| | | command_code: Option<String>, |
| | | params: Option<serde_json::Value>, |
| | | } |
| | | |
| | | fn require_safe_segment(value: &str) -> Result<()> { |
| | | if value.is_empty() |
| | | || value.contains('/') |
| | | || value.contains("..") |
| | | || !value |
| | | .chars() |
| | | .all(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '_' | '-' | '.')) |
| | | { |
| | | return Err(anyhow!("unsafe path segment")); |
| | | } |
| | | Ok(()) |
| | | } |
| | | |
| | | fn normalize_path_lexically(path: &Path) -> PathBuf { |
| | | let mut normalized = PathBuf::new(); |
| | | for component in path.components() { |
| | | match component { |
| | | std::path::Component::CurDir => {} |
| | | std::path::Component::ParentDir => { |
| | | normalized.pop(); |
| | | } |
| | | _ => normalized.push(component.as_os_str()), |
| | | } |
| | | } |
| | | normalized |
| | | } |
| | | |
| | | fn spawn_user_audio_frame_observer( |
| | | track: RemoteAudioTrack, |
| | | call_id: String, |
| | | trace_id: String, |
| | | participant_alias: String, |
| | | track_sid_alias: String, |
| | | simple_vad_enabled: bool, |
| | | simple_vad_config: SimpleVadConfig, |
| | | vad_enabled_gate: Arc<AtomicBool>, |
| | | turn_bridge_config: TurnBridgeConfig, |
| | | http: Client, |
| | | sink: Arc<BotAudioOutputSink>, |
| | | ) -> JoinHandle<()> { |
| | | tokio::spawn(async move { |
| | | let mut stream = NativeAudioStream::new( |
| | | track.rtc_track(), |
| | | TARGET_SAMPLE_RATE_HZ as i32, |
| | | i32::from(TARGET_NUM_CHANNELS), |
| | | ); |
| | | let started_at = Instant::now(); |
| | | let mut frame_count: u64 = 0; |
| | | let mut sample_count: u64 = 0; |
| | | let mut simple_vad = if simple_vad_enabled { |
| | | Some(SimpleVad::new(simple_vad_config)) |
| | | } else { |
| | | None |
| | | }; |
| | | |
| | | while let Some(frame) = stream.next().await { |
| | | frame_count += 1; |
| | | sample_count += u64::from(frame.samples_per_channel) * u64::from(frame.num_channels); |
| | | let elapsed_ms = started_at.elapsed().as_millis() as u64; |
| | | if frame_count == 1 { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | participant_alias = %participant_alias, |
| | | track_sid_alias = %track_sid_alias, |
| | | sample_rate_hz = frame.sample_rate, |
| | | num_channels = frame.num_channels, |
| | | samples_per_channel = frame.samples_per_channel, |
| | | first_frame_elapsed_ms = elapsed_ms, |
| | | "runtime helper user_audio_frame_received" |
| | | ); |
| | | } else if frame_count % 250 == 0 { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | participant_alias = %participant_alias, |
| | | track_sid_alias = %track_sid_alias, |
| | | frame_count, |
| | | sample_count, |
| | | observed_wall_ms = elapsed_ms, |
| | | "runtime helper user_audio_frame_summary" |
| | | ); |
| | | } |
| | | |
| | | if let Some(vad) = simple_vad.as_mut() { |
| | | if vad_enabled_gate.load(Ordering::Acquire) { |
| | | if let Some(turn) = vad.observe_frame( |
| | | &call_id, |
| | | &trace_id, |
| | | &participant_alias, |
| | | &track_sid_alias, |
| | | frame_count, |
| | | elapsed_ms, |
| | | &frame, |
| | | ) { |
| | | handle_finished_turn( |
| | | &http, |
| | | &turn_bridge_config, |
| | | &sink, |
| | | &call_id, |
| | | &trace_id, |
| | | turn, |
| | | ) |
| | | .await; |
| | | } |
| | | } else { |
| | | vad.observe_disabled_frame( |
| | | &call_id, |
| | | &trace_id, |
| | | &participant_alias, |
| | | &track_sid_alias, |
| | | frame_count, |
| | | elapsed_ms, |
| | | ); |
| | | } |
| | | } |
| | | } |
| | | |
| | | if let Some(vad) = simple_vad.as_mut() { |
| | | if let Some(turn) = vad.finish_stream( |
| | | &call_id, |
| | | &trace_id, |
| | | &participant_alias, |
| | | &track_sid_alias, |
| | | started_at.elapsed().as_millis() as u64, |
| | | ) { |
| | | handle_finished_turn(&http, &turn_bridge_config, &sink, &call_id, &trace_id, turn) |
| | | .await; |
| | | } |
| | | } |
| | | |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | participant_alias = %participant_alias, |
| | | track_sid_alias = %track_sid_alias, |
| | | frame_count, |
| | | sample_count, |
| | | observed_wall_ms = started_at.elapsed().as_millis() as u64, |
| | | "runtime helper user_audio_stream_ended" |
| | | ); |
| | | }) |
| | | } |
| | | |
| | | #[derive(Clone)] |
| | | struct SimpleVadConfig { |
| | | rms_threshold: f64, |
| | | peak_threshold: f64, |
| | | start_frames: u32, |
| | | end_silence_ms: u64, |
| | | min_speech_ms: u64, |
| | | max_turn_ms: u64, |
| | | initial_ignore_ms: u64, |
| | | } |
| | | |
| | | impl SimpleVadConfig { |
| | | fn from_env() -> Self { |
| | | Self { |
| | | rms_threshold: f64_env("CV_VAD_RMS_THRESHOLD", 0.012), |
| | | peak_threshold: f64_env("CV_VAD_PEAK_THRESHOLD", 0.08), |
| | | start_frames: u32_env("CV_VAD_START_FRAMES", 5).max(1), |
| | | end_silence_ms: u64_env("CV_VAD_END_SILENCE_MS", 700).max(100), |
| | | min_speech_ms: u64_env("CV_VAD_MIN_SPEECH_MS", 300).max(1), |
| | | max_turn_ms: u64_env("CV_VAD_MAX_TURN_MS", 10_000).max(1_000), |
| | | initial_ignore_ms: u64_env("CV_VAD_INITIAL_IGNORE_MS", 500), |
| | | } |
| | | } |
| | | |
| | | fn end_silence_frames(&self, frame_duration_ms: u64) -> u32 { |
| | | let frame_duration_ms = frame_duration_ms.max(1); |
| | | self.end_silence_ms.div_ceil(frame_duration_ms) as u32 |
| | | } |
| | | } |
| | | |
| | | struct SimpleVad { |
| | | config: SimpleVadConfig, |
| | | turn_index: u64, |
| | | in_speech: bool, |
| | | voiced_run_frames: u32, |
| | | silence_run_frames: u32, |
| | | speech_start_elapsed_ms: u64, |
| | | speech_frame_count: u64, |
| | | speech_sample_count: u64, |
| | | speech_rms_sum: f64, |
| | | speech_peak: f64, |
| | | ignored_before_enabled_frames: u64, |
| | | pre_speech_frames: Vec<Vec<i16>>, |
| | | speech_samples: Vec<i16>, |
| | | } |
| | | |
| | | struct FinishedSpeechTurn { |
| | | turn_index: u64, |
| | | turn_id: String, |
| | | samples: Vec<i16>, |
| | | duration_ms: u64, |
| | | frame_count: u64, |
| | | sample_count: u64, |
| | | end_reason: String, |
| | | } |
| | | |
| | | impl SimpleVad { |
| | | fn new(config: SimpleVadConfig) -> Self { |
| | | Self { |
| | | config, |
| | | turn_index: 0, |
| | | in_speech: false, |
| | | voiced_run_frames: 0, |
| | | silence_run_frames: 0, |
| | | speech_start_elapsed_ms: 0, |
| | | speech_frame_count: 0, |
| | | speech_sample_count: 0, |
| | | speech_rms_sum: 0.0, |
| | | speech_peak: 0.0, |
| | | ignored_before_enabled_frames: 0, |
| | | pre_speech_frames: Vec::new(), |
| | | speech_samples: Vec::new(), |
| | | } |
| | | } |
| | | |
| | | fn observe_disabled_frame( |
| | | &mut self, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | participant_alias: &str, |
| | | track_sid_alias: &str, |
| | | frame_count: u64, |
| | | elapsed_ms: u64, |
| | | ) { |
| | | self.reset_current_turn(); |
| | | self.ignored_before_enabled_frames = self.ignored_before_enabled_frames.saturating_add(1); |
| | | if self.ignored_before_enabled_frames == 1 || self.ignored_before_enabled_frames % 250 == 0 |
| | | { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | participant_alias = %participant_alias, |
| | | track_sid_alias = %track_sid_alias, |
| | | frame_count, |
| | | ignored_frame_count = self.ignored_before_enabled_frames, |
| | | elapsed_ms, |
| | | "runtime helper vad_ignored_before_enabled" |
| | | ); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | None, |
| | | "ignored_user_audio_before_vad_enabled", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "frameCount": frame_count, |
| | | "ignoredFrameCount": self.ignored_before_enabled_frames, |
| | | "elapsedMs": elapsed_ms, |
| | | }), |
| | | ); |
| | | } |
| | | } |
| | | |
| | | fn observe_frame( |
| | | &mut self, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | participant_alias: &str, |
| | | track_sid_alias: &str, |
| | | frame_count: u64, |
| | | elapsed_ms: u64, |
| | | frame: &AudioFrame<'_>, |
| | | ) -> Option<FinishedSpeechTurn> { |
| | | let frame_duration_ms = frame_duration_ms(frame); |
| | | let (rms, peak) = pcm_energy_stats(frame.data.as_ref()); |
| | | if elapsed_ms < self.config.initial_ignore_ms { |
| | | return None; |
| | | } |
| | | |
| | | let voiced = rms >= self.config.rms_threshold || peak >= self.config.peak_threshold; |
| | | if !self.in_speech { |
| | | self.remember_pre_speech_frame(frame); |
| | | } |
| | | if voiced { |
| | | self.voiced_run_frames = self.voiced_run_frames.saturating_add(1); |
| | | self.silence_run_frames = 0; |
| | | } else { |
| | | self.voiced_run_frames = 0; |
| | | self.silence_run_frames = self.silence_run_frames.saturating_add(1); |
| | | } |
| | | |
| | | if !self.in_speech { |
| | | if self.voiced_run_frames >= self.config.start_frames { |
| | | self.turn_index += 1; |
| | | self.in_speech = true; |
| | | self.speech_start_elapsed_ms = elapsed_ms |
| | | .saturating_sub(u64::from(self.config.start_frames) * frame_duration_ms); |
| | | self.speech_frame_count = u64::from(self.config.start_frames); |
| | | self.speech_sample_count = u64::from(frame.samples_per_channel) |
| | | * u64::from(frame.num_channels) |
| | | * u64::from(self.config.start_frames); |
| | | self.speech_rms_sum = rms * f64::from(self.config.start_frames); |
| | | self.speech_peak = peak; |
| | | self.speech_samples = self |
| | | .pre_speech_frames |
| | | .iter() |
| | | .flat_map(|samples| samples.iter().copied()) |
| | | .collect(); |
| | | let turn_id = format!("turn-{:04}", self.turn_index); |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | participant_alias = %participant_alias, |
| | | track_sid_alias = %track_sid_alias, |
| | | turn_index = self.turn_index, |
| | | frame_count, |
| | | speech_start_elapsed_ms = self.speech_start_elapsed_ms, |
| | | rms = round4(rms), |
| | | peak = round4(peak), |
| | | "runtime helper vad_speech_start" |
| | | ); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn_id), |
| | | "vad_speech_start", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "turnIndex": self.turn_index, |
| | | "frameCount": frame_count, |
| | | "speechStartElapsedMs": self.speech_start_elapsed_ms, |
| | | "rms": round4(rms), |
| | | "peak": round4(peak), |
| | | }), |
| | | ); |
| | | } |
| | | return None; |
| | | } |
| | | |
| | | self.speech_samples.extend_from_slice(frame.data.as_ref()); |
| | | self.speech_frame_count += 1; |
| | | self.speech_sample_count += |
| | | u64::from(frame.samples_per_channel) * u64::from(frame.num_channels); |
| | | self.speech_rms_sum += rms; |
| | | self.speech_peak = self.speech_peak.max(peak); |
| | | |
| | | let speech_duration_ms = elapsed_ms.saturating_sub(self.speech_start_elapsed_ms); |
| | | if speech_duration_ms >= self.config.max_turn_ms { |
| | | return self.finish_turn( |
| | | call_id, |
| | | trace_id, |
| | | participant_alias, |
| | | track_sid_alias, |
| | | elapsed_ms, |
| | | "max_turn_ms", |
| | | ); |
| | | } |
| | | |
| | | if !voiced && self.silence_run_frames >= self.config.end_silence_frames(frame_duration_ms) { |
| | | return self.finish_turn( |
| | | call_id, |
| | | trace_id, |
| | | participant_alias, |
| | | track_sid_alias, |
| | | elapsed_ms, |
| | | "silence", |
| | | ); |
| | | } |
| | | None |
| | | } |
| | | |
| | | fn finish_stream( |
| | | &mut self, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | participant_alias: &str, |
| | | track_sid_alias: &str, |
| | | elapsed_ms: u64, |
| | | ) -> Option<FinishedSpeechTurn> { |
| | | if self.in_speech { |
| | | return self.finish_turn( |
| | | call_id, |
| | | trace_id, |
| | | participant_alias, |
| | | track_sid_alias, |
| | | elapsed_ms, |
| | | "stream_end", |
| | | ); |
| | | } else if self.turn_index == 0 { |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | participant_alias = %participant_alias, |
| | | track_sid_alias = %track_sid_alias, |
| | | observed_wall_ms = elapsed_ms, |
| | | "runtime helper vad_no_speech_summary" |
| | | ); |
| | | } |
| | | None |
| | | } |
| | | |
| | | fn finish_turn( |
| | | &mut self, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | participant_alias: &str, |
| | | track_sid_alias: &str, |
| | | elapsed_ms: u64, |
| | | end_reason: &str, |
| | | ) -> Option<FinishedSpeechTurn> { |
| | | let speech_duration_ms = elapsed_ms.saturating_sub(self.speech_start_elapsed_ms); |
| | | let event_name = if speech_duration_ms < self.config.min_speech_ms { |
| | | "runtime helper vad_speech_too_short" |
| | | } else { |
| | | "runtime helper vad_speech_end" |
| | | }; |
| | | let avg_rms = if self.speech_frame_count == 0 { |
| | | 0.0 |
| | | } else { |
| | | self.speech_rms_sum / self.speech_frame_count as f64 |
| | | }; |
| | | |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | participant_alias = %participant_alias, |
| | | track_sid_alias = %track_sid_alias, |
| | | turn_index = self.turn_index, |
| | | speech_start_elapsed_ms = self.speech_start_elapsed_ms, |
| | | speech_end_elapsed_ms = elapsed_ms, |
| | | speech_duration_ms, |
| | | speech_frame_count = self.speech_frame_count, |
| | | speech_sample_count = self.speech_sample_count, |
| | | avg_rms = round4(avg_rms), |
| | | peak = round4(self.speech_peak), |
| | | end_reason = %end_reason, |
| | | "{}", event_name |
| | | ); |
| | | |
| | | let finished_turn = if speech_duration_ms >= self.config.min_speech_ms { |
| | | let turn_id = format!("turn-{:04}", self.turn_index); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(&turn_id), |
| | | "vad_speech_end", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "turnIndex": self.turn_index, |
| | | "speechStartElapsedMs": self.speech_start_elapsed_ms, |
| | | "speechEndElapsedMs": elapsed_ms, |
| | | "speechDurationMs": speech_duration_ms, |
| | | "speechFrameCount": self.speech_frame_count, |
| | | "speechSampleCount": self.speech_sample_count, |
| | | "avgRms": round4(avg_rms), |
| | | "peak": round4(self.speech_peak), |
| | | "endReason": end_reason, |
| | | }), |
| | | ); |
| | | Some(FinishedSpeechTurn { |
| | | turn_index: self.turn_index, |
| | | turn_id, |
| | | samples: std::mem::take(&mut self.speech_samples), |
| | | duration_ms: speech_duration_ms, |
| | | frame_count: self.speech_frame_count, |
| | | sample_count: self.speech_sample_count, |
| | | end_reason: end_reason.to_string(), |
| | | }) |
| | | } else { |
| | | None |
| | | }; |
| | | self.reset_current_turn(); |
| | | finished_turn |
| | | } |
| | | |
| | | fn reset_current_turn(&mut self) { |
| | | self.in_speech = false; |
| | | self.voiced_run_frames = 0; |
| | | self.silence_run_frames = 0; |
| | | self.speech_start_elapsed_ms = 0; |
| | | self.speech_frame_count = 0; |
| | | self.speech_sample_count = 0; |
| | | self.speech_rms_sum = 0.0; |
| | | self.speech_peak = 0.0; |
| | | self.speech_samples.clear(); |
| | | self.pre_speech_frames.clear(); |
| | | } |
| | | |
| | | fn remember_pre_speech_frame(&mut self, frame: &AudioFrame<'_>) { |
| | | self.pre_speech_frames.push(frame.data.as_ref().to_vec()); |
| | | let max_frames = self.config.start_frames as usize; |
| | | if self.pre_speech_frames.len() > max_frames { |
| | | let remove_count = self.pre_speech_frames.len() - max_frames; |
| | | self.pre_speech_frames.drain(0..remove_count); |
| | | } |
| | | } |
| | | } |
| | | |
| | | fn frame_duration_ms(frame: &AudioFrame<'_>) -> u64 { |
| | | if frame.sample_rate == 0 { |
| | | return 10; |
| | | } |
| | | ((u64::from(frame.samples_per_channel) * 1000) / u64::from(frame.sample_rate)).max(1) |
| | | } |
| | | |
| | | fn pcm_energy_stats(samples: &[i16]) -> (f64, f64) { |
| | | if samples.is_empty() { |
| | | return (0.0, 0.0); |
| | | } |
| | | let mut square_sum = 0.0; |
| | | let mut peak = 0.0; |
| | | for sample in samples { |
| | | let normalized = f64::from(*sample) / f64::from(i16::MAX); |
| | | square_sum += normalized * normalized; |
| | | let abs = normalized.abs(); |
| | | if abs > peak { |
| | | peak = abs; |
| | | } |
| | | } |
| | | ((square_sum / samples.len() as f64).sqrt(), peak) |
| | | } |
| | | |
| | | fn round4(value: f64) -> f64 { |
| | | (value * 10_000.0).round() / 10_000.0 |
| | | } |
| | | |
| | | struct BotAudioOutputSink { |
| | | room: Arc<Room>, |
| | | rtc_source: NativeAudioSource, |
| | | track: LocalAudioTrack, |
| | | device_output_destination_identity: Option<String>, |
| | | } |
| | | |
| | | impl BotAudioOutputSink { |
| | |
| | | room: Arc<Room>, |
| | | room_name: &str, |
| | | participant_identity: &str, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | track_name: &str, |
| | | sample_rate: u32, |
| | | num_channels: u32, |
| | | device_output_destination_identity: Option<String>, |
| | | ) -> Result<Self> { |
| | | let rtc_source = NativeAudioSource::new( |
| | | AudioSourceOptions::default(), |
| | |
| | | track_name, |
| | | RtcAudioSource::Native(rtc_source.clone()), |
| | | ); |
| | | let room_alias = redact(room_name); |
| | | let participant_alias = redact(participant_identity); |
| | | |
| | | room.local_participant() |
| | | .publish_track( |
| | |
| | | .await |
| | | .map_err(|error| { |
| | | anyhow!( |
| | | "failed to publish bot audio track in room {room_name} for participant {participant_identity}: {error}" |
| | | "failed to publish bot audio track in room {room_alias} for participant {participant_alias}: {error}" |
| | | ) |
| | | })?; |
| | | |
| | | info!( |
| | | room_id = %room_name, |
| | | participant_alias = %redact(participant_identity), |
| | | room_alias = %room_alias, |
| | | participant_alias = %participant_alias, |
| | | track_name = %track_name, |
| | | sample_rate, |
| | | num_channels, |
| | | "runtime helper published bot audio track" |
| | | ); |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | None, |
| | | "bot_track_ready", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "trackName": track_name, |
| | | "sampleRate": sample_rate, |
| | | "numChannels": num_channels, |
| | | }), |
| | | ); |
| | | |
| | | Ok(Self { |
| | | room, |
| | | rtc_source, |
| | | track, |
| | | device_output_destination_identity, |
| | | }) |
| | | } |
| | | |
| | | async fn publish_device_output( |
| | | &self, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | turn_id: &str, |
| | | output: &RuntimeTurnDeviceOutput, |
| | | ) -> Result<()> { |
| | | let command_id = output |
| | | .command_id |
| | | .as_deref() |
| | | .map(str::trim) |
| | | .filter(|value| !value.is_empty()) |
| | | .ok_or_else(|| anyhow!("device output commandId missing"))?; |
| | | let command_code = output |
| | | .command_code |
| | | .as_deref() |
| | | .map(str::trim) |
| | | .filter(|value| !value.is_empty()) |
| | | .ok_or_else(|| anyhow!("device output commandCode missing"))?; |
| | | let payload = serde_json::json!({ |
| | | "type": "device_output", |
| | | "schemaVersion": "1.0", |
| | | "callId": call_id, |
| | | "traceId": trace_id, |
| | | "turnId": turn_id, |
| | | "commandId": command_id, |
| | | "commandCode": command_code, |
| | | "params": output.params.clone().unwrap_or_else(|| serde_json::json!({})), |
| | | "source": { |
| | | "kind": "voice_command" |
| | | }, |
| | | }); |
| | | let destinations = self |
| | | .device_output_destination_identity |
| | | .as_ref() |
| | | .filter(|value| !value.trim().is_empty()) |
| | | .map(|value| vec![ParticipantIdentity(value.trim().to_string())]) |
| | | .unwrap_or_default(); |
| | | let destination_count = destinations.len(); |
| | | self.room |
| | | .local_participant() |
| | | .publish_data(DataPacket { |
| | | payload: serde_json::to_vec(&payload) |
| | | .context("failed to encode device output data message")?, |
| | | topic: Some("device_output".to_string()), |
| | | reliable: true, |
| | | destination_identities: destinations, |
| | | }) |
| | | .await |
| | | .map_err(|error| anyhow!("failed to publish device output data message: {error}"))?; |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn_id, |
| | | command_id_alias = %redact(command_id), |
| | | command_code = %command_code, |
| | | destination_count, |
| | | "runtime helper device_output_sent" |
| | | ); |
| | | Ok(()) |
| | | } |
| | | |
| | | async fn publish_reply_state( |
| | | &self, |
| | | call_id: &str, |
| | | trace_id: &str, |
| | | turn_id: &str, |
| | | reply_playback_mode: &str, |
| | | state: &str, |
| | | seq: u64, |
| | | ) -> Result<()> { |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(turn_id), |
| | | "reply_state_send_started", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "state": state, |
| | | "seq": seq, |
| | | "replyPlaybackMode": reply_playback_mode, |
| | | }), |
| | | ); |
| | | let payload = serde_json::json!({ |
| | | "type": "reply_state", |
| | | "schemaVersion": "1.0", |
| | | "callId": call_id, |
| | | "traceId": trace_id, |
| | | "turnId": turn_id, |
| | | "replyPlaybackMode": reply_playback_mode, |
| | | "state": state, |
| | | "seq": seq, |
| | | "tsMs": current_time_millis(), |
| | | }); |
| | | let destinations = self |
| | | .device_output_destination_identity |
| | | .as_ref() |
| | | .filter(|value| !value.trim().is_empty()) |
| | | .map(|value| vec![ParticipantIdentity(value.trim().to_string())]) |
| | | .unwrap_or_default(); |
| | | let destination_count = destinations.len(); |
| | | let publish_result = self |
| | | .room |
| | | .local_participant() |
| | | .publish_data(DataPacket { |
| | | payload: serde_json::to_vec(&payload) |
| | | .context("failed to encode reply state data message")?, |
| | | topic: Some("combrabo_voice.reply_state".to_string()), |
| | | reliable: true, |
| | | destination_identities: destinations, |
| | | }) |
| | | .await |
| | | .map_err(|error| anyhow!("failed to publish reply state data message: {error}")); |
| | | match publish_result { |
| | | Ok(_) => { |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(turn_id), |
| | | "reply_state_send_finished", |
| | | "ok", |
| | | None, |
| | | None, |
| | | json!({ |
| | | "state": state, |
| | | "seq": seq, |
| | | "replyPlaybackMode": reply_playback_mode, |
| | | "destinationCount": destination_count, |
| | | }), |
| | | ); |
| | | info!( |
| | | call_id = %call_id, |
| | | trace_id = %trace_id, |
| | | turn_id = %turn_id, |
| | | state = %state, |
| | | seq, |
| | | reply_playback_mode = %reply_playback_mode, |
| | | destination_count, |
| | | "runtime helper reply_state_sent" |
| | | ); |
| | | Ok(()) |
| | | } |
| | | Err(error) => { |
| | | emit_activity( |
| | | call_id, |
| | | trace_id, |
| | | Some(turn_id), |
| | | "reply_state_send_failed", |
| | | "failed", |
| | | Some("REPLY_STATE_SEND_FAILED"), |
| | | Some(true), |
| | | json!({ |
| | | "state": state, |
| | | "seq": seq, |
| | | "replyPlaybackMode": reply_playback_mode, |
| | | }), |
| | | ); |
| | | Err(error) |
| | | } |
| | | } |
| | | } |
| | | |
| | | async fn write_pcm_frame(&self, frame: &audio::PcmFrame) -> Result<()> { |
| | | let audio_frame = AudioFrame { |
| | | data: frame.data.as_slice().into(), |