cai
2026-07-09 57753e1b1e169369bdc83f3945ce5246c844b316
feat: callback runtime events
3 files modified
250 ■■■■ changed files
README.md 6 ●●●● patch | view | raw | blame | history
src/main.rs 8 ●●●●● patch | view | raw | blame | history
src/service.rs 236 ●●●● patch | view | raw | blame | history
README.md
@@ -235,7 +235,7 @@
export CV_ENABLE_SIMPLE_VAD=true
```
首版 VAD 不引入外部模型,只基于上行 PCM frame 的 RMS / peak 做保守阈值判断。配置了 Java runtime turn bridge 后,helper 会在有效 `vad_speech_end` 后把本轮 PCM 写成 `user.wav` artifact,并调用 Java 内部 bridge;helper 仍不做 ASR / LLM / TTS / 消息落库。
首版 VAD 不引入外部模型,只基于上行 PCM frame 的 RMS 做保守阈值判断,peak 只作为诊断字段记录。这个口径对齐 cb-sdk 的能量阈值分段思路,避免单个尖峰噪声反复打断静音窗口,使用户 turn 拖到 `max_turn_ms`。配置了 Java runtime turn bridge 后,helper 会在有效 `vad_speech_end` 后把本轮 PCM 写成 `user.wav` artifact,并调用 Java 内部 bridge;helper 仍不做 ASR / LLM / TTS / 消息落库。
可调参数:
@@ -243,8 +243,8 @@
export CV_VAD_RMS_THRESHOLD=0.012
export CV_VAD_PEAK_THRESHOLD=0.08
export CV_VAD_START_FRAMES=5
export CV_VAD_END_SILENCE_MS=700
export CV_VAD_MIN_SPEECH_MS=300
export CV_VAD_END_SILENCE_MS=400
export CV_VAD_MIN_SPEECH_MS=250
export CV_VAD_MAX_TURN_MS=10000
export CV_VAD_INITIAL_IGNORE_MS=500
export CV_VAD_GATE_UNTIL_GREETING_DONE=true
src/main.rs
@@ -2258,8 +2258,8 @@
            rms_threshold: f64_env("CV_VAD_RMS_THRESHOLD", 0.012),
            peak_threshold: f64_env("CV_VAD_PEAK_THRESHOLD", 0.08),
            start_frames: u32_env("CV_VAD_START_FRAMES", 5).max(1),
            end_silence_ms: u64_env("CV_VAD_END_SILENCE_MS", 700).max(100),
            min_speech_ms: u64_env("CV_VAD_MIN_SPEECH_MS", 300).max(1),
            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),
        }
@@ -2372,7 +2372,9 @@
            return None;
        }
        let voiced = rms >= self.config.rms_threshold || peak >= self.config.peak_threshold;
        // 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);
        }
src/service.rs
@@ -219,6 +219,12 @@
        trace_id: trace_id.clone(),
        runtime_session_nonce: start_req.runtime_session_nonce.clone(),
        runtime_session_id,
        event_callback_url: start_req
            .event_callback
            .as_ref()
            .and_then(|value| value.url.clone())
            .filter(|value| !value.trim().is_empty()),
        event_callback_token: Some(profile.event_callback_token.clone()),
        status: "STARTING".to_string(),
        bot_participant_joined: false,
        bot_track_ready: false,
@@ -374,6 +380,30 @@
        if let Some(value) = &runtime.turn_artifact_dir {
            command.env("CV_RUNTIME_TURN_ARTIFACT_DIR", value);
        }
        if let Some(value) = runtime.vad_rms_threshold {
            command.env("CV_VAD_RMS_THRESHOLD", value.to_string());
        }
        if let Some(value) = runtime.vad_peak_threshold {
            command.env("CV_VAD_PEAK_THRESHOLD", value.to_string());
        }
        if let Some(value) = runtime.vad_start_frames {
            command.env("CV_VAD_START_FRAMES", value.to_string());
        }
        if let Some(value) = runtime.vad_end_silence_ms {
            command.env("CV_VAD_END_SILENCE_MS", value.to_string());
        }
        if let Some(value) = runtime.vad_min_speech_ms {
            command.env("CV_VAD_MIN_SPEECH_MS", value.to_string());
        }
        if let Some(value) = runtime.vad_max_turn_ms {
            command.env("CV_VAD_MAX_TURN_MS", value.to_string());
        }
        if let Some(value) = runtime.vad_initial_ignore_ms {
            command.env("CV_VAD_INITIAL_IGNORE_MS", value.to_string());
        }
        if let Some(value) = runtime.vad_post_greeting_delay_ms {
            command.env("CV_VAD_POST_GREETING_DELAY_MS", value.to_string());
        }
        if let Some(true) = runtime.audio_debug_dump_enabled {
            if let Some(value) = &config.audio_debug_dump_dir {
                command.env("CV_AUDIO_DEBUG_DUMP_DIR", value);
@@ -434,34 +464,159 @@
        .and_then(Value::as_str)
        .unwrap_or("worker_event");
    let result = value.get("result").and_then(Value::as_str).unwrap_or("ok");
    let Ok(mut sessions) = state.sessions.lock() else {
        return;
    let callback = {
        let Ok(mut sessions) = state.sessions.lock() else {
            return;
        };
        let Some(session) = sessions.get_mut(call_id) else {
            return;
        };
        session.last_event_type = Some(event_name.to_string());
        session.last_event_at = now_millis();
        let callback = build_event_callback_dispatch(session, value);
        if result != "ok" {
            session.status = "FAILED".to_string();
            session.bot_participant_joined = false;
            session.bot_track_ready = false;
        } else {
            match event_name {
                "helper_worker_started" => {
                    session.status = "STARTING".to_string();
                }
                "bot_participant_joined" => {
                    session.bot_participant_joined = true;
                }
                "bot_track_ready" => {
                    session.bot_participant_joined = true;
                    session.bot_track_ready = true;
                    session.status = "STARTED".to_string();
                }
                _ => {}
            }
        }
        callback
    };
    let Some(session) = sessions.get_mut(call_id) else {
        return;
    if let Some(callback) = callback {
        dispatch_event_callback(callback);
    }
}
fn build_event_callback_dispatch(
    session: &HelperSession,
    value: &Value,
) -> Option<EventCallbackDispatch> {
    let url = session
        .event_callback_url
        .as_ref()
        .filter(|value| !value.trim().is_empty())?
        .clone();
    let token = session
        .event_callback_token
        .as_ref()
        .filter(|value| !value.trim().is_empty())?
        .clone();
    Some(EventCallbackDispatch {
        url,
        token,
        call_id: session.call_id.clone(),
        trace_id: session.trace_id.clone(),
        runtime_session_nonce: session.runtime_session_nonce.clone(),
        payload: value.clone(),
    })
}
fn dispatch_event_callback(callback: EventCallbackDispatch) {
    thread::spawn(move || {
        let event_name = callback
            .payload
            .get("eventName")
            .and_then(Value::as_str)
            .unwrap_or("worker_event")
            .to_string();
        match post_event_callback(&callback) {
            Ok(status) if (200..300).contains(&status) => {
                info!(
                    call_id = %callback.call_id,
                    trace_id = ?callback.trace_id,
                    event_name = %event_name,
                    status = status,
                    "helper event callback delivered"
                );
            }
            Ok(status) => {
                warn!(
                    call_id = %callback.call_id,
                    trace_id = ?callback.trace_id,
                    event_name = %event_name,
                    status = status,
                    "helper event callback rejected"
                );
            }
            Err(error) => {
                warn!(
                    call_id = %callback.call_id,
                    trace_id = ?callback.trace_id,
                    event_name = %event_name,
                    error = %error,
                    "helper event callback failed"
                );
            }
        }
    });
}
fn post_event_callback(callback: &EventCallbackDispatch) -> Result<u16> {
    let (host, path) = parse_http_url(&callback.url)?;
    let body = serde_json::to_vec(&callback.payload)?;
    let mut stream = TcpStream::connect(&host)
        .with_context(|| format!("failed to connect callback host {host}"))?;
    stream.set_read_timeout(Some(Duration::from_secs(3)))?;
    stream.set_write_timeout(Some(Duration::from_secs(3)))?;
    let trace_header = callback.trace_id.clone().unwrap_or_default();
    let request = format!(
        "POST {path} HTTP/1.1\r\n\
         Host: {host}\r\n\
         Authorization: Bearer {}\r\n\
         X-CV-Runtime-Token: {}\r\n\
         X-CV-Call-Id: {}\r\n\
         X-CV-Trace-Id: {}\r\n\
         X-CV-Runtime-Session-Nonce: {}\r\n\
         Content-Type: application/json\r\n\
         Content-Length: {}\r\n\
         Connection: close\r\n\
         \r\n",
        callback.token,
        callback.token,
        callback.call_id,
        trace_header,
        callback.runtime_session_nonce,
        body.len()
    );
    stream.write_all(request.as_bytes())?;
    stream.write_all(&body)?;
    stream.flush()?;
    let mut response = String::new();
    stream.read_to_string(&mut response)?;
    response
        .lines()
        .next()
        .and_then(|line| line.split_whitespace().nth(1))
        .and_then(|value| value.parse::<u16>().ok())
        .ok_or_else(|| anyhow!("callback response status missing"))
}
fn parse_http_url(url: &str) -> Result<(String, String)> {
    let Some(rest) = url.strip_prefix("http://") else {
        return Err(anyhow!("only http callback url is supported"));
    };
    session.last_event_type = Some(event_name.to_string());
    session.last_event_at = now_millis();
    if result != "ok" {
        session.status = "FAILED".to_string();
        session.bot_participant_joined = false;
        session.bot_track_ready = false;
        return;
    let (host, path) = match rest.split_once('/') {
        Some((host, path)) => (host, format!("/{path}")),
        None => (rest, "/".to_string()),
    };
    if host.trim().is_empty() {
        return Err(anyhow!("callback host missing"));
    }
    match event_name {
        "helper_worker_started" => {
            session.status = "STARTING".to_string();
        }
        "bot_participant_joined" => {
            session.bot_participant_joined = true;
        }
        "bot_track_ready" => {
            session.bot_participant_joined = true;
            session.bot_track_ready = true;
            session.status = "STARTED".to_string();
        }
        _ => {}
    }
    Ok((host.to_string(), path))
}
enum WaitReadyResult {
@@ -692,6 +847,7 @@
struct AuthProfile {
    helper_auth_token: String,
    turn_bridge_token: String,
    event_callback_token: String,
}
impl AuthProfile {
@@ -707,11 +863,24 @@
        let turn_bridge_token = env::var(format!("CV_TURN_BRIDGE_TOKEN_{suffix}"))
            .or_else(|_| env::var("CV_RUNTIME_TURN_BRIDGE_TOKEN"))
            .ok()?;
        let event_callback_token = env::var(format!("CV_EVENT_CALLBACK_TOKEN_{suffix}"))
            .or_else(|_| env::var("CV_RUNTIME_EVENT_CALLBACK_TOKEN"))
            .unwrap_or_else(|_| turn_bridge_token.clone());
        Some(Self {
            helper_auth_token,
            turn_bridge_token,
            event_callback_token,
        })
    }
}
struct EventCallbackDispatch {
    url: String,
    token: String,
    call_id: String,
    trace_id: Option<String>,
    runtime_session_nonce: String,
    payload: Value,
}
struct ServiceState {
@@ -724,6 +893,8 @@
    trace_id: Option<String>,
    runtime_session_nonce: String,
    runtime_session_id: String,
    event_callback_url: Option<String>,
    event_callback_token: Option<String>,
    status: String,
    bot_participant_joined: bool,
    bot_track_ready: bool,
@@ -742,6 +913,7 @@
    runtime_session_nonce: String,
    livekit: LiveKitStart,
    turn_bridge: Option<TurnBridgeStart>,
    event_callback: Option<EventCallbackStart>,
    auth_profile: Option<String>,
    audio: Option<AudioStart>,
    runtime: Option<RuntimeStart>,
@@ -765,6 +937,12 @@
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct EventCallbackStart {
    url: Option<String>,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct AudioStart {
    first_audio_source: Option<String>,
    greeting_audio: Option<GreetingAudioStart>,
@@ -781,6 +959,14 @@
#[serde(rename_all = "camelCase")]
struct RuntimeStart {
    turn_artifact_dir: Option<String>,
    vad_rms_threshold: Option<f64>,
    vad_peak_threshold: Option<f64>,
    vad_start_frames: Option<u32>,
    vad_end_silence_ms: Option<u64>,
    vad_min_speech_ms: Option<u64>,
    vad_max_turn_ms: Option<u64>,
    vad_initial_ignore_ms: Option<u64>,
    vad_post_greeting_delay_ms: Option<u64>,
    audio_debug_dump_enabled: Option<bool>,
}