From a590c735cfa4ec55ba6b32e297600d83e1a3e46a Mon Sep 17 00:00:00 2001
From: cai <cai@nbcai.cc>
Date: Fri, 10 Jul 2026 18:03:29 +0800
Subject: [PATCH] feat: stream realtime asr from helper

---
 src/main.rs | 3253 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
 1 files changed, 3,106 insertions(+), 147 deletions(-)

diff --git a/src/main.rs b/src/main.rs
index 07f4e9f..62107a9 100644
--- a/src/main.rs
+++ b/src/main.rs
@@ -1,30 +1,59 @@
+mod asr_realtime;
 mod audio;
+mod service;
 
-use std::{env, sync::Arc, time::Duration};
+use std::{
+    env, fs,
+    path::{Path, PathBuf},
+    sync::{
+        Arc,
+        atomic::{AtomicBool, Ordering},
+    },
+    time::{Duration, Instant, SystemTime, UNIX_EPOCH},
+};
 
 use anyhow::{Context, Result, anyhow};
-use audio::load_pre_recorded_frames;
+use asr_realtime::{RealtimeAsrConfig, RealtimeAsrOutcome, RealtimeAsrUpload};
+use audio::{AudioDiagnostics, load_pre_recorded_frames};
+use base64::{Engine as _, engine::general_purpose};
+use futures_util::StreamExt;
 use libwebrtc::{
     audio_source::native::NativeAudioSource,
+    audio_stream::native::NativeAudioStream,
     prelude::{AudioFrame, AudioSourceOptions, RtcAudioSource},
 };
 use livekit::{
     options::TrackPublishOptions,
-    prelude::{DataPacket, LocalAudioTrack, LocalTrack, ParticipantIdentity, Room, RoomOptions},
+    prelude::{
+        DataPacket, LocalAudioTrack, LocalTrack, ParticipantIdentity, RemoteAudioTrack,
+        RemoteTrack, Room, RoomEvent, RoomOptions,
+    },
 };
 use reqwest::Client;
-use serde::Serialize;
-use tokio::time::sleep;
+use serde::{Deserialize, Serialize};
+use serde_json::json;
+use tokio::time::{sleep, sleep_until};
+use tokio::{sync::mpsc::UnboundedReceiver, task::JoinHandle};
 use tracing::{info, warn};
 
-const TARGET_SAMPLE_RATE_HZ: u32 = 48_000;
-const TARGET_NUM_CHANNELS: u16 = 1;
+const USER_AUDIO_SAMPLE_RATE_HZ: u32 = 48_000;
+const USER_AUDIO_NUM_CHANNELS: u16 = 1;
+const DEFAULT_BOT_AUDIO_PROFILE: &str = "pcm-16k";
+const LIVEKIT_48K_SAMPLE_RATE_HZ: u32 = 48_000;
+const PCM_16K_SAMPLE_RATE_HZ: u32 = 16_000;
+const BOT_NUM_CHANNELS: u16 = 1;
 const TRACK_NAME: &str = "bot-main-audio";
-const DEVICE_OUTPUT_TOPIC: &str = "device_output";
 
 #[tokio::main(flavor = "multi_thread")]
 async fn main() -> Result<()> {
     init_tracing();
+    if service::service_mode_enabled() {
+        return service::run_service().await;
+    }
+    run_worker().await
+}
+
+async fn run_worker() -> Result<()> {
     let config = Config::from_env()?;
     let http = Client::builder()
         .use_rustls_tls()
@@ -35,12 +64,15 @@
         &http,
         config.greeting_audio_file.as_deref(),
         config.greeting_audio_url.as_deref(),
-        TARGET_SAMPLE_RATE_HZ,
-        TARGET_NUM_CHANNELS,
+        config.bot_audio_profile.sample_rate_hz,
+        config.bot_audio_profile.num_channels,
+        config.audio_debug_dump_dir.as_deref(),
+        &config.call_id,
+        "greeting",
     )
     .await?;
 
-    let (room, _events) = Room::connect(
+    let (room, events) = Room::connect(
         config.livekit_url.as_str(),
         config.bot_token.as_str(),
         RoomOptions::default(),
@@ -52,42 +84,138 @@
     info!(
         call_id = %config.call_id,
         trace_id = %config.trace_id,
-        room_id = %config.room_id,
+        room_alias = %redact(&config.room_id),
         participant_alias = %redact(&config.bot_participant_identity),
         greeting_source = %config.greeting_source,
+        bot_audio_profile = %config.bot_audio_profile.profile,
+        bot_sample_rate_hz = config.bot_audio_profile.sample_rate_hz,
+        bot_num_channels = config.bot_audio_profile.num_channels,
         "combrabo voice runtime helper connected"
     );
+    emit_activity(
+        &config.call_id,
+        &config.trace_id,
+        None,
+        "bot_participant_joined",
+        "ok",
+        None,
+        None,
+        json!({"participantAlias": redact(&config.bot_participant_identity)}),
+    );
+
+    let vad_enabled_gate = Arc::new(AtomicBool::new(
+        !config.simple_vad_enabled || !config.simple_vad_gate_until_greeting_done,
+    ));
+    if config.simple_vad_enabled && config.simple_vad_gate_until_greeting_done {
+        info!(
+            call_id = %config.call_id,
+            trace_id = %config.trace_id,
+            post_greeting_delay_ms = config.simple_vad_post_greeting_delay_ms,
+            "runtime helper vad_disabled_greeting"
+        );
+    } else if config.simple_vad_enabled {
+        info!(
+            call_id = %config.call_id,
+            trace_id = %config.trace_id,
+            "runtime helper vad_enabled"
+        );
+        emit_activity(
+            &config.call_id,
+            &config.trace_id,
+            None,
+            "vad_enabled",
+            "ok",
+            None,
+            None,
+            json!({"reason": "gate_disabled"}),
+        );
+    }
 
     let sink = BotAudioOutputSink::publish(
         room.clone(),
         &config.room_id,
         &config.bot_participant_identity,
+        &config.call_id,
+        &config.trace_id,
         TRACK_NAME,
-        TARGET_SAMPLE_RATE_HZ,
-        u32::from(TARGET_NUM_CHANNELS),
+        config.bot_audio_profile.profile.clone(),
+        config.bot_audio_profile.sample_rate_hz,
+        u32::from(config.bot_audio_profile.num_channels),
+        config.user_participant_identity.clone(),
     )
     .await?;
+    let sink = Arc::new(sink);
 
-    if config.device_output_smoke_enabled {
-        publish_device_output_smoke(room.as_ref(), &config).await?;
-    }
+    let user_audio_observer = spawn_user_audio_observer(
+        events,
+        &config,
+        vad_enabled_gate.clone(),
+        http.clone(),
+        sink.clone(),
+    );
 
-    if let Some(frames) = greeting_frames {
+    if let Some(loaded_audio) = greeting_frames {
+        let diagnostics = loaded_audio.diagnostics;
+        let frames = loaded_audio.frames;
+        log_audio_diagnostics(&config, &diagnostics);
         info!(
             call_id = %config.call_id,
             frame_count = frames.len(),
+            theoretical_duration_ms = diagnostics.target_duration_ms,
             greeting_source = %config.greeting_source,
             "runtime helper starting greeting playback"
         );
-        for frame in &frames {
+        emit_activity(
+            &config.call_id,
+            &config.trace_id,
+            None,
+            "greeting_write_started",
+            "ok",
+            None,
+            None,
+            json!({
+                "frameCount": frames.len(),
+                "greetingTheoreticalDurationMs": diagnostics.target_duration_ms,
+            }),
+        );
+        let playback_started_at = Instant::now();
+        let pacing_started_at = tokio::time::Instant::now();
+        for (index, frame) in frames.iter().enumerate() {
             sink.write_pcm_frame(frame).await?;
-            sleep(Duration::from_millis(20)).await;
+            sleep_until(pacing_started_at + Duration::from_millis(((index + 1) as u64) * 20)).await;
         }
+        let playback_wall_ms = playback_started_at.elapsed().as_millis() as i64;
+        let drift_ms = playback_wall_ms - diagnostics.target_duration_ms as i64;
         sink.clear_buffer();
         info!(
             call_id = %config.call_id,
             greeting_source = %config.greeting_source,
+            frame_count = frames.len(),
+            theoretical_duration_ms = diagnostics.target_duration_ms,
+            push_wall_duration_ms = playback_wall_ms,
+            push_drift_ms = drift_ms,
             "runtime helper finished greeting playback"
+        );
+        emit_activity(
+            &config.call_id,
+            &config.trace_id,
+            None,
+            "greeting_write_finished",
+            "ok",
+            None,
+            None,
+            json!({
+                "frameCount": frames.len(),
+                "greetingTheoreticalDurationMs": diagnostics.target_duration_ms,
+                "greetingPushWallDurationMs": playback_wall_ms,
+                "greetingPushDriftMs": drift_ms,
+            }),
+        );
+        schedule_vad_gate_enable(
+            vad_enabled_gate,
+            &config,
+            config.simple_vad_post_greeting_delay_ms,
+            "greeting_finished",
         );
     } else {
         warn!(
@@ -95,9 +223,12 @@
             greeting_source = %config.greeting_source,
             "runtime helper started without greeting audio; keeping published track alive"
         );
+        schedule_vad_gate_enable(vad_enabled_gate, &config, 0, "no_greeting_audio");
     }
 
     wait_for_shutdown_signal().await?;
+    user_audio_observer.abort();
+    let _ = user_audio_observer.await;
 
     if let Err(error) = sink.close().await {
         warn!(
@@ -133,12 +264,159 @@
     room_id: String,
     bot_token: String,
     bot_participant_identity: String,
+    user_participant_identity: Option<String>,
     greeting_source: String,
     greeting_audio_file: Option<String>,
     greeting_audio_url: Option<String>,
-    role_id: String,
-    device_output_smoke_enabled: bool,
-    device_output_destination_identities: Vec<ParticipantIdentity>,
+    audio_debug_dump_dir: Option<String>,
+    runtime_turn_bridge_url: Option<String>,
+    runtime_turn_bridge_token: Option<String>,
+    runtime_turn_bridge_mode: String,
+    runtime_asr_stream_enabled: bool,
+    runtime_asr_stream_url: Option<String>,
+    runtime_asr_realtime_enabled: bool,
+    runtime_asr_realtime_url: Option<String>,
+    runtime_asr_realtime_chunk_duration_ms: u64,
+    runtime_turn_artifact_dir: Option<String>,
+    runtime_session_nonce: Option<String>,
+    user_audio_observer_enabled: bool,
+    simple_vad_enabled: bool,
+    simple_vad_gate_until_greeting_done: bool,
+    simple_vad_post_greeting_delay_ms: u64,
+    simple_vad_config: SimpleVadConfig,
+    bot_audio_profile: BotAudioProfile,
+}
+
+#[derive(Clone)]
+struct BotAudioProfile {
+    profile: String,
+    sample_rate_hz: u32,
+    num_channels: u16,
+}
+
+impl BotAudioProfile {
+    fn from_env() -> Result<Self> {
+        let profile = env::var("CV_BOT_AUDIO_PROFILE")
+            .unwrap_or_else(|_| DEFAULT_BOT_AUDIO_PROFILE.to_string())
+            .trim()
+            .to_ascii_lowercase();
+        match profile.as_str() {
+            "livekit-48k" | "48k" => Ok(Self {
+                profile: "livekit-48k".to_string(),
+                sample_rate_hz: LIVEKIT_48K_SAMPLE_RATE_HZ,
+                num_channels: BOT_NUM_CHANNELS,
+            }),
+            "pcm-16k" | "16k" => Ok(Self {
+                profile: "pcm-16k".to_string(),
+                sample_rate_hz: PCM_16K_SAMPLE_RATE_HZ,
+                num_channels: BOT_NUM_CHANNELS,
+            }),
+            "custom" => {
+                let sample_rate_hz = u32_env("CV_BOT_SAMPLE_RATE_HZ", LIVEKIT_48K_SAMPLE_RATE_HZ);
+                let num_channels = u16_env("CV_BOT_NUM_CHANNELS", BOT_NUM_CHANNELS);
+                if sample_rate_hz == 0 {
+                    return Err(anyhow!("CV_BOT_SAMPLE_RATE_HZ must be positive"));
+                }
+                if num_channels == 0 {
+                    return Err(anyhow!("CV_BOT_NUM_CHANNELS must be positive"));
+                }
+                Ok(Self {
+                    profile,
+                    sample_rate_hz,
+                    num_channels,
+                })
+            }
+            _ => Err(anyhow!(
+                "unsupported CV_BOT_AUDIO_PROFILE {}; expected livekit-48k, pcm-16k or custom",
+                profile
+            )),
+        }
+    }
+}
+
+#[derive(Clone)]
+struct TurnBridgeConfig {
+    bridge_url: Option<String>,
+    bridge_token: Option<String>,
+    bridge_mode: String,
+    asr_stream_enabled: bool,
+    asr_stream_url: Option<String>,
+    asr_realtime_enabled: bool,
+    asr_realtime_url: Option<String>,
+    asr_realtime_chunk_duration_ms: u64,
+    artifact_dir: Option<String>,
+    runtime_session_nonce: Option<String>,
+    audio_debug_dump_dir: Option<String>,
+}
+
+impl TurnBridgeConfig {
+    fn from_config(config: &Config) -> Self {
+        Self {
+            bridge_url: config.runtime_turn_bridge_url.clone(),
+            bridge_token: config.runtime_turn_bridge_token.clone(),
+            bridge_mode: config.runtime_turn_bridge_mode.clone(),
+            asr_stream_enabled: config.runtime_asr_stream_enabled,
+            asr_stream_url: config.runtime_asr_stream_url.clone(),
+            asr_realtime_enabled: config.runtime_asr_realtime_enabled,
+            asr_realtime_url: config.runtime_asr_realtime_url.clone(),
+            asr_realtime_chunk_duration_ms: config.runtime_asr_realtime_chunk_duration_ms,
+            artifact_dir: config.runtime_turn_artifact_dir.clone(),
+            runtime_session_nonce: config.runtime_session_nonce.clone(),
+            audio_debug_dump_dir: config.audio_debug_dump_dir.clone(),
+        }
+    }
+
+    fn is_ready(&self) -> bool {
+        self.bridge_url
+            .as_ref()
+            .is_some_and(|value| !value.is_empty())
+            && self
+                .bridge_token
+                .as_ref()
+                .is_some_and(|value| !value.is_empty())
+            && self
+                .artifact_dir
+                .as_ref()
+                .is_some_and(|value| !value.is_empty())
+            && self
+                .runtime_session_nonce
+                .as_ref()
+                .is_some_and(|value| !value.is_empty())
+    }
+
+    fn is_stream_mode(&self) -> bool {
+        self.bridge_mode.eq_ignore_ascii_case("stream")
+            || self
+                .bridge_url
+                .as_deref()
+                .is_some_and(|value| value.trim_end_matches('/').ends_with("/stream"))
+    }
+
+    fn is_asr_stream_ready(&self) -> bool {
+        self.asr_stream_enabled
+            && self
+                .asr_stream_url
+                .as_ref()
+                .is_some_and(|value| !value.is_empty())
+            && self
+                .bridge_token
+                .as_ref()
+                .is_some_and(|value| !value.is_empty())
+            && self
+                .runtime_session_nonce
+                .as_ref()
+                .is_some_and(|value| !value.is_empty())
+    }
+
+    fn realtime_asr_config(&self) -> RealtimeAsrConfig {
+        RealtimeAsrConfig {
+            enabled: self.asr_realtime_enabled,
+            url: self.asr_realtime_url.clone(),
+            runtime_token: self.bridge_token.clone(),
+            runtime_session_nonce: self.runtime_session_nonce.clone(),
+            chunk_duration_ms: self.asr_realtime_chunk_duration_ms,
+        }
+    }
 }
 
 impl Config {
@@ -150,18 +428,124 @@
             room_id: required_env("CV_LIVEKIT_ROOM_ID")?,
             bot_token: required_env("CV_LIVEKIT_BOT_TOKEN")?,
             bot_participant_identity: required_env("CV_LIVEKIT_BOT_PARTICIPANT_IDENTITY")?,
+            user_participant_identity: optional_env("CV_LIVEKIT_USER_PARTICIPANT_IDENTITY"),
             greeting_source: env::var("CV_GREETING_SOURCE")
                 .unwrap_or_else(|_| "no_audio".to_string()),
             greeting_audio_file: optional_env("CV_GREETING_AUDIO_FILE"),
             greeting_audio_url: optional_env("CV_GREETING_AUDIO_URL"),
-            role_id: env::var("CV_ROLE_ID").unwrap_or_else(|_| "90".to_string()),
-            device_output_smoke_enabled: bool_env("CV_DEVICE_OUTPUT_SMOKE_ENABLED"),
-            device_output_destination_identities: device_output_destinations(
-                optional_env("CV_DEVICE_OUTPUT_DESTINATION_IDENTITIES"),
-                optional_env("CV_LIVEKIT_USER_PARTICIPANT_IDENTITY"),
+            audio_debug_dump_dir: optional_env("CV_AUDIO_DEBUG_DUMP_DIR"),
+            runtime_turn_bridge_url: optional_env("CV_RUNTIME_TURN_BRIDGE_URL"),
+            runtime_turn_bridge_token: optional_env("CV_RUNTIME_TURN_BRIDGE_TOKEN"),
+            runtime_turn_bridge_mode: env::var("CV_RUNTIME_TURN_BRIDGE_MODE")
+                .unwrap_or_else(|_| "json".to_string()),
+            runtime_asr_stream_enabled: bool_env("CV_RUNTIME_ASR_STREAM_ENABLED", false),
+            runtime_asr_stream_url: optional_env("CV_RUNTIME_ASR_STREAM_URL"),
+            runtime_asr_realtime_enabled: bool_env("CV_RUNTIME_ASR_REALTIME_ENABLED", false),
+            runtime_asr_realtime_url: optional_env("CV_RUNTIME_ASR_REALTIME_URL"),
+            runtime_asr_realtime_chunk_duration_ms: u64_env(
+                "CV_RUNTIME_ASR_REALTIME_CHUNK_DURATION_MS",
+                200,
             ),
+            runtime_turn_artifact_dir: optional_env("CV_RUNTIME_TURN_ARTIFACT_DIR"),
+            runtime_session_nonce: optional_env("CV_RUNTIME_SESSION_NONCE"),
+            user_audio_observer_enabled: bool_env("CV_ENABLE_USER_AUDIO_OBSERVER", true),
+            simple_vad_enabled: bool_env("CV_ENABLE_SIMPLE_VAD", true),
+            simple_vad_gate_until_greeting_done: bool_env("CV_VAD_GATE_UNTIL_GREETING_DONE", true),
+            simple_vad_post_greeting_delay_ms: u64_env("CV_VAD_POST_GREETING_DELAY_MS", 800),
+            simple_vad_config: SimpleVadConfig::from_env(),
+            bot_audio_profile: BotAudioProfile::from_env()?,
         })
     }
+}
+
+fn schedule_vad_gate_enable(
+    gate: Arc<AtomicBool>,
+    config: &Config,
+    delay_ms: u64,
+    reason: &'static str,
+) {
+    if !config.simple_vad_enabled || !config.simple_vad_gate_until_greeting_done {
+        return;
+    }
+
+    let call_id = config.call_id.clone();
+    let trace_id = config.trace_id.clone();
+    info!(
+        call_id = %call_id,
+        trace_id = %trace_id,
+        delay_ms,
+        reason,
+        "runtime helper vad_enable_scheduled"
+    );
+    emit_activity(
+        &call_id,
+        &trace_id,
+        None,
+        "vad_enable_scheduled",
+        "ok",
+        None,
+        None,
+        json!({
+            "vadEnableDelayMs": delay_ms,
+            "reason": reason,
+        }),
+    );
+
+    tokio::spawn(async move {
+        if delay_ms > 0 {
+            sleep(Duration::from_millis(delay_ms)).await;
+        }
+        gate.store(true, Ordering::Release);
+        info!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            delay_ms,
+            reason,
+            "runtime helper vad_enabled"
+        );
+        emit_activity(
+            &call_id,
+            &trace_id,
+            None,
+            "vad_enabled",
+            "ok",
+            None,
+            None,
+            json!({
+                "vadEnableDelayMs": delay_ms,
+                "reason": reason,
+            }),
+        );
+    });
+}
+
+fn log_audio_diagnostics(config: &Config, diagnostics: &AudioDiagnostics) {
+    info!(
+        call_id = %config.call_id,
+        trace_id = %config.trace_id,
+        greeting_source = %config.greeting_source,
+        source_kind = diagnostics.source_kind,
+        source_format = diagnostics.source_format,
+        source_bytes = diagnostics.source_bytes,
+        source_sample_rate_hz = diagnostics.source_sample_rate_hz,
+        source_num_channels = diagnostics.source_num_channels,
+        decoded_sample_count = diagnostics.decoded_sample_count,
+        decoded_duration_ms = diagnostics.decoded_duration_ms,
+        target_sample_rate_hz = diagnostics.target_sample_rate_hz,
+        target_num_channels = diagnostics.target_num_channels,
+        target_sample_count = diagnostics.target_sample_count,
+        target_duration_ms = diagnostics.target_duration_ms,
+        frame_count = diagnostics.frame_count,
+        rms = diagnostics.rms,
+        peak = diagnostics.peak,
+        clipped_sample_count = diagnostics.clipped_sample_count,
+        silence_ratio = diagnostics.silence_ratio,
+        mp3_skipped_data_count = diagnostics.mp3_skipped_data_count,
+        mp3_insufficient_data_count = diagnostics.mp3_insufficient_data_count,
+        debug_source_path = diagnostics.debug_source_path.as_deref().unwrap_or(""),
+        debug_pcm_wav_path = diagnostics.debug_pcm_wav_path.as_deref().unwrap_or(""),
+        "runtime helper greeting audio quality diagnostics"
+    );
 }
 
 fn required_env(key: &str) -> Result<String> {
@@ -183,34 +567,44 @@
     })
 }
 
