From b0a4e2ea93fafc4d65d474f3ada616a0ad5e19d5 Mon Sep 17 00:00:00 2001
From: Ariver <ar@Arm1.local>
Date: Sun, 12 Jul 2026 11:39:35 +0800
Subject: [PATCH] feat: bind helper stream timing markers
---
src/main.rs | 3849 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
1 files changed, 3,702 insertions(+), 147 deletions(-)
diff --git a/src/main.rs b/src/main.rs
index 07f4e9f..611993a 100644
--- a/src/main.rs
+++ b/src/main.rs
@@ -1,30 +1,73 @@
+mod asr_realtime;
mod audio;
+mod service;
-use std::{env, sync::Arc, time::Duration};
+use std::{
+ borrow::Cow,
+ collections::HashSet,
+ 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 asr_realtime::{RealtimeAsrConfig, RealtimeAsrOutcome, RealtimeAsrUpload};
+use audio::{AudioDiagnostics, 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 sha2::{Digest, Sha256};
+use tokio::time::{sleep, sleep_until, timeout};
+use tokio::{
+ sync::{mpsc, mpsc::UnboundedReceiver, watch},
+ task::JoinHandle,
+};
use tracing::{info, warn};
-const TARGET_SAMPLE_RATE_HZ: u32 = 48_000;
-const TARGET_NUM_CHANNELS: u16 = 1;
+const USER_AUDIO_SAMPLE_RATE_HZ: u32 = 48_000;
+const USER_AUDIO_NUM_CHANNELS: u16 = 1;
+const INBOUND_AUDIO_QUEUE_CAPACITY: usize = 100;
+const INBOUND_AUDIO_DROP_LOG_INTERVAL: u64 = 100;
+const AUDIO_DRAIN_STOP_GRACE: Duration = Duration::from_secs(2);
+const DEFAULT_BOT_AUDIO_PROFILE: &str = "pcm-16k";
+const LIVEKIT_48K_SAMPLE_RATE_HZ: u32 = 48_000;
+const PCM_16K_SAMPLE_RATE_HZ: u32 = 16_000;
+const BOT_NUM_CHANNELS: u16 = 1;
const TRACK_NAME: &str = "bot-main-audio";
-const DEVICE_OUTPUT_TOPIC: &str = "device_output";
+const STREAM_TIMING_VERSION: u32 = 1;
+const STREAM_TIMING_FIRST_SEGMENT: u64 = 1;
+const STREAM_TIMING_FIRST_CHUNK: u64 = 1;
+const STREAM_TIMING_MAX_ELAPSED_MS: u64 = 5_000;
+const STREAM_TIMING_SOURCE: &str = "stream_anchor_monotonic";
#[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()
@@ -35,12 +78,15 @@
&http,
config.greeting_audio_file.as_deref(),
config.greeting_audio_url.as_deref(),
- TARGET_SAMPLE_RATE_HZ,
- TARGET_NUM_CHANNELS,
+ config.bot_audio_profile.sample_rate_hz,
+ config.bot_audio_profile.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(),
@@ -52,42 +98,138 @@
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,
+ bot_audio_profile = %config.bot_audio_profile.profile,
+ bot_sample_rate_hz = config.bot_audio_profile.sample_rate_hz,
+ bot_num_channels = config.bot_audio_profile.num_channels,
"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.bot_audio_profile.profile.clone(),
+ config.bot_audio_profile.sample_rate_hz,
+ u32::from(config.bot_audio_profile.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!(
@@ -95,9 +237,12 @@
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!(
@@ -133,12 +278,159 @@
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_asr_stream_enabled: bool,
+ runtime_asr_stream_url: Option<String>,
+ runtime_asr_realtime_enabled: bool,
+ runtime_asr_realtime_url: Option<String>,
+ runtime_asr_realtime_chunk_duration_ms: u64,
+ 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,
+ bot_audio_profile: BotAudioProfile,
+}
+
+#[derive(Clone)]
+struct BotAudioProfile {
+ profile: String,
+ sample_rate_hz: u32,
+ num_channels: u16,
+}
+
+impl BotAudioProfile {
+ fn from_env() -> Result<Self> {
+ let profile = env::var("CV_BOT_AUDIO_PROFILE")
+ .unwrap_or_else(|_| DEFAULT_BOT_AUDIO_PROFILE.to_string())
+ .trim()
+ .to_ascii_lowercase();
+ match profile.as_str() {
+ "livekit-48k" | "48k" => Ok(Self {
+ profile: "livekit-48k".to_string(),
+ sample_rate_hz: LIVEKIT_48K_SAMPLE_RATE_HZ,
+ num_channels: BOT_NUM_CHANNELS,
+ }),
+ "pcm-16k" | "16k" => Ok(Self {
+ profile: "pcm-16k".to_string(),
+ sample_rate_hz: PCM_16K_SAMPLE_RATE_HZ,
+ num_channels: BOT_NUM_CHANNELS,
+ }),
+ "custom" => {
+ let sample_rate_hz = u32_env("CV_BOT_SAMPLE_RATE_HZ", LIVEKIT_48K_SAMPLE_RATE_HZ);
+ let num_channels = u16_env("CV_BOT_NUM_CHANNELS", BOT_NUM_CHANNELS);
+ if sample_rate_hz == 0 {
+ return Err(anyhow!("CV_BOT_SAMPLE_RATE_HZ must be positive"));
+ }
+ if num_channels == 0 {
+ return Err(anyhow!("CV_BOT_NUM_CHANNELS must be positive"));
+ }
+ Ok(Self {
+ profile,
+ sample_rate_hz,
+ num_channels,
+ })
+ }
+ _ => Err(anyhow!(
+ "unsupported CV_BOT_AUDIO_PROFILE {}; expected livekit-48k, pcm-16k or custom",
+ profile
+ )),
+ }
+ }
+}
+
+#[derive(Clone)]
+struct TurnBridgeConfig {
+ bridge_url: Option<String>,
+ bridge_token: Option<String>,
+ bridge_mode: String,
+ asr_stream_enabled: bool,
+ asr_stream_url: Option<String>,
+ asr_realtime_enabled: bool,
+ asr_realtime_url: Option<String>,
+ asr_realtime_chunk_duration_ms: u64,
+ 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(),
+ asr_stream_enabled: config.runtime_asr_stream_enabled,
+ asr_stream_url: config.runtime_asr_stream_url.clone(),
+ asr_realtime_enabled: config.runtime_asr_realtime_enabled,
+ asr_realtime_url: config.runtime_asr_realtime_url.clone(),
+ asr_realtime_chunk_duration_ms: config.runtime_asr_realtime_chunk_duration_ms,
+ 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"))
+ }
+
+ fn is_asr_stream_ready(&self) -> bool {
+ self.asr_stream_enabled
+ && self
+ .asr_stream_url
+ .as_ref()
+ .is_some_and(|value| !value.is_empty())
+ && self
+ .bridge_token
+ .as_ref()
+ .is_some_and(|value| !value.is_empty())
+ && self
+ .runtime_session_nonce
+ .as_ref()
+ .is_some_and(|value| !value.is_empty())
+ }
+
+ fn realtime_asr_config(&self) -> RealtimeAsrConfig {
+ RealtimeAsrConfig {
+ enabled: self.asr_realtime_enabled,
+ url: self.asr_realtime_url.clone(),
+ runtime_token: self.bridge_token.clone(),
+ runtime_session_nonce: self.runtime_session_nonce.clone(),
+ chunk_duration_ms: self.asr_realtime_chunk_duration_ms,
+ }
+ }
}
impl Config {
@@ -150,18 +442,124 @@
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_asr_stream_enabled: bool_env("CV_RUNTIME_ASR_STREAM_ENABLED", false),
+ runtime_asr_stream_url: optional_env("CV_RUNTIME_ASR_STREAM_URL"),
+ runtime_asr_realtime_enabled: bool_env("CV_RUNTIME_ASR_REALTIME_ENABLED", false),
+ runtime_asr_realtime_url: optional_env("CV_RUNTIME_ASR_REALTIME_URL"),
+ runtime_asr_realtime_chunk_duration_ms: u64_env(
+ "CV_RUNTIME_ASR_REALTIME_CHUNK_DURATION_MS",
+ 200,
),
+ 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(),
+ bot_audio_profile: BotAudioProfile::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> {
@@ -183,34 +581,44 @@
})
}
-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 u16_env(key: &str, default_value: u16) -> u16 {
+ env::var(key)
+ .ok()
+ .and_then(|value| value.trim().parse::<u16>().ok())
+ .unwrap_or(default_value)
}
fn redact(value: &str) -> String {
@@ -220,118 +628,2969 @@
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,
+ "eventWallTimeMs": current_time_millis(),
+ "result": result,
+ "reasonCode": reason_code,
+ "retryable": retryable,
+ "extension": extension,
+ });
+ println!("{payload}");
+}
+
+fn emit_anchored_activity(
+ call_id: &str,
+ trace_id: &str,
+ turn_id: &str,
+ event_name: &str,
+ marker: &RuntimeTurnStreamTimingMarker,
+ extension: serde_json::Value,
+) {
+ let payload = json!({
+ "type": "cv_activity",
+ "callId": call_id,
+ "traceId": trace_id,
+ "turnId": turn_id,
+ "eventName": event_name,
+ "eventWallTimeMs": current_time_millis(),
+ "serverDeltaMs": marker.server_delta_ms,
+ "serverDeltaSource": STREAM_TIMING_SOURCE,
+ "result": "ok",
+ "reasonCode": null,
+ "retryable": null,
+ "extension": marker.extension_with(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,
+ realtime_asr_result_ref: Option<String>,
+) {
+ 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.as_str(),
+ }),
+ );
+ let asr_result_ref = match realtime_asr_result_ref {
+ Some(value) => Some(value),
+ None => request_asr_result_ref(http, bridge_config, call_id, trace_id, &turn).await,
+ };
+ if bridge_config.is_stream_mode() {
+ match request_turn_bridge_stream(
+ http,
+ bridge_config,
+ sink,
+ call_id,
+ trace_id,
+ &turn,
+ &path_ref,
+ byte_size,
+ asr_result_ref.as_deref(),
+ 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,
+ asr_result_ref.as_deref(),
+ 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,
+ USER_AUDIO_SAMPLE_RATE_HZ,
+ USER_AUDIO_NUM_CHANNELS,
+ )?;
+ let byte_size = fs::metadata(&output_path)
+ .context("failed to stat turn artifact")?
+ .len();
+ Ok((path_ref, byte_size))
+}
+
+async fn request_asr_result_ref(
+ http: &Client,
+ bridge_config: &TurnBridgeConfig,
+ call_id: &str,
+ trace_id: &str,
+ turn: &FinishedSpeechTurn,
+) -> Option<String> {
+ if !bridge_config.is_asr_stream_ready() {
+ return None;
+ }
+ match request_asr_stream(http, bridge_config, call_id, trace_id, turn).await {
+ Ok(Some(asr_result_ref)) => {
+ emit_activity(
+ call_id,
+ trace_id,
+ Some(&turn.turn_id),
+ "asr_stream_ref_ready",
+ "ok",
+ None,
+ None,
+ json!({
+ "asrResultRefPresent": true,
+ "format": "pcm_s16le",
+ "sampleRate": 16000,
+ "channels": 1,
+ }),
+ );
+ Some(asr_result_ref)
+ }
+ Ok(None) => None,
+ Err(error) => {
+ warn!(
+ call_id = %call_id,
+ trace_id = %trace_id,
+ turn_id = %turn.turn_id,
+ error = %safe_error(&error.to_string()),
+ "runtime helper asr_stream_failed_fallback"
+ );
+ emit_activity(
+ call_id,
+ trace_id,
+ Some(&turn.turn_id),
+ "asr_stream_fallback",
+ "skipped",
+ Some("ASR_STREAM_INTERRUPTED"),
+ Some(true),
+ json!({
+ "fallbackReason": "asr_stream_request_failed",
+ "fallbackStage": "asr_stream",
+ }),
+ );
+ None
+ }
+ }
+}
+
+async fn request_asr_stream(
+ http: &Client,
+ bridge_config: &TurnBridgeConfig,
+ call_id: &str,
+ trace_id: &str,
+ turn: &FinishedSpeechTurn,
+) -> Result<Option<String>> {
+ let started_at = Instant::now();
+ let ndjson = build_asr_stream_ndjson(call_id, trace_id, turn, bridge_config)?;
+ let response = http
+ .post(bridge_config.asr_stream_url.as_deref().unwrap_or_default())
+ .header("Content-Type", "application/x-ndjson")
+ .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(),
+ )
+ .body(ndjson)
+ .send()
.await
- .map_err(|error| anyhow!("failed to publish livekit device output data: {error}"))?;
+ .context("failed to post asr stream")?;
+ 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 asr_stream_http_failed"
+ );
+ return Ok(None);
+ }
+ let body: RuntimeTurnCommonResult<RuntimeAsrStreamResp> = response
+ .json()
+ .await
+ .context("failed to decode asr stream response")?;
+ if body.code != 0 {
+ warn!(
+ call_id = %call_id,
+ trace_id = %trace_id,
+ turn_id = %turn.turn_id,
+ code = body.code,
+ msg_len = body.msg.as_deref().unwrap_or_default().len(),
+ "runtime helper asr_stream_common_result_failed"
+ );
+ return Ok(None);
+ }
+ let Some(data) = body.data else {
+ return Ok(None);
+ };
+ if data.status.as_deref() == Some("final") {
+ info!(
+ call_id = %call_id,
+ trace_id = %trace_id,
+ turn_id = %turn.turn_id,
+ chunk_count = data.chunk_count.unwrap_or_default(),
+ audio_bytes = data.audio_bytes.unwrap_or_default(),
+ asr_duration_ms = data.asr_duration_ms.unwrap_or_default(),
+ wall_ms = started_at.elapsed().as_millis() as u64,
+ provider = %data.provider_alias.as_deref().unwrap_or("unknown"),
+ text_len = data.text_len.unwrap_or_default(),
+ "runtime helper asr_stream_final"
+ );
+ return Ok(data.asr_result_ref);
+ }
+ emit_activity(
+ call_id,
+ trace_id,
+ Some(&turn.turn_id),
+ "asr_stream_fallback",
+ "skipped",
+ None,
+ Some(true),
+ json!({
+ "fallbackReason": data.fallback_reason,
+ "fallbackStage": data.fallback_stage,
+ "status": data.status,
+ }),
+ );
+ Ok(None)
+}
+
+fn build_asr_stream_ndjson(
+ call_id: &str,
+ trace_id: &str,
+ turn: &FinishedSpeechTurn,
+ bridge_config: &TurnBridgeConfig,
+) -> Result<String> {
+ let chunks = asr_pcm_16k_chunks(turn)?;
+ let mut seq = 1u64;
+ let mut lines = Vec::with_capacity(chunks.len() + 3);
+ lines.push(serde_json::to_string(&json!({
+ "event": "asr_stream_started",
+ "seq": seq,
+ "callId": call_id,
+ "traceId": trace_id,
+ "turnId": turn.turn_id.as_str(),
+ "tsMs": current_time_millis(),
+ "payload": {
+ "format": "pcm_s16le",
+ "sampleRate": 16000,
+ "channels": 1,
+ "runtimeSessionNonce": bridge_config.runtime_session_nonce.as_deref().unwrap_or_default(),
+ "providerHint": "volcengine",
+ }
+ }))?);
+ for (index, samples) in chunks.iter().enumerate() {
+ seq += 1;
+ let bytes = pcm_i16_to_le_bytes(samples);
+ let duration_ms = ((samples.len() as u64) * 1000 / 16_000).max(1);
+ lines.push(serde_json::to_string(&json!({
+ "event": "asr_audio_chunk",
+ "seq": seq,
+ "callId": call_id,
+ "traceId": trace_id,
+ "turnId": turn.turn_id.as_str(),
+ "tsMs": current_time_millis(),
+ "payload": {
+ "chunkSeq": index + 1,
+ "format": "pcm_s16le",
+ "sampleRate": 16000,
+ "channels": 1,
+ "durationMs": duration_ms,
+ "payloadBase64": general_purpose::STANDARD.encode(bytes),
+ }
+ }))?);
+ }
+ seq += 1;
+ lines.push(serde_json::to_string(&json!({
+ "event": "vad_speech_end",
+ "seq": seq,
+ "callId": call_id,
+ "traceId": trace_id,
+ "turnId": turn.turn_id.as_str(),
+ "tsMs": current_time_millis(),
+ "payload": {
+ "endReason": turn.end_reason.as_str(),
+ "speechDurationMs": turn.duration_ms,
+ }
+ }))?);
+ seq += 1;
+ lines.push(serde_json::to_string(&json!({
+ "event": "asr_stream_finish",
+ "seq": seq,
+ "callId": call_id,
+ "traceId": trace_id,
+ "turnId": turn.turn_id.as_str(),
+ "tsMs": current_time_millis(),
+ "payload": {
+ "finalChunkSeq": chunks.len(),
+ "audioDurationMs": turn.duration_ms,
+ }
+ }))?);
+ Ok(lines.join("\n") + "\n")
+}
+
+fn asr_pcm_16k_chunks(turn: &FinishedSpeechTurn) -> Result<Vec<Vec<i16>>> {
+ if turn.samples.is_empty() {
+ return Err(anyhow!("empty turn samples"));
+ }
+ let samples_16k: Vec<i16> = turn.samples.iter().step_by(3).copied().collect();
+ if samples_16k.is_empty() {
+ return Err(anyhow!("empty 16k asr samples"));
+ }
+ let samples_per_chunk = 320usize;
+ Ok(samples_16k
+ .chunks(samples_per_chunk)
+ .map(|chunk| chunk.to_vec())
+ .collect())
+}
+
+fn pcm_i16_to_le_bytes(samples: &[i16]) -> Vec<u8> {
+ let mut bytes = Vec::with_capacity(samples.len() * 2);
+ for sample in samples {
+ bytes.extend_from_slice(&sample.to_le_bytes());
+ }
+ bytes
+}
+
+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,
+ asr_result_ref: Option<&str>,
+ 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: Some(RuntimeTurnAudioArtifact {
+ artifact_type: "local_file".to_string(),
+ path_ref: path_ref.to_string(),
+ format: "wav".to_string(),
+ sample_rate: USER_AUDIO_SAMPLE_RATE_HZ,
+ channels: u32::from(USER_AUDIO_NUM_CHANNELS),
+ duration_ms: turn.duration_ms,
+ byte_size,
+ }),
+ asr_result_ref: asr_result_ref.map(str::to_string),
+ };
+ 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 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?;
+ }
+ }
+ state.close_timing();
+ 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,
+ &event,
+ 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;
+ state.close_timing();
+ 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") => {
+ state.close_timing();
+ 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;
+ state.close_timing();
+ 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() {
+ if activity.event_type.as_deref() == Some("tts_first_audio_chunk_ready") {
+ if let Some(runtime_session_nonce) =
+ bridge_config.runtime_session_nonce.as_deref()
+ {
+ let _ = state.arm_timing_anchor(
+ call_id,
+ trace_id,
+ &turn.turn_id,
+ runtime_session_nonce,
+ &event,
+ );
+ }
+ }
+ 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,
+ event: &RuntimeTurnStreamEvent,
+ 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();
+ if !matches!(format.as_str(), "pcm_s16le" | "mp3" | "mpeg" | "wav") {
+ return Err(anyhow!("unsupported reply_audio_chunk format {format}"));
+ }
+ match state.reply_chunk_markers.observe(audio_chunk.segment_seq) {
+ ReplyChunkMarker::FirstReply => {
+ let extension = json!({
+ "segmentSeq": audio_chunk.segment_seq,
+ "chunkSeq": audio_chunk.chunk_seq,
+ "format": format.as_str(),
+ "bytes": payload.len(),
+ });
+ if let Some(marker) =
+ state.record_m6(call_id, trace_id, &turn.turn_id, event, audio_chunk)
+ {
+ emit_anchored_activity(
+ call_id,
+ trace_id,
+ &turn.turn_id,
+ "helper_first_reply_audio_chunk_received",
+ &marker,
+ extension,
+ );
+ } else {
+ emit_activity(
+ call_id,
+ trace_id,
+ Some(&turn.turn_id),
+ "helper_first_reply_audio_chunk_received",
+ "ok",
+ None,
+ None,
+ extension,
+ );
+ }
+ }
+ ReplyChunkMarker::SegmentFirst => emit_activity(
+ call_id,
+ trace_id,
+ Some(&turn.turn_id),
+ "helper_segment_first_audio_chunk_received",
+ "ok",
+ None,
+ None,
+ json!({
+ "segmentSeq": audio_chunk.segment_seq,
+ "chunkSeq": audio_chunk.chunk_seq,
+ "format": format.as_str(),
+ "bytes": payload.len(),
+ }),
+ ),
+ ReplyChunkMarker::None => {}
+ }
+ let frames = if format == "pcm_s16le" {
+ let sample_rate = audio_chunk.sample_rate.unwrap_or(sink.sample_rate_hz);
+ let channels = audio_chunk.channels.unwrap_or(u32::from(sink.num_channels));
+ if state.pcm_stream_decoder.is_none() {
+ state.pcm_stream_decoder = Some(audio::PcmS16leStreamDecoder::new(
+ sample_rate,
+ channels,
+ sink.sample_rate_hz,
+ sink.num_channels,
+ )?);
+ }
+ state.pcm_stream_network_chunk_count =
+ state.pcm_stream_network_chunk_count.saturating_add(1);
+ let stream_result = state
+ .pcm_stream_decoder
+ .as_mut()
+ .expect("pcm stream decoder initialized")
+ .push_bytes(
+ &payload,
+ sample_rate,
+ channels,
+ audio_chunk.last.unwrap_or(false),
+ )?;
+ if stream_result.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 = stream_result.dropped_tail_bytes,
+ "runtime helper stream_audio_pcm_unaligned_tail_dropped"
+ );
+ }
+ if stream_result.frames.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_source_bytes = stream_result.buffered_source_bytes,
+ buffered_source_samples = stream_result.buffered_source_samples,
+ network_chunk_count = state.pcm_stream_network_chunk_count,
+ "runtime helper stream_audio_pcm_waiting_for_20ms_frame"
+ );
+ return Ok(0);
+ }
+ stream_result.frames
+ } 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",
+ sink.sample_rate_hz,
+ sink.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 {
+ unreachable!("supported encoded format checked above")
+ };
+ 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;
+ let extension = 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,
+ });
+ if let Some(marker) = state.record_m7() {
+ emit_anchored_activity(
+ call_id,
+ trace_id,
+ &turn.turn_id,
+ "bot_reply_first_audio_frame_written",
+ &marker,
+ extension,
+ );
+ } else {
+ emit_activity(
+ call_id,
+ trace_id,
+ Some(&turn.turn_id),
+ "bot_reply_first_audio_frame_written",
+ "ok",
+ None,
+ None,
+ extension,
+ );
+ }
+ }
+ sleep_until(pacing_started_at + Duration::from_millis(((index + 1) as u64) * 20)).await;
+ }
+ if audio_chunk.last.unwrap_or(false) {
+ let mut debug_source_path = None;
+ let mut debug_pcm_wav_path = None;
+ let mut debug_pcm_wav_size_bytes = None;
+ if format == "pcm_s16le" {
+ if let (Some(debug_dump_dir), Some(decoder)) = (
+ bridge_config.audio_debug_dump_dir.as_deref(),
+ state.pcm_stream_decoder.as_ref(),
+ ) {
+ match decoder.write_debug_dump(
+ debug_dump_dir,
+ call_id,
+ &format!("stream-reply-{}", turn.turn_id),
+ ) {
+ Ok(debug_dump) => {
+ debug_source_path = debug_dump.debug_source_path;
+ debug_pcm_wav_path = debug_dump.debug_pcm_wav_path;
+ debug_pcm_wav_size_bytes = debug_dump.debug_pcm_wav_size_bytes;
+ }
+ Err(error) => warn!(
+ call_id = %call_id,
+ trace_id = %trace_id,
+ turn_id = %turn.turn_id,
+ error = %safe_error(&error.to_string()),
+ "runtime helper stream_audio_pcm_debug_dump_failed"
+ ),
+ }
+ }
+ }
+ 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,
+ "networkChunkCount": state.pcm_stream_network_chunk_count,
+ "debugSourcePath": debug_source_path,
+ "debugPcmWavPath": debug_pcm_wav_path,
+ "debugPcmWavSizeBytes": debug_pcm_wav_size_bytes,
+ "sourceSampleRate": audio_chunk.sample_rate,
+ "sourceChannels": audio_chunk.channels,
+ "targetAudioProfile": sink.profile.as_str(),
+ "targetSampleRate": sink.sample_rate_hz,
+ "targetChannels": sink.num_channels,
+ "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64,
+ }),
+ );
+ }
+ Ok(frames.len())
+}
+
+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,
+ asr_result_ref: Option<&str>,
+ 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: Some(RuntimeTurnAudioArtifact {
+ artifact_type: "local_file".to_string(),
+ path_ref: path_ref.to_string(),
+ format: "wav".to_string(),
+ sample_rate: USER_AUDIO_SAMPLE_RATE_HZ,
+ channels: u32::from(USER_AUDIO_NUM_CHANNELS),
+ duration_ms: turn.duration_ms,
+ byte_size,
+ }),
+ asr_result_ref: asr_result_ref.map(str::to_string),
+ };
+ 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,
+ sink.sample_rate_hz,
+ sink.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")]
+ #[serde(skip_serializing_if = "Option::is_none")]
+ audio_artifact: Option<RuntimeTurnAudioArtifact>,
+ #[serde(rename = "asrResultRef", skip_serializing_if = "Option::is_none")]
+ asr_result_ref: Option<String>,
+}
+
+#[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(Deserialize)]
+#[serde(rename_all = "camelCase")]
+struct RuntimeAsrStreamResp {
+ status: Option<String>,
+ asr_result_ref: Option<String>,
+ chunk_count: Option<u64>,
+ audio_bytes: Option<u64>,
+ asr_duration_ms: Option<u64>,
+ provider_alias: Option<String>,
+ text_len: Option<u64>,
+ fallback_reason: Option<String>,
+ fallback_stage: Option<String>,
+}
+
+#[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_stream_decoder: Option<audio::PcmS16leStreamDecoder>,
+ pcm_stream_network_chunk_count: u64,
+ reply_chunk_markers: ReplyChunkMarkerState,
+ timing: RuntimeTurnStreamTimingState,
+}
+
+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_stream_decoder: None,
+ pcm_stream_network_chunk_count: 0,
+ reply_chunk_markers: ReplyChunkMarkerState::default(),
+ timing: RuntimeTurnStreamTimingState::default(),
+ }
}
}
-fn command_id(call_id: &str, sequence: u8) -> String {
- format!("{call_id}-device-smoke-{sequence:03}")
+impl RuntimeTurnStreamState {
+ fn arm_timing_anchor(
+ &mut self,
+ call_id: &str,
+ trace_id: &str,
+ turn_id: &str,
+ runtime_session_nonce: &str,
+ event: &RuntimeTurnStreamEvent,
+ ) -> bool {
+ if self.timing.phase != RuntimeTurnStreamTimingPhase::Empty
+ || event.call_id.as_deref() != Some(call_id)
+ || event.trace_id.as_deref() != Some(trace_id)
+ || event.turn_id.as_deref() != Some(turn_id)
+ {
+ return false;
+ }
+ let Some(activity) = event.activity.as_ref() else {
+ return false;
+ };
+ if activity.event_type.as_deref() != Some("tts_first_audio_chunk_ready") {
+ return false;
+ }
+ let Some(extension) = activity.extension.as_ref() else {
+ return false;
+ };
+ let Some(anchor_id) = extension.stream_anchor_id.as_deref() else {
+ return false;
+ };
+ let valid_anchor_id = (16..=64).contains(&anchor_id.len()) && anchor_id.is_ascii();
+ let expected_nonce_hash = runtime_session_nonce_hash(runtime_session_nonce);
+ if extension.stream_timing_version != Some(STREAM_TIMING_VERSION)
+ || !valid_anchor_id
+ || extension.runtime_session_nonce_hash.as_deref() != Some(expected_nonce_hash.as_str())
+ || extension.segment_seq != Some(STREAM_TIMING_FIRST_SEGMENT)
+ || extension.stream_timing_validation.as_deref() != Some("bound")
+ {
+ return false;
+ }
+ let Some(server_delta_ms) = extension.stream_anchor_server_delta_ms else {
+ return false;
+ };
+ self.timing.anchor = Some(RuntimeTurnStreamTimingAnchor {
+ call_id: call_id.to_string(),
+ trace_id: trace_id.to_string(),
+ turn_id: turn_id.to_string(),
+ anchor_id: anchor_id.to_string(),
+ server_delta_ms,
+ runtime_session_nonce_hash: expected_nonce_hash,
+ segment_seq: STREAM_TIMING_FIRST_SEGMENT,
+ received_at: Instant::now(),
+ });
+ self.timing.phase = RuntimeTurnStreamTimingPhase::Armed;
+ true
+ }
+
+ fn record_m6(
+ &mut self,
+ call_id: &str,
+ trace_id: &str,
+ turn_id: &str,
+ event: &RuntimeTurnStreamEvent,
+ audio_chunk: &RuntimeTurnStreamAudioChunk,
+ ) -> Option<RuntimeTurnStreamTimingMarker> {
+ if self.timing.phase != RuntimeTurnStreamTimingPhase::Armed {
+ return None;
+ }
+ let anchor = self.timing.anchor.as_ref()?;
+ if anchor.call_id != call_id
+ || anchor.trace_id != trace_id
+ || anchor.turn_id != turn_id
+ || event.call_id.as_deref() != Some(call_id)
+ || event.trace_id.as_deref() != Some(trace_id)
+ || event.turn_id.as_deref() != Some(turn_id)
+ || audio_chunk.segment_seq != Some(anchor.segment_seq)
+ || audio_chunk.chunk_seq != Some(STREAM_TIMING_FIRST_CHUNK)
+ || audio_chunk.stream_timing_version != Some(STREAM_TIMING_VERSION)
+ || audio_chunk.stream_anchor_id.as_deref() != Some(anchor.anchor_id.as_str())
+ {
+ return None;
+ }
+ let Some(marker) = RuntimeTurnStreamTimingMarker::from_anchor(
+ anchor,
+ audio_chunk.chunk_seq.unwrap_or(STREAM_TIMING_FIRST_CHUNK),
+ ) else {
+ self.close_timing();
+ return None;
+ };
+ self.timing.chunk_seq = Some(marker.chunk_seq);
+ self.timing.phase = RuntimeTurnStreamTimingPhase::M6Recorded;
+ Some(marker)
+ }
+
+ fn record_m7(&mut self) -> Option<RuntimeTurnStreamTimingMarker> {
+ if self.timing.phase != RuntimeTurnStreamTimingPhase::M6Recorded {
+ return None;
+ }
+ let anchor = self.timing.anchor.as_ref()?;
+ let Some(marker) = RuntimeTurnStreamTimingMarker::from_anchor(
+ anchor,
+ self.timing.chunk_seq.unwrap_or(STREAM_TIMING_FIRST_CHUNK),
+ ) else {
+ self.close_timing();
+ return None;
+ };
+ self.timing.phase = RuntimeTurnStreamTimingPhase::M7Recorded;
+ Some(marker)
+ }
+
+ fn close_timing(&mut self) {
+ self.timing.close();
+ }
+}
+
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+enum RuntimeTurnStreamTimingPhase {
+ Empty,
+ Armed,
+ M6Recorded,
+ M7Recorded,
+ Closed,
+}
+
+struct RuntimeTurnStreamTimingState {
+ phase: RuntimeTurnStreamTimingPhase,
+ anchor: Option<RuntimeTurnStreamTimingAnchor>,
+ chunk_seq: Option<u64>,
+}
+
+impl Default for RuntimeTurnStreamTimingState {
+ fn default() -> Self {
+ Self {
+ phase: RuntimeTurnStreamTimingPhase::Empty,
+ anchor: None,
+ chunk_seq: None,
+ }
+ }
+}
+
+impl RuntimeTurnStreamTimingState {
+ fn close(&mut self) {
+ self.anchor = None;
+ self.chunk_seq = None;
+ self.phase = RuntimeTurnStreamTimingPhase::Closed;
+ }
+}
+
+impl Drop for RuntimeTurnStreamTimingState {
+ fn drop(&mut self) {
+ self.anchor = None;
+ self.chunk_seq = None;
+ }
+}
+
+struct RuntimeTurnStreamTimingAnchor {
+ call_id: String,
+ trace_id: String,
+ turn_id: String,
+ anchor_id: String,
+ server_delta_ms: u64,
+ runtime_session_nonce_hash: String,
+ segment_seq: u64,
+ received_at: Instant,
+}
+
+struct RuntimeTurnStreamTimingMarker {
+ version: u32,
+ anchor_id: String,
+ anchor_server_delta_ms: u64,
+ anchor_elapsed_ms: u64,
+ server_delta_ms: u64,
+ runtime_session_nonce_hash: String,
+ segment_seq: u64,
+ chunk_seq: u64,
+}
+
+impl RuntimeTurnStreamTimingMarker {
+ fn from_anchor(anchor: &RuntimeTurnStreamTimingAnchor, chunk_seq: u64) -> Option<Self> {
+ let elapsed_ms = u64::try_from(anchor.received_at.elapsed().as_millis()).ok()?;
+ if elapsed_ms > STREAM_TIMING_MAX_ELAPSED_MS {
+ return None;
+ }
+ Some(Self {
+ version: STREAM_TIMING_VERSION,
+ anchor_id: anchor.anchor_id.clone(),
+ anchor_server_delta_ms: anchor.server_delta_ms,
+ anchor_elapsed_ms: elapsed_ms,
+ server_delta_ms: anchor.server_delta_ms.checked_add(elapsed_ms)?,
+ runtime_session_nonce_hash: anchor.runtime_session_nonce_hash.clone(),
+ segment_seq: anchor.segment_seq,
+ chunk_seq,
+ })
+ }
+
+ fn extension_with(&self, extra: serde_json::Value) -> serde_json::Value {
+ let mut extension = match extra {
+ serde_json::Value::Object(value) => value,
+ _ => serde_json::Map::new(),
+ };
+ extension.insert("streamTimingVersion".to_string(), json!(self.version));
+ extension.insert("streamAnchorId".to_string(), json!(self.anchor_id));
+ extension.insert(
+ "streamAnchorServerDeltaMs".to_string(),
+ json!(self.anchor_server_delta_ms),
+ );
+ extension.insert("anchorElapsedMs".to_string(), json!(self.anchor_elapsed_ms));
+ extension.insert(
+ "runtimeSessionNonceHash".to_string(),
+ json!(self.runtime_session_nonce_hash),
+ );
+ extension.insert("segmentSeq".to_string(), json!(self.segment_seq));
+ extension.insert("chunkSeq".to_string(), json!(self.chunk_seq));
+ extension.insert("streamTimingValidation".to_string(), json!("bound"));
+ serde_json::Value::Object(extension)
+ }
+}
+
+fn runtime_session_nonce_hash(value: &str) -> String {
+ let digest = Sha256::digest(value.as_bytes());
+ digest[..6]
+ .iter()
+ .map(|byte| format!("{byte:02x}"))
+ .collect()
+}
+
+#[derive(Debug, PartialEq, Eq)]
+enum ReplyChunkMarker {
+ FirstReply,
+ SegmentFirst,
+ None,
+}
+
+#[derive(Default)]
+struct ReplyChunkMarkerState {
+ first_reply_seen: bool,
+ seen_segments: HashSet<u64>,
+}
+
+impl ReplyChunkMarkerState {
+ fn observe(&mut self, segment_seq: Option<u64>) -> ReplyChunkMarker {
+ let first_for_segment = segment_seq
+ .map(|value| self.seen_segments.insert(value))
+ .unwrap_or(false);
+ if !self.first_reply_seen {
+ self.first_reply_seen = true;
+ return ReplyChunkMarker::FirstReply;
+ }
+ if first_for_segment {
+ return ReplyChunkMarker::SegmentFirst;
+ }
+ ReplyChunkMarker::None
+ }
+}
+
+#[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 = "callId")]
+ call_id: Option<String>,
+ #[serde(rename = "traceId")]
+ trace_id: Option<String>,
+ #[serde(rename = "turnId")]
+ turn_id: 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>,
+ #[serde(rename = "segmentSeq")]
+ segment_seq: Option<u64>,
+ #[serde(rename = "streamTimingVersion")]
+ stream_timing_version: Option<u32>,
+ #[serde(rename = "streamAnchorId")]
+ stream_anchor_id: Option<String>,
+ 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>,
+ extension: Option<RuntimeTurnStreamTimingExtension>,
+}
+
+#[derive(Clone, Deserialize)]
+struct RuntimeTurnStreamTimingExtension {
+ #[serde(rename = "streamTimingVersion")]
+ stream_timing_version: Option<u32>,
+ #[serde(rename = "streamAnchorId")]
+ stream_anchor_id: Option<String>,
+ #[serde(rename = "streamAnchorServerDeltaMs")]
+ stream_anchor_server_delta_ms: Option<u64>,
+ #[serde(rename = "runtimeSessionNonceHash")]
+ runtime_session_nonce_hash: Option<String>,
+ #[serde(rename = "segmentSeq")]
+ segment_seq: Option<u64>,
+ #[serde(rename = "streamTimingValidation")]
+ stream_timing_validation: 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(),
+ USER_AUDIO_SAMPLE_RATE_HZ as i32,
+ i32::from(USER_AUDIO_NUM_CHANNELS),
+ );
+ let started_at = Instant::now();
+ let (frame_tx, mut frame_rx) =
+ mpsc::channel::<DrainedUserAudioFrame>(INBOUND_AUDIO_QUEUE_CAPACITY);
+ let (drain_shutdown_tx, mut drain_shutdown_rx) = watch::channel(false);
+ let drain_call_id = call_id.clone();
+ let drain_trace_id = trace_id.clone();
+ let drain_task = tokio::spawn(async move {
+ let mut received_frame_count: u64 = 0;
+ let mut dropped_frame_count: u64 = 0;
+ loop {
+ tokio::select! {
+ changed = drain_shutdown_rx.changed() => {
+ match changed {
+ Ok(()) if *drain_shutdown_rx.borrow() => break,
+ Ok(()) => {}
+ Err(_) => break,
+ }
+ }
+ maybe_frame = stream.next() => {
+ let Some(frame) = maybe_frame else {
+ break;
+ };
+ received_frame_count = received_frame_count.saturating_add(1);
+ let drained = DrainedUserAudioFrame {
+ frame_index: received_frame_count,
+ captured_elapsed_ms: started_at.elapsed().as_millis() as u64,
+ frame: AudioFrame {
+ data: Cow::Owned(frame.data.as_ref().to_vec()),
+ sample_rate: frame.sample_rate,
+ num_channels: frame.num_channels,
+ samples_per_channel: frame.samples_per_channel,
+ },
+ };
+ match frame_tx.try_send(drained) {
+ Ok(()) => {}
+ Err(mpsc::error::TrySendError::Full(_)) => {
+ dropped_frame_count = dropped_frame_count.saturating_add(1);
+ if dropped_frame_count % INBOUND_AUDIO_DROP_LOG_INTERVAL == 1 {
+ warn!(
+ call_id = %drain_call_id,
+ trace_id = %drain_trace_id,
+ dropped_frame_count,
+ received_frame_count,
+ queue_capacity = INBOUND_AUDIO_QUEUE_CAPACITY,
+ "runtime helper inbound_audio_queue_full_dropping_newest"
+ );
+ }
+ }
+ Err(mpsc::error::TrySendError::Closed(_)) => break,
+ }
+ }
+ }
+ }
+ stream.close();
+ info!(
+ call_id = %drain_call_id,
+ trace_id = %drain_trace_id,
+ received_frame_count,
+ dropped_frame_count,
+ queue_capacity = INBOUND_AUDIO_QUEUE_CAPACITY,
+ "runtime helper user_audio_drain_ended"
+ );
+ });
+ 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
+ };
+ let mut realtime_asr_upload: Option<RealtimeAsrUpload> = None;
+
+ while let Some(drained) = frame_rx.recv().await {
+ let frame = drained.frame;
+ frame_count = drained.frame_index;
+ sample_count += u64::from(frame.samples_per_channel) * u64::from(frame.num_channels);
+ let elapsed_ms = drained.captured_elapsed_ms;
+ 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) {
+ let was_in_speech = vad.in_speech;
+ let turn = vad.observe_frame(
+ &call_id,
+ &trace_id,
+ &participant_alias,
+ &track_sid_alias,
+ frame_count,
+ elapsed_ms,
+ &frame,
+ );
+ let is_in_speech = vad.in_speech;
+
+ if !was_in_speech && is_in_speech {
+ let turn_id = format!("turn-{:04}", vad.turn_index);
+ match RealtimeAsrUpload::start(
+ http.clone(),
+ turn_bridge_config.realtime_asr_config(),
+ &call_id,
+ &trace_id,
+ &turn_id,
+ &vad.speech_samples,
+ ) {
+ Ok(upload) => {
+ info!(
+ call_id = %call_id,
+ trace_id = %trace_id,
+ turn_id = %turn_id,
+ "runtime helper asr_realtime_session_started"
+ );
+ realtime_asr_upload = Some(upload);
+ }
+ Err(error) if turn_bridge_config.asr_realtime_enabled => {
+ warn!(
+ call_id = %call_id,
+ trace_id = %trace_id,
+ turn_id = %turn_id,
+ error = %safe_error(&error.to_string()),
+ "runtime helper asr_realtime_start_failed_fallback"
+ );
+ }
+ Err(_) => {}
+ }
+ } else if was_in_speech {
+ let push_failed = realtime_asr_upload
+ .as_mut()
+ .and_then(|upload| upload.push_48k_samples(frame.data.as_ref()).err());
+ if let Some(error) = push_failed {
+ warn!(
+ call_id = %call_id,
+ trace_id = %trace_id,
+ error = %safe_error(&error.to_string()),
+ "runtime helper asr_realtime_upload_failed_fallback"
+ );
+ if let Some(upload) = realtime_asr_upload.take() {
+ tokio::spawn(async move {
+ upload.cancel("upload_backpressure").await;
+ });
+ }
+ }
+ }
+
+ if let Some(turn) = turn {
+ let realtime_asr_result_ref = match realtime_asr_upload.take() {
+ Some(upload) => {
+ finish_realtime_asr_upload(upload, &call_id, &trace_id, &turn).await
+ }
+ None => None,
+ };
+ handle_finished_turn(
+ &http,
+ &turn_bridge_config,
+ &sink,
+ &call_id,
+ &trace_id,
+ turn,
+ realtime_asr_result_ref,
+ )
+ .await;
+ } else if was_in_speech && !is_in_speech {
+ if let Some(upload) = realtime_asr_upload.take() {
+ tokio::spawn(async move {
+ upload.cancel("speech_too_short").await;
+ });
+ }
+ }
+ } else {
+ if let Some(upload) = realtime_asr_upload.take() {
+ tokio::spawn(async move {
+ upload.cancel("vad_disabled").await;
+ });
+ }
+ 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,
+ ) {
+ let realtime_asr_result_ref = match realtime_asr_upload.take() {
+ Some(upload) => {
+ finish_realtime_asr_upload(upload, &call_id, &trace_id, &turn).await
+ }
+ None => None,
+ };
+ handle_finished_turn(
+ &http,
+ &turn_bridge_config,
+ &sink,
+ &call_id,
+ &trace_id,
+ turn,
+ realtime_asr_result_ref,
+ )
+ .await;
+ }
+ }
+ if let Some(upload) = realtime_asr_upload.take() {
+ upload.cancel("stream_end").await;
+ }
+ let _ = drain_shutdown_tx.send(true);
+ let mut drain_task = drain_task;
+ if timeout(AUDIO_DRAIN_STOP_GRACE, &mut drain_task)
+ .await
+ .is_err()
+ {
+ drain_task.abort();
+ let _ = drain_task.await;
+ warn!(
+ call_id = %call_id,
+ trace_id = %trace_id,
+ stop_grace_ms = AUDIO_DRAIN_STOP_GRACE.as_millis() as u64,
+ "runtime helper user_audio_drain_stop_timeout"
+ );
+ }
+
+ 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"
+ );
+ })
+}
+
+struct DrainedUserAudioFrame {
+ frame_index: u64,
+ captured_elapsed_ms: u64,
+ frame: AudioFrame<'static>,
+}
+
+async fn finish_realtime_asr_upload(
+ upload: RealtimeAsrUpload,
+ call_id: &str,
+ trace_id: &str,
+ turn: &FinishedSpeechTurn,
+) -> Option<String> {
+ match upload.finish(turn.duration_ms, &turn.end_reason).await {
+ Ok(RealtimeAsrOutcome {
+ status,
+ asr_result_ref,
+ provider_alias,
+ partial_count,
+ fallback_reason,
+ fallback_stage,
+ chunk_count,
+ audio_bytes,
+ wall_ms,
+ }) => {
+ info!(
+ call_id = %call_id,
+ trace_id = %trace_id,
+ turn_id = %turn.turn_id,
+ status = %status,
+ provider_alias = ?provider_alias,
+ partial_count,
+ chunk_count,
+ audio_bytes,
+ wall_ms,
+ asr_result_ref_present = asr_result_ref.is_some(),
+ fallback_reason = ?fallback_reason,
+ fallback_stage = ?fallback_stage,
+ "runtime helper asr_realtime_finished"
+ );
+ if status == "final" {
+ asr_result_ref
+ } else {
+ None
+ }
+ }
+ Err(error) => {
+ warn!(
+ call_id = %call_id,
+ trace_id = %trace_id,
+ turn_id = %turn.turn_id,
+ error = %safe_error(&error.to_string()),
+ "runtime helper asr_realtime_failed_fallback"
+ );
+ None
+ }
+ }
+}
+
+#[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", 400).max(100),
+ min_speech_ms: u64_env("CV_VAD_MIN_SPEECH_MS", 250).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;
+ }
+
+ // Align with cb-sdk's energy-based segmentation: peak is diagnostic only,
+ // otherwise isolated spikes can keep a turn open until max_turn_ms.
+ let voiced = rms >= self.config.rms_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>,
+ profile: String,
+ sample_rate_hz: u32,
+ num_channels: u16,
}
impl BotAudioOutputSink {
@@ -339,9 +3598,13 @@
room: Arc<Room>,
room_name: &str,
participant_identity: &str,
+ call_id: &str,
+ trace_id: &str,
track_name: &str,
+ profile: String,
sample_rate: u32,
num_channels: u32,
+ device_output_destination_identity: Option<String>,
) -> Result<Self> {
let rtc_source = NativeAudioSource::new(
AudioSourceOptions::default(),
@@ -353,6 +3616,8 @@
track_name,
RtcAudioSource::Native(rtc_source.clone()),
);
+ let room_alias = redact(room_name);
+ let participant_alias = redact(participant_identity);
room.local_participant()
.publish_track(
@@ -362,24 +3627,210 @@
.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}"
)
})?;
+ let num_channels_u16 = u16::try_from(num_channels)
+ .map_err(|_| anyhow!("unsupported bot audio channel count {num_channels}"))?;
info!(
- room_id = %room_name,
- participant_alias = %redact(participant_identity),
+ room_alias = %room_alias,
+ participant_alias = %participant_alias,
track_name = %track_name,
+ bot_audio_profile = %profile,
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,
+ "audioProfile": profile,
+ "sampleRate": sample_rate,
+ "numChannels": num_channels,
+ }),
);
Ok(Self {
room,
rtc_source,
track,
+ device_output_destination_identity,
+ profile,
+ sample_rate_hz: sample_rate,
+ num_channels: num_channels_u16,
})
+ }
+
+ 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<()> {
@@ -439,3 +3890,107 @@
Ok(())
}
}
+
+#[cfg(test)]
+mod tests {
+ use super::{
+ ReplyChunkMarker, ReplyChunkMarkerState, RuntimeTurnStreamEvent, RuntimeTurnStreamState,
+ RuntimeTurnStreamTimingPhase, runtime_session_nonce_hash,
+ };
+
+ #[test]
+ fn reply_chunk_marker_state_emits_turn_first_once_and_later_segment_first_once() {
+ let mut state = ReplyChunkMarkerState::default();
+
+ assert_eq!(ReplyChunkMarker::FirstReply, state.observe(Some(1)));
+ assert_eq!(ReplyChunkMarker::None, state.observe(Some(1)));
+ assert_eq!(ReplyChunkMarker::SegmentFirst, state.observe(Some(2)));
+ assert_eq!(ReplyChunkMarker::None, state.observe(Some(2)));
+ assert_eq!(ReplyChunkMarker::SegmentFirst, state.observe(Some(3)));
+ }
+
+ #[test]
+ fn reply_chunk_marker_state_without_segment_only_emits_turn_first() {
+ let mut state = ReplyChunkMarkerState::default();
+
+ assert_eq!(ReplyChunkMarker::FirstReply, state.observe(None));
+ assert_eq!(ReplyChunkMarker::None, state.observe(None));
+ }
+
+ #[test]
+ fn runtime_turn_stream_audio_chunk_reads_segment_seq() {
+ let event: RuntimeTurnStreamEvent = serde_json::from_str(
+ r#"{"type":"reply_audio_chunk","audioChunk":{"chunkSeq":4,"segmentSeq":2,"format":"pcm_s16le","payloadBase64":"AA==","last":false}}"#,
+ )
+ .expect("turn stream event");
+
+ assert_eq!(
+ Some(2),
+ event.audio_chunk.and_then(|chunk| chunk.segment_seq)
+ );
+ }
+
+ #[test]
+ fn runtime_turn_stream_reads_frozen_m5_and_audio_chunk_timing_contract() {
+ let activity: RuntimeTurnStreamEvent = serde_json::from_str(
+ r#"{"type":"activity","callId":"call-1","traceId":"trace-1","turnId":"turn-1","activity":{"eventType":"tts_first_audio_chunk_ready","extension":{"streamTimingVersion":1,"streamAnchorId":"0123456789abcdef","streamAnchorServerDeltaMs":1200,"runtimeSessionNonceHash":"abcdef012345","segmentSeq":1,"streamTimingValidation":"bound"}}}"#,
+ )
+ .expect("m5 activity event");
+ let audio: RuntimeTurnStreamEvent = serde_json::from_str(
+ r#"{"type":"reply_audio_chunk","callId":"call-1","traceId":"trace-1","turnId":"turn-1","audioChunk":{"chunkSeq":1,"segmentSeq":1,"streamTimingVersion":1,"streamAnchorId":"0123456789abcdef","format":"pcm_s16le","payloadBase64":"AA==","last":false}}"#,
+ )
+ .expect("timed audio chunk event");
+
+ assert_eq!(Some("call-1"), activity.call_id.as_deref());
+ assert_eq!(Some("trace-1"), activity.trace_id.as_deref());
+ assert_eq!(Some("turn-1"), activity.turn_id.as_deref());
+ let extension = activity.activity.unwrap().extension.unwrap();
+ assert_eq!(Some(1), extension.stream_timing_version);
+ assert_eq!(Some(1200), extension.stream_anchor_server_delta_ms);
+ assert_eq!(
+ Some("0123456789abcdef"),
+ audio.audio_chunk.unwrap().stream_anchor_id.as_deref()
+ );
+ }
+
+ #[test]
+ fn runtime_turn_stream_timing_records_m6_and_m7_once_then_rejects_terminal_late_events() {
+ let nonce = "runtime-nonce";
+ let nonce_hash = runtime_session_nonce_hash(nonce);
+ let activity: RuntimeTurnStreamEvent = serde_json::from_str(&format!(
+ r#"{{"type":"activity","callId":"call-1","traceId":"trace-1","turnId":"turn-1","activity":{{"eventType":"tts_first_audio_chunk_ready","extension":{{"streamTimingVersion":1,"streamAnchorId":"0123456789abcdef","streamAnchorServerDeltaMs":1200,"runtimeSessionNonceHash":"{nonce_hash}","segmentSeq":1,"streamTimingValidation":"bound"}}}}}}"#,
+ ))
+ .expect("m5 activity event");
+ let audio: RuntimeTurnStreamEvent = serde_json::from_str(
+ r#"{"type":"reply_audio_chunk","callId":"call-1","traceId":"trace-1","turnId":"turn-1","audioChunk":{"chunkSeq":1,"segmentSeq":1,"streamTimingVersion":1,"streamAnchorId":"0123456789abcdef","format":"pcm_s16le","payloadBase64":"AA==","last":false}}"#,
+ )
+ .expect("timed audio chunk event");
+ let mut state = RuntimeTurnStreamState::default();
+
+ assert!(state.arm_timing_anchor("call-1", "trace-1", "turn-1", nonce, &activity));
+ let chunk = audio.audio_chunk.as_ref().expect("audio chunk");
+ assert!(
+ state
+ .record_m6("call-1", "trace-1", "turn-1", &audio, chunk)
+ .is_some()
+ );
+ assert!(
+ state
+ .record_m6("call-1", "trace-1", "turn-1", &audio, chunk)
+ .is_none()
+ );
+ assert!(state.record_m7().is_some());
+ assert!(state.record_m7().is_none());
+
+ state.close_timing();
+ assert_eq!(RuntimeTurnStreamTimingPhase::Closed, state.timing.phase);
+ assert!(state.timing.anchor.is_none());
+ assert!(!state.arm_timing_anchor("call-1", "trace-1", "turn-1", nonce, &activity));
+ assert!(
+ state
+ .record_m6("call-1", "trace-1", "turn-1", &audio, chunk)
+ .is_none()
+ );
+ assert!(state.record_m7().is_none());
+ }
+}
--
Gitblit v1.9.3