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