//! Provider-neutral helpers for the OpenAI Realtime wire protocol. //! //! The model and voice are deliberately supplied by discovery/configuration; //! this module contains no catalogue or fallback identifiers. Keeping the //! wire messages here makes the streaming transport usable by HTTP, desktop, //! and channel surfaces without duplicating protocol details. use crate::error::LlmError; use futures::{SinkExt, StreamExt}; use serde_json::{Value, json}; use tokio_tungstenite::tungstenite::Message; use tokio_util::sync::CancellationToken; #[derive(Debug, Clone, PartialEq, Eq)] pub struct RealtimeConfig { pub model: String, pub voice: Option, pub input_format: String, pub output_format: String, } impl RealtimeConfig { pub fn validate(&self) -> Result<(), LlmError> { for (name, value) in [ ("model", self.model.as_str()), ("input_format", self.input_format.as_str()), ("output_format", self.output_format.as_str()), ] { if value.trim().is_empty() { return Err(LlmError::InvalidRequest(format!( "realtime {name} is required" ))); } } if self.voice.as_deref().is_some_and(|v| v.trim().is_empty()) { return Err(LlmError::InvalidRequest( "realtime voice cannot be empty".into(), )); } Ok(()) } } /// Build the initial session update. Empty optional voice is omitted so the /// provider can apply its configured default. pub fn build_session_update( config: &RealtimeConfig, instructions: Option<&str>, ) -> Result { config.validate()?; let mut session = json!({ "modalities": ["text", "audio"], "input_audio_format": config.input_format, "output_audio_format": config.output_format, }); if let Some(voice) = config.voice.as_deref().filter(|v| !v.trim().is_empty()) { session["voice"] = json!(voice); } if let Some(text) = instructions.filter(|v| !v.trim().is_empty()) { session["instructions"] = json!(text); } Ok(json!({"type":"session.update", "session": session})) } pub fn build_audio_append(audio: &[u8]) -> Result { if audio.is_empty() { return Err(LlmError::InvalidRequest( "realtime audio cannot be empty".into(), )); } let encoded = base64::engine::general_purpose::STANDARD.encode(audio); Ok(json!({"type":"input_audio_buffer.append", "audio": encoded})) } pub fn build_response_create() -> Value { json!({"type":"response.create", "response": {"modalities":["audio","text"]}}) } /// Execute one OpenAI Realtime turn over a websocket. The endpoint is /// supplied by configuration so compatible providers can use the same /// transport. Audio is returned as the concatenated `response.audio.delta` /// payload; all other provider events are ignored by this low-level adapter. pub async fn round_trip( api_key: &str, endpoint: &str, config: &RealtimeConfig, audio: &[u8], instructions: Option<&str>, cancel: &CancellationToken, ) -> Result, LlmError> { config.validate()?; if api_key.trim().is_empty() || endpoint.trim().is_empty() { return Err(LlmError::InvalidRequest( "realtime credentials and endpoint are required".into(), )); } if audio.is_empty() { return Err(LlmError::InvalidRequest( "realtime audio cannot be empty".into(), )); } let separator = if endpoint.contains('?') { '&' } else { '?' }; let url = format!( "{endpoint}{separator}model={}", percent_encoding::utf8_percent_encode(&config.model, percent_encoding::NON_ALPHANUMERIC) ); let request = tokio_tungstenite::tungstenite::http::Request::builder() .uri(url) .header("Authorization", format!("Bearer {api_key}")) .header("OpenAI-Beta", "realtime=v1") .body(()) .map_err(|e| LlmError::InvalidRequest(e.to_string()))?; let (mut socket, _) = tokio::select! { _ = cancel.cancelled() => return Err(LlmError::Aborted { partial: None }), result = tokio_tungstenite::connect_async(request) => result.map_err(|e| LlmError::Network(e.to_string()))?, }; let send = |value: Value| Message::Text(value.to_string()); for value in [ build_session_update(config, instructions)?, build_audio_append(audio)?, build_response_create(), ] { tokio::select! { _ = cancel.cancelled() => return Err(LlmError::Aborted { partial: None }), result = socket.send(send(value)) => result.map_err(|e| LlmError::Network(e.to_string()))?, } } const MAX_AUDIO_BYTES: usize = 16 * 1024 * 1024; let mut output = Vec::new(); while let Some(message) = tokio::select! { _ = cancel.cancelled() => return Err(LlmError::Aborted { // Realtime audio is returned as bytes, while the shared LLM // abort contract stores textual assistant messages. The caller // still owns the already-emitted audio buffer and can preserve it. partial: None, }), message = socket.next() => message, } { let message = message.map_err(|e| LlmError::Network(e.to_string()))?; let Message::Text(text) = message else { continue; }; let event: Value = serde_json::from_str(&text).map_err(|e| LlmError::Parse(e.to_string()))?; match event.get("type").and_then(Value::as_str) { Some("response.audio.delta") => { let Some(delta) = event.get("delta").and_then(Value::as_str) else { continue; }; let bytes = base64::engine::general_purpose::STANDARD .decode(delta) .map_err(|e| LlmError::Parse(e.to_string()))?; if output.len().saturating_add(bytes.len()) > MAX_AUDIO_BYTES { return Err(LlmError::InvalidRequest( "provider audio exceeds 16 MiB".into(), )); } output.extend(bytes); } Some("error") => return Err(LlmError::InvalidRequest(event.to_string())), Some("response.done") => break, _ => {} } } if output.is_empty() { return Err(LlmError::Parse( "provider returned empty realtime audio".into(), )); } Ok(output) } use base64::Engine; #[cfg(test)] #[allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)] mod tests { use super::*; #[test] fn session_payload_requires_discovered_model_and_omits_voice_default() { let cfg = RealtimeConfig { model: "discovered-model".into(), voice: None, input_format: "pcm16".into(), output_format: "pcm16".into(), }; let body = build_session_update(&cfg, Some("be concise")).unwrap(); assert_eq!(body["type"], "session.update"); assert!(body["session"].get("voice").is_none()); assert_eq!(body["session"]["instructions"], "be concise"); } #[test] fn audio_append_is_base64_and_rejects_empty() { let body = build_audio_append(&[1, 2, 3]).unwrap(); assert_eq!(body["audio"], "AQID"); assert!(build_audio_append(&[]).is_err()); } }