| | |
| | | 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, |
| | |
| | | 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(value) = runtime.asr_streaming_enabled { |
| | | command.env("CV_RUNTIME_ASR_STREAM_ENABLED", value.to_string()); |
| | | } |
| | | if let Some(value) = &runtime.asr_stream_url { |
| | | command.env("CV_RUNTIME_ASR_STREAM_URL", value); |
| | | } |
| | | if let Some(value) = runtime.asr_realtime_enabled { |
| | | command.env("CV_RUNTIME_ASR_REALTIME_ENABLED", value.to_string()); |
| | | } |
| | | if let Some(value) = &runtime.asr_realtime_url { |
| | | command.env("CV_RUNTIME_ASR_REALTIME_URL", value); |
| | | } |
| | | if let Some(value) = runtime.asr_realtime_chunk_duration_ms { |
| | | command.env( |
| | | "CV_RUNTIME_ASR_REALTIME_CHUNK_DURATION_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); |
| | |
| | | .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 { |
| | |
| | | struct AuthProfile { |
| | | helper_auth_token: String, |
| | | turn_bridge_token: String, |
| | | event_callback_token: String, |
| | | } |
| | | |
| | | impl AuthProfile { |
| | |
| | | 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 { |
| | |
| | | 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, |
| | |
| | | runtime_session_nonce: String, |
| | | livekit: LiveKitStart, |
| | | turn_bridge: Option<TurnBridgeStart>, |
| | | event_callback: Option<EventCallbackStart>, |
| | | auth_profile: Option<String>, |
| | | audio: Option<AudioStart>, |
| | | runtime: Option<RuntimeStart>, |
| | |
| | | |
| | | #[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>, |
| | |
| | | #[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>, |
| | | asr_streaming_enabled: Option<bool>, |
| | | asr_stream_url: Option<String>, |
| | | asr_realtime_enabled: Option<bool>, |
| | | asr_realtime_url: Option<String>, |
| | | asr_realtime_chunk_duration_ms: Option<u64>, |
| | | } |
| | | |
| | | #[derive(Default, Deserialize)] |