mod asr_realtime;
|
mod audio;
|
mod service;
|
|
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 asr_realtime::{
|
AudioIngressMetadata, 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, RemoteAudioTrack,
|
RemoteParticipant, RemoteTrack, Room, RoomEvent, RoomOptions,
|
},
|
};
|
use reqwest::Client;
|
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 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 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()
|
.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(),
|
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(
|
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,
|
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,
|
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);
|
|
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_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 {
|
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_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> {
|
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 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 {
|
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,
|
"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 user_participant_identity = config.user_participant_identity.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,
|
user_participant_identity,
|
)
|
.await;
|
})
|
}
|
|
fn is_bound_user_participant(identity: &str, expected: Option<&str>) -> bool {
|
expected.is_none_or(|value| identity == value)
|
}
|
|
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>,
|
user_participant_identity: Option<String>,
|
) {
|
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,
|
} => {
|
if !is_bound_user_participant(
|
&participant.identity().to_string(),
|
user_participant_identity.as_deref(),
|
) {
|
warn!(call_id = %call_id, trace_id = %trace_id,
|
metadata_status = "wrong_participant",
|
"runtime helper ignored non-user audio participant");
|
continue;
|
}
|
let participant_alias = redact(&participant.identity().to_string());
|
let track_sid_alias = redact(&track.sid().to_string());
|
let track_name = track.name();
|
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(),
|
participant,
|
);
|
}
|
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,
|
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) => {
|
let mut published_device_outputs = HashSet::new();
|
for output in &outcome.device_outputs {
|
if !should_publish_device_output(&mut published_device_outputs, output) {
|
continue;
|
}
|
if let Err(error) = sink
|
.publish_device_output(call_id, trace_id, &turn.turn_id, output)
|
.await
|
{
|
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,
|
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
|
.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 = %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() {
|
if !should_publish_device_output(&mut state.published_device_output_ids, output) {
|
return Ok(());
|
}
|
sink.publish_device_output(call_id, trace_id, &turn.turn_id, output)
|
.await?;
|
state.device_output_count = state.device_output_count.saturating_add(1);
|
}
|
}
|
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 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,
|
published_device_output_ids: HashSet<String>,
|
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,
|
published_device_output_ids: HashSet::new(),
|
encoded_audio_buffer: Vec::new(),
|
pcm_stream_decoder: None,
|
pcm_stream_network_chunk_count: 0,
|
reply_chunk_markers: ReplyChunkMarkerState::default(),
|
timing: RuntimeTurnStreamTimingState::default(),
|
}
|
}
|
}
|
|
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 should_publish_device_output(
|
published_ids: &mut HashSet<String>,
|
output: &RuntimeTurnDeviceOutput,
|
) -> bool {
|
let Some(command_id) = output
|
.command_id
|
.as_deref()
|
.map(str::trim)
|
.filter(|value| !value.is_empty())
|
else {
|
return false;
|
};
|
if output
|
.command_code
|
.as_deref()
|
.map(str::trim)
|
.filter(|value| !value.is_empty())
|
.is_none()
|
{
|
return false;
|
}
|
published_ids.insert(command_id.to_string())
|
}
|
|
fn require_safe_segment(value: &str) -> Result<()> {
|
if value.is_empty()
|
|| value.contains('/')
|
|| 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>,
|
participant: RemoteParticipant,
|
) -> 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;
|
let mut last_fixture_sequence: Option<String> = 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, is_in_speech, turn) = observe_frame_and_start_session(
|
vad,
|
&call_id,
|
&trace_id,
|
&participant_alias,
|
&track_sid_alias,
|
frame_count,
|
elapsed_ms,
|
&frame,
|
http.clone(),
|
turn_bridge_config.realtime_asr_config(),
|
|| participant.attributes(),
|
&mut realtime_asr_upload,
|
&mut last_fixture_sequence,
|
turn_bridge_config.asr_realtime_enabled,
|
);
|
|
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"
|
);
|
})
|
}
|
|
fn observe_frame_and_start_session<F>(
|
vad: &mut SimpleVad,
|
call_id: &str,
|
trace_id: &str,
|
participant_alias: &str,
|
track_sid_alias: &str,
|
frame_count: u64,
|
elapsed_ms: u64,
|
frame: &AudioFrame<'_>,
|
http: Client,
|
config: RealtimeAsrConfig,
|
read_attributes: F,
|
upload_slot: &mut Option<RealtimeAsrUpload>,
|
last_fixture_sequence: &mut Option<String>,
|
realtime_enabled: bool,
|
) -> (bool, bool, Option<FinishedSpeechTurn>)
|
where
|
F: FnOnce() -> std::collections::HashMap<String, String>,
|
{
|
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 {
|
start_realtime_session_for_new_speech(
|
http,
|
config,
|
call_id,
|
trace_id,
|
vad,
|
read_attributes,
|
upload_slot,
|
last_fixture_sequence,
|
realtime_enabled,
|
);
|
}
|
(was_in_speech, is_in_speech, turn)
|
}
|
|
fn start_realtime_session_for_new_speech(
|
http: Client,
|
config: RealtimeAsrConfig,
|
call_id: &str,
|
trace_id: &str,
|
vad: &SimpleVad,
|
read_attributes: impl FnOnce() -> std::collections::HashMap<String, String>,
|
upload_slot: &mut Option<RealtimeAsrUpload>,
|
last_fixture_sequence: &mut Option<String>,
|
realtime_enabled: bool,
|
) {
|
let turn_id = format!("turn-{:04}", vad.turn_index);
|
let metadata = match AudioIngressMetadata::from_participant(&read_attributes()) {
|
Ok(metadata) => metadata,
|
Err(reason) => {
|
warn!(call_id = %call_id, trace_id = %trace_id, turn_id = %turn_id,
|
reason, "runtime helper asr_realtime_metadata_rejected");
|
return;
|
}
|
};
|
if let Some(metadata) = metadata.as_ref() {
|
if !fixture_sequence_is_new(
|
last_fixture_sequence.as_deref(),
|
&metadata.client_fixture_sequence,
|
) {
|
warn!(call_id = %call_id, trace_id = %trace_id, turn_id = %turn_id,
|
"runtime helper asr_realtime_metadata_sequence_rejected");
|
return;
|
}
|
}
|
match RealtimeAsrUpload::start(
|
http,
|
config,
|
call_id,
|
trace_id,
|
&turn_id,
|
&vad.speech_samples,
|
metadata.as_ref(),
|
) {
|
Ok(upload) => {
|
if let Some(metadata) = metadata {
|
*last_fixture_sequence = Some(metadata.client_fixture_sequence);
|
}
|
info!(call_id = %call_id, trace_id = %trace_id, turn_id = %turn_id,
|
"runtime helper asr_realtime_session_started");
|
*upload_slot = Some(upload);
|
}
|
Err(error) if 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(_) => {}
|
}
|
}
|
|
fn fixture_sequence_is_new(previous: Option<&str>, current: &str) -> bool {
|
let Some(previous) = previous else {
|
return true;
|
};
|
let current_number = current
|
.rsplit_once('-')
|
.and_then(|(_, value)| value.parse::<u64>().ok());
|
let previous_number = previous
|
.rsplit_once('-')
|
.and_then(|(_, value)| value.parse::<u64>().ok());
|
match (previous_number, current_number) {
|
(Some(previous), Some(current)) => current > previous,
|
_ => previous != current,
|
}
|
}
|
|
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 {
|
async fn publish(
|
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(),
|
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}"
|
)
|
})?;
|
let num_channels_u16 = u16::try_from(num_channels)
|
.map_err(|_| anyhow!("unsupported bot audio channel count {num_channels}"))?;
|
|
info!(
|
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<()> {
|
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(())
|
}
|
}
|
|
#[cfg(test)]
|
mod tests {
|
use super::{
|
ReplyChunkMarker, ReplyChunkMarkerState, RuntimeTurnDeviceOutput, RuntimeTurnStreamEvent,
|
RuntimeTurnStreamState, RuntimeTurnStreamTimingPhase, runtime_session_nonce_hash,
|
should_publish_device_output,
|
};
|
use std::{
|
collections::HashSet,
|
io::{Read, Write},
|
net::TcpListener,
|
sync::mpsc,
|
thread,
|
time::Duration,
|
};
|
|
#[test]
|
fn production_vad_session_boundary_reads_updated_attributes() {
|
let config = SimpleVadConfig {
|
rms_threshold: 0.001,
|
peak_threshold: 0.01,
|
start_frames: 2,
|
end_silence_ms: 100,
|
min_speech_ms: 1,
|
max_turn_ms: 1_000,
|
initial_ignore_ms: 0,
|
};
|
let mut vad = SimpleVad::new(config);
|
let samples = vec![1_000i16; 160];
|
let frame = AudioFrame {
|
data: samples.as_slice().into(),
|
sample_rate: 16_000,
|
num_channels: 1,
|
samples_per_channel: 160,
|
};
|
let mut attributes = std::collections::HashMap::from([
|
(
|
"inputSourceCategory".to_string(),
|
"controlled_fixture".to_string(),
|
),
|
(
|
"clientFixtureSequence".to_string(),
|
"fixture-01".to_string(),
|
),
|
]);
|
let mut starts = Vec::new();
|
for (session_index, sequence) in [(1, "fixture-01"), (2, "fixture-02")] {
|
let was_in_speech = vad.in_speech;
|
vad.observe_frame(
|
"call-001",
|
"trace-001",
|
"participant",
|
"track",
|
session_index * 2 - 1,
|
1_000 * session_index,
|
&frame,
|
);
|
vad.observe_frame(
|
"call-001",
|
"trace-001",
|
"participant",
|
"track",
|
session_index * 2,
|
1_000 * session_index + 10,
|
&frame,
|
);
|
let is_in_speech = vad.in_speech;
|
assert!(!was_in_speech && is_in_speech);
|
attributes.insert("clientFixtureSequence".to_string(), sequence.to_string());
|
let metadata = AudioIngressMetadata::from_participant(&attributes)
|
.expect("valid participant attributes")
|
.expect("controlled fixture metadata");
|
let session_line = asr_realtime::session_start_line(
|
"call-001",
|
"trace-001",
|
&format!("turn-{session_index:04}"),
|
"nonce-001",
|
Some(&metadata),
|
)
|
.expect("session start line");
|
let session_json: serde_json::Value =
|
serde_json::from_slice(&session_line).expect("session start json");
|
assert_eq!(sequence, session_json["clientFixtureSequence"]);
|
starts.push(metadata.client_fixture_sequence);
|
vad.reset_current_turn();
|
}
|
assert_eq!(vec!["fixture-01", "fixture-02"], starts);
|
attributes.insert("inputSourceCategory".to_string(), "other".to_string());
|
assert!(AudioIngressMetadata::from_participant(&attributes).is_err());
|
assert!(
|
AudioIngressMetadata::from_participant(&std::collections::HashMap::new())
|
.expect("missing attributes is absent")
|
.is_none()
|
);
|
attributes.insert(
|
"clientFixtureSequence".to_string(),
|
"fixture-01".to_string(),
|
);
|
assert!(AudioIngressMetadata::from_participant(&attributes).is_err());
|
}
|
|
#[test]
|
fn production_observer_rejects_wrong_participant_before_vad_session() {
|
assert!(!is_bound_user_participant(
|
"participant-other",
|
Some("participant-user")
|
));
|
assert!(is_bound_user_participant(
|
"participant-user",
|
Some("participant-user")
|
));
|
assert!(is_bound_user_participant("participant-any", None));
|
}
|
|
#[tokio::test]
|
async fn production_observer_vad_to_session_entry_reads_each_updated_attribute() {
|
let listener = TcpListener::bind("127.0.0.1:0").expect("bind local ASR fixture");
|
let address = listener.local_addr().expect("fixture address");
|
let (request_tx, request_rx) = mpsc::channel::<String>();
|
let server = thread::spawn(move || {
|
for _ in 0..2 {
|
let (mut stream, _) = listener.accept().expect("accept ASR session");
|
stream
|
.set_read_timeout(Some(Duration::from_secs(2)))
|
.expect("set fixture timeout");
|
let mut bytes = Vec::new();
|
let mut buffer = [0_u8; 4096];
|
loop {
|
match stream.read(&mut buffer) {
|
Ok(0) => break,
|
Ok(size) => {
|
bytes.extend_from_slice(&buffer[..size]);
|
if bytes.windows(7).any(|window| window == b"0\r\n\r\n") {
|
break;
|
}
|
}
|
Err(_) => break,
|
}
|
}
|
request_tx
|
.send(String::from_utf8_lossy(&bytes).into_owned())
|
.expect("capture ASR request");
|
stream
|
.write_all(b"HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: 39\r\nconnection: close\r\n\r\n{\"code\":0,\"data\":{\"status\":\"ok\"}}")
|
.expect("write fixture response");
|
}
|
});
|
let mut vad = SimpleVad::new(SimpleVadConfig {
|
rms_threshold: 0.001,
|
peak_threshold: 0.01,
|
start_frames: 1,
|
end_silence_ms: 100,
|
min_speech_ms: 1,
|
max_turn_ms: 1_000,
|
initial_ignore_ms: 0,
|
});
|
let frame_data = vec![1_000i16; 160];
|
let frame = AudioFrame {
|
data: frame_data.as_slice().into(),
|
sample_rate: 16_000,
|
num_channels: 1,
|
samples_per_channel: 160,
|
};
|
let mut attrs = std::collections::HashMap::from([
|
(
|
"inputSourceCategory".to_string(),
|
"controlled_fixture".to_string(),
|
),
|
(
|
"clientFixtureSequence".to_string(),
|
"fixture-01".to_string(),
|
),
|
]);
|
let mut upload = None;
|
let mut last_fixture_sequence = None;
|
let config = RealtimeAsrConfig {
|
enabled: true,
|
url: Some(format!("http://{address}/runtime/asr/realtime")),
|
runtime_token: Some("test".to_string()),
|
runtime_session_nonce: Some("test".to_string()),
|
chunk_duration_ms: 200,
|
};
|
let (was, is, turn) = observe_frame_and_start_session(
|
&mut vad,
|
"call-001",
|
"trace-001",
|
"participant-user",
|
"track-001",
|
1,
|
1_000,
|
&frame,
|
Client::new(),
|
config.clone(),
|
|| attrs.clone(),
|
&mut upload,
|
&mut last_fixture_sequence,
|
true,
|
);
|
assert!(!was && is && turn.is_none());
|
assert!(upload.is_some());
|
upload.take().unwrap().cancel("test").await;
|
vad.reset_current_turn();
|
attrs.insert(
|
"clientFixtureSequence".to_string(),
|
"fixture-02".to_string(),
|
);
|
let (was, is, turn) = observe_frame_and_start_session(
|
&mut vad,
|
"call-001",
|
"trace-001",
|
"participant-user",
|
"track-001",
|
2,
|
2_000,
|
&frame,
|
Client::new(),
|
config,
|
|| attrs.clone(),
|
&mut upload,
|
&mut last_fixture_sequence,
|
true,
|
);
|
assert!(!was && is && turn.is_none());
|
assert!(upload.is_some());
|
upload.take().unwrap().cancel("test").await;
|
|
let first_request = request_rx
|
.recv_timeout(Duration::from_secs(2))
|
.expect("first session request");
|
let second_request = request_rx
|
.recv_timeout(Duration::from_secs(2))
|
.expect("second session request");
|
assert!(first_request.contains("\\\"clientFixtureSequence\\\":\\\"fixture-01\\\""));
|
assert!(second_request.contains("\\\"clientFixtureSequence\\\":\\\"fixture-02\\\""));
|
server.join().expect("fixture server");
|
|
// The same production boundary rejects a wrong participant before VAD/session creation.
|
assert!(!is_bound_user_participant(
|
"participant-other",
|
Some("participant-user")
|
));
|
assert_eq!(0, request_rx.try_iter().count());
|
vad.reset_current_turn();
|
attrs.insert(
|
"clientFixtureSequence".to_string(),
|
"fixture-01".to_string(),
|
);
|
let (_, _, _) = observe_frame_and_start_session(
|
&mut vad,
|
"call-001",
|
"trace-001",
|
"participant-user",
|
"track-001",
|
3,
|
3_000,
|
&frame,
|
Client::new(),
|
RealtimeAsrConfig {
|
enabled: true,
|
url: Some(format!("http://{address}/runtime/asr/realtime")),
|
runtime_token: Some("test".to_string()),
|
runtime_session_nonce: Some("test".to_string()),
|
chunk_duration_ms: 200,
|
},
|
|| attrs.clone(),
|
&mut upload,
|
&mut last_fixture_sequence,
|
true,
|
);
|
assert!(upload.is_none());
|
vad.reset_current_turn();
|
attrs.insert(
|
"clientFixtureSequence".to_string(),
|
"fixture-00".to_string(),
|
);
|
let (_, _, _) = observe_frame_and_start_session(
|
&mut vad,
|
"call-001",
|
"trace-001",
|
"participant-user",
|
"track-001",
|
4,
|
4_000,
|
&frame,
|
Client::new(),
|
RealtimeAsrConfig {
|
enabled: true,
|
url: Some(format!("http://{address}/runtime/asr/realtime")),
|
runtime_token: Some("test".to_string()),
|
runtime_session_nonce: Some("test".to_string()),
|
chunk_duration_ms: 200,
|
},
|
|| attrs.clone(),
|
&mut upload,
|
&mut last_fixture_sequence,
|
true,
|
);
|
assert!(upload.is_none());
|
vad.reset_current_turn();
|
attrs.insert("inputSourceCategory".to_string(), "other".to_string());
|
let mut invalid_upload = None;
|
let (_, _, invalid_turn) = observe_frame_and_start_session(
|
&mut vad,
|
"call-001",
|
"trace-001",
|
"participant-other",
|
"track-001",
|
3,
|
3_000,
|
&frame,
|
Client::new(),
|
RealtimeAsrConfig {
|
enabled: true,
|
url: Some("http://127.0.0.1:9".to_string()),
|
runtime_token: Some("test".to_string()),
|
runtime_session_nonce: Some("test".to_string()),
|
chunk_duration_ms: 200,
|
},
|
|| attrs,
|
&mut invalid_upload,
|
&mut last_fixture_sequence,
|
true,
|
);
|
assert!(invalid_turn.is_none());
|
assert!(invalid_upload.is_none());
|
}
|
|
#[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());
|
}
|
|
#[test]
|
fn device_output_contract_is_reliable_and_deduplicated() {
|
let output = RuntimeTurnDeviceOutput {
|
command_id: Some("cmd-1".to_string()),
|
command_code: Some("custom.app.DeviceLevelChange".to_string()),
|
params: Some(serde_json::json!({"level": 1})),
|
};
|
let mut published = HashSet::new();
|
assert!(should_publish_device_output(&mut published, &output));
|
assert!(!should_publish_device_output(&mut published, &output));
|
assert_eq!(published.len(), 1);
|
}
|
|
#[test]
|
fn device_output_invalid_or_missing_command_is_fail_closed() {
|
for output in [
|
RuntimeTurnDeviceOutput {
|
command_id: None,
|
command_code: Some("custom.app.DeviceLevelChange".to_string()),
|
params: None,
|
},
|
RuntimeTurnDeviceOutput {
|
command_id: Some("cmd-1".to_string()),
|
command_code: None,
|
params: None,
|
},
|
] {
|
let mut published = HashSet::new();
|
assert!(!should_publish_device_output(&mut published, &output));
|
assert!(published.is_empty());
|
}
|
}
|
}
|