-fn bool_env(key: &str) -> bool {
-    matches!(
-        env::var(key)
-            .unwrap_or_default()
-            .trim()
-            .to_ascii_lowercase()
-            .as_str(),
-        "1" | "true" | "yes" | "y" | "on"
-    )
+fn bool_env(key: &str, default_value: bool) -> bool {
+    match env::var(key) {
+        Ok(value) => match value.trim().to_ascii_lowercase().as_str() {
+            "1" | "true" | "yes" | "on" => true,
+            "0" | "false" | "no" | "off" => false,
+            _ => default_value,
+        },
+        Err(_) => default_value,
+    }
 }
 
-fn device_output_destinations(
-    configured: Option<String>,
-    user_identity: Option<String>,
-) -> Vec<ParticipantIdentity> {
-    let identities = configured
-        .filter(|value| !value.trim().is_empty())
-        .map(|value| {
-            value
-                .split(',')
-                .map(str::trim)
-                .filter(|identity| !identity.is_empty())
-                .map(ToOwned::to_owned)
-                .collect::<Vec<_>>()
-        })
-        .or_else(|| user_identity.map(|identity| vec![identity]))
-        .unwrap_or_default();
-    identities.into_iter().map(Into::into).collect()
+fn f64_env(key: &str, default_value: f64) -> f64 {
+    env::var(key)
+        .ok()
+        .and_then(|value| value.trim().parse::<f64>().ok())
+        .filter(|value| value.is_finite() && *value >= 0.0)
+        .unwrap_or(default_value)
+}
+
+fn u64_env(key: &str, default_value: u64) -> u64 {
+    env::var(key)
+        .ok()
+        .and_then(|value| value.trim().parse::<u64>().ok())
+        .unwrap_or(default_value)
+}
+
+fn u32_env(key: &str, default_value: u32) -> u32 {
+    env::var(key)
+        .ok()
+        .and_then(|value| value.trim().parse::<u32>().ok())
+        .unwrap_or(default_value)
+}
+
+fn u16_env(key: &str, default_value: u16) -> u16 {
+    env::var(key)
+        .ok()
+        .and_then(|value| value.trim().parse::<u16>().ok())
+        .unwrap_or(default_value)
 }
 
 fn redact(value: &str) -> String {
@@ -220,118 +614,2491 @@
     format!("{}***{}", &value[..4], &value[value.len() - 4..])
 }
 
-#[derive(Serialize)]
-struct DeviceOutputMessage {
-    #[serde(rename = "type")]
-    message_type: &'static str,
-    #[serde(rename = "schemaVersion")]
-    schema_version: &'static str,
-    #[serde(rename = "callId")]
+fn safe_error(value: &str) -> String {
+    let sanitized = value.replace(['\r', '\n'], " ");
+    let trimmed = sanitized.trim();
+    let mut output: String = trimmed.chars().take(180).collect();
+    if trimmed.chars().count() > 180 {
+        output.push_str("...");
+    }
+    output
+}
+
+fn current_time_millis() -> u64 {
+    SystemTime::now()
+        .duration_since(UNIX_EPOCH)
+        .map(|value| value.as_millis() as u64)
+        .unwrap_or_default()
+}
+
+fn emit_activity(
+    call_id: &str,
+    trace_id: &str,
+    turn_id: Option<&str>,
+    event_name: &str,
+    result: &str,
+    reason_code: Option<&str>,
+    retryable: Option<bool>,
+    extension: serde_json::Value,
+) {
+    let payload = json!({
+        "type": "cv_activity",
+        "callId": call_id,
+        "traceId": trace_id,
+        "turnId": turn_id,
+        "eventName": event_name,
+        "eventWallTimeMs": current_time_millis(),
+        "result": result,
+        "reasonCode": reason_code,
+        "retryable": retryable,
+        "extension": extension,
+    });
+    println!("{payload}");
+}
+
+fn spawn_user_audio_observer(
+    events: UnboundedReceiver<RoomEvent>,
+    config: &Config,
+    vad_enabled_gate: Arc<AtomicBool>,
+    http: Client,
+    sink: Arc<BotAudioOutputSink>,
+) -> JoinHandle<()> {
+    let call_id = config.call_id.clone();
+    let trace_id = config.trace_id.clone();
+    let enabled = config.user_audio_observer_enabled;
+    let simple_vad_enabled = config.simple_vad_enabled;
+    let simple_vad_config = config.simple_vad_config.clone();
+    let turn_bridge_config = TurnBridgeConfig::from_config(config);
+
+    tokio::spawn(async move {
+        if !enabled {
+            info!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                "runtime helper user audio observer disabled"
+            );
+            return;
+        }
+        observe_user_audio_events(
+            events,
+            call_id,
+            trace_id,
+            simple_vad_enabled,
+            simple_vad_config,
+            vad_enabled_gate,
+            turn_bridge_config,
+            http,
+            sink,
+        )
+        .await;
+    })
+}
+
+async fn observe_user_audio_events(
+    mut events: UnboundedReceiver<RoomEvent>,
     call_id: String,
-    #[serde(rename = "roleId")]
-    role_id: String,
-    #[serde(rename = "traceId")]
     trace_id: String,
-    #[serde(rename = "ackMode")]
-    ack_mode: &'static str,
-    #[serde(rename = "sensorInstructions")]
-    sensor_instructions: Vec<SensorInstruction>,
+    simple_vad_enabled: bool,
+    simple_vad_config: SimpleVadConfig,
+    vad_enabled_gate: Arc<AtomicBool>,
+    turn_bridge_config: TurnBridgeConfig,
+    http: Client,
+    sink: Arc<BotAudioOutputSink>,
+) {
+    info!(
+        call_id = %call_id,
+        trace_id = %trace_id,
+        simple_vad_enabled,
+        vad_rms_threshold = simple_vad_config.rms_threshold,
+        vad_peak_threshold = simple_vad_config.peak_threshold,
+        vad_start_frames = simple_vad_config.start_frames,
+        vad_end_silence_ms = simple_vad_config.end_silence_ms,
+        vad_min_speech_ms = simple_vad_config.min_speech_ms,
+        vad_max_turn_ms = simple_vad_config.max_turn_ms,
+        vad_initial_ignore_ms = simple_vad_config.initial_ignore_ms,
+        vad_gate_enabled = vad_enabled_gate.load(Ordering::Acquire),
+        "runtime helper user_track_subscribe_requested"
+    );
+
+    while let Some(event) = events.recv().await {
+        match event {
+            RoomEvent::TrackSubscribed {
+                track: RemoteTrack::Audio(track),
+                publication: _,
+                participant,
+            } => {
+                let participant_alias = redact(&participant.identity().to_string());
+                let track_sid_alias = redact(&track.sid().to_string());
+                let track_name = track.name();
+                let track_source = format!("{:?}", track.source());
+                info!(
+                    call_id = %call_id,
+                    trace_id = %trace_id,
+                    participant_alias = %participant_alias,
+                    track_sid_alias = %track_sid_alias,
+                    track_name = %track_name,
+                    track_source = %track_source,
+                    "runtime helper user_track_subscribed"
+                );
+                spawn_user_audio_frame_observer(
+                    track,
+                    call_id.clone(),
+                    trace_id.clone(),
+                    participant_alias,
+                    track_sid_alias,
+                    simple_vad_enabled,
+                    simple_vad_config.clone(),
+                    vad_enabled_gate.clone(),
+                    turn_bridge_config.clone(),
+                    http.clone(),
+                    sink.clone(),
+                );
+            }
+            RoomEvent::TrackSubscribed {
+                track: RemoteTrack::Video(track),
+                publication: _,
+                participant,
+            } => {
+                info!(
+                    call_id = %call_id,
+                    trace_id = %trace_id,
+                    participant_alias = %redact(&participant.identity().to_string()),
+                    track_sid_alias = %redact(&track.sid().to_string()),
+                    "runtime helper ignored non-audio subscribed track"
+                );
+            }
+            RoomEvent::TrackSubscriptionFailed {
+                participant,
+                error,
+                track_sid,
+            } => {
+                warn!(
+                    call_id = %call_id,
+                    trace_id = %trace_id,
+                    participant_alias = %redact(&participant.identity().to_string()),
+                    track_sid_alias = %redact(&track_sid.to_string()),
+                    error = %error,
+                    "runtime helper user_track_subscription_failed"
+                );
+            }
+            RoomEvent::Disconnected { reason } => {
+                info!(
+                    call_id = %call_id,
+                    trace_id = %trace_id,
+                    reason = ?reason,
+                    "runtime helper room event stream disconnected"
+                );
+                break;
+            }
+            _ => {}
+        }
+    }
 }
 
-#[derive(Serialize)]
-struct SensorInstruction {
-    #[serde(rename = "commandId")]
-    command_id: String,
-    #[serde(rename = "sensorType")]
-    sensor_type: &'static str,
-    #[serde(rename = "operationType")]
-    operation_type: &'static str,
-    step: i32,
-    #[serde(rename = "durationSec", skip_serializing_if = "Option::is_none")]
-    duration_sec: Option<i32>,
-    extension: String,
+async fn handle_finished_turn(
+    http: &Client,
+    bridge_config: &TurnBridgeConfig,
+    sink: &BotAudioOutputSink,
+    call_id: &str,
+    trace_id: &str,
+    turn: FinishedSpeechTurn,
+    realtime_asr_result_ref: Option<String>,
+) {
+    let turn_pipeline_started_at = Instant::now();
+    if !bridge_config.is_ready() {
+        warn!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            turn_id = %turn.turn_id,
+            turn_index = turn.turn_index,
+            duration_ms = turn.duration_ms,
+            "runtime helper turn_bridge_skipped"
+        );
+        emit_activity(
+            call_id,
+            trace_id,
+            Some(&turn.turn_id),
+            "turn_bridge_skipped",
+            "skipped",
+            Some("TURN_BRIDGE_NOT_CONFIGURED"),
+            Some(true),
+            json!({
+                "durationMs": turn.duration_ms,
+                "frameCount": turn.frame_count,
+                "sampleCount": turn.sample_count,
+            }),
+        );
+        return;
+    }
+
+    let artifact_root = PathBuf::from(bridge_config.artifact_dir.as_deref().unwrap_or_default());
+    match write_user_turn_artifact(&artifact_root, call_id, &turn) {
+        Ok((path_ref, byte_size)) => {
+            info!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                turn_id = %turn.turn_id,
+                turn_index = turn.turn_index,
+                duration_ms = turn.duration_ms,
+                frame_count = turn.frame_count,
+                sample_count = turn.sample_count,
+                byte_size,
+                end_reason = %turn.end_reason,
+                "runtime helper turn_artifact_written"
+            );
+            emit_activity(
+                call_id,
+                trace_id,
+                Some(&turn.turn_id),
+                "turn_bridge_requested",
+                "ok",
+                None,
+                None,
+                json!({
+                    "turnDurationMs": turn.duration_ms,
+                    "turnArtifactBytes": byte_size,
+                    "frameCount": turn.frame_count,
+                    "sampleCount": turn.sample_count,
+                    "endReason": turn.end_reason.as_str(),
+                }),
+            );
+            let asr_result_ref = match realtime_asr_result_ref {
+                Some(value) => Some(value),
+                None => request_asr_result_ref(http, bridge_config, call_id, trace_id, &turn).await,
+            };
+            if bridge_config.is_stream_mode() {
+                match request_turn_bridge_stream(
+                    http,
+                    bridge_config,
+                    sink,
+                    call_id,
+                    trace_id,
+                    &turn,
+                    &path_ref,
+                    byte_size,
+                    asr_result_ref.as_deref(),
+                    turn_pipeline_started_at,
+                )
+                .await
+                {
+                    Ok(outcome) => {
+                        info!(
+                            call_id = %call_id,
+                            trace_id = %trace_id,
+                            turn_id = %turn.turn_id,
+                            audio_chunk_count = outcome.audio_chunk_count,
+                            device_output_count = outcome.device_output_count,
+                            "runtime helper turn_stream_completed"
+                        );
+                        emit_activity(
+                            call_id,
+                            trace_id,
+                            Some(&turn.turn_id),
+                            "turn_completed",
+                            "ok",
+                            None,
+                            None,
+                            json!({
+                                "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64,
+                                "audioChunkCount": outcome.audio_chunk_count,
+                                "deviceOutputCount": outcome.device_output_count,
+                                "replyPlaybackMode": outcome.reply_playback_mode,
+                            }),
+                        );
+                    }
+                    Err(error) => {
+                        let safe = safe_error(&error.to_string());
+                        warn!(
+                            call_id = %call_id,
+                            trace_id = %trace_id,
+                            turn_id = %turn.turn_id,
+                            error = %safe,
+                            "runtime helper turn_stream_failed"
+                        );
+                        emit_activity(
+                            call_id,
+                            trace_id,
+                            Some(&turn.turn_id),
+                            "turn_failed",
+                            "failed",
+                            Some("TURN_STREAM_FAILED"),
+                            Some(true),
+                            json!({
+                                "stage": "turn_bridge_stream",
+                                "error": safe,
+                            }),
+                        );
+                    }
+                }
+                return;
+            }
+            match request_turn_bridge(
+                http,
+                bridge_config,
+                call_id,
+                trace_id,
+                &turn,
+                &path_ref,
+                byte_size,
+                asr_result_ref.as_deref(),
+                turn_pipeline_started_at,
+            )
+            .await
+            {
+                Ok(outcome) => {
+                    for output in &outcome.device_outputs {
+                        if let Err(error) = sink
+                            .publish_device_output(call_id, trace_id, &turn.turn_id, output)
+                            .await
+                        {
+                            warn!(
+                                call_id = %call_id,
+                                trace_id = %trace_id,
+                                turn_id = %turn.turn_id,
+                                reason_code = "BOT_DATA_WRITE_FAILED",
+                                error = %safe_error(&error.to_string()),
+                                "runtime helper device_output_failed"
+                            );
+                        }
+                    }
+                    let Some(reply_audio_artifact) = outcome.reply_audio_artifact else {
+                        warn!(
+                            call_id = %call_id,
+                            trace_id = %trace_id,
+                            turn_id = %turn.turn_id,
+                            reason_code = "REPLY_AUDIO_MISSING",
+                            "runtime helper turn_failed"
+                        );
+                        return;
+                    };
+                    if let Err(error) = write_reply_audio_artifact(
+                        http,
+                        bridge_config,
+                        sink,
+                        call_id,
+                        trace_id,
+                        &turn,
+                        &reply_audio_artifact,
+                        turn_pipeline_started_at,
+                    )
+                    .await
+                    {
+                        warn!(
+                            call_id = %call_id,
+                            trace_id = %trace_id,
+                            turn_id = %turn.turn_id,
+                            reason_code = "BOT_AUDIO_WRITE_FAILED",
+                            error = %safe_error(&error.to_string()),
+                            "runtime helper turn_failed"
+                        );
+                        return;
+                    }
+                    info!(
+                        call_id = %call_id,
+                        trace_id = %trace_id,
+                        turn_id = %turn.turn_id,
+                        "runtime helper turn_completed"
+                    );
+                    emit_activity(
+                        call_id,
+                        trace_id,
+                        Some(&turn.turn_id),
+                        "turn_completed",
+                        "ok",
+                        None,
+                        None,
+                        json!({
+                            "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64,
+                        }),
+                    );
+                }
+                Err(error) => warn!(
+                    call_id = %call_id,
+                    trace_id = %trace_id,
+                    turn_id = %turn.turn_id,
+                    error = %safe_error(&error.to_string()),
+                    "runtime helper turn_bridge_failed"
+                ),
+            }
+        }
+        Err(error) => warn!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            turn_id = %turn.turn_id,
+            error = %safe_error(&error.to_string()),
+            "runtime helper turn_artifact_write_failed"
+        ),
+    }
 }
 
