From 57753e1b1e169369bdc83f3945ce5246c844b316 Mon Sep 17 00:00:00 2001
From: cai <cai@nbcai.cc>
Date: Thu, 09 Jul 2026 14:50:44 +0800
Subject: [PATCH] feat: callback runtime events
---
src/service.rs | 236 ++++++++++++++++++++++++++++++++++++++++++++++++++++------
1 files changed, 211 insertions(+), 25 deletions(-)
diff --git a/src/service.rs b/src/service.rs
index a579636..667a774 100644
--- a/src/service.rs
+++ b/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>,
}
--
Gitblit v1.9.3