diff --git a/CHANGELOG.md b/CHANGELOG.md index 50da47b..0c5d7ca 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +## [0.45.4] - 2026-07-09 + +### Fixed +- LLM — `parse_sse_line` now emits every tool_call in a single OpenAI-compatible streaming chunk, instead of silently dropping indices 1+. ([#247](https://github.com/Fullstop000/ignis/pull/247)) +- Tools — `web_fetch` and `web_search` no longer buffer arbitrary HTTP response bodies into memory; a hard byte cap now bounds peak memory usage. ([#247](https://github.com/Fullstop000/ignis/pull/247)) + ## [0.45.3] - 2026-07-07 ### Fixed diff --git a/Cargo.lock b/Cargo.lock index bfa45e9..ff9f75e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1051,7 +1051,7 @@ dependencies = [ [[package]] name = "ignis" -version = "0.45.3" +version = "0.45.4" dependencies = [ "anyhow", "async-trait", diff --git a/ignis/Cargo.toml b/ignis/Cargo.toml index 10261c3..5739070 100644 --- a/ignis/Cargo.toml +++ b/ignis/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "ignis" -version = "0.45.3" +version = "0.45.4" edition = "2021" description = "A single-binary, multi-provider AI coding agent for your terminal." license = "Apache-2.0" diff --git a/ignis/src/llm/protocols/mod.rs b/ignis/src/llm/protocols/mod.rs index 1c6d472..d379824 100644 --- a/ignis/src/llm/protocols/mod.rs +++ b/ignis/src/llm/protocols/mod.rs @@ -293,50 +293,66 @@ where ) } -/// Parse one SSE line's payload into a response delta. Returns `None` for lines -/// that carry no delta: blank lines, comments / non-`data:` lines, the terminal -/// `[DONE]`, unparseable JSON, and empty-content deltas. Shared by every -/// OpenAI-compatible provider so the mapping lives — and is tested — once. -pub(crate) fn parse_sse_line(line: &str) -> Option { +/// Parse one SSE line's payload into one or more response deltas. Returns an +/// empty vec for lines that carry no delta: blank lines, comments / +/// non-`data:` lines, the terminal `[DONE]`, unparseable JSON, and +/// empty-content deltas. Shared by every OpenAI-compatible provider so the +/// mapping lives — and is tested — once. +pub(crate) fn parse_sse_line(line: &str) -> Vec { let line = line.trim(); if line.is_empty() { - return None; + return Vec::new(); } - let data_part = line.strip_prefix("data:")?.trim(); + let data_part = match line.strip_prefix("data:").map(str::trim) { + Some(d) => d, + None => return Vec::new(), + }; if data_part == "[DONE]" { - return None; + return Vec::new(); } - let chunk: Chunk = serde_json::from_str(data_part).ok()?; + let chunk: Chunk = match serde_json::from_str(data_part).ok() { + Some(c) => c, + None => return Vec::new(), + }; + + let mut out = Vec::new(); if let Some(choice) = chunk.choices.as_ref().and_then(|c| c.first()) { if let Some(content) = &choice.delta.content { if !content.is_empty() { - return Some(LlmResponseDelta::Text(content.clone())); + out.push(LlmResponseDelta::Text(content.clone())); } } - if let Some(reasoning) = &choice.delta.reasoning_content { - if !reasoning.is_empty() { - return Some(LlmResponseDelta::Reasoning(reasoning.clone())); + // Legacy single-delta behavior: when both content and reasoning are + // present in the same chunk, content "wins". Keep that ordering so a + // reasoning-only chunk is only emitted when there is no content. + if out.is_empty() { + if let Some(reasoning) = &choice.delta.reasoning_content { + if !reasoning.is_empty() { + out.push(LlmResponseDelta::Reasoning(reasoning.clone())); + } } } - if let Some(tc) = choice.delta.tool_calls.as_ref().and_then(|t| t.first()) { - let name = tc.function.as_ref().and_then(|f| f.name.clone()); - let arguments = tc - .function - .as_ref() - .and_then(|f| f.arguments.clone()) - .unwrap_or_default(); - return Some(LlmResponseDelta::ToolCall { - index: tc.index, - id: tc.id.clone(), - name, - arguments, - }); + if let Some(tool_calls) = &choice.delta.tool_calls { + for tc in tool_calls { + let name = tc.function.as_ref().and_then(|f| f.name.clone()); + let arguments = tc + .function + .as_ref() + .and_then(|f| f.arguments.clone()) + .unwrap_or_default(); + out.push(LlmResponseDelta::ToolCall { + index: tc.index, + id: tc.id.clone(), + name, + arguments, + }); + } } } if let Some(u) = &chunk.usage { - return Some(LlmResponseDelta::Usage(u.to_usage())); + out.push(LlmResponseDelta::Usage(u.to_usage())); } - None + out } /// Policy for trimming prior-turn material from a `history` before serializing @@ -543,10 +559,9 @@ pub(crate) async fn openai_compatible_chat_stream( let status = res.status(); if !status.is_success() { - let error_text = res - .text() + let (error_text, _) = crate::tools::util::read_body_with_cap(res, 64 * 1024) .await - .unwrap_or_else(|_| "Unknown error".to_string()); + .unwrap_or_else(|e| (format!("Unknown error: {e}"), false)); return Err(LlmHttpError { status, body: error_text, @@ -555,11 +570,12 @@ pub(crate) async fn openai_compatible_chat_stream( } let line_stream = bytes_to_lines(res.bytes_stream()); - let delta_stream = line_stream.filter_map(|line_result| async move { - match line_result { - Err(err) => Some(Err(err)), - Ok(line) => parse_sse_line(&line).map(Ok), - } + let delta_stream = line_stream.flat_map(|line_result| { + let deltas = match line_result { + Err(err) => vec![Err(err)], + Ok(line) => parse_sse_line(&line).into_iter().map(Ok).collect(), + }; + futures_util::stream::iter(deltas) }); Ok(delta_stream.boxed()) } @@ -600,13 +616,13 @@ mod tests { #[test] fn parse_sse_text_delta() { let d = parse_sse_line(r#"data: {"choices":[{"delta":{"content":"hello"}}]}"#); - assert!(matches!(d, Some(LlmResponseDelta::Text(t)) if t == "hello")); + assert!(matches!(d.as_slice(), [LlmResponseDelta::Text(t)] if t == "hello")); } #[test] fn parse_sse_reasoning_delta() { let d = parse_sse_line(r#"data: {"choices":[{"delta":{"reasoning_content":"hmm"}}]}"#); - assert!(matches!(d, Some(LlmResponseDelta::Reasoning(t)) if t == "hmm")); + assert!(matches!(d.as_slice(), [LlmResponseDelta::Reasoning(t)] if t == "hmm")); } #[test] @@ -614,14 +630,14 @@ mod tests { let d = parse_sse_line( r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_1","function":{"name":"bash","arguments":"{\"command\":"}}]}}]}"#, ); - match d { - Some(LlmResponseDelta::ToolCall { + match d.as_slice() { + [LlmResponseDelta::ToolCall { index, id, name, arguments, - }) => { - assert_eq!(index, 0); + }] => { + assert_eq!(*index, 0); assert_eq!(id.as_deref(), Some("call_1")); assert_eq!(name.as_deref(), Some("bash")); assert_eq!(arguments, r#"{"command":"#); @@ -636,7 +652,21 @@ mod tests { r#"data: {"choices":[{"delta":{"tool_calls":[{"index":1,"function":{"name":"x"}}]}}]}"#, ); assert!( - matches!(d, Some(LlmResponseDelta::ToolCall { arguments, .. }) if arguments.is_empty()) + matches!(d.as_slice(), [LlmResponseDelta::ToolCall { arguments, .. }] if arguments.is_empty()) + ); + } + + #[test] + fn parse_sse_emits_all_tool_calls_in_chunk() { + let d = parse_sse_line( + r#"data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_1","function":{"name":"bash","arguments":"{\"command\":\"ls\"}"}},{"index":1,"id":"call_2","function":{"name":"read_file","arguments":"{\"path\":\"/tmp/a\"}"}}]}}]}"#, + ); + assert_eq!(d.len(), 2); + assert!( + matches!(&d[0], LlmResponseDelta::ToolCall { index: 0, id, name, .. } if id.as_deref() == Some("call_1") && name.as_deref() == Some("bash")) + ); + assert!( + matches!(&d[1], LlmResponseDelta::ToolCall { index: 1, id, name, .. } if id.as_deref() == Some("call_2") && name.as_deref() == Some("read_file")) ); } @@ -645,18 +675,18 @@ mod tests { let d = parse_sse_line( r#"data: {"choices":[],"usage":{"prompt_tokens":10,"completion_tokens":3}}"#, ); - assert!(matches!(d, Some(LlmResponseDelta::Usage(u)) if u.input_tokens == 10)); + assert!(matches!(d.as_slice(), [LlmResponseDelta::Usage(u)] if u.input_tokens == 10)); } #[test] - fn parse_sse_done_and_noise_yield_none() { - assert!(parse_sse_line("data: [DONE]").is_none()); - assert!(parse_sse_line("").is_none()); - assert!(parse_sse_line(" ").is_none()); - assert!(parse_sse_line(": keep-alive comment").is_none()); // not a data: line - assert!(parse_sse_line("data: {not json}").is_none()); + fn parse_sse_done_and_noise_yield_empty() { + assert!(parse_sse_line("data: [DONE]").is_empty()); + assert!(parse_sse_line("").is_empty()); + assert!(parse_sse_line(" ").is_empty()); + assert!(parse_sse_line(": keep-alive comment").is_empty()); // not a data: line + assert!(parse_sse_line("data: {not json}").is_empty()); // Empty content delta carries no signal. - assert!(parse_sse_line(r#"data: {"choices":[{"delta":{"content":""}}]}"#).is_none()); + assert!(parse_sse_line(r#"data: {"choices":[{"delta":{"content":""}}]}"#).is_empty()); } #[test] @@ -665,7 +695,7 @@ mod tests { let d = parse_sse_line( r#"data: {"choices":[{"delta":{"content":"a","reasoning_content":"b"}}]}"#, ); - assert!(matches!(d, Some(LlmResponseDelta::Text(t)) if t == "a")); + assert!(matches!(d.as_slice(), [LlmResponseDelta::Text(t)] if t == "a")); } // ---- prep_outbound_history: stale tool-output masking + reasoning strip ---- diff --git a/ignis/src/tools/mod.rs b/ignis/src/tools/mod.rs index 5c074ba..4886567 100644 --- a/ignis/src/tools/mod.rs +++ b/ignis/src/tools/mod.rs @@ -13,7 +13,7 @@ mod list_dir; mod read_file; mod skill; mod todo_write; -mod util; +pub(crate) mod util; mod web_fetch; mod web_search; mod worktree; diff --git a/ignis/src/tools/util.rs b/ignis/src/tools/util.rs index 445a2c3..d0eabef 100644 --- a/ignis/src/tools/util.rs +++ b/ignis/src/tools/util.rs @@ -1,5 +1,6 @@ //! Shared helpers for the native tool implementations. +use futures_util::StreamExt; use std::time::Duration; /// The single marker every tool appends when it cuts its output, so the driving @@ -18,6 +19,55 @@ pub(crate) fn truncate_chars(s: &str, max: usize) -> String { format!("{kept}\n{TRUNCATION_MARKER}") } +/// Read the HTTP body with a hard byte cap, returning the prefix plus a +/// truncation marker if the server sent more than `limit` bytes. Uses a +/// streaming read so a huge response cannot OOM the process before truncation. +/// +/// The returned `bool` is `true` when the body was larger than `limit` and was +/// cut; callers that need strict completeness (e.g. JSON parsers) can turn +/// this into their own error instead of attempting a parse. +pub(crate) async fn read_body_with_cap( + resp: reqwest::Response, + limit: usize, +) -> Result<(String, bool), String> { + let mut bytes = Vec::with_capacity(limit.min(4096)); + let mut stream = resp.bytes_stream(); + while let Some(chunk) = stream.next().await { + let chunk = chunk.map_err(|e| format!("Failed to read body: {e}"))?; + if bytes.len() + chunk.len() > limit { + let take = limit.saturating_sub(bytes.len()); + bytes.extend_from_slice(&chunk[..take]); + break; + } + bytes.extend_from_slice(&chunk); + } + Ok(truncate_bytes_with_marker(bytes, limit)) +} + +/// Apply the same truncation/marker logic that `read_body_with_cap` uses, +/// but synchronously on an already-collected byte buffer. Kept separate so +/// UTF-8 boundary handling can be unit-tested without constructing a live +/// HTTP response. +pub(crate) fn truncate_bytes_with_marker(bytes: Vec, limit: usize) -> (String, bool) { + if bytes.len() <= limit { + return ( + String::from_utf8(bytes).unwrap_or_else(|e| format!("Invalid UTF-8: {e}")), + false, + ); + } + let mut trimmed = bytes; + trimmed.truncate(limit); + let mut end = trimmed.len(); + while end > 0 && std::str::from_utf8(&trimmed[..end]).is_err() { + end -= 1; + } + trimmed.truncate(end); + let mut text = String::from_utf8(trimmed).expect("truncated at a valid UTF-8 boundary"); + text.push('\n'); + text.push_str(TRUNCATION_MARKER); + (text, true) +} + const MAX_RETRIES: u32 = 3; const BASE_BACKOFF_MS: u64 = 500; const MAX_BACKOFF_MS: u64 = 5_000; @@ -93,4 +143,37 @@ mod tests { assert!(out.ends_with(TRUNCATION_MARKER)); assert_eq!(out.chars().filter(|&c| c == '中').count(), 4); } + + #[test] + fn truncate_bytes_with_marker_under_limit_passes_through() { + let (text, truncated) = truncate_bytes_with_marker(b"hello".to_vec(), 10); + assert_eq!(text, "hello"); + assert!(!truncated); + } + + #[test] + fn truncate_bytes_with_marker_at_exact_limit_has_no_marker() { + let (text, truncated) = truncate_bytes_with_marker(b"hello".to_vec(), 5); + assert_eq!(text, "hello"); + assert!(!truncated); + } + + #[test] + fn truncate_bytes_with_marker_over_limit_appends_marker() { + let (text, truncated) = truncate_bytes_with_marker(b"hello world".to_vec(), 5); + assert!(truncated); + assert!(text.starts_with("hello")); + assert!(text.ends_with(TRUNCATION_MARKER)); + } + + #[test] + fn truncate_bytes_with_marker_respects_cjk_boundary() { + // Each '中' is 3 bytes. A cap of 5 bytes lands inside the second CJK + // character; the function must trim back to the last valid codepoint. + let bytes = "中中".as_bytes().to_vec(); + let (text, truncated) = truncate_bytes_with_marker(bytes, 5); + assert!(truncated); + assert_eq!(text.chars().filter(|&c| c == '中').count(), 1); + assert!(text.ends_with(TRUNCATION_MARKER)); + } } diff --git a/ignis/src/tools/web_fetch.rs b/ignis/src/tools/web_fetch.rs index fac8b95..a4292b7 100644 --- a/ignis/src/tools/web_fetch.rs +++ b/ignis/src/tools/web_fetch.rs @@ -4,6 +4,10 @@ use std::time::Duration; /// Cap on returned text so a large page can't blow up the context. const MAX_OUTPUT_CHARS: usize = 20_000; +/// Hard cap on bytes read before the HTML→text pass. This bounds peak memory +/// regardless of how large the server response is; output is still truncated to +/// `MAX_OUTPUT_CHARS` afterward. +const MAX_BODY_BYTES: usize = 4 * 1024 * 1024; /// Fetch a URL and return its readable text. HTML is stripped to plain text. /// Pairs with `web_search` (search finds URLs, fetch reads them). @@ -60,10 +64,8 @@ impl StaticTool for WebFetchTool { .and_then(|v| v.to_str().ok()) .map(|ct| ct.contains("html")) .unwrap_or(false); - let body = resp - .text() - .await - .map_err(|e| format!("Failed to read body: {e}"))?; + let (body, truncated) = + crate::tools::util::read_body_with_cap(resp, MAX_BODY_BYTES).await?; let text = if is_html || looks_like_html(&body) { html_to_text(&body) @@ -74,7 +76,14 @@ impl StaticTool for WebFetchTool { if text.is_empty() { return Ok("(empty response)".to_string()); } - Ok(crate::tools::util::truncate_chars(text, MAX_OUTPUT_CHARS)) + // If the body hit the byte cap, preserve the marker even after the + // character-based output truncation that follows. + let mut out = crate::tools::util::truncate_chars(text, MAX_OUTPUT_CHARS); + if truncated && !out.ends_with(crate::tools::util::TRUNCATION_MARKER) { + out.push('\n'); + out.push_str(crate::tools::util::TRUNCATION_MARKER); + } + Ok(out) } } diff --git a/ignis/src/tools/web_search.rs b/ignis/src/tools/web_search.rs index cb5df33..f51b5c2 100644 --- a/ignis/src/tools/web_search.rs +++ b/ignis/src/tools/web_search.rs @@ -3,6 +3,12 @@ use async_trait::async_trait; use serde_json::json; const RESULT_COUNT: u32 = 5; +/// Cap on bytes read from a search API error response; keeps the error message +/// small without buffering an arbitrary server body into memory. +const MAX_ERROR_BODY_BYTES: usize = 64 * 1024; +/// Cap on bytes read from a search API success response. Search results are +/// JSON and should be small, but this bounds memory against a misbehaving API. +const MAX_JSON_BODY_BYTES: usize = 2 * 1024 * 1024; /// A normalized search hit, independent of which backend produced it. struct SearchResult { @@ -55,13 +61,19 @@ impl WebSearchTool { .await?; if !resp.status().is_success() { let status = resp.status(); - let body = resp.text().await.unwrap_or_default(); + let (body, _) = + crate::tools::util::read_body_with_cap(resp, MAX_ERROR_BODY_BYTES).await?; return Err(format!("Brave API error {status}: {}", truncate(&body))); } - let json: serde_json::Value = resp - .json() - .await - .map_err(|e| format!("Failed to parse response: {e}"))?; + let (body, truncated) = + crate::tools::util::read_body_with_cap(resp, MAX_JSON_BODY_BYTES).await?; + if truncated { + return Err(format!( + "Brave API response body exceeded maximum allowed size (>{MAX_JSON_BODY_BYTES} bytes)" + )); + } + let json: serde_json::Value = + serde_json::from_str(&body).map_err(|e| format!("Failed to parse response: {e}"))?; Ok(parse_brave(&json)) } @@ -78,13 +90,19 @@ impl WebSearchTool { .await?; if !resp.status().is_success() { let status = resp.status(); - let body = resp.text().await.unwrap_or_default(); + let (body, _) = + crate::tools::util::read_body_with_cap(resp, MAX_ERROR_BODY_BYTES).await?; return Err(format!("Tavily API error {status}: {}", truncate(&body))); } - let json: serde_json::Value = resp - .json() - .await - .map_err(|e| format!("Failed to parse response: {e}"))?; + let (body, truncated) = + crate::tools::util::read_body_with_cap(resp, MAX_JSON_BODY_BYTES).await?; + if truncated { + return Err(format!( + "Tavily API response body exceeded maximum allowed size (>{MAX_JSON_BODY_BYTES} bytes)" + )); + } + let json: serde_json::Value = + serde_json::from_str(&body).map_err(|e| format!("Failed to parse response: {e}"))?; Ok(parse_tavily(&json)) } }