-async fn publish_device_output_smoke(room: &Room, config: &Config) -> Result<()> {
-    let payload = build_device_output_smoke_payload(config);
-    let payload_json =
-        serde_json::to_vec(&payload).context("failed to serialize device output smoke payload")?;
-    let destination_count = config.device_output_destination_identities.len();
-    room.local_participant()
-        .publish_data(DataPacket {
-            reliable: true,
-            payload: payload_json,
-            topic: Some(DEVICE_OUTPUT_TOPIC.to_string()),
-            destination_identities: config.device_output_destination_identities.clone(),
-        })
+fn write_user_turn_artifact(
+    artifact_root: &Path,
+    call_id: &str,
+    turn: &FinishedSpeechTurn,
+) -> Result<(String, u64)> {
+    require_safe_segment(call_id)?;
+    require_safe_segment(&turn.turn_id)?;
+    if turn.samples.is_empty() {
+        return Err(anyhow!("empty turn samples"));
+    }
+    let path_ref = format!("{}/{}/user.wav", call_id, turn.turn_id);
+    let output_path = normalize_path_lexically(&artifact_root.join(&path_ref));
+    let root = normalize_path_lexically(artifact_root);
+    if !output_path.starts_with(&root) {
+        return Err(anyhow!("turn artifact path escapes root"));
+    }
+    if let Some(parent) = output_path.parent() {
+        fs::create_dir_all(parent).context("failed to create turn artifact dir")?;
+    }
+    audio::write_pcm_wav(
+        &output_path,
+        &turn.samples,
+        USER_AUDIO_SAMPLE_RATE_HZ,
+        USER_AUDIO_NUM_CHANNELS,
+    )?;
+    let byte_size = fs::metadata(&output_path)
+        .context("failed to stat turn artifact")?
+        .len();
+    Ok((path_ref, byte_size))
+}
+
+async fn request_asr_result_ref(
+    http: &Client,
+    bridge_config: &TurnBridgeConfig,
+    call_id: &str,
+    trace_id: &str,
+    turn: &FinishedSpeechTurn,
+) -> Option<String> {
+    if !bridge_config.is_asr_stream_ready() {
+        return None;
+    }
+    match request_asr_stream(http, bridge_config, call_id, trace_id, turn).await {
+        Ok(Some(asr_result_ref)) => {
+            emit_activity(
+                call_id,
+                trace_id,
+                Some(&turn.turn_id),
+                "asr_stream_ref_ready",
+                "ok",
+                None,
+                None,
+                json!({
+                    "asrResultRefPresent": true,
+                    "format": "pcm_s16le",
+                    "sampleRate": 16000,
+                    "channels": 1,
+                }),
+            );
+            Some(asr_result_ref)
+        }
+        Ok(None) => None,
+        Err(error) => {
+            warn!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                turn_id = %turn.turn_id,
+                error = %safe_error(&error.to_string()),
+                "runtime helper asr_stream_failed_fallback"
+            );
+            emit_activity(
+                call_id,
+                trace_id,
+                Some(&turn.turn_id),
+                "asr_stream_fallback",
+                "skipped",
+                Some("ASR_STREAM_INTERRUPTED"),
+                Some(true),
+                json!({
+                    "fallbackReason": "asr_stream_request_failed",
+                    "fallbackStage": "asr_stream",
+                }),
+            );
+            None
+        }
+    }
+}
+
+async fn request_asr_stream(
+    http: &Client,
+    bridge_config: &TurnBridgeConfig,
+    call_id: &str,
+    trace_id: &str,
+    turn: &FinishedSpeechTurn,
+) -> Result<Option<String>> {
+    let started_at = Instant::now();
+    let ndjson = build_asr_stream_ndjson(call_id, trace_id, turn, bridge_config)?;
+    let response = http
+        .post(bridge_config.asr_stream_url.as_deref().unwrap_or_default())
+        .header("Content-Type", "application/x-ndjson")
+        .header(
+            "X-CV-Runtime-Token",
+            bridge_config.bridge_token.as_deref().unwrap_or_default(),
+        )
+        .header("X-CV-Call-Id", call_id)
+        .header("X-CV-Trace-Id", trace_id)
+        .header(
+            "X-CV-Runtime-Session-Nonce",
+            bridge_config
+                .runtime_session_nonce
+                .as_deref()
+                .unwrap_or_default(),
+        )
+        .body(ndjson)
+        .send()
         .await
-        .map_err(|error| anyhow!("failed to publish livekit device output data: {error}"))?;
+        .context("failed to post asr stream")?;
+    let status = response.status();
+    if !status.is_success() {
+        let body_len = response.text().await.map(|body| body.len()).unwrap_or(0);
+        warn!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            turn_id = %turn.turn_id,
+            http_status = status.as_u16(),
+            body_len,
+            "runtime helper asr_stream_http_failed"
+        );
+        return Ok(None);
+    }
+    let body: RuntimeTurnCommonResult<RuntimeAsrStreamResp> = response
+        .json()
+        .await
+        .context("failed to decode asr stream response")?;
+    if body.code != 0 {
+        warn!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            turn_id = %turn.turn_id,
+            code = body.code,
+            msg_len = body.msg.as_deref().unwrap_or_default().len(),
+            "runtime helper asr_stream_common_result_failed"
+        );
+        return Ok(None);
+    }
+    let Some(data) = body.data else {
+        return Ok(None);
+    };
+    if data.status.as_deref() == Some("final") {
+        info!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            turn_id = %turn.turn_id,
+            chunk_count = data.chunk_count.unwrap_or_default(),
+            audio_bytes = data.audio_bytes.unwrap_or_default(),
+            asr_duration_ms = data.asr_duration_ms.unwrap_or_default(),
+            wall_ms = started_at.elapsed().as_millis() as u64,
+            provider = %data.provider_alias.as_deref().unwrap_or("unknown"),
+            text_len = data.text_len.unwrap_or_default(),
+            "runtime helper asr_stream_final"
+        );
+        return Ok(data.asr_result_ref);
+    }
+    emit_activity(
+        call_id,
+        trace_id,
+        Some(&turn.turn_id),
+        "asr_stream_fallback",
+        "skipped",
+        None,
+        Some(true),
+        json!({
+            "fallbackReason": data.fallback_reason,
+            "fallbackStage": data.fallback_stage,
+            "status": data.status,
+        }),
+    );
+    Ok(None)
+}
+
+fn build_asr_stream_ndjson(
+    call_id: &str,
+    trace_id: &str,
+    turn: &FinishedSpeechTurn,
+    bridge_config: &TurnBridgeConfig,
+) -> Result<String> {
+    let chunks = asr_pcm_16k_chunks(turn)?;
+    let mut seq = 1u64;
+    let mut lines = Vec::with_capacity(chunks.len() + 3);
+    lines.push(serde_json::to_string(&json!({
+        "event": "asr_stream_started",
+        "seq": seq,
+        "callId": call_id,
+        "traceId": trace_id,
+        "turnId": turn.turn_id.as_str(),
+        "tsMs": current_time_millis(),
+        "payload": {
+            "format": "pcm_s16le",
+            "sampleRate": 16000,
+            "channels": 1,
+            "runtimeSessionNonce": bridge_config.runtime_session_nonce.as_deref().unwrap_or_default(),
+            "providerHint": "volcengine",
+        }
+    }))?);
+    for (index, samples) in chunks.iter().enumerate() {
+        seq += 1;
+        let bytes = pcm_i16_to_le_bytes(samples);
+        let duration_ms = ((samples.len() as u64) * 1000 / 16_000).max(1);
+        lines.push(serde_json::to_string(&json!({
+            "event": "asr_audio_chunk",
+            "seq": seq,
+            "callId": call_id,
+            "traceId": trace_id,
+            "turnId": turn.turn_id.as_str(),
+            "tsMs": current_time_millis(),
+            "payload": {
+                "chunkSeq": index + 1,
+                "format": "pcm_s16le",
+                "sampleRate": 16000,
+                "channels": 1,
+                "durationMs": duration_ms,
+                "payloadBase64": general_purpose::STANDARD.encode(bytes),
+            }
+        }))?);
+    }
+    seq += 1;
+    lines.push(serde_json::to_string(&json!({
+        "event": "vad_speech_end",
+        "seq": seq,
+        "callId": call_id,
+        "traceId": trace_id,
+        "turnId": turn.turn_id.as_str(),
+        "tsMs": current_time_millis(),
+        "payload": {
+            "endReason": turn.end_reason.as_str(),
+            "speechDurationMs": turn.duration_ms,
+        }
+    }))?);
+    seq += 1;
+    lines.push(serde_json::to_string(&json!({
+        "event": "asr_stream_finish",
+        "seq": seq,
+        "callId": call_id,
+        "traceId": trace_id,
+        "turnId": turn.turn_id.as_str(),
+        "tsMs": current_time_millis(),
+        "payload": {
+            "finalChunkSeq": chunks.len(),
+            "audioDurationMs": turn.duration_ms,
+        }
+    }))?);
+    Ok(lines.join("\n") + "\n")
+}
+
+fn asr_pcm_16k_chunks(turn: &FinishedSpeechTurn) -> Result<Vec<Vec<i16>>> {
+    if turn.samples.is_empty() {
+        return Err(anyhow!("empty turn samples"));
+    }
+    let samples_16k: Vec<i16> = turn.samples.iter().step_by(3).copied().collect();
+    if samples_16k.is_empty() {
+        return Err(anyhow!("empty 16k asr samples"));
+    }
+    let samples_per_chunk = 320usize;
+    Ok(samples_16k
+        .chunks(samples_per_chunk)
+        .map(|chunk| chunk.to_vec())
+        .collect())
+}
+
+fn pcm_i16_to_le_bytes(samples: &[i16]) -> Vec<u8> {
+    let mut bytes = Vec::with_capacity(samples.len() * 2);
+    for sample in samples {
+        bytes.extend_from_slice(&sample.to_le_bytes());
+    }
+    bytes
+}
+
+async fn request_turn_bridge_stream(
+    http: &Client,
+    bridge_config: &TurnBridgeConfig,
+    sink: &BotAudioOutputSink,
+    call_id: &str,
+    trace_id: &str,
+    turn: &FinishedSpeechTurn,
+    path_ref: &str,
+    byte_size: u64,
+    asr_result_ref: Option<&str>,
+    turn_pipeline_started_at: Instant,
+) -> Result<RuntimeTurnStreamOutcome> {
+    let bridge_started_at = Instant::now();
+    let request = RuntimeTurnRequest {
+        call_id: call_id.to_string(),
+        trace_id: trace_id.to_string(),
+        turn_id: turn.turn_id.clone(),
+        audio_artifact: Some(RuntimeTurnAudioArtifact {
+            artifact_type: "local_file".to_string(),
+            path_ref: path_ref.to_string(),
+            format: "wav".to_string(),
+            sample_rate: USER_AUDIO_SAMPLE_RATE_HZ,
+            channels: u32::from(USER_AUDIO_NUM_CHANNELS),
+            duration_ms: turn.duration_ms,
+            byte_size,
+        }),
+        asr_result_ref: asr_result_ref.map(str::to_string),
+    };
+    let response = http
+        .post(bridge_config.bridge_url.as_deref().unwrap_or_default())
+        .header(
+            "X-CV-Runtime-Token",
+            bridge_config.bridge_token.as_deref().unwrap_or_default(),
+        )
+        .header("X-CV-Call-Id", call_id)
+        .header("X-CV-Trace-Id", trace_id)
+        .header(
+            "X-CV-Runtime-Session-Nonce",
+            bridge_config
+                .runtime_session_nonce
+                .as_deref()
+                .unwrap_or_default(),
+        )
+        .json(&request)
+        .send()
+        .await
+        .context("failed to post turn stream bridge")?;
+    let status = response.status();
+    if !status.is_success() {
+        let body_len = response.text().await.map(|body| body.len()).unwrap_or(0);
+        warn!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            turn_id = %turn.turn_id,
+            http_status = status.as_u16(),
+            body_len,
+            "runtime helper turn_stream_http_failed"
+        );
+        return Err(anyhow!("turn stream bridge http failed"));
+    }
 
     info!(
-        call_id = %config.call_id,
-        trace_id = %config.trace_id,
-        topic = DEVICE_OUTPUT_TOPIC,
-        reliable = true,
-        command_count = payload.sensor_instructions.len(),
-        destination_count,
-        "runtime helper published device output smoke data"
+        call_id = %call_id,
+        trace_id = %trace_id,
+        turn_id = %turn.turn_id,
+        "runtime helper turn_stream_connected"
+    );
+    let mut byte_stream = response.bytes_stream();
+    let mut line_buffer: Vec<u8> = Vec::new();
+    let mut state = RuntimeTurnStreamState::default();
+    while let Some(chunk) = byte_stream.next().await {
+        let chunk = chunk.context("failed to read turn stream chunk")?;
+        line_buffer.extend_from_slice(&chunk);
+        while let Some(newline_index) = line_buffer.iter().position(|value| *value == b'\n') {
+            let line: Vec<u8> = line_buffer.drain(..=newline_index).collect();
+            if let Some(event) = parse_turn_stream_event_line(&line)? {
+                handle_turn_stream_event(
+                    bridge_config,
+                    sink,
+                    call_id,
+                    trace_id,
+                    turn,
+                    event,
+                    &mut state,
+                    turn_pipeline_started_at,
+                )
+                .await?;
+            }
+        }
+    }
+    if !line_buffer.is_empty() {
+        if let Some(event) = parse_turn_stream_event_line(&line_buffer)? {
+            handle_turn_stream_event(
+                bridge_config,
+                sink,
+                call_id,
+                trace_id,
+                turn,
+                event,
+                &mut state,
+                turn_pipeline_started_at,
+            )
+            .await?;
+        }
+    }
+    if !state.completed {
+        warn!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            turn_id = %turn.turn_id,
+            audio_chunk_count = state.audio_chunk_count,
+            "runtime helper turn_stream_completed_without_final_event"
+        );
+        return Err(anyhow!("turn stream ended without turn_completed"));
+    }
+    emit_activity(
+        call_id,
+        trace_id,
+        Some(&turn.turn_id),
+        "turn_bridge_completed",
+        "ok",
+        None,
+        None,
+        json!({
+            "bridgeWallDurationMs": bridge_started_at.elapsed().as_millis() as u64,
+            "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64,
+            "audioChunkCount": state.audio_chunk_count,
+            "deviceOutputCount": state.device_output_count,
+            "replyPlaybackMode": state.reply_playback_mode.as_str(),
+        }),
+    );
+    Ok(RuntimeTurnStreamOutcome {
+        reply_playback_mode: state.reply_playback_mode,
+        audio_chunk_count: state.audio_chunk_count,
+        device_output_count: state.device_output_count,
+    })
+}
+
+fn parse_turn_stream_event_line(line: &[u8]) -> Result<Option<RuntimeTurnStreamEvent>> {
+    let line = trim_ascii_whitespace(line);
+    if line.is_empty() {
+        return Ok(None);
+    }
+    serde_json::from_slice(line)
+        .context("failed to parse turn stream event")
+        .map(Some)
+}
+
+fn diagnostic_str<'a>(diagnostics: Option<&'a serde_json::Value>, key: &str) -> Option<&'a str> {
+    diagnostics?.get(key)?.as_str()
+}
+
+fn diagnostic_bool(diagnostics: Option<&serde_json::Value>, key: &str) -> Option<bool> {
+    diagnostics?.get(key)?.as_bool()
+}
+
+fn diagnostic_u64(diagnostics: Option<&serde_json::Value>, key: &str) -> Option<u64> {
+    diagnostics?.get(key)?.as_u64()
+}
+
+async fn handle_turn_stream_event(
+    bridge_config: &TurnBridgeConfig,
+    sink: &BotAudioOutputSink,
+    call_id: &str,
+    trace_id: &str,
+    turn: &FinishedSpeechTurn,
+    event: RuntimeTurnStreamEvent,
+    state: &mut RuntimeTurnStreamState,
+    turn_pipeline_started_at: Instant,
+) -> Result<()> {
+    let event_type = event.event_type();
+    match event_type.as_deref() {
+        Some("reply_playback_mode_selected") => {
+            if let Some(reply_playback_mode) = event.reply_playback_mode.as_deref() {
+                state.reply_playback_mode = reply_playback_mode.to_string();
+            }
+            let diagnostics = event.diagnostics.as_ref();
+            info!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                turn_id = %turn.turn_id,
+                reply_playback_mode = %state.reply_playback_mode,
+                streaming_enabled = ?diagnostic_bool(diagnostics, "streamingEnabled"),
+                provider_streaming_supported = ?diagnostic_bool(diagnostics, "providerStreamingSupported"),
+                first_chunk_received = ?diagnostic_bool(diagnostics, "firstChunkReceived"),
+                stream_chunk_count = ?diagnostic_u64(diagnostics, "streamChunkCount"),
+                fallback_reason = %diagnostic_str(diagnostics, "fallbackReason").unwrap_or("none"),
+                fallback_stage = %diagnostic_str(diagnostics, "fallbackStage").unwrap_or("none"),
+                stream_bridge_mode = %diagnostic_str(diagnostics, "streamBridgeMode").unwrap_or("unknown"),
+                tts_provider = %diagnostic_str(diagnostics, "ttsProvider").unwrap_or("unknown"),
+                "runtime helper reply_playback_mode_selected"
+            );
+        }
+        Some("reply_state") => {
+            if let Some(reply_state) = event.state.as_deref() {
+                publish_reply_state_from_stream_event(
+                    sink,
+                    call_id,
+                    trace_id,
+                    &turn.turn_id,
+                    &state.reply_playback_mode,
+                    reply_state,
+                    event.seq,
+                )
+                .await?;
+                if reply_state == "reply_playback_started" {
+                    state.playback_started_sent = true;
+                }
+            }
+        }
+        Some("reply_audio_chunk") => {
+            let audio_chunk = event
+                .audio_chunk
+                .as_ref()
+                .ok_or_else(|| anyhow!("reply_audio_chunk event missing audioChunk"))?;
+            if !state.playback_started_sent {
+                state.reply_state_seq = state.reply_state_seq.saturating_add(1);
+                sink.publish_reply_state(
+                    call_id,
+                    trace_id,
+                    &turn.turn_id,
+                    &state.reply_playback_mode,
+                    "reply_playback_started",
+                    state.reply_state_seq,
+                )
+                .await?;
+                state.playback_started_sent = true;
+            }
+            let written_frames = write_stream_audio_chunk(
+                bridge_config,
+                sink,
+                call_id,
+                trace_id,
+                turn,
+                audio_chunk,
+                state,
+                turn_pipeline_started_at,
+            )
+            .await?;
+            if written_frames > 0 {
+                state.audio_chunk_count = state.audio_chunk_count.saturating_add(1);
+            }
+        }
+        Some("device_output") => {
+            if let Some(output) = event.device_output.as_ref() {
+                sink.publish_device_output(call_id, trace_id, &turn.turn_id, output)
+                    .await?;
+                state.device_output_count = state.device_output_count.saturating_add(1);
+            }
+        }
+        Some("turn_completed") => {
+            state.completed = true;
+            info!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                turn_id = %turn.turn_id,
+                message_id_alias = %event.completion.as_ref().and_then(|value| value.message_id.as_deref()).map(redact).unwrap_or_else(|| "none".to_string()),
+                audio_chunk_count = state.audio_chunk_count,
+                "runtime helper turn_stream_final_received"
+            );
+        }
+        Some("turn_failed") => {
+            let error = event.error.as_ref();
+            let reason_code = error
+                .and_then(|value| value.reason_code.as_deref())
+                .unwrap_or("TURN_STREAM_FAILED");
+            let stage = error
+                .and_then(|value| value.stage.as_deref())
+                .unwrap_or("turn_bridge_stream");
+            let retryable = error.and_then(|value| value.retryable).unwrap_or(false);
+            warn!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                turn_id = %turn.turn_id,
+                reason_code = %reason_code,
+                stage = %stage,
+                retryable = retryable,
+                "runtime helper turn_stream_failed_event"
+            );
+            emit_activity(
+                call_id,
+                trace_id,
+                Some(&turn.turn_id),
+                "turn_failed",
+                "failed",
+                Some(reason_code),
+                Some(retryable),
+                json!({
+                    "stage": stage,
+                    "replyPlaybackMode": state.reply_playback_mode.as_str(),
+                }),
+            );
+            return Err(anyhow!("turn stream failed event"));
+        }
+        Some("turn_cancelled") => {
+            state.completed = true;
+            info!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                turn_id = %turn.turn_id,
+                "runtime helper turn_stream_cancelled_event"
+            );
+            publish_reply_state_from_stream_event(
+                sink,
+                call_id,
+                trace_id,
+                &turn.turn_id,
+                &state.reply_playback_mode,
+                "reply_playback_cancelled",
+                event.seq,
+            )
+            .await?;
+            emit_activity(
+                call_id,
+                trace_id,
+                Some(&turn.turn_id),
+                "turn_cancelled",
+                "ok",
+                None,
+                Some(false),
+                json!({
+                    "replyPlaybackMode": state.reply_playback_mode.as_str(),
+                    "audioChunkCount": state.audio_chunk_count,
+                }),
+            );
+        }
+        Some("activity") => {
+            if let Some(activity) = event.activity.as_ref() {
+                info!(
+                    call_id = %call_id,
+                    trace_id = %trace_id,
+                    turn_id = %turn.turn_id,
+                    activity_event = %activity.event_type.as_deref().unwrap_or("unknown"),
+                    stage = %activity.stage.as_deref().unwrap_or("unknown"),
+                    reason_code = %activity.reason_code.as_deref().unwrap_or("none"),
+                    "runtime helper turn_stream_activity"
+                );
+            }
+        }
+        Some(other) => {
+            info!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                turn_id = %turn.turn_id,
+                event_type = %other,
+                "runtime helper ignored turn stream event"
+            );
+        }
+        None => {
+            warn!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                turn_id = %turn.turn_id,
+                "runtime helper ignored turn stream event without type"
+            );
+        }
+    }
+    Ok(())
+}
+
+async fn publish_reply_state_from_stream_event(
+    sink: &BotAudioOutputSink,
+    call_id: &str,
+    trace_id: &str,
+    turn_id: &str,
+    reply_playback_mode: &str,
+    state: &str,
+    seq: Option<u64>,
+) -> Result<()> {
+    sink.publish_reply_state(
+        call_id,
+        trace_id,
+        turn_id,
+        reply_playback_mode,
+        state,
+        seq.unwrap_or(0),
+    )
+    .await
+}
+
+async fn write_stream_audio_chunk(
+    bridge_config: &TurnBridgeConfig,
+    sink: &BotAudioOutputSink,
+    call_id: &str,
+    trace_id: &str,
+    turn: &FinishedSpeechTurn,
+    audio_chunk: &RuntimeTurnStreamAudioChunk,
+    state: &mut RuntimeTurnStreamState,
+    turn_pipeline_started_at: Instant,
+) -> Result<usize> {
+    let payload_base64 = audio_chunk
+        .payload_base64
+        .as_deref()
+        .map(str::trim)
+        .filter(|value| !value.is_empty())
+        .ok_or_else(|| anyhow!("reply_audio_chunk payloadBase64 missing"))?;
+    let payload = general_purpose::STANDARD
+        .decode(payload_base64)
+        .context("failed to decode reply_audio_chunk payloadBase64")?;
+    let format = audio_chunk
+        .format
+        .as_deref()
+        .unwrap_or("pcm_s16le")
+        .trim()
+        .to_ascii_lowercase();
+    let frames = if format == "pcm_s16le" {
+        let 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 {
+        return Err(anyhow!("unsupported reply_audio_chunk format {format}"));
+    };
+    if frames.is_empty() {
+        return Ok(0);
+    }
+
+    if !state.first_audio_frame_written {
+        emit_activity(
+            call_id,
+            trace_id,
+            Some(&turn.turn_id),
+            "bot_reply_audio_write_started",
+            "ok",
+            None,
+            None,
+            json!({
+                "replyPlaybackMode": state.reply_playback_mode.as_str(),
+                "format": format.as_str(),
+                "chunkSeq": audio_chunk.chunk_seq,
+                "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64,
+            }),
+        );
+    }
+    let pacing_started_at = tokio::time::Instant::now();
+    for (index, frame) in frames.iter().enumerate() {
+        sink.write_pcm_frame(frame).await?;
+        if !state.first_audio_frame_written {
+            state.first_audio_frame_written = true;
+            emit_activity(
+                call_id,
+                trace_id,
+                Some(&turn.turn_id),
+                "bot_reply_first_audio_frame_written",
+                "ok",
+                None,
+                None,
+                json!({
+                    "replyPlaybackMode": state.reply_playback_mode.as_str(),
+                    "format": format.as_str(),
+                    "chunkSeq": audio_chunk.chunk_seq,
+                    "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64,
+                }),
+            );
+        }
+        sleep_until(pacing_started_at + Duration::from_millis(((index + 1) as u64) * 20)).await;
+    }
+    if audio_chunk.last.unwrap_or(false) {
+        let mut debug_source_path = None;
+        let mut debug_pcm_wav_path = None;
+        let mut debug_pcm_wav_size_bytes = None;
+        if format == "pcm_s16le" {
+            if let (Some(debug_dump_dir), Some(decoder)) = (
+                bridge_config.audio_debug_dump_dir.as_deref(),
+                state.pcm_stream_decoder.as_ref(),
+            ) {
+                match decoder.write_debug_dump(
+                    debug_dump_dir,
+                    call_id,
+                    &format!("stream-reply-{}", turn.turn_id),
+                ) {
+                    Ok(debug_dump) => {
+                        debug_source_path = debug_dump.debug_source_path;
+                        debug_pcm_wav_path = debug_dump.debug_pcm_wav_path;
+                        debug_pcm_wav_size_bytes = debug_dump.debug_pcm_wav_size_bytes;
+                    }
+                    Err(error) => warn!(
+                        call_id = %call_id,
+                        trace_id = %trace_id,
+                        turn_id = %turn.turn_id,
+                        error = %safe_error(&error.to_string()),
+                        "runtime helper stream_audio_pcm_debug_dump_failed"
+                    ),
+                }
+            }
+        }
+        emit_activity(
+            call_id,
+            trace_id,
+            Some(&turn.turn_id),
+            "bot_reply_audio_write_finished",
+            "ok",
+            None,
+            None,
+            json!({
+                "replyPlaybackMode": state.reply_playback_mode.as_str(),
+                "format": format.as_str(),
+                "chunkSeq": audio_chunk.chunk_seq,
+                "networkChunkCount": state.pcm_stream_network_chunk_count,
+                "debugSourcePath": debug_source_path,
+                "debugPcmWavPath": debug_pcm_wav_path,
+                "debugPcmWavSizeBytes": debug_pcm_wav_size_bytes,
+                "sourceSampleRate": audio_chunk.sample_rate,
+                "sourceChannels": audio_chunk.channels,
+                "targetAudioProfile": sink.profile.as_str(),
+                "targetSampleRate": sink.sample_rate_hz,
+                "targetChannels": sink.num_channels,
+                "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64,
+            }),
+        );
+    }
+    Ok(frames.len())
+}
+
+fn trim_ascii_whitespace(value: &[u8]) -> &[u8] {
+    let mut start = 0;
+    let mut end = value.len();
+    while start < end && value[start].is_ascii_whitespace() {
+        start += 1;
+    }
+    while end > start && value[end - 1].is_ascii_whitespace() {
+        end -= 1;
+    }
+    &value[start..end]
+}
+
+async fn request_turn_bridge(
+    http: &Client,
+    bridge_config: &TurnBridgeConfig,
+    call_id: &str,
+    trace_id: &str,
+    turn: &FinishedSpeechTurn,
+    path_ref: &str,
+    byte_size: u64,
+    asr_result_ref: Option<&str>,
+    turn_pipeline_started_at: Instant,
+) -> Result<RuntimeTurnBridgeOutcome> {
+    let bridge_started_at = Instant::now();
+    let request = RuntimeTurnRequest {
+        call_id: call_id.to_string(),
+        trace_id: trace_id.to_string(),
+        turn_id: turn.turn_id.clone(),
+        audio_artifact: Some(RuntimeTurnAudioArtifact {
+            artifact_type: "local_file".to_string(),
+            path_ref: path_ref.to_string(),
+            format: "wav".to_string(),
+            sample_rate: USER_AUDIO_SAMPLE_RATE_HZ,
+            channels: u32::from(USER_AUDIO_NUM_CHANNELS),
+            duration_ms: turn.duration_ms,
+            byte_size,
+        }),
+        asr_result_ref: asr_result_ref.map(str::to_string),
+    };
+    let response = http
+        .post(bridge_config.bridge_url.as_deref().unwrap_or_default())
+        .header(
+            "X-CV-Runtime-Token",
+            bridge_config.bridge_token.as_deref().unwrap_or_default(),
+        )
+        .header("X-CV-Call-Id", call_id)
+        .header("X-CV-Trace-Id", trace_id)
+        .header(
+            "X-CV-Runtime-Session-Nonce",
+            bridge_config
+                .runtime_session_nonce
+                .as_deref()
+                .unwrap_or_default(),
+        )
+        .json(&request)
+        .send()
+        .await
+        .context("failed to post turn bridge")?;
+    let status = response.status();
+    let body = response
+        .text()
+        .await
+        .context("failed to read turn bridge response")?;
+    let body_len = body.len();
+    let parsed: Option<RuntimeTurnCommonResult<RuntimeTurnResponseData>> =
+        serde_json::from_str(&body).ok();
+    if !status.is_success() {
+        warn!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            turn_id = %turn.turn_id,
+            http_status = status.as_u16(),
+            body_len,
+            code = parsed.as_ref().map(|value| value.code).unwrap_or_default(),
+            reason_code = %parsed.as_ref().and_then(|value| value.msg.as_deref()).unwrap_or("unknown"),
+            "runtime helper turn_bridge_http_failed"
+        );
+        return Err(anyhow!("turn bridge http failed"));
+    }
+    let Some(result) = parsed else {
+        warn!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            turn_id = %turn.turn_id,
+            body_len,
+            "runtime helper turn_bridge_response_invalid"
+        );
+        return Err(anyhow!("turn bridge response invalid"));
+    };
+    if result.code != 0 {
+        warn!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            turn_id = %turn.turn_id,
+            code = result.code,
+            reason_code = %result.msg.as_deref().unwrap_or("unknown"),
+            retryable = result.retryable.unwrap_or(false),
+            stage = %result.stage.as_deref().unwrap_or("unknown"),
+            "runtime helper turn_bridge_business_failed"
+        );
+        return Err(anyhow!("turn bridge business failed"));
+    }
+    let data = result.data;
+    let reply_audio_artifact = data
+        .as_ref()
+        .and_then(|value| value.reply_audio_artifact.as_ref());
+    let device_outputs = data
+        .as_ref()
+        .and_then(|value| value.device_outputs.clone())
+        .unwrap_or_default();
+    info!(
+        call_id = %call_id,
+        trace_id = %trace_id,
+        turn_id = %turn.turn_id,
+        response_turn_id = %data.as_ref().and_then(|value| value.turn_id.as_deref()).unwrap_or("unknown"),
+        message_id_alias = %data.as_ref().and_then(|value| value.message_id.as_deref()).map(redact).unwrap_or_else(|| "none".to_string()),
+        reply_audio_present = reply_audio_artifact.is_some(),
+        reply_audio_type = %reply_audio_artifact.and_then(|value| value.artifact_type.as_deref()).unwrap_or("none"),
+        reply_audio_path_alias = %reply_audio_artifact.and_then(|value| value.path_ref.as_deref()).map(redact).unwrap_or_else(|| "none".to_string()),
+        reply_audio_format = %reply_audio_artifact.and_then(|value| value.format.as_deref()).unwrap_or("none"),
+        device_output_count = device_outputs.len(),
+        "runtime helper turn_bridge_completed"
+    );
+    emit_activity(
+        call_id,
+        trace_id,
+        Some(&turn.turn_id),
+        "turn_bridge_completed",
+        "ok",
+        None,
+        None,
+        json!({
+            "bridgeWallDurationMs": bridge_started_at.elapsed().as_millis() as u64,
+            "replyAudioPresent": reply_audio_artifact.is_some(),
+            "deviceOutputCount": device_outputs.len(),
+            "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64,
+        }),
+    );
+    Ok(RuntimeTurnBridgeOutcome {
+        reply_audio_artifact: reply_audio_artifact.cloned(),
+        device_outputs,
+    })
+}
+
+async fn write_reply_audio_artifact(
+    http: &Client,
+    bridge_config: &TurnBridgeConfig,
+    sink: &BotAudioOutputSink,
+    call_id: &str,
+    trace_id: &str,
+    turn: &FinishedSpeechTurn,
+    artifact: &RuntimeTurnReplyAudioArtifact,
+    turn_pipeline_started_at: Instant,
+) -> Result<()> {
+    let artifact_type = artifact.artifact_type.as_deref().unwrap_or_default().trim();
+    if artifact_type != "local_file" {
+        return Err(anyhow!("unsupported reply audio artifact type"));
+    }
+    let path_ref = artifact
+        .path_ref
+        .as_deref()
+        .map(str::trim)
+        .filter(|value| !value.is_empty())
+        .ok_or_else(|| anyhow!("reply audio pathRef missing"))?;
+    let artifact_root = PathBuf::from(bridge_config.artifact_dir.as_deref().unwrap_or_default());
+    let audio_path = resolve_artifact_path(&artifact_root, path_ref)?;
+    let audio_path_string = audio_path.to_string_lossy().to_string();
+    let reply_audio_format = artifact.format.as_deref().unwrap_or("unknown");
+
+    info!(
+        call_id = %call_id,
+        trace_id = %trace_id,
+        turn_id = %turn.turn_id,
+        reply_audio_type = %artifact_type,
+        reply_audio_path_alias = %redact(path_ref),
+        reply_audio_format = %reply_audio_format,
+        "runtime helper bot_reply_audio_write_started"
+    );
+    emit_activity(
+        call_id,
+        trace_id,
+        Some(&turn.turn_id),
+        "bot_reply_audio_write_started",
+        "ok",
+        None,
+        None,
+        json!({
+            "replyAudioFormat": reply_audio_format,
+            "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64,
+        }),
+    );
+
+    let reply_debug_label = format!("reply-{}", turn.turn_id);
+    let loaded_audio = load_pre_recorded_frames(
+        http,
+        Some(&audio_path_string),
+        None,
+        sink.sample_rate_hz,
+        sink.num_channels,
+        bridge_config.audio_debug_dump_dir.as_deref(),
+        call_id,
+        &reply_debug_label,
+    )
+    .await?
+    .ok_or_else(|| anyhow!("reply audio artifact decode returned empty"))?;
+    let AudioDiagnostics {
+        source_format,
+        source_bytes,
+        source_sample_rate_hz,
+        source_num_channels,
+        target_duration_ms,
+        frame_count,
+        rms,
+        peak,
+        clipped_sample_count,
+        silence_ratio,
+        mp3_skipped_data_count,
+        mp3_insufficient_data_count,
+        debug_source_path,
+        debug_pcm_wav_path,
+        ..
+    } = loaded_audio.diagnostics;
+    let frames = loaded_audio.frames;
+
+    let playback_started_at = Instant::now();
+    let pacing_started_at = tokio::time::Instant::now();
+    for (index, frame) in frames.iter().enumerate() {
+        sink.write_pcm_frame(frame).await?;
+        sleep_until(pacing_started_at + Duration::from_millis(((index + 1) as u64) * 20)).await;
+    }
+    let playback_wall_ms = playback_started_at.elapsed().as_millis() as i64;
+    let drift_ms = playback_wall_ms - target_duration_ms as i64;
+    sink.clear_buffer();
+
+    info!(
+        call_id = %call_id,
+        trace_id = %trace_id,
+        turn_id = %turn.turn_id,
+        source_format,
+        source_bytes,
+        source_sample_rate_hz,
+        source_num_channels,
+        frame_count,
+        theoretical_duration_ms = target_duration_ms,
+        push_wall_duration_ms = playback_wall_ms,
+        push_drift_ms = drift_ms,
+        rms = round4(rms),
+        peak = round4(peak),
+        clipped_sample_count,
+        silence_ratio = round4(silence_ratio),
+        mp3_skipped_data_count,
+        mp3_insufficient_data_count,
+        debug_source_path = debug_source_path.as_deref().unwrap_or(""),
+        debug_pcm_wav_path = debug_pcm_wav_path.as_deref().unwrap_or(""),
+        "runtime helper bot_reply_audio_write_finished"
+    );
+    emit_activity(
+        call_id,
+        trace_id,
+        Some(&turn.turn_id),
+        "bot_reply_audio_write_finished",
+        "ok",
+        None,
+        None,
+        json!({
+            "sourceFormat": source_format,
+            "sourceBytes": source_bytes,
+            "sourceSampleRateHz": source_sample_rate_hz,
+            "sourceNumChannels": source_num_channels,
+            "frameCount": frame_count,
+            "theoreticalDurationMs": target_duration_ms,
+            "pushWallDurationMs": playback_wall_ms,
+            "pushDriftMs": drift_ms,
+            "rms": round4(rms),
+            "peak": round4(peak),
+            "clippedSampleCount": clipped_sample_count,
+            "silenceRatio": round4(silence_ratio),
+            "mp3SkippedDataCount": mp3_skipped_data_count,
+            "mp3InsufficientDataCount": mp3_insufficient_data_count,
+            "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64,
+        }),
     );
     Ok(())
 }
 
