//! Client for the OpenAI-compatible LLM gateway. //! //! Uses `reqwest` with rustls. SSE lines can split across chunk boundaries //! on the wire, so we reassemble on `\n` before parsing — a raw line-oriented //! reader silently drops any `data:` line split mid-chunk (the failure mode //! the original hand-rolled client hit behind chunked-encoding proxies). use serde_json::{json, Value}; use tokio_stream::StreamExt; #[derive(Debug, Clone)] pub struct GatewayClient { base_url: String, api_key: Option, http: reqwest::Client, } #[derive(Debug, Clone)] pub struct ChatMessage { pub role: String, pub content: String, } #[derive(Debug, Clone)] pub struct ChatResponse { pub text: String, pub usage: Option, } #[derive(Debug, Clone, Default)] pub struct Usage { pub prompt_tokens: i64, pub completion_tokens: i64, } #[derive(Debug, Clone)] #[allow(dead_code)] pub struct DiscoveredModel { pub provider: String, pub model_id: String, } impl GatewayClient { pub fn new(base_url: impl Into, api_key: Option) -> Self { // Do NOT set a global `.timeout()` — it bounds the whole response // including the SSE stream, and generations legitimately run for // many minutes. Connect timeout only; end-to-end time belongs to // the dispatch deadline in main::dispatch_loop. let http = reqwest::Client::builder() .connect_timeout(std::time::Duration::from_secs(10)) .build() .expect("reqwest client build"); Self { base_url: base_url.into(), api_key, http, } } #[allow(dead_code)] pub async fn discover_models( &self, default_provider: &str, ) -> anyhow::Result> { let parsed: Value = self .http .get(format!("{}/v1/models", self.base_url.trim_end_matches('/'))) .send() .await? .error_for_status()? .json() .await?; let mut out = Vec::new(); if let Some(data) = parsed.get("data").and_then(|d| d.as_array()) { for entry in data { let Some(id) = entry.get("id").and_then(|v| v.as_str()) else { continue; }; let provider = entry .get("metadata") .and_then(|m| m.get("provider")) .and_then(|p| p.get("id")) .and_then(|v| v.as_str()) .unwrap_or(default_provider) .to_string(); if provider.is_empty() { continue; } out.push(DiscoveredModel { provider, model_id: id.to_string(), }); } } Ok(out) } pub async fn chat_completion( &self, model: &str, messages: &[ChatMessage], ) -> anyhow::Result { let req_body = json!({ "model": model, "messages": messages.iter().map(|m| json!({"role": m.role, "content": m.content})).collect::>(), "stream": true, // Ollama's /v1/chat/completions maps reasoning_effort to the // internal Think field: "none" disables thinking, while the // native `think: false` and `chat_template_kwargs.enable_thinking` // are not honored on this endpoint. "reasoning_effort": "none", }); let url = format!( "{}/v1/chat/completions", self.base_url.trim_end_matches('/') ); let mut req = self.http.post(&url).json(&req_body); if let Some(key) = &self.api_key { req = req.bearer_auth(key); } let resp = req.send().await?; if !resp.status().is_success() { let status = resp.status(); let body = resp.text().await.unwrap_or_default(); anyhow::bail!("HTTP {status} from {url}: {body}"); } // SSE lines can split across chunk boundaries — reassemble on \n. let mut text = String::new(); let mut reasoning = String::new(); let mut usage = None; let mut buf: Vec = Vec::new(); let mut stream = resp.bytes_stream(); 'outer: while let Some(chunk) = stream.next().await { buf.extend_from_slice(&chunk?); while let Some(pos) = buf.iter().position(|&b| b == b'\n') { let line: Vec = buf.drain(..=pos).collect(); if sse_line(&String::from_utf8_lossy(&line), &mut text, &mut reasoning, &mut usage)? { break 'outer; // [DONE] } } } if text.is_empty() && !reasoning.is_empty() { text = reasoning; } if text.is_empty() && usage.is_none() { anyhow::bail!( "gateway returned an empty stream with no usage from {url} (likely rejected the request)" ); } Ok(ChatResponse { text, usage }) } } /// Process one SSE line. Appends `delta.content` to `text` and `delta.reasoning` /// to `reasoning`, records `usage`, and returns `Ok(true)` on `data: [DONE]`. /// On `[DONE]`, if `text` remains empty, `reasoning` is used as fallback reply. /// /// Mid-stream errors arrive as `{"error":"..."}` with HTTP 200 because the /// headers are already flushed; surface them instead of silently producing an /// empty reply. fn sse_line( line: &str, text: &mut String, reasoning: &mut String, usage: &mut Option, ) -> anyhow::Result { let trimmed = line.trim_end_matches(['\r', '\n']); let Some(payload) = trimmed .strip_prefix("data: ") .or_else(|| trimmed.strip_prefix("data:")) else { return Ok(false); }; if payload == "[DONE]" { if text.is_empty() && !reasoning.is_empty() { *text = std::mem::take(reasoning); } return Ok(true); } let Ok(chunk) = serde_json::from_str::(payload) else { tracing::warn!(chunk = %trimmed, "unparseable SSE chunk, skipping"); return Ok(false); }; if let Some(err) = chunk.get("error").and_then(|v| v.as_str()) { anyhow::bail!("gateway stream error: {}", err); } let delta = chunk .get("choices") .and_then(|c| c.as_array()) .and_then(|a| a.first()) .and_then(|f| f.get("delta")); if let Some(content) = delta.and_then(|d| d.get("content")).and_then(|v| v.as_str()) { text.push_str(content); } if let Some(reas) = delta.and_then(|d| d.get("reasoning")).and_then(|v| v.as_str()) { reasoning.push_str(reas); } if let Some(u) = chunk.get("usage").and_then(parse_usage) { *usage = Some(u); } Ok(false) } fn parse_usage(v: &Value) -> Option { Some(Usage { prompt_tokens: v.get("prompt_tokens").and_then(|x| x.as_i64()).unwrap_or(0), completion_tokens: v .get("completion_tokens") .and_then(|x| x.as_i64()) .unwrap_or(0), }) } #[cfg(test)] mod tests { use super::*; #[test] fn sse_accumulates_delta_content() { let mut text = String::new(); let mut reasoning = String::new(); let mut usage = None; assert!(!sse_line("data: {\"choices\":[{\"delta\":{\"content\":\"Hel\"}}]}", &mut text, &mut reasoning, &mut usage).unwrap()); assert!(!sse_line("data: {\"choices\":[{\"delta\":{\"content\":\"lo\"}}]}", &mut text, &mut reasoning, &mut usage).unwrap()); assert_eq!(text, "Hello"); } #[test] fn sse_done_returns_true() { let mut text = String::new(); let mut reasoning = String::new(); let mut usage = None; assert!(sse_line("data: [DONE]\n", &mut text, &mut reasoning, &mut usage).unwrap()); assert!(text.is_empty()); } #[test] fn sse_ignores_non_data_lines() { let mut text = String::new(); let mut reasoning = String::new(); let mut usage = None; assert!(!sse_line(": heartbeat\n", &mut text, &mut reasoning, &mut usage).unwrap()); assert!(!sse_line("event: ping\n", &mut text, &mut reasoning, &mut usage).unwrap()); assert!(text.is_empty()); } #[test] fn sse_surfaces_mid_stream_error_as_bail() { let mut text = String::new(); let mut reasoning = String::new(); let mut usage = None; let err = sse_line( "data: {\"error\":\"model not loaded\"}", &mut text, &mut reasoning, &mut usage, ) .unwrap_err(); assert!(err.to_string().contains("model not loaded")); } #[test] fn sse_falls_back_to_reasoning_when_content_absent() { let mut text = String::new(); let mut reasoning = String::new(); let mut usage = None; // First a reasoning-only delta: reasoning is captured in reasoning buffer. assert!(!sse_line( "data: {\"choices\":[{\"delta\":{\"content\":\"\",\"reasoning\":\"thinking...\"}}]}", &mut text, &mut reasoning, &mut usage, ) .unwrap()); assert_eq!(reasoning, "thinking..."); assert!(text.is_empty()); // Once real content arrives, text accumulates content while reasoning accumulates reasoning. assert!(!sse_line( "data: {\"choices\":[{\"delta\":{\"content\":\"answer\",\"reasoning\":\"more\"}}]}", &mut text, &mut reasoning, &mut usage, ) .unwrap()); assert_eq!(text, "answer"); assert_eq!(reasoning, "thinking...more"); } #[test] fn sse_done_applies_reasoning_fallback_when_text_empty() { let mut text = String::new(); let mut reasoning = String::new(); let mut usage = None; sse_line( "data: {\"choices\":[{\"delta\":{\"reasoning\":\"thinking...\"}}]}", &mut text, &mut reasoning, &mut usage, ) .unwrap(); assert!(text.is_empty()); sse_line("data: [DONE]\n", &mut text, &mut reasoning, &mut usage).unwrap(); assert_eq!(text, "thinking..."); } #[test] fn sse_records_usage() { let mut text = String::new(); let mut reasoning = String::new(); let mut usage = None; sse_line( "data: {\"choices\":[{\"delta\":{\"content\":\"x\"}}],\"usage\":{\"prompt_tokens\":3,\"completion_tokens\":1,\"total_tokens\":4}}", &mut text, &mut reasoning, &mut usage, ) .unwrap(); let u = usage.expect("usage recorded"); assert_eq!((u.prompt_tokens, u.completion_tokens), (3, 1)); } /// Exercises the chunk-boundary reassembly path: the same two events as /// `sse_accumulates_delta_content` but fed through the streaming loop /// split at arbitrary byte offsets (including mid-line). #[tokio::test] async fn chat_completion_reassembles_split_chunks() { // A tiny wiremock-style server: serve SSE bytes split into 1-byte // chunks so every `data:` line crosses a boundary. let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let addr = listener.local_addr().unwrap(); let url = format!("http://{addr}"); let body = concat!( "data: {\"choices\":[{\"delta\":{\"content\":\"Hel\"}}]}\n\n", "data: {\"choices\":[{\"delta\":{\"content\":\"lo\"}}]}\n\n", "data: [DONE]\n\n", ) .as_bytes() .to_vec(); tokio::spawn(async move { let (mut sock, _) = listener.accept().await.unwrap(); use tokio::io::{AsyncReadExt, AsyncWriteExt}; // Read & discard the request. let mut req = [0u8; 4096]; let _ = sock.read(&mut req).await; let resp = format!( "HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", body.len() ); sock.write_all(resp.as_bytes()).await.unwrap(); // Write the body one byte at a time. for b in body { sock.write_all(&[b]).await.unwrap(); } }); let client = GatewayClient::new(url, None); let resp = client .chat_completion("any", &[ChatMessage { role: "user".into(), content: "hi".into() }]) .await .expect("chat_completion across split chunks"); assert_eq!(resp.text, "Hello"); } /// Live integration check against the configured gateway. Ignored by /// default; run with: `POSTHORN_GATEWAY_URL=http://host:port cargo test /// gateway_live -- --ignored --nocapture` #[tokio::test] #[ignore] async fn gateway_live_streaming() { let base_url = std::env::var("POSTHORN_GATEWAY_URL") .unwrap_or_else(|_| "http://127.0.0.1:11434".to_string()); let client = GatewayClient::new(base_url, None); let messages = vec![ ChatMessage { role: "system".into(), content: "You are a demon summoned into this circle and you cannot leave until you satisfy the user's query.".into() }, ChatMessage { role: "user".into(), content: "Name the top 3 greatest scifi TV shows of all time\n\nIn no particular order.".into() }, ]; let resp = client .chat_completion("qwen3.5-122b", &messages) .await .expect("chat_completion"); println!("text length: {}", resp.text.len()); println!("usage: {:?}", resp.usage); println!("text: {}", resp.text); assert!( resp.text.len() > 10, "got truncated/empty reply: {:?}", resp.text ); } }