diff --git a/crates/tui/src/config.rs b/crates/tui/src/config.rs index e4abf074df..1bc2337570 100644 --- a/crates/tui/src/config.rs +++ b/crates/tui/src/config.rs @@ -2967,6 +2967,19 @@ pub struct VisionModelConfig { /// Base URL for the vision model API. Defaults to OpenAI. #[serde(default)] pub base_url: Option, + /// Total request budget in seconds, including retries and response body + /// consumption. `None` or `0` preserves the 120-second default; larger + /// values are capped at one hour. + #[serde(default)] + pub request_timeout_secs: Option, + /// Whether to request an SSE response. `None` preserves the existing + /// non-streaming request shape; callers must opt in with `Some(true)`. + #[serde(default)] + pub stream: Option, + /// Whether the existing transient-error retry policy is enabled. `None` + /// preserves the current default; `Some(false)` makes one attempt. + #[serde(default)] + pub retry_on_transient_errors: Option, } /// `[runtime_api]` table — knobs for the local HTTP/SSE daemon. diff --git a/crates/tui/src/vision/tools.rs b/crates/tui/src/vision/tools.rs index de06240c21..f02364a47c 100644 --- a/crates/tui/src/vision/tools.rs +++ b/crates/tui/src/vision/tools.rs @@ -5,15 +5,97 @@ use std::time::Duration; use async_trait::async_trait; use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64}; +use futures_util::StreamExt; use serde_json::{Value, json}; +use crate::client::{ERROR_BODY_MAX_BYTES, bounded_error_text}; use crate::config::VisionModelConfig; -use crate::llm_client::{LlmError, RetryConfig, sanitize_http_error_body, with_retry}; +use crate::llm_client::{ + LlmError, RetryConfig, extract_retry_after, sanitize_http_error_body, with_retry, +}; use crate::tools::spec::{ ToolCapability, ToolContext, ToolError, ToolResult, ToolSpec, required_str, }; const DEFAULT_VISION_MAX_OUTPUT_TOKENS: u32 = 4096; +const DEFAULT_VISION_REQUEST_TIMEOUT_SECS: u64 = 120; +const MAX_VISION_REQUEST_TIMEOUT_SECS: u64 = 3600; +const MAX_VISION_RESPONSE_BYTES: usize = 8 * 1024 * 1024; + +#[derive(Default)] +struct SseLineBuffer { + pending: Vec, + scan_from: usize, +} + +impl SseLineBuffer { + fn push(&mut self, chunk: &[u8]) -> Result, ToolError> { + if self.pending.len().saturating_add(chunk.len()) > MAX_VISION_RESPONSE_BYTES { + return Err(ToolError::execution_failed(format!( + "Vision SSE frame exceeded {MAX_VISION_RESPONSE_BYTES} bytes" + ))); + } + self.pending.extend_from_slice(chunk); + + let mut lines = Vec::new(); + let mut line_start = 0; + for (index, byte) in self.pending.iter().enumerate().skip(self.scan_from) { + if *byte == b'\n' { + lines.push(String::from_utf8_lossy(&self.pending[line_start..index]).into_owned()); + line_start = index + 1; + } + } + if line_start > 0 { + self.pending.drain(..line_start); + } + self.scan_from = self.pending.len(); + Ok(lines) + } + + fn finish(&mut self) -> Option { + (!self.pending.is_empty()) + .then(|| String::from_utf8_lossy(&std::mem::take(&mut self.pending)).into_owned()) + } +} + +#[derive(Default)] +struct VisionStreamState { + content: String, + saw_data_event: bool, + completed: bool, + truncated: bool, +} + +fn process_sse_line(line: &str, state: &mut VisionStreamState) -> bool { + let Some(data) = line.trim().strip_prefix("data:").map(str::trim) else { + return false; + }; + state.saw_data_event = true; + if data == "[DONE]" { + state.completed = true; + return true; + } + + let Ok(event) = serde_json::from_str::(data) else { + tracing::warn!("Vision SSE frame is not valid JSON; marking the stream truncated"); + state.truncated = true; + return false; + }; + if let Some(delta) = event + .pointer("/choices/0/delta/content") + .and_then(Value::as_str) + { + state.content.push_str(delta); + } + if let Some(reason) = event + .pointer("/choices/0/finish_reason") + .and_then(Value::as_str) + { + state.completed = true; + state.truncated |= reason == "length"; + } + false +} pub struct ImageAnalyzeTool { config: VisionModelConfig, @@ -24,7 +106,7 @@ impl ImageAnalyzeTool { #[must_use] pub fn new(config: VisionModelConfig) -> Self { let client = crate::tls::reqwest_client_builder() - .timeout(Duration::from_secs(120)) + .connect_timeout(Duration::from_secs(30)) .build() .expect("Failed to build HTTP client"); Self { config, client } @@ -147,9 +229,165 @@ impl ImageAnalyzeTool { "max_tokens" }; payload[token_limit_field] = json!(DEFAULT_VISION_MAX_OUTPUT_TOKENS); + if let Some(stream) = self.config.stream { + payload["stream"] = json!(stream); + } payload } + + fn request_timeout(&self) -> Duration { + let seconds = self + .config + .request_timeout_secs + .filter(|seconds| *seconds > 0) + .unwrap_or(DEFAULT_VISION_REQUEST_TIMEOUT_SECS) + .min(MAX_VISION_REQUEST_TIMEOUT_SECS); + Duration::from_secs(seconds) + } + + fn retry_config(&self) -> RetryConfig { + RetryConfig { + max_retries: 3, + initial_delay: 1.0, + max_delay: 30.0, + enabled: self.config.retry_on_transient_errors.unwrap_or(true), + ..Default::default() + } + } + + fn parse_non_streaming_body(&self, body: &[u8]) -> Result<(String, String, bool), ToolError> { + let json: Value = serde_json::from_slice(body).map_err(|error| { + ToolError::execution_failed(format!("Vision API returned invalid JSON: {error}")) + })?; + let content = json + .pointer("/choices/0/message/content") + .and_then(Value::as_str) + .unwrap_or_default() + .to_string(); + let model = json + .get("model") + .and_then(Value::as_str) + .unwrap_or(&self.config.model) + .to_string(); + let truncated = json + .pointer("/choices/0/finish_reason") + .and_then(Value::as_str) + == Some("length"); + Ok((content, model, truncated)) + } + + fn vision_result( + &self, + content: String, + model: String, + truncated: bool, + ) -> Result { + if content.trim().is_empty() { + return Err(ToolError::execution_failed( + "Vision API returned no usable content", + )); + } + + let mut result = json!({ + "analysis": content, + "model": model, + }); + if truncated { + result["truncated"] = json!(true); + } + ToolResult::json(&result).map_err(|error| { + ToolError::execution_failed(format!("Failed to serialize result: {error}")) + }) + } + + async fn read_bounded_body( + response: reqwest::Response, + deadline: tokio::time::Instant, + ) -> Result, ToolError> { + let mut body = Vec::new(); + let mut stream = response.bytes_stream(); + loop { + let chunk = tokio::time::timeout_at(deadline, stream.next()) + .await + .map_err(|_| ToolError::execution_failed("Vision API response timed out"))?; + let Some(chunk) = chunk else { break }; + let chunk = chunk.map_err(|error| { + ToolError::execution_failed(format!("Failed to read Vision API response: {error}")) + })?; + if body.len().saturating_add(chunk.len()) > MAX_VISION_RESPONSE_BYTES { + return Err(ToolError::execution_failed(format!( + "Vision API response exceeded {MAX_VISION_RESPONSE_BYTES} bytes" + ))); + } + body.extend_from_slice(&chunk); + } + Ok(body) + } + + async fn execute_non_streaming( + &self, + response: reqwest::Response, + deadline: tokio::time::Instant, + ) -> Result { + let body = Self::read_bounded_body(response, deadline).await?; + let (content, model, truncated) = self.parse_non_streaming_body(&body)?; + self.vision_result(content, model, truncated) + } + + async fn execute_streaming( + &self, + response: reqwest::Response, + deadline: tokio::time::Instant, + ) -> Result { + let mut raw_body = Vec::new(); + let mut line_buffer = SseLineBuffer::default(); + let mut state = VisionStreamState::default(); + let mut stream = response.bytes_stream(); + + 'response: loop { + let chunk = match tokio::time::timeout_at(deadline, stream.next()).await { + Ok(Some(Ok(chunk))) => chunk, + // A mid-stream read failure is not a clean truncation: the + // request phase already succeeded, so the shared retry policy + // no longer applies and the caller needs the real cause. + Ok(Some(Err(error))) => { + return Err(ToolError::execution_failed(format!( + "Vision API stream failed mid-response: {error}" + ))); + } + Err(_) => { + return Err(ToolError::execution_failed("Vision API response timed out")); + } + Ok(None) => break, + }; + if raw_body.len().saturating_add(chunk.len()) > MAX_VISION_RESPONSE_BYTES { + return Err(ToolError::execution_failed(format!( + "Vision API response exceeded {MAX_VISION_RESPONSE_BYTES} bytes" + ))); + } + raw_body.extend_from_slice(&chunk); + for line in line_buffer.push(&chunk)? { + if process_sse_line(&line, &mut state) { + break 'response; + } + } + } + + if !state.completed + && let Some(line) = line_buffer.finish() + { + let _ = process_sse_line(&line, &mut state); + } + + if !state.saw_data_event { + let (content, model, truncated) = self.parse_non_streaming_body(&raw_body)?; + return self.vision_result(content, model, truncated || state.truncated); + } + + state.truncated |= !state.completed; + self.vision_result(state.content, self.config.model.clone(), state.truncated) + } } #[async_trait] @@ -199,79 +437,59 @@ impl ToolSpec for ImageAnalyzeTool { let url = format!("{}/chat/completions", self.base_url()); let api_key = self.api_key(); - let retry_config = RetryConfig { - max_retries: 3, - initial_delay: 1.0, - max_delay: 30.0, - enabled: true, - ..Default::default() - }; - - let response = with_retry( - &retry_config, - || { - let client = self.client.clone(); - let url = url.clone(); - let api_key = api_key.clone(); - let payload = payload.clone(); - async move { - let response = client - .post(&url) - .header("Content-Type", "application/json") - .header("Authorization", format!("Bearer {api_key}")) - .json(&payload) - .send() - .await - .map_err(|e| LlmError::from_reqwest(&e))?; - - let status = response.status(); - if !status.is_success() { - let error_text = response - .text() + let retry_config = self.retry_config(); + let deadline = tokio::time::Instant::now() + self.request_timeout(); + + let response = tokio::time::timeout_at( + deadline, + with_retry( + &retry_config, + || { + let client = self.client.clone(); + let url = url.clone(); + let api_key = api_key.clone(); + let payload = payload.clone(); + async move { + let response = client + .post(&url) + .header("Content-Type", "application/json") + .header("Authorization", format!("Bearer {api_key}")) + .json(&payload) + .send() .await - .unwrap_or_else(|_| "Unknown error".to_string()); - let error_text = sanitize_http_error_body( - Some("Vision provider"), - status.as_u16(), - &error_text, - ); - return Err(LlmError::from_http_response(status.as_u16(), &error_text)); + .map_err(|e| LlmError::from_reqwest(&e))?; + + let status = response.status(); + if !status.is_success() { + let retry_after = extract_retry_after(response.headers()); + let error_text = + bounded_error_text(response, ERROR_BODY_MAX_BYTES).await; + let error_text = sanitize_http_error_body( + Some("Vision provider"), + status.as_u16(), + &error_text, + ); + return Err(LlmError::from_http_response_with_retry_after( + status.as_u16(), + &error_text, + retry_after, + )); + } + Ok(response) } - Ok(response) - } - }, - None, + }, + None, + ), ) .await + .map_err(|_| ToolError::execution_failed("Vision API request timed out"))? .map_err(|e| ToolError::execution_failed(format!("Vision API request failed: {e}")))?; - let json: Value = response - .json() - .await - .map_err(|e| ToolError::execution_failed(format!("Failed to parse response: {e}")))?; - - let content = json - .get("choices") - .and_then(|c| c.get(0)) - .and_then(|c| c.get("message")) - .and_then(|m| m.get("content")) - .and_then(|c| c.as_str()) - .unwrap_or("") - .to_string(); - - let model = json - .get("model") - .and_then(|m| m.as_str()) - .unwrap_or(&self.config.model) - .to_string(); - - let result = json!({ - "analysis": content, - "model": model, - }); - - ToolResult::json(&result) - .map_err(|e| ToolError::execution_failed(format!("Failed to serialize result: {e}"))) + if self.config.stream == Some(true) { + self.execute_streaming(response, deadline).await + } else { + self.execute_non_streaming(response, deadline).await + } } } @@ -279,6 +497,10 @@ impl ToolSpec for ImageAnalyzeTool { mod tests { use super::*; use tempfile::tempdir; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::TcpListener; + use wiremock::matchers::{method, path}; + use wiremock::{Mock, MockServer, ResponseTemplate}; #[cfg(unix)] fn create_file_symlink( @@ -301,9 +523,120 @@ mod tests { model: "test-vision-model".to_string(), api_key: Some("test-key".to_string()), base_url: Some("https://example.invalid/v1".to_string()), + request_timeout_secs: None, + stream: None, + retry_on_transient_errors: None, } } + async fn serve_slow_error_body() -> String { + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("bind slow test server"); + let address = listener.local_addr().expect("slow test server address"); + + tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.expect("accept test request"); + let mut buffer = [0_u8; 4096]; + let _ = socket.read(&mut buffer).await.expect("read test request"); + socket + .write_all( + b"HTTP/1.1 500 Internal Server Error\r\nContent-Length: 100\r\nConnection: close\r\n\r\n", + ) + .await + .expect("write slow response headers"); + tokio::time::sleep(Duration::from_secs(2)).await; + }); + + format!("http://{address}/v1") + } + + async fn serve_broken_stream_body() -> String { + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("bind broken-stream test server"); + let address = listener + .local_addr() + .expect("broken-stream test server address"); + + tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.expect("accept test request"); + let mut buffer = [0_u8; 4096]; + let _ = socket.read(&mut buffer).await.expect("read test request"); + socket + .write_all( + b"HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nContent-Length: 500\r\nConnection: close\r\n\r\n", + ) + .await + .expect("write broken-stream response headers"); + socket + .write_all(b"data: {\"choices\":[{\"delta\":{\"content\":\"partial\"}}]}\n\n") + .await + .expect("write partial stream body"); + // Close before Content-Length is satisfied so the client sees a + // mid-stream read error, not a clean EOF. + }); + + format!("http://{address}/v1") + } + + async fn serve_hanging_stream_body() -> String { + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("bind hanging-stream test server"); + let address = listener + .local_addr() + .expect("hanging-stream test server address"); + + tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.expect("accept test request"); + let mut buffer = [0_u8; 4096]; + let _ = socket.read(&mut buffer).await.expect("read test request"); + socket + .write_all( + b"HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nContent-Length: 500\r\nConnection: close\r\n\r\n", + ) + .await + .expect("write hanging-stream response headers"); + socket + .write_all(b"data: {\"choices\":[{\"delta\":{\"content\":\"partial\"}}]}\n\n") + .await + .expect("write partial stream body"); + tokio::time::sleep(Duration::from_secs(30)).await; + }); + + format!("http://{address}/v1") + } + + async fn execute_against( + mut config: VisionModelConfig, + status: u16, + content_type: &str, + body: &str, + ) -> (Result, Value) { + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path("/v1/chat/completions")) + .respond_with(ResponseTemplate::new(status).set_body_raw(body, content_type)) + .mount(&server) + .await; + config.base_url = Some(format!("{}/v1", server.uri())); + let workspace = tempdir().expect("workspace tempdir"); + std::fs::write(workspace.path().join("image.png"), b"test image") + .expect("write test image"); + let context = ToolContext::new(workspace.path().to_path_buf()); + let result = ImageAnalyzeTool::new(config) + .execute(json!({"image_path": "image.png"}), &context) + .await; + let requests = server.received_requests().await.expect("captured request"); + let request = serde_json::from_slice(&requests[0].body).expect("JSON request body"); + (result, request) + } + + fn result_json(result: ToolResult) -> Value { + serde_json::from_str(&result.content).expect("tool result JSON") + } + #[test] fn tool_metadata_is_read_only_and_named_image_analyze() { let tool = ImageAnalyzeTool::new(fake_config()); @@ -348,6 +681,227 @@ mod tests { Some(u64::from(DEFAULT_VISION_MAX_OUTPUT_TOKENS)) ); assert!(payload.get("max_completion_tokens").is_none()); + assert!( + payload.get("stream").is_none(), + "the default request shape must remain non-streaming" + ); + } + + #[test] + fn vision_payload_only_emits_stream_when_explicitly_configured() { + for (configured, expected) in [(Some(true), Some(true)), (Some(false), Some(false))] { + let mut config = fake_config(); + config.stream = configured; + let payload = + ImageAnalyzeTool::new(config).request_payload("describe", "abc123", "image/png"); + + assert_eq!(payload.get("stream").and_then(Value::as_bool), expected); + } + } + + #[test] + fn request_controls_preserve_defaults_and_apply_safe_bounds() { + let default_tool = ImageAnalyzeTool::new(fake_config()); + assert_eq!( + default_tool.request_timeout(), + Duration::from_secs(DEFAULT_VISION_REQUEST_TIMEOUT_SECS) + ); + assert!(default_tool.retry_config().enabled); + + let mut config = fake_config(); + config.request_timeout_secs = Some(u64::MAX); + config.retry_on_transient_errors = Some(false); + let configured_tool = ImageAnalyzeTool::new(config); + assert_eq!( + configured_tool.request_timeout(), + Duration::from_secs(MAX_VISION_REQUEST_TIMEOUT_SECS) + ); + assert!(!configured_tool.retry_config().enabled); + } + + #[test] + fn sse_buffer_preserves_utf8_split_across_chunks() { + let event = "data: {\"choices\":[{\"delta\":{\"content\":\"图\"}}]}\n"; + let bytes = event.as_bytes(); + let split = event.find('图').expect("multibyte character") + 1; + let mut buffer = SseLineBuffer::default(); + let mut lines = buffer.push(&bytes[..split]).expect("first chunk"); + lines.extend(buffer.push(&bytes[split..]).expect("second chunk")); + + assert_eq!(lines, vec![event.trim_end()]); + let mut state = VisionStreamState::default(); + assert!(!process_sse_line(&lines[0], &mut state)); + assert_eq!(state.content, "图"); + } + + #[test] + fn empty_vision_content_is_always_an_error() { + let tool = ImageAnalyzeTool::new(fake_config()); + for truncated in [false, true] { + let error = tool + .vision_result(String::new(), "test-model".to_string(), truncated) + .expect_err("empty content cannot be a successful analysis"); + assert!(error.to_string().contains("no usable content")); + } + } + + #[test] + fn done_only_stream_is_not_misclassified_as_json() { + let mut state = VisionStreamState::default(); + + assert!(process_sse_line("data: [DONE]", &mut state)); + assert!(state.saw_data_event); + assert!(state.completed); + } + + #[tokio::test] + async fn execute_streaming_accumulates_deltas() { + let mut config = fake_config(); + config.stream = Some(true); + let body = concat!( + "data: {\"choices\":[{\"delta\":{\"content\":\"hello \"}}]}\n\n", + "data: {\"choices\":[{\"delta\":{\"content\":\"world\"},\"finish_reason\":\"stop\"}]}\n\n" + ); + + let (result, request) = execute_against(config, 200, "text/event-stream", body).await; + let result = result_json(result.expect("streaming response succeeds")); + + assert_eq!(request.get("stream").and_then(Value::as_bool), Some(true)); + assert_eq!( + result.get("analysis").and_then(Value::as_str), + Some("hello world") + ); + assert!(result.get("truncated").is_none()); + } + + #[tokio::test] + async fn execute_streaming_marks_clean_eof_without_terminal_event_as_truncated() { + let mut config = fake_config(); + config.stream = Some(true); + let body = "data: {\"choices\":[{\"delta\":{\"content\":\"partial\"}}]}\n\n"; + + let (result, _) = execute_against(config, 200, "text/event-stream", body).await; + let result = result_json(result.expect("partial content remains usable")); + + assert_eq!( + result.get("analysis").and_then(Value::as_str), + Some("partial") + ); + assert_eq!(result.get("truncated").and_then(Value::as_bool), Some(true)); + } + + #[tokio::test] + async fn execute_streaming_mid_stream_read_error_fails_with_cause() { + let workspace = tempdir().expect("workspace tempdir"); + std::fs::write(workspace.path().join("image.png"), b"test image") + .expect("write test image"); + let context = ToolContext::new(workspace.path().to_path_buf()); + let mut config = fake_config(); + config.stream = Some(true); + config.base_url = Some(serve_broken_stream_body().await); + + let error = ImageAnalyzeTool::new(config) + .execute(json!({"image_path": "image.png"}), &context) + .await + .expect_err("a broken stream must fail, not report truncated success"); + + assert!( + error.to_string().contains("stream failed mid-response"), + "error must carry the read failure cause; got {error}" + ); + } + + #[tokio::test] + async fn execute_streaming_mid_stream_timeout_fails() { + let workspace = tempdir().expect("workspace tempdir"); + std::fs::write(workspace.path().join("image.png"), b"test image") + .expect("write test image"); + let context = ToolContext::new(workspace.path().to_path_buf()); + let mut config = fake_config(); + config.stream = Some(true); + config.base_url = Some(serve_hanging_stream_body().await); + config.request_timeout_secs = Some(1); + + let error = tokio::time::timeout( + Duration::from_secs(5), + ImageAnalyzeTool::new(config).execute(json!({"image_path": "image.png"}), &context), + ) + .await + .expect("the tool must enforce its own deadline") + .expect_err("a stalled stream must fail, not report truncated success"); + + assert!( + error.to_string().contains("timed out"), + "error must call out the deadline; got {error}" + ); + } + + #[tokio::test] + async fn execute_streaming_marks_malformed_frame_as_truncated() { + let mut config = fake_config(); + config.stream = Some(true); + let body = concat!( + "data: {not json}\n\n", + "data: {\"choices\":[{\"delta\":{\"content\":\"kept\"},\"finish_reason\":\"stop\"}]}\n\n" + ); + + let (result, _) = execute_against(config, 200, "text/event-stream", body).await; + let result = result_json(result.expect("valid frames survive a malformed one")); + + assert_eq!(result.get("analysis").and_then(Value::as_str), Some("kept")); + assert_eq!(result.get("truncated").and_then(Value::as_bool), Some(true)); + } + + #[tokio::test] + async fn execute_streaming_falls_back_to_ordinary_json() { + let mut config = fake_config(); + config.stream = Some(true); + let body = r#"{"model":"server-model","choices":[{"message":{"content":"fallback"},"finish_reason":"stop"}]}"#; + + let (result, _) = execute_against(config, 200, "application/json", body).await; + let result = result_json(result.expect("ordinary JSON fallback succeeds")); + + assert_eq!( + result.get("analysis").and_then(Value::as_str), + Some("fallback") + ); + assert_eq!( + result.get("model").and_then(Value::as_str), + Some("server-model") + ); + } + + #[tokio::test] + async fn execute_non_streaming_rejects_empty_content() { + let body = r#"{"choices":[{"message":{"content":""},"finish_reason":"length"}]}"#; + + let (result, request) = execute_against(fake_config(), 200, "application/json", body).await; + let error = result.expect_err("empty response must fail"); + + assert!(request.get("stream").is_none()); + assert!(error.to_string().contains("no usable content")); + } + + #[tokio::test] + async fn request_deadline_covers_error_response_body() { + let workspace = tempdir().expect("workspace tempdir"); + std::fs::write(workspace.path().join("image.png"), b"test image") + .expect("write test image"); + let context = ToolContext::new(workspace.path().to_path_buf()); + let mut config = fake_config(); + config.base_url = Some(serve_slow_error_body().await); + config.request_timeout_secs = Some(1); + config.retry_on_transient_errors = Some(false); + + let result = tokio::time::timeout( + Duration::from_secs(2), + ImageAnalyzeTool::new(config).execute(json!({"image_path": "image.png"}), &context), + ) + .await + .expect("the tool must enforce its own deadline") + .expect_err("an incomplete error response must fail"); + + assert!(result.to_string().contains("request timed out")); } #[test] diff --git a/docs/CONFIGURATION.md b/docs/CONFIGURATION.md index a4f4995156..56982a2095 100644 --- a/docs/CONFIGURATION.md +++ b/docs/CONFIGURATION.md @@ -555,6 +555,28 @@ api_key = "YOUR_XIAOMI_KEY" base_url = "https://api.xiaomimimo.com/v1" ``` +The vision request also accepts three optional compatibility controls: + +```toml +[vision_model] +# Opt in only when the endpoint supports OpenAI-compatible SSE responses. +stream = true +# Total budget for retries and response consumption (default: 120, maximum: 3600). +request_timeout_secs = 300 +# Set to false when a local endpoint should receive exactly one attempt. +retry_on_transient_errors = false +``` + +Omitting `stream` preserves the existing non-streaming request shape. When +streaming is enabled but an endpoint returns an ordinary JSON response, +Codewhale parses that response as a compatibility fallback. Omitting +`retry_on_transient_errors` keeps the existing transient-error retry policy. +Empty responses fail instead of producing a successful result with no usable +analysis. A mid-stream read failure or a stalled response fails the tool with +the underlying cause; only a clean end of stream without a terminal event, a +malformed SSE frame, or a `finish_reason` of `length` is reported as a +truncated result. + The example above uses Xiaomi MiMo's pay-as-you-go OpenAI-compatible endpoint. If you are using a Token Plan key (`tp-...`) for `[vision_model]`, you must set `base_url` explicitly because this generic OpenAI-compatible block does not