-fn build_device_output_smoke_payload(config: &Config) -> DeviceOutputMessage {
-    DeviceOutputMessage {
-        message_type: DEVICE_OUTPUT_TOPIC,
-        schema_version: "1.0",
-        call_id: config.call_id.clone(),
-        role_id: config.role_id.clone(),
-        trace_id: config.trace_id.clone(),
-        ack_mode: "http",
-        sensor_instructions: vec![
-            SensorInstruction {
-                command_id: command_id(&config.call_id, 1),
-                sensor_type: "Vibrator",
-                operation_type: "VibratorStart",
-                step: 1,
-                duration_sec: None,
-                extension: r#"{"levelList":[{"level":1,"percent":"0.6"}]}"#.to_string(),
-            },
-            SensorInstruction {
-                command_id: command_id(&config.call_id, 2),
-                sensor_type: "Vibrator",
-                operation_type: "VibratorUp",
-                step: 1,
-                duration_sec: None,
-                extension: r#"{"delta":1}"#.to_string(),
-            },
-            SensorInstruction {
-                command_id: command_id(&config.call_id, 3),
-                sensor_type: "Pump",
-                operation_type: "JiaStart",
-                step: 1,
-                duration_sec: None,
-                extension: r#"{"levelList":[{"level":1,"percent":"0.5"}]}"#.to_string(),
-            },
-            SensorInstruction {
-                command_id: command_id(&config.call_id, 4),
-                sensor_type: "Heating",
-                operation_type: "HeatingStart",
-                step: 1,
-                duration_sec: Some(3),
-                extension: r#"{"target":"warm"}"#.to_string(),
-            },
-        ],
+fn resolve_artifact_path(artifact_root: &Path, path_ref: &str) -> Result<PathBuf> {
+    let path_ref = path_ref.trim();
+    if path_ref.is_empty() {
+        return Err(anyhow!("empty artifact pathRef"));
+    }
+    let relative = Path::new(path_ref);
+    if relative.is_absolute() {
+        return Err(anyhow!("absolute artifact pathRef is not allowed"));
+    }
+    let mut safe_relative = PathBuf::new();
+    for component in relative.components() {
+        match component {
+            std::path::Component::Normal(segment) => {
+                let segment = segment
+                    .to_str()
+                    .ok_or_else(|| anyhow!("non-utf8 artifact pathRef segment"))?;
+                require_safe_segment(segment)?;
+                safe_relative.push(segment);
+            }
+            std::path::Component::CurDir => {}
+            _ => return Err(anyhow!("unsafe artifact pathRef component")),
+        }
+    }
+    if safe_relative.as_os_str().is_empty() {
+        return Err(anyhow!("artifact pathRef has no safe components"));
+    }
+    let root = normalize_path_lexically(artifact_root);
+    let output_path = normalize_path_lexically(&root.join(safe_relative));
+    if !output_path.starts_with(&root) {
+        return Err(anyhow!("reply artifact path escapes root"));
+    }
+    Ok(output_path)
+}
+
+#[derive(Serialize)]
+struct RuntimeTurnRequest {
+    #[serde(rename = "callId")]
+    call_id: String,
+    #[serde(rename = "traceId")]
+    trace_id: String,
+    #[serde(rename = "turnId")]
+    turn_id: String,
+    #[serde(rename = "audioArtifact")]
+    #[serde(skip_serializing_if = "Option::is_none")]
+    audio_artifact: Option<RuntimeTurnAudioArtifact>,
+    #[serde(rename = "asrResultRef", skip_serializing_if = "Option::is_none")]
+    asr_result_ref: Option<String>,
+}
+
+#[derive(Serialize)]
+struct RuntimeTurnAudioArtifact {
+    #[serde(rename = "type")]
+    artifact_type: String,
+    #[serde(rename = "pathRef")]
+    path_ref: String,
+    format: String,
+    #[serde(rename = "sampleRate")]
+    sample_rate: u32,
+    channels: u32,
+    #[serde(rename = "durationMs")]
+    duration_ms: u64,
+    #[serde(rename = "byteSize")]
+    byte_size: u64,
+}
+
+#[derive(Deserialize)]
+struct RuntimeTurnCommonResult<T> {
+    code: i64,
+    msg: Option<String>,
+    data: Option<T>,
+    stage: Option<String>,
+    retryable: Option<bool>,
+}
+
+#[derive(Deserialize)]
+#[serde(rename_all = "camelCase")]
+struct RuntimeAsrStreamResp {
+    status: Option<String>,
+    asr_result_ref: Option<String>,
+    chunk_count: Option<u64>,
+    audio_bytes: Option<u64>,
+    asr_duration_ms: Option<u64>,
+    provider_alias: Option<String>,
+    text_len: Option<u64>,
+    fallback_reason: Option<String>,
+    fallback_stage: Option<String>,
+}
+
+#[derive(Default)]
+struct RuntimeTurnBridgeOutcome {
+    reply_audio_artifact: Option<RuntimeTurnReplyAudioArtifact>,
+    device_outputs: Vec<RuntimeTurnDeviceOutput>,
+}
+
+struct RuntimeTurnStreamOutcome {
+    reply_playback_mode: String,
+    audio_chunk_count: u64,
+    device_output_count: u64,
+}
+
+struct RuntimeTurnStreamState {
+    reply_playback_mode: String,
+    reply_state_seq: u64,
+    playback_started_sent: bool,
+    first_audio_frame_written: bool,
+    completed: bool,
+    audio_chunk_count: u64,
+    device_output_count: u64,
+    encoded_audio_buffer: Vec<u8>,
+    pcm_stream_decoder: Option<audio::PcmS16leStreamDecoder>,
+    pcm_stream_network_chunk_count: u64,
+}
+
+impl Default for RuntimeTurnStreamState {
+    fn default() -> Self {
+        Self {
+            reply_playback_mode: "full_tts_fallback".to_string(),
+            reply_state_seq: 0,
+            playback_started_sent: false,
+            first_audio_frame_written: false,
+            completed: false,
+            audio_chunk_count: 0,
+            device_output_count: 0,
+            encoded_audio_buffer: Vec::new(),
+            pcm_stream_decoder: None,
+            pcm_stream_network_chunk_count: 0,
+        }
     }
 }
 
-fn command_id(call_id: &str, sequence: u8) -> String {
-    format!("{call_id}-device-smoke-{sequence:03}")
+#[derive(Deserialize)]
+struct RuntimeTurnResponseData {
+    #[serde(rename = "turnId")]
+    turn_id: Option<String>,
+    #[serde(rename = "messageId")]
+    message_id: Option<String>,
+    #[serde(rename = "replyAudioArtifact")]
+    reply_audio_artifact: Option<RuntimeTurnReplyAudioArtifact>,
+    #[serde(rename = "deviceOutputs")]
+    device_outputs: Option<Vec<RuntimeTurnDeviceOutput>>,
+}
+
+#[derive(Deserialize)]
+struct RuntimeTurnStreamEvent {
+    #[serde(rename = "type", alias = "event")]
+    event_type: Option<String>,
+    #[serde(rename = "seq")]
+    seq: Option<u64>,
+    #[serde(rename = "replyPlaybackMode")]
+    reply_playback_mode: Option<String>,
+    state: Option<String>,
+    #[serde(rename = "audioChunk")]
+    audio_chunk: Option<RuntimeTurnStreamAudioChunk>,
+    activity: Option<RuntimeTurnStreamActivity>,
+    #[serde(rename = "deviceOutput")]
+    device_output: Option<RuntimeTurnDeviceOutput>,
+    error: Option<RuntimeTurnStreamError>,
+    completion: Option<RuntimeTurnStreamCompletion>,
+    diagnostics: Option<serde_json::Value>,
+}
+
+impl RuntimeTurnStreamEvent {
+    fn event_type(&self) -> Option<String> {
+        self.event_type.as_deref().map(|value| match value {
+            "reply_playback_mode_selected" => "reply_playback_mode_selected".to_string(),
+            "reply_state" => "reply_state".to_string(),
+            "reply_audio_chunk" => "reply_audio_chunk".to_string(),
+            "device_output" => "device_output".to_string(),
+            "turn_completed" => "turn_completed".to_string(),
+            "turn_failed" => "turn_failed".to_string(),
+            "turn_cancelled" => "turn_cancelled".to_string(),
+            "activity" => "activity".to_string(),
+            other => other.to_string(),
+        })
+    }
+}
+
+#[derive(Deserialize)]
+struct RuntimeTurnStreamAudioChunk {
+    #[serde(rename = "chunkSeq", alias = "seq")]
+    chunk_seq: Option<u64>,
+    format: Option<String>,
+    #[serde(rename = "sampleRate")]
+    sample_rate: Option<u32>,
+    channels: Option<u32>,
+    #[serde(rename = "payloadBase64", alias = "audioBase64")]
+    payload_base64: Option<String>,
+    last: Option<bool>,
+}
+
+#[derive(Deserialize)]
+struct RuntimeTurnStreamActivity {
+    #[serde(rename = "eventType", alias = "event")]
+    event_type: Option<String>,
+    stage: Option<String>,
+    #[serde(rename = "reasonCode")]
+    reason_code: Option<String>,
+}
+
+#[derive(Deserialize)]
+struct RuntimeTurnStreamError {
+    #[serde(rename = "reasonCode")]
+    reason_code: Option<String>,
+    stage: Option<String>,
+    retryable: Option<bool>,
+}
+
+#[derive(Deserialize)]
+struct RuntimeTurnStreamCompletion {
+    #[serde(rename = "messageId")]
+    message_id: Option<String>,
+}
+
+#[derive(Clone, Deserialize)]
+struct RuntimeTurnReplyAudioArtifact {
+    #[serde(rename = "type")]
+    artifact_type: Option<String>,
+    #[serde(rename = "pathRef")]
+    path_ref: Option<String>,
+    format: Option<String>,
+}
+
+#[derive(Clone, Deserialize)]
+struct RuntimeTurnDeviceOutput {
+    #[serde(rename = "commandId")]
+    command_id: Option<String>,
+    #[serde(rename = "commandCode")]
+    command_code: Option<String>,
+    params: Option<serde_json::Value>,
+}
+
+fn require_safe_segment(value: &str) -> Result<()> {
+    if value.is_empty()
+        || value.contains('/')
+        || value.contains("..")
+        || !value
+            .chars()
+            .all(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '_' | '-' | '.'))
+    {
+        return Err(anyhow!("unsafe path segment"));
+    }
+    Ok(())
+}
+
+fn normalize_path_lexically(path: &Path) -> PathBuf {
+    let mut normalized = PathBuf::new();
+    for component in path.components() {
+        match component {
+            std::path::Component::CurDir => {}
+            std::path::Component::ParentDir => {
+                normalized.pop();
+            }
+            _ => normalized.push(component.as_os_str()),
+        }
+    }
+    normalized
+}
+
+fn spawn_user_audio_frame_observer(
+    track: RemoteAudioTrack,
+    call_id: String,
+    trace_id: String,
+    participant_alias: String,
+    track_sid_alias: String,
+    simple_vad_enabled: bool,
+    simple_vad_config: SimpleVadConfig,
+    vad_enabled_gate: Arc<AtomicBool>,
+    turn_bridge_config: TurnBridgeConfig,
+    http: Client,
+    sink: Arc<BotAudioOutputSink>,
+) -> JoinHandle<()> {
+    tokio::spawn(async move {
+        let mut stream = NativeAudioStream::new(
+            track.rtc_track(),
+            USER_AUDIO_SAMPLE_RATE_HZ as i32,
+            i32::from(USER_AUDIO_NUM_CHANNELS),
+        );
+        let started_at = Instant::now();
+        let mut frame_count: u64 = 0;
+        let mut sample_count: u64 = 0;
+        let mut simple_vad = if simple_vad_enabled {
+            Some(SimpleVad::new(simple_vad_config))
+        } else {
+            None
+        };
+        let mut realtime_asr_upload: Option<RealtimeAsrUpload> = None;
+
+        while let Some(frame) = stream.next().await {
+            frame_count += 1;
+            sample_count += u64::from(frame.samples_per_channel) * u64::from(frame.num_channels);
+            let elapsed_ms = started_at.elapsed().as_millis() as u64;
+            if frame_count == 1 {
+                info!(
+                    call_id = %call_id,
+                    trace_id = %trace_id,
+                    participant_alias = %participant_alias,
+                    track_sid_alias = %track_sid_alias,
+                    sample_rate_hz = frame.sample_rate,
+                    num_channels = frame.num_channels,
+                    samples_per_channel = frame.samples_per_channel,
+                    first_frame_elapsed_ms = elapsed_ms,
+                    "runtime helper user_audio_frame_received"
+                );
+            } else if frame_count % 250 == 0 {
+                info!(
+                    call_id = %call_id,
+                    trace_id = %trace_id,
+                    participant_alias = %participant_alias,
+                    track_sid_alias = %track_sid_alias,
+                    frame_count,
+                    sample_count,
+                    observed_wall_ms = elapsed_ms,
+                    "runtime helper user_audio_frame_summary"
+                );
+            }
+
+            if let Some(vad) = simple_vad.as_mut() {
+                if vad_enabled_gate.load(Ordering::Acquire) {
+                    let was_in_speech = vad.in_speech;
+                    let turn = vad.observe_frame(
+                        &call_id,
+                        &trace_id,
+                        &participant_alias,
+                        &track_sid_alias,
+                        frame_count,
+                        elapsed_ms,
+                        &frame,
+                    );
+                    let is_in_speech = vad.in_speech;
+
+                    if !was_in_speech && is_in_speech {
+                        let turn_id = format!("turn-{:04}", vad.turn_index);
+                        match RealtimeAsrUpload::start(
+                            http.clone(),
+                            turn_bridge_config.realtime_asr_config(),
+                            &call_id,
+                            &trace_id,
+                            &turn_id,
+                            &vad.speech_samples,
+                        ) {
+                            Ok(upload) => {
+                                info!(
+                                    call_id = %call_id,
+                                    trace_id = %trace_id,
+                                    turn_id = %turn_id,
+                                    "runtime helper asr_realtime_session_started"
+                                );
+                                realtime_asr_upload = Some(upload);
+                            }
+                            Err(error) if turn_bridge_config.asr_realtime_enabled => {
+                                warn!(
+                                    call_id = %call_id,
+                                    trace_id = %trace_id,
+                                    turn_id = %turn_id,
+                                    error = %safe_error(&error.to_string()),
+                                    "runtime helper asr_realtime_start_failed_fallback"
+                                );
+                            }
+                            Err(_) => {}
+                        }
+                    } else if was_in_speech {
+                        let push_failed = realtime_asr_upload
+                            .as_mut()
+                            .and_then(|upload| upload.push_48k_samples(frame.data.as_ref()).err());
+                        if let Some(error) = push_failed {
+                            warn!(
+                                call_id = %call_id,
+                                trace_id = %trace_id,
+                                error = %safe_error(&error.to_string()),
+                                "runtime helper asr_realtime_upload_failed_fallback"
+                            );
+                            if let Some(upload) = realtime_asr_upload.take() {
+                                tokio::spawn(async move {
+                                    upload.cancel("upload_backpressure").await;
+                                });
+                            }
+                        }
+                    }
+
+                    if let Some(turn) = turn {
+                        let realtime_asr_result_ref = match realtime_asr_upload.take() {
+                            Some(upload) => {
+                                finish_realtime_asr_upload(upload, &call_id, &trace_id, &turn).await
+                            }
+                            None => None,
+                        };
+                        handle_finished_turn(
+                            &http,
+                            &turn_bridge_config,
+                            &sink,
+                            &call_id,
+                            &trace_id,
+                            turn,
+                            realtime_asr_result_ref,
+                        )
+                        .await;
+                    } else if was_in_speech && !is_in_speech {
+                        if let Some(upload) = realtime_asr_upload.take() {
+                            tokio::spawn(async move {
+                                upload.cancel("speech_too_short").await;
+                            });
+                        }
+                    }
+                } else {
+                    if let Some(upload) = realtime_asr_upload.take() {
+                        tokio::spawn(async move {
+                            upload.cancel("vad_disabled").await;
+                        });
+                    }
+                    vad.observe_disabled_frame(
+                        &call_id,
+                        &trace_id,
+                        &participant_alias,
+                        &track_sid_alias,
+                        frame_count,
+                        elapsed_ms,
+                    );
+                }
+            }
+        }
+
+        if let Some(vad) = simple_vad.as_mut() {
+            if let Some(turn) = vad.finish_stream(
+                &call_id,
+                &trace_id,
+                &participant_alias,
+                &track_sid_alias,
+                started_at.elapsed().as_millis() as u64,
+            ) {
+                let realtime_asr_result_ref = match realtime_asr_upload.take() {
+                    Some(upload) => {
+                        finish_realtime_asr_upload(upload, &call_id, &trace_id, &turn).await
+                    }
+                    None => None,
+                };
+                handle_finished_turn(
+                    &http,
+                    &turn_bridge_config,
+                    &sink,
+                    &call_id,
+                    &trace_id,
+                    turn,
+                    realtime_asr_result_ref,
+                )
+                .await;
+            }
+        }
+        if let Some(upload) = realtime_asr_upload.take() {
+            upload.cancel("stream_end").await;
+        }
+
+        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"
+        );
+    })
+}
+
+async fn finish_realtime_asr_upload(
+    upload: RealtimeAsrUpload,
+    call_id: &str,
+    trace_id: &str,
+    turn: &FinishedSpeechTurn,
+) -> Option<String> {
+    match upload.finish(turn.duration_ms, &turn.end_reason).await {
+        Ok(RealtimeAsrOutcome {
+            status,
+            asr_result_ref,
+            provider_alias,
+            partial_count,
+            fallback_reason,
+            fallback_stage,
+            chunk_count,
+            audio_bytes,
+            wall_ms,
+        }) => {
+            info!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                turn_id = %turn.turn_id,
+                status = %status,
+                provider_alias = ?provider_alias,
+                partial_count,
+                chunk_count,
+                audio_bytes,
+                wall_ms,
+                asr_result_ref_present = asr_result_ref.is_some(),
+                fallback_reason = ?fallback_reason,
+                fallback_stage = ?fallback_stage,
+                "runtime helper asr_realtime_finished"
+            );
+            if status == "final" {
+                asr_result_ref
+            } else {
+                None
+            }
+        }
+        Err(error) => {
+            warn!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                turn_id = %turn.turn_id,
+                error = %safe_error(&error.to_string()),
+                "runtime helper asr_realtime_failed_fallback"
+            );
+            None
+        }
+    }
+}
+
+#[derive(Clone)]
+struct SimpleVadConfig {
+    rms_threshold: f64,
+    peak_threshold: f64,
+    start_frames: u32,
+    end_silence_ms: u64,
+    min_speech_ms: u64,
+    max_turn_ms: u64,
+    initial_ignore_ms: u64,
+}
+
+impl SimpleVadConfig {
+    fn from_env() -> Self {
+        Self {
+            rms_threshold: f64_env("CV_VAD_RMS_THRESHOLD", 0.012),
+            peak_threshold: f64_env("CV_VAD_PEAK_THRESHOLD", 0.08),
+            start_frames: u32_env("CV_VAD_START_FRAMES", 5).max(1),
+            end_silence_ms: u64_env("CV_VAD_END_SILENCE_MS", 400).max(100),
+            min_speech_ms: u64_env("CV_VAD_MIN_SPEECH_MS", 250).max(1),
+            max_turn_ms: u64_env("CV_VAD_MAX_TURN_MS", 10_000).max(1_000),
+            initial_ignore_ms: u64_env("CV_VAD_INITIAL_IGNORE_MS", 500),
+        }
+    }
+
+    fn end_silence_frames(&self, frame_duration_ms: u64) -> u32 {
+        let frame_duration_ms = frame_duration_ms.max(1);
+        self.end_silence_ms.div_ceil(frame_duration_ms) as u32
+    }
+}
+
+struct SimpleVad {
+    config: SimpleVadConfig,
+    turn_index: u64,
+    in_speech: bool,
+    voiced_run_frames: u32,
+    silence_run_frames: u32,
+    speech_start_elapsed_ms: u64,
+    speech_frame_count: u64,
+    speech_sample_count: u64,
+    speech_rms_sum: f64,
+    speech_peak: f64,
+    ignored_before_enabled_frames: u64,
+    pre_speech_frames: Vec<Vec<i16>>,
+    speech_samples: Vec<i16>,
+}
+
+struct FinishedSpeechTurn {
+    turn_index: u64,
+    turn_id: String,
+    samples: Vec<i16>,
+    duration_ms: u64,
+    frame_count: u64,
+    sample_count: u64,
+    end_reason: String,
+}
+
+impl SimpleVad {
+    fn new(config: SimpleVadConfig) -> Self {
+        Self {
+            config,
+            turn_index: 0,
+            in_speech: false,
+            voiced_run_frames: 0,
+            silence_run_frames: 0,
+            speech_start_elapsed_ms: 0,
+            speech_frame_count: 0,
+            speech_sample_count: 0,
+            speech_rms_sum: 0.0,
+            speech_peak: 0.0,
+            ignored_before_enabled_frames: 0,
+            pre_speech_frames: Vec::new(),
+            speech_samples: Vec::new(),
+        }
+    }
+
+    fn observe_disabled_frame(
+        &mut self,
+        call_id: &str,
+        trace_id: &str,
+        participant_alias: &str,
+        track_sid_alias: &str,
+        frame_count: u64,
+        elapsed_ms: u64,
+    ) {
+        self.reset_current_turn();
+        self.ignored_before_enabled_frames = self.ignored_before_enabled_frames.saturating_add(1);
+        if self.ignored_before_enabled_frames == 1 || self.ignored_before_enabled_frames % 250 == 0
+        {
+            info!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                participant_alias = %participant_alias,
+                track_sid_alias = %track_sid_alias,
+                frame_count,
+                ignored_frame_count = self.ignored_before_enabled_frames,
+                elapsed_ms,
+                "runtime helper vad_ignored_before_enabled"
+            );
+            emit_activity(
+                call_id,
+                trace_id,
+                None,
+                "ignored_user_audio_before_vad_enabled",
+                "ok",
+                None,
+                None,
+                json!({
+                    "frameCount": frame_count,
+                    "ignoredFrameCount": self.ignored_before_enabled_frames,
+                    "elapsedMs": elapsed_ms,
+                }),
+            );
+        }
+    }
+
+    fn observe_frame(
+        &mut self,
+        call_id: &str,
+        trace_id: &str,
+        participant_alias: &str,
+        track_sid_alias: &str,
+        frame_count: u64,
+        elapsed_ms: u64,
+        frame: &AudioFrame<'_>,
+    ) -> Option<FinishedSpeechTurn> {
+        let frame_duration_ms = frame_duration_ms(frame);
+        let (rms, peak) = pcm_energy_stats(frame.data.as_ref());
+        if elapsed_ms < self.config.initial_ignore_ms {
+            return None;
+        }
+
+        // Align with cb-sdk's energy-based segmentation: peak is diagnostic only,
+        // otherwise isolated spikes can keep a turn open until max_turn_ms.
+        let voiced = rms >= self.config.rms_threshold;
+        if !self.in_speech {
+            self.remember_pre_speech_frame(frame);
+        }
+        if voiced {
+            self.voiced_run_frames = self.voiced_run_frames.saturating_add(1);
+            self.silence_run_frames = 0;
+        } else {
+            self.voiced_run_frames = 0;
+            self.silence_run_frames = self.silence_run_frames.saturating_add(1);
+        }
+
+        if !self.in_speech {
+            if self.voiced_run_frames >= self.config.start_frames {
+                self.turn_index += 1;
+                self.in_speech = true;
+                self.speech_start_elapsed_ms = elapsed_ms
+                    .saturating_sub(u64::from(self.config.start_frames) * frame_duration_ms);
+                self.speech_frame_count = u64::from(self.config.start_frames);
+                self.speech_sample_count = u64::from(frame.samples_per_channel)
+                    * u64::from(frame.num_channels)
+                    * u64::from(self.config.start_frames);
+                self.speech_rms_sum = rms * f64::from(self.config.start_frames);
+                self.speech_peak = peak;
+                self.speech_samples = self
+                    .pre_speech_frames
+                    .iter()
+                    .flat_map(|samples| samples.iter().copied())
+                    .collect();
+                let turn_id = format!("turn-{:04}", self.turn_index);
+                info!(
+                    call_id = %call_id,
+                    trace_id = %trace_id,
+                    participant_alias = %participant_alias,
+                    track_sid_alias = %track_sid_alias,
+                    turn_index = self.turn_index,
+                    frame_count,
+                    speech_start_elapsed_ms = self.speech_start_elapsed_ms,
+                    rms = round4(rms),
+                    peak = round4(peak),
+                    "runtime helper vad_speech_start"
+                );
+                emit_activity(
+                    call_id,
+                    trace_id,
+                    Some(&turn_id),
+                    "vad_speech_start",
+                    "ok",
+                    None,
+                    None,
+                    json!({
+                        "turnIndex": self.turn_index,
+                        "frameCount": frame_count,
+                        "speechStartElapsedMs": self.speech_start_elapsed_ms,
+                        "rms": round4(rms),
+                        "peak": round4(peak),
+                    }),
+                );
+            }
+            return None;
+        }
+
+        self.speech_samples.extend_from_slice(frame.data.as_ref());
+        self.speech_frame_count += 1;
+        self.speech_sample_count +=
+            u64::from(frame.samples_per_channel) * u64::from(frame.num_channels);
+        self.speech_rms_sum += rms;
+        self.speech_peak = self.speech_peak.max(peak);
+
+        let speech_duration_ms = elapsed_ms.saturating_sub(self.speech_start_elapsed_ms);
+        if speech_duration_ms >= self.config.max_turn_ms {
+            return self.finish_turn(
+                call_id,
+                trace_id,
+                participant_alias,
+                track_sid_alias,
+                elapsed_ms,
+                "max_turn_ms",
+            );
+        }
+
+        if !voiced && self.silence_run_frames >= self.config.end_silence_frames(frame_duration_ms) {
+            return self.finish_turn(
+                call_id,
+                trace_id,
+                participant_alias,
+                track_sid_alias,
+                elapsed_ms,
+                "silence",
+            );
+        }
+        None
+    }
+
+    fn finish_stream(
+        &mut self,
+        call_id: &str,
+        trace_id: &str,
+        participant_alias: &str,
+        track_sid_alias: &str,
+        elapsed_ms: u64,
+    ) -> Option<FinishedSpeechTurn> {
+        if self.in_speech {
+            return self.finish_turn(
+                call_id,
+                trace_id,
+                participant_alias,
+                track_sid_alias,
+                elapsed_ms,
+                "stream_end",
+            );
+        } else if self.turn_index == 0 {
+            info!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                participant_alias = %participant_alias,
+                track_sid_alias = %track_sid_alias,
+                observed_wall_ms = elapsed_ms,
+                "runtime helper vad_no_speech_summary"
+            );
+        }
+        None
+    }
+
+    fn finish_turn(
+        &mut self,
+        call_id: &str,
+        trace_id: &str,
+        participant_alias: &str,
+        track_sid_alias: &str,
+        elapsed_ms: u64,
+        end_reason: &str,
+    ) -> Option<FinishedSpeechTurn> {
+        let speech_duration_ms = elapsed_ms.saturating_sub(self.speech_start_elapsed_ms);
+        let event_name = if speech_duration_ms < self.config.min_speech_ms {
+            "runtime helper vad_speech_too_short"
+        } else {
+            "runtime helper vad_speech_end"
+        };
+        let avg_rms = if self.speech_frame_count == 0 {
+            0.0
+        } else {
+            self.speech_rms_sum / self.speech_frame_count as f64
+        };
+
+        info!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            participant_alias = %participant_alias,
+            track_sid_alias = %track_sid_alias,
+            turn_index = self.turn_index,
+            speech_start_elapsed_ms = self.speech_start_elapsed_ms,
+            speech_end_elapsed_ms = elapsed_ms,
+            speech_duration_ms,
+            speech_frame_count = self.speech_frame_count,
+            speech_sample_count = self.speech_sample_count,
+            avg_rms = round4(avg_rms),
+            peak = round4(self.speech_peak),
+            end_reason = %end_reason,
+            "{}", event_name
+        );
+
+        let finished_turn = if speech_duration_ms >= self.config.min_speech_ms {
+            let turn_id = format!("turn-{:04}", self.turn_index);
+            emit_activity(
+                call_id,
+                trace_id,
+                Some(&turn_id),
+                "vad_speech_end",
+                "ok",
+                None,
+                None,
+                json!({
+                    "turnIndex": self.turn_index,
+                    "speechStartElapsedMs": self.speech_start_elapsed_ms,
+                    "speechEndElapsedMs": elapsed_ms,
+                    "speechDurationMs": speech_duration_ms,
+                    "speechFrameCount": self.speech_frame_count,
+                    "speechSampleCount": self.speech_sample_count,
+                    "avgRms": round4(avg_rms),
+                    "peak": round4(self.speech_peak),
+                    "endReason": end_reason,
+                }),
+            );
+            Some(FinishedSpeechTurn {
+                turn_index: self.turn_index,
+                turn_id,
+                samples: std::mem::take(&mut self.speech_samples),
+                duration_ms: speech_duration_ms,
+                frame_count: self.speech_frame_count,
+                sample_count: self.speech_sample_count,
+                end_reason: end_reason.to_string(),
+            })
+        } else {
+            None
+        };
+        self.reset_current_turn();
+        finished_turn
+    }
+
+    fn reset_current_turn(&mut self) {
+        self.in_speech = false;
+        self.voiced_run_frames = 0;
+        self.silence_run_frames = 0;
+        self.speech_start_elapsed_ms = 0;
+        self.speech_frame_count = 0;
+        self.speech_sample_count = 0;
+        self.speech_rms_sum = 0.0;
+        self.speech_peak = 0.0;
+        self.speech_samples.clear();
+        self.pre_speech_frames.clear();
+    }
+
+    fn remember_pre_speech_frame(&mut self, frame: &AudioFrame<'_>) {
+        self.pre_speech_frames.push(frame.data.as_ref().to_vec());
+        let max_frames = self.config.start_frames as usize;
+        if self.pre_speech_frames.len() > max_frames {
+            let remove_count = self.pre_speech_frames.len() - max_frames;
+            self.pre_speech_frames.drain(0..remove_count);
+        }
+    }
+}
+
+fn frame_duration_ms(frame: &AudioFrame<'_>) -> u64 {
+    if frame.sample_rate == 0 {
+        return 10;
+    }
+    ((u64::from(frame.samples_per_channel) * 1000) / u64::from(frame.sample_rate)).max(1)
+}
+
+fn pcm_energy_stats(samples: &[i16]) -> (f64, f64) {
+    if samples.is_empty() {
+        return (0.0, 0.0);
+    }
+    let mut square_sum = 0.0;
+    let mut peak = 0.0;
+    for sample in samples {
+        let normalized = f64::from(*sample) / f64::from(i16::MAX);
+        square_sum += normalized * normalized;
+        let abs = normalized.abs();
+        if abs > peak {
+            peak = abs;
+        }
+    }
+    ((square_sum / samples.len() as f64).sqrt(), peak)
+}
+
+fn round4(value: f64) -> f64 {
+    (value * 10_000.0).round() / 10_000.0
 }
 
 struct BotAudioOutputSink {
     room: Arc<Room>,
     rtc_source: NativeAudioSource,
     track: LocalAudioTrack,
+    device_output_destination_identity: Option<String>,
+    profile: String,
+    sample_rate_hz: u32,
+    num_channels: u16,
 }
 
 impl BotAudioOutputSink {
@@ -339,9 +3106,13 @@
         room: Arc<Room>,
         room_name: &str,
         participant_identity: &str,
+        call_id: &str,
+        trace_id: &str,
         track_name: &str,
+        profile: String,
         sample_rate: u32,
         num_channels: u32,
+        device_output_destination_identity: Option<String>,
     ) -> Result<Self> {
         let rtc_source = NativeAudioSource::new(
             AudioSourceOptions::default(),
@@ -353,6 +3124,8 @@
             track_name,
             RtcAudioSource::Native(rtc_source.clone()),
         );
+        let room_alias = redact(room_name);
+        let participant_alias = redact(participant_identity);
 
         room.local_participant()
             .publish_track(
@@ -362,26 +3135,212 @@
             .await
             .map_err(|error| {
                 anyhow!(
-                    "failed to publish bot audio track in room {room_name} for participant {participant_identity}: {error}"
+                    "failed to publish bot audio track in room {room_alias} for participant {participant_alias}: {error}"
                 )
             })?;
+        let num_channels_u16 = u16::try_from(num_channels)
+            .map_err(|_| anyhow!("unsupported bot audio channel count {num_channels}"))?;
 
         info!(
-            room_id = %room_name,
-            participant_alias = %redact(participant_identity),
+            room_alias = %room_alias,
+            participant_alias = %participant_alias,
             track_name = %track_name,
+            bot_audio_profile = %profile,
             sample_rate,
             num_channels,
             "runtime helper published bot audio track"
+        );
+        emit_activity(
+            call_id,
+            trace_id,
+            None,
+            "bot_track_ready",
+            "ok",
+            None,
+            None,
+            json!({
+                "trackName": track_name,
+                "audioProfile": profile,
+                "sampleRate": sample_rate,
+                "numChannels": num_channels,
+            }),
         );
 
         Ok(Self {
             room,
             rtc_source,
             track,
+            device_output_destination_identity,
+            profile,
+            sample_rate_hz: sample_rate,
+            num_channels: num_channels_u16,
         })
     }
 
+    async fn publish_device_output(
+        &self,
+        call_id: &str,
+        trace_id: &str,
+        turn_id: &str,
+        output: &RuntimeTurnDeviceOutput,
+    ) -> Result<()> {
+        let command_id = output
+            .command_id
+            .as_deref()
+            .map(str::trim)
+            .filter(|value| !value.is_empty())
+            .ok_or_else(|| anyhow!("device output commandId missing"))?;
+        let command_code = output
+            .command_code
+            .as_deref()
+            .map(str::trim)
+            .filter(|value| !value.is_empty())
+            .ok_or_else(|| anyhow!("device output commandCode missing"))?;
+        let payload = serde_json::json!({
+            "type": "device_output",
+            "schemaVersion": "1.0",
+            "callId": call_id,
+            "traceId": trace_id,
+            "turnId": turn_id,
+            "commandId": command_id,
+            "commandCode": command_code,
+            "params": output.params.clone().unwrap_or_else(|| serde_json::json!({})),
+            "source": {
+                "kind": "voice_command"
+            },
+        });
+        let destinations = self
+            .device_output_destination_identity
+            .as_ref()
+            .filter(|value| !value.trim().is_empty())
+            .map(|value| vec![ParticipantIdentity(value.trim().to_string())])
+            .unwrap_or_default();
+        let destination_count = destinations.len();
+        self.room
+            .local_participant()
+            .publish_data(DataPacket {
+                payload: serde_json::to_vec(&payload)
+                    .context("failed to encode device output data message")?,
+                topic: Some("device_output".to_string()),
+                reliable: true,
+                destination_identities: destinations,
+            })
+            .await
+            .map_err(|error| anyhow!("failed to publish device output data message: {error}"))?;
+        info!(
+            call_id = %call_id,
+            trace_id = %trace_id,
+            turn_id = %turn_id,
+            command_id_alias = %redact(command_id),
+            command_code = %command_code,
+            destination_count,
+            "runtime helper device_output_sent"
+        );
+        Ok(())
+    }
+
+    async fn publish_reply_state(
+        &self,
+        call_id: &str,
+        trace_id: &str,
+        turn_id: &str,
+        reply_playback_mode: &str,
+        state: &str,
+        seq: u64,
+    ) -> Result<()> {
+        emit_activity(
+            call_id,
+            trace_id,
+            Some(turn_id),
+            "reply_state_send_started",
+            "ok",
+            None,
+            None,
+            json!({
+                "state": state,
+                "seq": seq,
+                "replyPlaybackMode": reply_playback_mode,
+            }),
+        );
+        let payload = serde_json::json!({
+            "type": "reply_state",
+            "schemaVersion": "1.0",
+            "callId": call_id,
+            "traceId": trace_id,
+            "turnId": turn_id,
+            "replyPlaybackMode": reply_playback_mode,
+            "state": state,
+            "seq": seq,
+            "tsMs": current_time_millis(),
+        });
+        let destinations = self
+            .device_output_destination_identity
+            .as_ref()
+            .filter(|value| !value.trim().is_empty())
+            .map(|value| vec![ParticipantIdentity(value.trim().to_string())])
+            .unwrap_or_default();
+        let destination_count = destinations.len();
+        let publish_result = self
+            .room
+            .local_participant()
+            .publish_data(DataPacket {
+                payload: serde_json::to_vec(&payload)
+                    .context("failed to encode reply state data message")?,
+                topic: Some("combrabo_voice.reply_state".to_string()),
+                reliable: true,
+                destination_identities: destinations,
+            })
+            .await
+            .map_err(|error| anyhow!("failed to publish reply state data message: {error}"));
+        match publish_result {
+            Ok(_) => {
+                emit_activity(
+                    call_id,
+                    trace_id,
+                    Some(turn_id),
+                    "reply_state_send_finished",
+                    "ok",
+                    None,
+                    None,
+                    json!({
+                        "state": state,
+                        "seq": seq,
+                        "replyPlaybackMode": reply_playback_mode,
+                        "destinationCount": destination_count,
+                    }),
+                );
+                info!(
+                    call_id = %call_id,
+                    trace_id = %trace_id,
+                    turn_id = %turn_id,
+                    state = %state,
+                    seq,
+                    reply_playback_mode = %reply_playback_mode,
+                    destination_count,
+                    "runtime helper reply_state_sent"
+                );
+                Ok(())
+            }
+            Err(error) => {
+                emit_activity(
+                    call_id,
+                    trace_id,
+                    Some(turn_id),
+                    "reply_state_send_failed",
+                    "failed",
+                    Some("REPLY_STATE_SEND_FAILED"),
+                    Some(true),
+                    json!({
+                        "state": state,
+                        "seq": seq,
+                        "replyPlaybackMode": reply_playback_mode,
+                    }),
+                );
+                Err(error)
+            }
+        }
+    }
+
     async fn write_pcm_frame(&self, frame: &audio::PcmFrame) -> Result<()> {
         let audio_frame = AudioFrame {
             data: frame.data.as_slice().into(),

--
Gitblit v1.9.3