From 1227a74ebe964ee19c7ae3d6b3e3d970cd091db8 Mon Sep 17 00:00:00 2001 From: Illia Polosukhin Date: Sat, 1 Aug 2026 05:03:09 +0000 Subject: [PATCH 1/4] feat(llm): explicit Anthropic cache_control breakpoints on both transports MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Closes #6984 (P0 of the pi-harness adoption program, docs/research/ pi-agent-deep-dive.md §7.3). The rig transport previously relied solely on Anthropic automatic caching via a top-level cache_control field, and the OAuth transport emitted no cache markers at all. Now both place explicit breakpoints so the tool/system prefix and the growing conversation cache independently: - OAuth transport: apply_cache_breakpoints marks the system prompt block, the last tool definition, and the last content block of the last message, all carrying the retention TTL. Retention None keeps the legacy wire shape (plain-string system, no markers). - rig transport: build_rig_request marks the last tool by moving it into rig's raw additional_params.tools (appended after typed tools, order preserved, Anthropic-native input_schema shape) and keeps the top-level automatic marker; Short retention additionally enables rig's typed system/last-message breakpoints. Long must not enable the typed breakpoints: rig markers cannot carry a TTL and a 5m block marker beside a 1h automatic marker is an API error. All markers in a request share one TTL, satisfying Anthropic's longer-TTL-first ordering rule. Unsupported models downgrade to None via supports_prompt_cache on both paths. Wire shape is pinned by loopback capture-server tests in both files (three per transport: short, long, none), plus build_rig_request seam tests for the tool move. Co-Authored-By: Claude Fable 5 --- crates/ironclaw_llm/CLAUDE.md | 23 ++ crates/ironclaw_llm/src/anthropic_oauth.rs | 349 +++++++++++++++++++- crates/ironclaw_llm/src/config.rs | 20 +- crates/ironclaw_llm/src/lib.rs | 23 +- crates/ironclaw_llm/src/rig_adapter.rs | 355 +++++++++++++++++++-- 5 files changed, 732 insertions(+), 38 deletions(-) diff --git a/crates/ironclaw_llm/CLAUDE.md b/crates/ironclaw_llm/CLAUDE.md index b857d695939..574e02d8fd3 100644 --- a/crates/ironclaw_llm/CLAUDE.md +++ b/crates/ironclaw_llm/CLAUDE.md @@ -285,6 +285,29 @@ Providers in this crate import it as `use ironclaw_common::llm_costs as costs;` (a plain import alias, **not** a re-export — see the relocation note in `.claude/rules/type-placement.md`). +## Anthropic Prompt Caching + +Both Anthropic transports emit explicit `cache_control` breakpoints when +`cache_retention` (env: `ANTHROPIC_CACHE_RETENTION`) is not `none` (#6984): + +- **OAuth transport** (`anthropic_oauth.rs`, `apply_cache_breakpoints`): system + prompt block, last tool definition, and the last content block of the last + message, all carrying the retention TTL (`{"type":"ephemeral"}` for short, + `+ "ttl":"1h"` for long). +- **API-key transport** (`rig_adapter.rs`, `build_rig_request` + + `create_anthropic_from_registry`): the top-level automatic-caching marker, + an explicit marker on the last tool (moved into rig's raw + `additional_params.tools`, which rig appends after typed tools), and — for + `short` only — rig's typed system/last-message breakpoints + (`CompletionModel::prompt_caching`). `long` must not enable the typed + breakpoints: rig's markers cannot carry a TTL, and a 5m block marker with a + 1h automatic marker is an API error (TTL conflict on the last block). + +All markers within a request share one TTL, satisfying Anthropic's +longer-TTL-first ordering rule. Models without cache support (claude-2 era) +downgrade to `none` via `supports_prompt_cache`. Wire shape is pinned by +capture-server tests in both files. + ## rig_adapter.rs Details `RigAdapter` bridges any rig-core `CompletionModel` to `LlmProvider`. It is actively used in production for all non-NEAR AI providers (OpenAI, Anthropic, Ollama, Tinfoil, OpenAI-compatible). Key behaviors: diff --git a/crates/ironclaw_llm/src/anthropic_oauth.rs b/crates/ironclaw_llm/src/anthropic_oauth.rs index 06367b8a2cd..1f812fb52af 100644 --- a/crates/ironclaw_llm/src/anthropic_oauth.rs +++ b/crates/ironclaw_llm/src/anthropic_oauth.rs @@ -98,6 +98,9 @@ pub(crate) struct AnthropicOAuthProvider { active_model: std::sync::RwLock, /// Parameter names that this provider does not support. unsupported_params: HashSet, + /// Anthropic prompt-cache retention; drives the explicit `cache_control` + /// breakpoints (system prompt, last tool, last message block). See #6984. + cache_retention: crate::config::CacheRetention, } impl AnthropicOAuthProvider { @@ -127,6 +130,14 @@ impl AnthropicOAuthProvider { let unsupported_params: HashSet = config.unsupported_params.iter().cloned().collect(); + let cache_retention = if config.cache_retention != crate::config::CacheRetention::None + && !crate::rig_adapter::supports_prompt_cache(&config.model) + { + crate::config::CacheRetention::None + } else { + config.cache_retention + }; + Ok(Self { client, token: std::sync::RwLock::new(token), @@ -134,6 +145,7 @@ impl AnthropicOAuthProvider { base_url, active_model, unsupported_params, + cache_retention, }) } @@ -326,17 +338,19 @@ impl LlmProvider for AnthropicOAuthProvider { let (system, messages) = convert_messages(req.messages); let max_tokens = req.max_tokens.unwrap_or(DEFAULT_MAX_TOKENS); - let request = AnthropicRequest { + let mut request = AnthropicRequest { thinking: thinking_for_request(&model, max_tokens, req.temperature, false), model, messages, - system, + system: system.map(AnthropicSystem::Text), max_tokens, temperature: req.temperature, tools: None, tool_choice: None, }; + apply_cache_breakpoints(&mut request, self.cache_retention); + let response: AnthropicResponse = self.send_request(&request).await?; let extracted = extract_response_content(&response); @@ -375,6 +389,7 @@ impl LlmProvider for AnthropicOAuthProvider { name: t.name, description: t.description, input_schema: t.parameters, + cache_control: None, }) .collect(); @@ -406,7 +421,7 @@ impl LlmProvider for AnthropicOAuthProvider { let has_tools = !tools.is_empty(); let opt_tools = if has_tools { Some(tools) } else { None }; - let request = AnthropicRequest { + let mut request = AnthropicRequest { thinking: if has_tools { None } else { @@ -414,13 +429,15 @@ impl LlmProvider for AnthropicOAuthProvider { }, model, messages, - system, + system: system.map(AnthropicSystem::Text), max_tokens, temperature: req.temperature, tools: opt_tools, tool_choice, }; + apply_cache_breakpoints(&mut request, self.cache_retention); + let response: AnthropicResponse = self.send_request(&request).await?; let extracted = extract_response_content(&response); @@ -486,7 +503,7 @@ struct AnthropicRequest { model: String, messages: Vec, #[serde(skip_serializing_if = "Option::is_none")] - system: Option, + system: Option, max_tokens: u32, #[serde(skip_serializing_if = "Option::is_none")] temperature: Option, @@ -504,6 +521,25 @@ struct AnthropicMessage { content: AnthropicContent, } +/// Anthropic system prompt: a plain string (legacy wire shape, kept when +/// caching is off) or content blocks so the last block can carry a +/// `cache_control` breakpoint. +#[derive(Debug, Serialize)] +#[serde(untagged)] +enum AnthropicSystem { + Text(String), + Blocks(Vec), +} + +#[derive(Debug, Serialize)] +struct AnthropicSystemBlock { + #[serde(rename = "type")] + block_type: &'static str, + text: String, + #[serde(skip_serializing_if = "Option::is_none")] + cache_control: Option, +} + /// Anthropic content can be a simple string or a list of content blocks. #[derive(Debug, Serialize)] #[serde(untagged)] @@ -516,22 +552,46 @@ enum AnthropicContent { #[serde(tag = "type")] enum AnthropicContentBlock { #[serde(rename = "text")] - Text { text: String }, + Text { + text: String, + #[serde(skip_serializing_if = "Option::is_none")] + cache_control: Option, + }, #[serde(rename = "image")] - Image { source: AnthropicImageSource }, + Image { + source: AnthropicImageSource, + #[serde(skip_serializing_if = "Option::is_none")] + cache_control: Option, + }, #[serde(rename = "tool_use")] ToolUse { id: String, name: String, input: serde_json::Value, + #[serde(skip_serializing_if = "Option::is_none")] + cache_control: Option, }, #[serde(rename = "tool_result")] ToolResult { tool_use_id: String, content: String, + #[serde(skip_serializing_if = "Option::is_none")] + cache_control: Option, }, } +impl AnthropicContentBlock { + /// Stamp a `cache_control` breakpoint on this block. + fn set_cache_control(&mut self, marker: serde_json::Value) { + match self { + Self::Text { cache_control, .. } + | Self::Image { cache_control, .. } + | Self::ToolUse { cache_control, .. } + | Self::ToolResult { cache_control, .. } => *cache_control = Some(marker), + } + } +} + /// Inline base64 image source for an Anthropic `image` content block. #[derive(Debug, Serialize)] struct AnthropicImageSource { @@ -546,6 +606,8 @@ struct AnthropicTool { name: String, description: String, input_schema: serde_json::Value, + #[serde(skip_serializing_if = "Option::is_none")] + cache_control: Option, } #[derive(Debug, Serialize)] @@ -613,6 +675,7 @@ fn user_image_blocks(parts: &[ContentPart]) -> Vec { ContentPart::ImageUrl { image_url } => { let (media_type, data) = image_url.decode_data_url()?; Some(AnthropicContentBlock::Image { + cache_control: None, source: AnthropicImageSource { source_type: "base64", media_type: media_type.to_string(), @@ -625,6 +688,56 @@ fn user_image_blocks(parts: &[ContentPart]) -> Vec { .collect() } +/// Place the explicit Anthropic `cache_control` breakpoints (issue #6984): +/// the system prompt, the last tool definition, and the last content block of +/// the last message. Mirrors pi's placement so the tool/system prefix and the +/// growing conversation each cache independently. All markers carry the same +/// TTL, satisfying Anthropic's longer-TTL-first ordering rule. No-op when +/// retention is `None`, preserving the legacy wire shape (plain-string +/// system, no markers). +fn apply_cache_breakpoints( + request: &mut AnthropicRequest, + retention: crate::config::CacheRetention, +) { + let Some(marker) = retention.cache_control_json() else { + return; + }; + + if let Some(AnthropicSystem::Text(text)) = request.system.take() { + request.system = Some(AnthropicSystem::Blocks(vec![AnthropicSystemBlock { + block_type: "text", + text, + cache_control: Some(marker.clone()), + }])); + } + + if let Some(tools) = request.tools.as_mut() + && let Some(last) = tools.last_mut() + { + last.cache_control = Some(marker.clone()); + } + + if let Some(last_message) = request.messages.last_mut() { + match &mut last_message.content { + // Empty text blocks cannot carry cache_control (API rejects + // them), so an empty trailing message keeps the string form. + AnthropicContent::Text(text) if !text.is_empty() => { + last_message.content = + AnthropicContent::Blocks(vec![AnthropicContentBlock::Text { + text: std::mem::take(text), + cache_control: Some(marker), + }]); + } + AnthropicContent::Text(_) => {} + AnthropicContent::Blocks(blocks) => { + if let Some(last_block) = blocks.last_mut() { + last_block.set_cache_control(marker); + } + } + } + } +} + /// Convert ChatMessage list to Anthropic format. /// /// Extracts system messages to the top-level `system` parameter (Anthropic @@ -649,7 +762,10 @@ fn convert_messages(messages: Vec) -> (Option, Vec { let mut blocks = Vec::with_capacity(1 + image_blocks.len()); if !msg.content.is_empty() { - blocks.push(AnthropicContentBlock::Text { text: msg.content }); + blocks.push(AnthropicContentBlock::Text { + text: msg.content, + cache_control: None, + }); } blocks.extend(image_blocks); AnthropicContent::Blocks(blocks) @@ -665,13 +781,17 @@ fn convert_messages(messages: Vec) -> (Option, Vec = Vec::new(); if !msg.content.is_empty() { - blocks.push(AnthropicContentBlock::Text { text: msg.content }); + blocks.push(AnthropicContentBlock::Text { + text: msg.content, + cache_control: None, + }); } for tc in tool_calls { blocks.push(AnthropicContentBlock::ToolUse { id: tc.id, name: tc.name, input: tc.arguments, + cache_control: None, }); } anthropic_msgs.push(AnthropicMessage { @@ -694,6 +814,7 @@ fn convert_messages(messages: Vec) -> (Option, Vec ExtractedAnthropicR #[cfg(test)] mod tests { use super::*; + use crate::config::CacheRetention; + use crate::provider::ToolDefinition; #[tokio::test] async fn complete_preserves_missing_retry_after_on_headerless_502() { @@ -853,6 +976,214 @@ mod tests { )); } + /// One-shot loopback capture server: returns the base URL and a handle + /// resolving to the captured request body. Replies 400 — these tests + /// assert the request wire shape, not response handling. + async fn capture_one_request() -> (String, tokio::sync::oneshot::Receiver) { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::TcpListener; + + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("loopback listener"); + let base_url = format!( + "http://{}", + listener.local_addr().expect("loopback address") + ); + let (tx, rx) = tokio::sync::oneshot::channel(); + tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.expect("accept request"); + let mut request = Vec::new(); + let mut buffer = [0_u8; 4096]; + loop { + let read = socket.read(&mut buffer).await.expect("read request"); + if read == 0 { + return; + } + request.extend_from_slice(&buffer[..read]); + let Some(header_end) = request.windows(4).position(|bytes| bytes == b"\r\n\r\n") + else { + continue; + }; + let headers = + std::str::from_utf8(&request[..header_end]).expect("request headers are UTF-8"); + let content_length = headers + .lines() + .find_map(|line| { + let (name, value) = line.split_once(':')?; + name.eq_ignore_ascii_case("content-length") + .then_some(value.trim()) + }) + .expect("content length header") + .parse::() + .expect("content length is numeric"); + let body_start = header_end + 4; + if request.len() < body_start + content_length { + continue; + } + tx.send( + String::from_utf8(request[body_start..body_start + content_length].to_vec()) + .expect("request body is UTF-8"), + ) + .expect("test receives captured request"); + socket + .write_all(b"HTTP/1.1 400 Bad Request\r\nContent-Length: 0\r\n\r\n") + .await + .expect("write test response"); + return; + } + }); + (base_url, rx) + } + + fn provider_with_retention( + base_url: &str, + retention: CacheRetention, + ) -> AnthropicOAuthProvider { + let mut config = RegistryProviderConfig::generic( + crate::registry::ProviderProtocol::Anthropic, + "anthropic_oauth", + None, + base_url, + "claude-opus-4-6", + ); + config.oauth_token = Some(SecretString::from("test-token".to_string())); + config.cache_retention = retention; + AnthropicOAuthProvider::new(&config).expect("provider") + } + + async fn captured_json(rx: tokio::sync::oneshot::Receiver) -> serde_json::Value { + let body = tokio::time::timeout(std::time::Duration::from_secs(5), rx) + .await + .expect("request capture timed out") + .expect("captured request body"); + serde_json::from_str(&body).expect("captured body is JSON") + } + + /// Wire-level pin for issue #6984: under `Short` retention the OAuth + /// transport emits the three explicit cache breakpoints — system prompt + /// block, last tool definition, and the last content block of the last + /// message (here a tool_result, the common agent-loop tail). + #[tokio::test] + async fn oauth_short_retention_places_explicit_cache_breakpoints() { + let (base_url, captured) = capture_one_request().await; + let provider = provider_with_retention(&base_url, CacheRetention::Short); + + let request = ToolCompletionRequest::new( + vec![ + ChatMessage::system("You are helpful."), + ChatMessage::user("Run the tool."), + ChatMessage::assistant_with_tool_calls( + None, + vec![ToolCall { + id: "call_1".to_string(), + name: "alpha".to_string(), + arguments: serde_json::json!({}), + ..ToolCall::default() + }], + ), + ChatMessage::tool_result("call_1", "alpha", "tool says hi"), + ], + vec![ + ToolDefinition { + name: "alpha".to_string(), + description: "First tool".to_string(), + parameters: serde_json::json!({"type": "object", "properties": {}}), + }, + ToolDefinition { + name: "beta".to_string(), + description: "Second tool".to_string(), + parameters: serde_json::json!({"type": "object", "properties": {}}), + }, + ], + ); + let _ = provider.complete_with_tools(request).await; + let body = captured_json(captured).await; + + let system = body["system"] + .as_array() + .expect("system serialized as blocks when caching is on"); + assert_eq!( + system.last().expect("system block")["cache_control"]["type"], + "ephemeral" + ); + assert!(system.last().unwrap()["cache_control"].get("ttl").is_none()); + + let tools = body["tools"].as_array().expect("tools array"); + assert_eq!(tools.len(), 2); + assert!(tools[0].get("cache_control").is_none()); + assert_eq!(tools[1]["name"], "beta"); + assert_eq!(tools[1]["cache_control"]["type"], "ephemeral"); + + let messages = body["messages"].as_array().expect("messages"); + let last_content = messages.last().expect("last message")["content"] + .as_array() + .expect("last message content blocks"); + let last_block = last_content.last().expect("content block"); + assert_eq!(last_block["type"], "tool_result"); + assert_eq!(last_block["cache_control"]["type"], "ephemeral"); + // Only the final block carries a marker. + for message in &messages[..messages.len() - 1] { + if let Some(blocks) = message["content"].as_array() { + for block in blocks { + assert!(block.get("cache_control").is_none(), "{body}"); + } + } + } + } + + /// Wire-level pin: `Long` retention stamps a 1h TTL on every breakpoint + /// (uniform TTLs satisfy Anthropic's longer-before-shorter ordering rule). + #[tokio::test] + async fn oauth_long_retention_uses_1h_ttl_markers() { + let (base_url, captured) = capture_one_request().await; + let provider = provider_with_retention(&base_url, CacheRetention::Long); + + let _ = provider + .complete(CompletionRequest::new(vec![ + ChatMessage::system("You are helpful."), + ChatMessage::user("Question"), + ])) + .await; + let body = captured_json(captured).await; + + let system = body["system"].as_array().expect("system blocks"); + assert_eq!(system.last().unwrap()["cache_control"]["ttl"], "1h"); + + let messages = body["messages"].as_array().expect("messages"); + let last_content = messages.last().unwrap()["content"] + .as_array() + .expect("last message content blocks"); + assert_eq!(last_content.last().unwrap()["cache_control"]["ttl"], "1h"); + } + + /// Wire-level pin: retention `None` keeps the legacy wire shape — system + /// as a plain string, and no cache_control anywhere. + #[tokio::test] + async fn oauth_no_retention_keeps_legacy_wire_shape() { + let (base_url, captured) = capture_one_request().await; + let provider = provider_with_retention(&base_url, CacheRetention::None); + + let _ = provider + .complete(CompletionRequest::new(vec![ + ChatMessage::system("You are helpful."), + ChatMessage::user("Question"), + ])) + .await; + let body = captured_json(captured).await; + + assert!( + body["system"].is_string(), + "system stays a plain string when caching is off: {body}" + ); + assert!( + !serde_json::to_string(&body) + .unwrap() + .contains("cache_control"), + "no cache_control may be emitted when caching is off: {body}" + ); + } + #[test] fn context_overflow_413_maps_to_context_length_exceeded() { // A raw HTTP 413 (payload too large) must become ContextLengthExceeded diff --git a/crates/ironclaw_llm/src/config.rs b/crates/ironclaw_llm/src/config.rs index 343d53d50aa..3626de35e30 100644 --- a/crates/ironclaw_llm/src/config.rs +++ b/crates/ironclaw_llm/src/config.rs @@ -23,9 +23,10 @@ pub const OAUTH_PLACEHOLDER: &str = "oauth-placeholder"; /// Prompt cache retention policy for Anthropic. /// -/// Controls Anthropic's automatic prompt caching via a top-level -/// `cache_control` field injected through rig-core's `additional_params`. -/// - `None` — caching disabled, no `cache_control` injected. +/// Controls Anthropic prompt caching — both the explicit per-block +/// `cache_control` breakpoints (system prompt, last tool definition, last +/// message block) and the top-level automatic-caching marker. +/// - `None` — caching disabled, no `cache_control` emitted anywhere. /// - `Short` — 5-minute TTL (default), `{"type": "ephemeral"}`, 1.25× write surcharge. /// - `Long` — 1-hour TTL, `{"type": "ephemeral", "ttl": "1h"}`, 2× write surcharge. #[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] @@ -39,6 +40,19 @@ pub enum CacheRetention { Long, } +impl CacheRetention { + /// The Anthropic `cache_control` marker for this retention, usable both + /// as a per-block breakpoint and as the request-level automatic-caching + /// field. `None` when caching is disabled. + pub(crate) fn cache_control_json(&self) -> Option { + match self { + Self::None => Option::None, + Self::Short => Some(serde_json::json!({"type": "ephemeral"})), + Self::Long => Some(serde_json::json!({"type": "ephemeral", "ttl": "1h"})), + } + } +} + impl std::str::FromStr for CacheRetention { type Err = String; diff --git a/crates/ironclaw_llm/src/lib.rs b/crates/ironclaw_llm/src/lib.rs index a6c88461835..e29e1e570e1 100644 --- a/crates/ironclaw_llm/src/lib.rs +++ b/crates/ironclaw_llm/src/lib.rs @@ -502,15 +502,32 @@ fn create_anthropic_from_registry( reason: format!("Failed to create Anthropic client: {e}"), })?; - let cache_retention = config.cache_retention; + // Downgrade retention up front for models without prompt-cache support so + // the rig `prompt_caching` flag below agrees with the adapter's own + // `with_cache_retention` validation. + let cache_retention = if config.cache_retention != CacheRetention::None + && !rig_adapter::supports_prompt_cache(&config.model) + { + CacheRetention::None + } else { + config.cache_retention + }; - let model = client.completion_model(&config.model); + let mut model = client.completion_model(&config.model); + + // Short retention: rig's typed breakpoints (system prompt + last message + // block, plain 5m ephemeral) complement the request-level automatic + // marker and the last-tool marker added in `build_rig_request`. Long + // retention must NOT set this — rig's markers cannot carry a TTL, and a + // 5m block marker alongside a 1h automatic marker is an API error + // (TTL conflict on the last block). See issue #6984. + model.prompt_caching = cache_retention == CacheRetention::Short; if cache_retention != CacheRetention::None { tracing::debug!( model = %config.model, retention = %cache_retention, - "Anthropic automatic prompt caching enabled" + "Anthropic prompt caching enabled (explicit breakpoints + automatic marker)" ); } diff --git a/crates/ironclaw_llm/src/rig_adapter.rs b/crates/ironclaw_llm/src/rig_adapter.rs index fdc4d11e46a..c8d0430038a 100644 --- a/crates/ironclaw_llm/src/rig_adapter.rs +++ b/crates/ironclaw_llm/src/rig_adapter.rs @@ -48,9 +48,11 @@ pub struct RigAdapter { input_cost: Decimal, output_cost: Decimal, /// Prompt cache retention policy (Anthropic only). - /// When not `CacheRetention::None`, injects top-level `cache_control` - /// via `additional_params` for Anthropic automatic caching. Also controls - /// the cost multiplier for cache-creation tokens. + /// When not `CacheRetention::None`, emits explicit `cache_control` + /// breakpoints (last tool via `additional_params.tools`, plus rig's + /// system/last-message markers under `Short`) and the top-level + /// automatic-caching marker — see `build_rig_request` and issue #6984. + /// Also controls the cost multiplier for cache-creation tokens. cache_retention: CacheRetention, /// Parameter names that this provider does not support (e.g., `"temperature"`). /// These are stripped from requests before sending to avoid 400 errors. @@ -288,9 +290,12 @@ impl RigAdapter { /// - `Short` — 5-minute TTL via `{"type": "ephemeral"}`, 1.25× write surcharge. /// - `Long` — 1-hour TTL via `{"type": "ephemeral", "ttl": "1h"}`, 2.0× write surcharge. /// - /// Cache injection uses Anthropic's **automatic caching** — a top-level - /// `cache_control` field in `additional_params` that gets `#[serde(flatten)]`'d - /// into the request body by rig-core. + /// Cache injection combines explicit breakpoints with Anthropic's + /// automatic caching: `build_rig_request` marks the last tool definition + /// (via `additional_params.tools`) and injects the top-level + /// `cache_control` field, while `Short` retention additionally enables + /// rig's typed system/last-message breakpoints at provider construction + /// (`create_anthropic_from_registry`). See issue #6984. /// /// If the configured model does not support caching (e.g. claude-2), /// a warning is logged once at construction and caching is disabled. @@ -911,7 +916,7 @@ fn saturate_u32(val: u64) -> u32 { /// /// Per Anthropic docs, only Claude 3+ models support prompt caching. /// Unsupported: claude-2, claude-2.1, claude-instant-*. -fn supports_prompt_cache(name: &str) -> bool { +pub(crate) fn supports_prompt_cache(name: &str) -> bool { let lower = name.to_lowercase(); // Strip optional provider prefix (e.g. "anthropic/claude-...") let model = lower.strip_prefix("anthropic/").unwrap_or(&lower); @@ -1057,17 +1062,32 @@ fn merge_additional_params(rig_req: &mut RigRequest, defaults: Option<&serde_jso /// Build a rig-core CompletionRequest from our internal types. /// -/// When `cache_retention` is not `None`, injects a top-level `cache_control` -/// field via `additional_params`. Rig-core's `AnthropicCompletionRequest` -/// uses `#[serde(flatten)]` on `additional_params`, so the field lands at -/// the request root — which is exactly what Anthropic's **automatic caching** -/// expects. The API auto-places the cache breakpoint at the last cacheable -/// block and moves it forward as conversations grow. +/// When `cache_retention` is not `None`, emits Anthropic cache markers +/// (issue #6984): +/// +/// - A top-level `cache_control` field via `additional_params` (rig-core's +/// `AnthropicCompletionRequest` uses `#[serde(flatten)]` on +/// `additional_params`, so the field lands at the request root). This is +/// Anthropic's **automatic caching**: the API places a breakpoint at the +/// last cacheable block and moves it forward as conversations grow. +/// - An explicit breakpoint on the **last tool definition**, so the tool +/// prefix stays cached even when later parts of the prompt change. rig's +/// typed `ToolDefinition` cannot carry `cache_control`, so the last tool +/// is moved into `additional_params.tools` in Anthropic's native shape — +/// rig appends raw `additional_params.tools` entries after the typed +/// tools, preserving order. +/// +/// The system-prompt and last-message-block breakpoints are rig's own +/// `prompt_caching` flag, set at provider construction for `Short` retention +/// only: rig's typed markers are always plain 5m ephemeral, and a 5m block +/// marker combined with a 1h automatic marker is rejected by the API (TTL +/// conflict on the last block), so `Long` relies on the automatic marker for +/// the conversation tail. #[allow(clippy::too_many_arguments)] fn build_rig_request( preamble: Option, mut history: Vec, - tools: Vec, + mut tools: Vec, tool_choice: Option, temperature: Option, max_tokens: Option, @@ -1083,15 +1103,23 @@ fn build_rig_request( reason: format!("Failed to build chat history: {}", e), })?; - // Inject top-level cache_control for Anthropic automatic prompt caching. - let additional_params = match cache_retention { - CacheRetention::None => None, - CacheRetention::Short => Some(serde_json::json!({ - "cache_control": {"type": "ephemeral"} - })), - CacheRetention::Long => Some(serde_json::json!({ - "cache_control": {"type": "ephemeral", "ttl": "1h"} - })), + let additional_params = match cache_retention.cache_control_json() { + None => None, + Some(marker) => { + let mut params = serde_json::json!({ "cache_control": marker.clone() }); + if let Some(last_tool) = tools.pop() { + // Anthropic-native tool shape: rig maps `parameters` → + // `input_schema` only for typed tools; raw entries pass + // through verbatim. + params["tools"] = serde_json::json!([{ + "name": last_tool.name, + "description": last_tool.description, + "input_schema": last_tool.parameters, + "cache_control": marker, + }]); + } + Some(params) + } }; Ok(RigRequest { @@ -3486,6 +3514,287 @@ mod tests { ); } + fn two_rig_tools() -> Vec { + vec![ + RigToolDefinition { + name: "alpha".to_string(), + description: "First tool".to_string(), + parameters: serde_json::json!({ + "type": "object", + "properties": { "a": { "type": "string" } } + }), + }, + RigToolDefinition { + name: "beta".to_string(), + description: "Second tool".to_string(), + parameters: serde_json::json!({ + "type": "object", + "properties": { "b": { "type": "string" } } + }), + }, + ] + } + + #[test] + fn test_build_rig_request_marks_last_tool_short() { + let req = build_rig_request( + Some("You are helpful.".to_string()), + vec![RigMessage::user("Hello")], + two_rig_tools(), + None, + None, + None, + CacheRetention::Short, + ) + .unwrap(); + + // The last tool moves into additional_params.tools (rig appends those + // after the typed tools, preserving order) carrying the explicit + // cache breakpoint in Anthropic's native shape. + assert_eq!(req.tools.len(), 1, "last tool moves to additional_params"); + assert_eq!(req.tools[0].name, "alpha"); + let params = req.additional_params.expect("additional_params"); + let moved = params["tools"].as_array().expect("raw tools array"); + assert_eq!(moved.len(), 1); + assert_eq!(moved[0]["name"], "beta"); + assert_eq!(moved[0]["description"], "Second tool"); + assert!( + moved[0].get("input_schema").is_some(), + "moved tool must use Anthropic's input_schema key: {moved:?}" + ); + assert!(moved[0].get("parameters").is_none()); + assert_eq!(moved[0]["cache_control"]["type"], "ephemeral"); + assert!(moved[0]["cache_control"].get("ttl").is_none()); + } + + #[test] + fn test_build_rig_request_marks_last_tool_long_with_ttl() { + let req = build_rig_request( + Some("You are helpful.".to_string()), + vec![RigMessage::user("Hello")], + two_rig_tools(), + None, + None, + None, + CacheRetention::Long, + ) + .unwrap(); + + let params = req.additional_params.expect("additional_params"); + let moved = params["tools"].as_array().expect("raw tools array"); + assert_eq!(moved[0]["cache_control"]["type"], "ephemeral"); + assert_eq!(moved[0]["cache_control"]["ttl"], "1h"); + } + + #[test] + fn test_build_rig_request_keeps_tools_typed_when_none() { + let req = build_rig_request( + Some("You are helpful.".to_string()), + vec![RigMessage::user("Hello")], + two_rig_tools(), + None, + None, + None, + CacheRetention::None, + ) + .unwrap(); + + assert_eq!(req.tools.len(), 2, "no tool moves when caching is off"); + assert!(req.additional_params.is_none()); + } + + /// Wire-level pin for the Anthropic explicit-breakpoint layout under + /// `Short` retention (issue #6984): system prompt, last tool, and last + /// message block each carry an ephemeral marker, and the request-level + /// automatic-caching marker is retained (it no-ops on the already-marked + /// last block per Anthropic's compatibility rule). + #[tokio::test] + async fn anthropic_short_retention_places_explicit_cache_breakpoints() { + use rig::client::CompletionClient; + use rig::providers::anthropic; + + let (base_url, captured_body) = capture_one_http_request().await; + let client = anthropic::Client::builder() + .api_key("test-key") + .base_url(&base_url) + .build() + .expect("build Anthropic client"); + let mut model = client.completion_model("claude-opus-4-6"); + // Mirrors create_anthropic_from_registry: Short retention enables + // rig's typed system/last-message breakpoints. + model.prompt_caching = true; + let adapter = + RigAdapter::new(model, "claude-opus-4-6").with_cache_retention(CacheRetention::Short); + + let request = ToolCompletionRequest::new( + vec![ + ChatMessage::system("You are helpful."), + ChatMessage::user("First question"), + ChatMessage::assistant("First answer"), + ChatMessage::user("Second question"), + ], + vec![ + IronToolDefinition { + name: "alpha".to_string(), + description: "First tool".to_string(), + parameters: serde_json::json!({ + "type": "object", + "properties": { "a": { "type": "string" } } + }), + }, + IronToolDefinition { + name: "beta".to_string(), + description: "Second tool".to_string(), + parameters: serde_json::json!({ + "type": "object", + "properties": { "b": { "type": "string" } } + }), + }, + ], + ); + let _ = adapter.complete_with_tools(request).await; + + let body: serde_json::Value = + serde_json::from_str(&captured_request_body(captured_body).await) + .expect("captured body is JSON"); + + assert_eq!(body["cache_control"]["type"], "ephemeral"); + assert!(body["cache_control"].get("ttl").is_none()); + + let system = body["system"] + .as_array() + .expect("system serialized as blocks"); + assert_eq!( + system.last().expect("system block")["cache_control"]["type"], + "ephemeral" + ); + + let tools = body["tools"].as_array().expect("tools array"); + assert_eq!(tools.len(), 2, "both tools reach the wire: {body}"); + assert_eq!(tools[0]["name"], "alpha"); + assert!(tools[0].get("cache_control").is_none()); + assert_eq!(tools[1]["name"], "beta", "order preserved after the move"); + assert_eq!(tools[1]["cache_control"]["type"], "ephemeral"); + assert!( + tools[1].get("input_schema").is_some(), + "moved tool keeps Anthropic's input_schema key: {tools:?}" + ); + + let messages = body["messages"].as_array().expect("messages"); + let last_content = messages.last().expect("last message")["content"] + .as_array() + .expect("last message content blocks"); + assert_eq!( + last_content.last().expect("content block")["cache_control"]["type"], + "ephemeral" + ); + } + + /// Wire-level pin for `Long` retention: rig's typed markers cannot carry + /// a TTL, so only the request-level automatic marker and the raw last-tool + /// marker are emitted, both at 1h. Mixing a 5m block marker with a 1h + /// automatic marker would be rejected by the API (TTL conflict on the + /// last block), so system/message blocks must stay unmarked. + #[tokio::test] + async fn anthropic_long_retention_uses_1h_markers_without_block_conflicts() { + use rig::client::CompletionClient; + use rig::providers::anthropic; + + let (base_url, captured_body) = capture_one_http_request().await; + let client = anthropic::Client::builder() + .api_key("test-key") + .base_url(&base_url) + .build() + .expect("build Anthropic client"); + let model = client.completion_model("claude-opus-4-6"); + let adapter = + RigAdapter::new(model, "claude-opus-4-6").with_cache_retention(CacheRetention::Long); + + let request = ToolCompletionRequest::new( + vec![ + ChatMessage::system("You are helpful."), + ChatMessage::user("Question"), + ], + vec![IronToolDefinition { + name: "alpha".to_string(), + description: "First tool".to_string(), + parameters: serde_json::json!({ + "type": "object", + "properties": { "a": { "type": "string" } } + }), + }], + ); + let _ = adapter.complete_with_tools(request).await; + + let body: serde_json::Value = + serde_json::from_str(&captured_request_body(captured_body).await) + .expect("captured body is JSON"); + + assert_eq!(body["cache_control"]["type"], "ephemeral"); + assert_eq!(body["cache_control"]["ttl"], "1h"); + + let tools = body["tools"].as_array().expect("tools array"); + assert_eq!(tools.len(), 1); + assert_eq!(tools[0]["cache_control"]["ttl"], "1h"); + + for block in body["system"].as_array().expect("system blocks") { + assert!( + block.get("cache_control").is_none(), + "no 5m block marker may precede the 1h automatic marker: {body}" + ); + } + for message in body["messages"].as_array().expect("messages") { + if let Some(blocks) = message["content"].as_array() { + for block in blocks { + assert!( + block.get("cache_control").is_none(), + "message blocks must stay unmarked under Long retention: {body}" + ); + } + } + } + } + + /// Wire-level pin: retention `None` emits no cache_control anywhere. + #[tokio::test] + async fn anthropic_no_retention_emits_no_cache_control() { + use rig::client::CompletionClient; + use rig::providers::anthropic; + + let (base_url, captured_body) = capture_one_http_request().await; + let client = anthropic::Client::builder() + .api_key("test-key") + .base_url(&base_url) + .build() + .expect("build Anthropic client"); + let adapter = RigAdapter::new( + client.completion_model("claude-opus-4-6"), + "claude-opus-4-6", + ); + + let request = ToolCompletionRequest::new( + vec![ + ChatMessage::system("You are helpful."), + ChatMessage::user("Question"), + ], + vec![IronToolDefinition { + name: "alpha".to_string(), + description: "First tool".to_string(), + parameters: serde_json::json!({ + "type": "object", + "properties": { "a": { "type": "string" } } + }), + }], + ); + let _ = adapter.complete_with_tools(request).await; + + let body = captured_request_body(captured_body).await; + assert!( + !body.contains("cache_control"), + "no cache_control may be emitted when caching is off: {body}" + ); + } + /// Verify that the multiplier match arms in `RigAdapter::cache_write_multiplier` /// produce the expected values. We use a standalone helper because constructing /// a real `RigAdapter` requires a rig `Model` (which needs network/provider setup). From 6f7047e351af8380dea7234520b74f4467693159 Mon Sep 17 00:00:00 2001 From: Illia Polosukhin Date: Sat, 1 Aug 2026 05:37:50 +0000 Subject: [PATCH 2/4] test(llm): close changed-line coverage gaps on cache breakpoints The Reborn integration-tier changed-coverage gate flagged uncovered branches in #6997: the unsupported-model retention downgrade on both transports, the Image/ToolUse marker arms, the empty-text guard, and the create_anthropic_from_registry wiring. - Extract the duplicated downgrade logic into rig_adapter::effective_cache_retention, shared by lib.rs and the OAuth constructor, with a direct unit test over all branches. - Wire test: an unsupported model (claude-2.1) with Short retention keeps the legacy no-caching shape end-to-end. - Direct apply_cache_breakpoints tests: tool_use tail without system/tools, image tail, empty-text tail, empty transcript. - Construction test driving create_anthropic_from_registry across all retention modes including the downgrade path. - Move the OAuth transport test suite to src/anthropic_oauth/tests.rs (same idiom as rig_adapter/tests/) to stay inside the file-size budget, with the matching coverage exemption entry. Co-Authored-By: Claude Fable 5 --- crates/ironclaw_llm/src/anthropic_oauth.rs | 519 +-------------- .../ironclaw_llm/src/anthropic_oauth/tests.rs | 623 ++++++++++++++++++ crates/ironclaw_llm/src/lib.rs | 36 +- crates/ironclaw_llm/src/rig_adapter.rs | 41 +- tests/integration/coverage-exemptions.toml | 5 + 5 files changed, 702 insertions(+), 522 deletions(-) create mode 100644 crates/ironclaw_llm/src/anthropic_oauth/tests.rs diff --git a/crates/ironclaw_llm/src/anthropic_oauth.rs b/crates/ironclaw_llm/src/anthropic_oauth.rs index 1f812fb52af..a4d13a82983 100644 --- a/crates/ironclaw_llm/src/anthropic_oauth.rs +++ b/crates/ironclaw_llm/src/anthropic_oauth.rs @@ -130,13 +130,8 @@ impl AnthropicOAuthProvider { let unsupported_params: HashSet = config.unsupported_params.iter().cloned().collect(); - let cache_retention = if config.cache_retention != crate::config::CacheRetention::None - && !crate::rig_adapter::supports_prompt_cache(&config.model) - { - crate::config::CacheRetention::None - } else { - config.cache_retention - }; + let cache_retention = + crate::rig_adapter::effective_cache_retention(config.cache_retention, &config.model); Ok(Self { client, @@ -917,511 +912,7 @@ fn extract_response_content(response: &AnthropicResponse) -> ExtractedAnthropicR } } +// The transport/cache-breakpoint test suite lives in its own file so this +// one stays inside the file-size budget: `src/anthropic_oauth/tests.rs`. #[cfg(test)] -mod tests { - use super::*; - use crate::config::CacheRetention; - use crate::provider::ToolDefinition; - - #[tokio::test] - async fn complete_preserves_missing_retry_after_on_headerless_502() { - use tokio::io::{AsyncReadExt, AsyncWriteExt}; - use tokio::net::TcpListener; - - let listener = TcpListener::bind("127.0.0.1:0") - .await - .expect("loopback listener"); - let base_url = format!( - "http://{}", - listener.local_addr().expect("loopback address") - ); - let server = tokio::spawn(async move { - let (mut socket, _) = listener.accept().await.expect("accept request"); - let mut request = vec![0_u8; 4096]; - let _ = socket.read(&mut request).await.expect("read request"); - let body = r#"{"error":{"message":"upstream unavailable"}}"#; - let response = format!( - "HTTP/1.1 502 Bad Gateway\r\ncontent-type: application/json\r\n\ - content-length: {}\r\n\r\n{body}", - body.len() - ); - socket - .write_all(response.as_bytes()) - .await - .expect("write error response"); - }); - - let mut config = RegistryProviderConfig::generic( - crate::registry::ProviderProtocol::Anthropic, - "anthropic_oauth", - None, - base_url, - "claude-test", - ); - config.oauth_token = Some(SecretString::from("test-token".to_string())); - let provider = AnthropicOAuthProvider::new(&config).expect("provider"); - let error = provider - .complete(CompletionRequest::new(vec![ChatMessage::user("hello")])) - .await - .expect_err("scripted provider error"); - server.await.expect("loopback server"); - - assert!(matches!( - error, - LlmError::BadGateway { - provider, - status: 502, - retry_after: None, - } if provider == "anthropic_oauth" - )); - } - - /// One-shot loopback capture server: returns the base URL and a handle - /// resolving to the captured request body. Replies 400 — these tests - /// assert the request wire shape, not response handling. - async fn capture_one_request() -> (String, tokio::sync::oneshot::Receiver) { - use tokio::io::{AsyncReadExt, AsyncWriteExt}; - use tokio::net::TcpListener; - - let listener = TcpListener::bind("127.0.0.1:0") - .await - .expect("loopback listener"); - let base_url = format!( - "http://{}", - listener.local_addr().expect("loopback address") - ); - let (tx, rx) = tokio::sync::oneshot::channel(); - tokio::spawn(async move { - let (mut socket, _) = listener.accept().await.expect("accept request"); - let mut request = Vec::new(); - let mut buffer = [0_u8; 4096]; - loop { - let read = socket.read(&mut buffer).await.expect("read request"); - if read == 0 { - return; - } - request.extend_from_slice(&buffer[..read]); - let Some(header_end) = request.windows(4).position(|bytes| bytes == b"\r\n\r\n") - else { - continue; - }; - let headers = - std::str::from_utf8(&request[..header_end]).expect("request headers are UTF-8"); - let content_length = headers - .lines() - .find_map(|line| { - let (name, value) = line.split_once(':')?; - name.eq_ignore_ascii_case("content-length") - .then_some(value.trim()) - }) - .expect("content length header") - .parse::() - .expect("content length is numeric"); - let body_start = header_end + 4; - if request.len() < body_start + content_length { - continue; - } - tx.send( - String::from_utf8(request[body_start..body_start + content_length].to_vec()) - .expect("request body is UTF-8"), - ) - .expect("test receives captured request"); - socket - .write_all(b"HTTP/1.1 400 Bad Request\r\nContent-Length: 0\r\n\r\n") - .await - .expect("write test response"); - return; - } - }); - (base_url, rx) - } - - fn provider_with_retention( - base_url: &str, - retention: CacheRetention, - ) -> AnthropicOAuthProvider { - let mut config = RegistryProviderConfig::generic( - crate::registry::ProviderProtocol::Anthropic, - "anthropic_oauth", - None, - base_url, - "claude-opus-4-6", - ); - config.oauth_token = Some(SecretString::from("test-token".to_string())); - config.cache_retention = retention; - AnthropicOAuthProvider::new(&config).expect("provider") - } - - async fn captured_json(rx: tokio::sync::oneshot::Receiver) -> serde_json::Value { - let body = tokio::time::timeout(std::time::Duration::from_secs(5), rx) - .await - .expect("request capture timed out") - .expect("captured request body"); - serde_json::from_str(&body).expect("captured body is JSON") - } - - /// Wire-level pin for issue #6984: under `Short` retention the OAuth - /// transport emits the three explicit cache breakpoints — system prompt - /// block, last tool definition, and the last content block of the last - /// message (here a tool_result, the common agent-loop tail). - #[tokio::test] - async fn oauth_short_retention_places_explicit_cache_breakpoints() { - let (base_url, captured) = capture_one_request().await; - let provider = provider_with_retention(&base_url, CacheRetention::Short); - - let request = ToolCompletionRequest::new( - vec![ - ChatMessage::system("You are helpful."), - ChatMessage::user("Run the tool."), - ChatMessage::assistant_with_tool_calls( - None, - vec![ToolCall { - id: "call_1".to_string(), - name: "alpha".to_string(), - arguments: serde_json::json!({}), - ..ToolCall::default() - }], - ), - ChatMessage::tool_result("call_1", "alpha", "tool says hi"), - ], - vec![ - ToolDefinition { - name: "alpha".to_string(), - description: "First tool".to_string(), - parameters: serde_json::json!({"type": "object", "properties": {}}), - }, - ToolDefinition { - name: "beta".to_string(), - description: "Second tool".to_string(), - parameters: serde_json::json!({"type": "object", "properties": {}}), - }, - ], - ); - let _ = provider.complete_with_tools(request).await; - let body = captured_json(captured).await; - - let system = body["system"] - .as_array() - .expect("system serialized as blocks when caching is on"); - assert_eq!( - system.last().expect("system block")["cache_control"]["type"], - "ephemeral" - ); - assert!(system.last().unwrap()["cache_control"].get("ttl").is_none()); - - let tools = body["tools"].as_array().expect("tools array"); - assert_eq!(tools.len(), 2); - assert!(tools[0].get("cache_control").is_none()); - assert_eq!(tools[1]["name"], "beta"); - assert_eq!(tools[1]["cache_control"]["type"], "ephemeral"); - - let messages = body["messages"].as_array().expect("messages"); - let last_content = messages.last().expect("last message")["content"] - .as_array() - .expect("last message content blocks"); - let last_block = last_content.last().expect("content block"); - assert_eq!(last_block["type"], "tool_result"); - assert_eq!(last_block["cache_control"]["type"], "ephemeral"); - // Only the final block carries a marker. - for message in &messages[..messages.len() - 1] { - if let Some(blocks) = message["content"].as_array() { - for block in blocks { - assert!(block.get("cache_control").is_none(), "{body}"); - } - } - } - } - - /// Wire-level pin: `Long` retention stamps a 1h TTL on every breakpoint - /// (uniform TTLs satisfy Anthropic's longer-before-shorter ordering rule). - #[tokio::test] - async fn oauth_long_retention_uses_1h_ttl_markers() { - let (base_url, captured) = capture_one_request().await; - let provider = provider_with_retention(&base_url, CacheRetention::Long); - - let _ = provider - .complete(CompletionRequest::new(vec![ - ChatMessage::system("You are helpful."), - ChatMessage::user("Question"), - ])) - .await; - let body = captured_json(captured).await; - - let system = body["system"].as_array().expect("system blocks"); - assert_eq!(system.last().unwrap()["cache_control"]["ttl"], "1h"); - - let messages = body["messages"].as_array().expect("messages"); - let last_content = messages.last().unwrap()["content"] - .as_array() - .expect("last message content blocks"); - assert_eq!(last_content.last().unwrap()["cache_control"]["ttl"], "1h"); - } - - /// Wire-level pin: retention `None` keeps the legacy wire shape — system - /// as a plain string, and no cache_control anywhere. - #[tokio::test] - async fn oauth_no_retention_keeps_legacy_wire_shape() { - let (base_url, captured) = capture_one_request().await; - let provider = provider_with_retention(&base_url, CacheRetention::None); - - let _ = provider - .complete(CompletionRequest::new(vec![ - ChatMessage::system("You are helpful."), - ChatMessage::user("Question"), - ])) - .await; - let body = captured_json(captured).await; - - assert!( - body["system"].is_string(), - "system stays a plain string when caching is off: {body}" - ); - assert!( - !serde_json::to_string(&body) - .unwrap() - .contains("cache_control"), - "no cache_control may be emitted when caching is off: {body}" - ); - } - - #[test] - fn context_overflow_413_maps_to_context_length_exceeded() { - // A raw HTTP 413 (payload too large) must become ContextLengthExceeded - // so the loop's context-shrink recovery fires. - match context_length_error_for_status(413, "Request Entity Too Large") { - Some(LlmError::ContextLengthExceeded { .. }) => {} - other => panic!("expected ContextLengthExceeded, got {other:?}"), - } - } - - #[test] - fn context_overflow_400_body_maps_to_context_length_exceeded() { - let body = r#"{"type":"error","error":{"type":"invalid_request_error","message":"prompt is too long: 234872 tokens > 200000 maximum"}}"#; - match context_length_error_for_status(400, body) { - Some(LlmError::ContextLengthExceeded { used, limit }) => { - assert_eq!(used, 234872); - assert_eq!(limit, 200000); - } - other => panic!("expected ContextLengthExceeded, got {other:?}"), - } - } - - #[test] - fn unrelated_400_is_not_context_overflow() { - // A plain bad-request (e.g. invalid request shape) must NOT be - // classified as context overflow — the caller falls through to - // RequestFailed. - assert!( - context_length_error_for_status(400, r#"{"error":{"message":"invalid request body"}}"#) - .is_none() - ); - } - - #[test] - fn unrelated_5xx_is_not_context_overflow() { - assert!(context_length_error_for_status(503, "service unavailable").is_none()); - } - - #[test] - fn test_convert_messages_extracts_system() { - let messages = vec![ - ChatMessage::system("You are helpful."), - ChatMessage::user("Hello"), - ]; - let (system, msgs) = convert_messages(messages); - assert_eq!(system, Some("You are helpful.".to_string())); - assert_eq!(msgs.len(), 1); - assert_eq!(msgs[0].role, "user"); - } - - #[test] - fn test_convert_messages_multiple_systems() { - let messages = vec![ - ChatMessage::system("System 1"), - ChatMessage::system("System 2"), - ChatMessage::user("Hello"), - ]; - let (system, msgs) = convert_messages(messages); - assert_eq!(system, Some("System 1\n\nSystem 2".to_string())); - assert_eq!(msgs.len(), 1); - } - - #[test] - fn test_convert_messages_user_image_becomes_base64_image_block() { - let messages = vec![ChatMessage::user_with_parts( - "what is this?", - vec![ContentPart::ImageUrl { - image_url: crate::provider::ImageUrl { - url: "data:image/png;base64,AQIDBA==".to_string(), - detail: None, - }, - }], - )]; - let (_system, msgs) = convert_messages(messages); - assert_eq!(msgs.len(), 1); - // Text rides as the first block, the image as a base64 `image` block. - let value = serde_json::to_value(&msgs[0]).expect("serialize"); - let blocks = value["content"].as_array().expect("content blocks"); - assert_eq!(blocks.len(), 2); - assert_eq!(blocks[0]["type"], "text"); - assert_eq!(blocks[0]["text"], "what is this?"); - assert_eq!(blocks[1]["type"], "image"); - assert_eq!(blocks[1]["source"]["type"], "base64"); - assert_eq!(blocks[1]["source"]["media_type"], "image/png"); - assert_eq!(blocks[1]["source"]["data"], "AQIDBA=="); - } - - #[test] - fn test_convert_messages_text_only_user_stays_a_string() { - let messages = vec![ChatMessage::user("just text")]; - let (_system, msgs) = convert_messages(messages); - let value = serde_json::to_value(&msgs[0]).expect("serialize"); - // No inline images → compact string content, not a blocks array. - assert_eq!(value["content"], "just text"); - } - - #[test] - fn test_convert_messages_tool_calls() { - let tool_calls = vec![ToolCall { - id: "call_1".to_string(), - name: "search".to_string(), - arguments: serde_json::json!({"q": "test"}), - reasoning: None, - signature: None, - arguments_parse_error: None, - }]; - let messages = vec![ - ChatMessage::user("Search for test"), - ChatMessage::assistant_with_tool_calls(Some("Let me search.".to_string()), tool_calls), - ChatMessage::tool_result("call_1", "search", "found it"), - ]; - let (system, msgs) = convert_messages(messages); - assert!(system.is_none()); - assert_eq!(msgs.len(), 3); - assert_eq!(msgs[0].role, "user"); - assert_eq!(msgs[1].role, "assistant"); - // Tool result should be a user message - assert_eq!(msgs[2].role, "user"); - } - - #[test] - fn test_extract_response_text_only() { - let response = AnthropicResponse { - content: vec![AnthropicResponseBlock::Text { - text: "Hello!".to_string(), - }], - stop_reason: Some("end_turn".to_string()), - usage: AnthropicUsage { - input_tokens: 10, - output_tokens: 5, - cache_creation_input_tokens: 0, - cache_read_input_tokens: 0, - }, - }; - let extracted = extract_response_content(&response); - assert_eq!(extracted.content, Some("Hello!".to_string())); - assert!(extracted.tool_calls.is_empty()); - } - - #[test] - fn test_extract_response_with_tool_use() { - let response = AnthropicResponse { - content: vec![ - AnthropicResponseBlock::Text { - text: "Let me search.".to_string(), - }, - AnthropicResponseBlock::ToolUse { - id: "call_1".to_string(), - name: "search".to_string(), - input: serde_json::json!({"q": "test"}), - }, - ], - stop_reason: Some("tool_use".to_string()), - usage: AnthropicUsage { - input_tokens: 20, - output_tokens: 15, - cache_creation_input_tokens: 0, - cache_read_input_tokens: 0, - }, - }; - let extracted = extract_response_content(&response); - assert_eq!(extracted.content, Some("Let me search.".to_string())); - assert_eq!(extracted.tool_calls.len(), 1); - assert_eq!(extracted.tool_calls[0].name, "search"); - } - - #[test] - fn test_extract_response_preserves_thinking_as_reasoning() { - let response = AnthropicResponse { - content: vec![ - AnthropicResponseBlock::Thinking { - thinking: Some("Raw thinking".to_string()), - summary: Some("Summarized thinking".to_string()), - _signature: Some("sig".to_string()), - }, - AnthropicResponseBlock::Text { - text: "Done.".to_string(), - }, - ], - stop_reason: Some("end_turn".to_string()), - usage: AnthropicUsage { - input_tokens: 20, - output_tokens: 15, - cache_creation_input_tokens: 0, - cache_read_input_tokens: 0, - }, - }; - let extracted = extract_response_content(&response); - assert_eq!(extracted.content, Some("Done.".to_string())); - assert_eq!(extracted.reasoning, Some("Summarized thinking".to_string())); - } - - #[test] - fn test_extract_response_uses_thinking_when_summary_absent() { - let response = AnthropicResponse { - content: vec![ - AnthropicResponseBlock::Thinking { - thinking: Some("Raw thinking fallback".to_string()), - summary: None, - _signature: Some("sig".to_string()), - }, - AnthropicResponseBlock::Text { - text: "Done.".to_string(), - }, - ], - stop_reason: Some("end_turn".to_string()), - usage: AnthropicUsage { - input_tokens: 20, - output_tokens: 15, - cache_creation_input_tokens: 0, - cache_read_input_tokens: 0, - }, - }; - - let extracted = extract_response_content(&response); - - assert_eq!(extracted.content, Some("Done.".to_string())); - assert_eq!( - extracted.reasoning, - Some("Raw thinking fallback".to_string()) - ); - } - - /// Regression test for #1136: token field must be mutable via RwLock - /// so that a refreshed token persists across subsequent requests. - #[test] - fn test_token_update_persists() { - let original = SecretString::from("old_token".to_string()); - let token = std::sync::RwLock::new(original); - - // Read the original - assert_eq!(token.read().unwrap().expose_secret(), "old_token"); - - // Simulate a successful refresh - let refreshed = SecretString::from("new_token".to_string()); - *token.write().unwrap() = refreshed; - - // Subsequent reads see the updated token - assert_eq!(token.read().unwrap().expose_secret(), "new_token"); - } -} +mod tests; diff --git a/crates/ironclaw_llm/src/anthropic_oauth/tests.rs b/crates/ironclaw_llm/src/anthropic_oauth/tests.rs new file mode 100644 index 00000000000..edf4b199410 --- /dev/null +++ b/crates/ironclaw_llm/src/anthropic_oauth/tests.rs @@ -0,0 +1,623 @@ +//! Wire-shape and breakpoint tests for the Anthropic OAuth transport. +//! +//! Split out of `anthropic_oauth.rs` (arch file-size budget): loopback +//! capture-server tests pinning the request JSON — including the #6984 +//! `cache_control` breakpoint layout — plus direct branch tests for +//! `apply_cache_breakpoints`. + +use super::*; +use crate::config::CacheRetention; +use crate::provider::ToolDefinition; + +#[tokio::test] +async fn complete_preserves_missing_retry_after_on_headerless_502() { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::TcpListener; + + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("loopback listener"); + let base_url = format!( + "http://{}", + listener.local_addr().expect("loopback address") + ); + let server = tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.expect("accept request"); + let mut request = vec![0_u8; 4096]; + let _ = socket.read(&mut request).await.expect("read request"); + let body = r#"{"error":{"message":"upstream unavailable"}}"#; + let response = format!( + "HTTP/1.1 502 Bad Gateway\r\ncontent-type: application/json\r\n\ + content-length: {}\r\n\r\n{body}", + body.len() + ); + socket + .write_all(response.as_bytes()) + .await + .expect("write error response"); + }); + + let mut config = RegistryProviderConfig::generic( + crate::registry::ProviderProtocol::Anthropic, + "anthropic_oauth", + None, + base_url, + "claude-test", + ); + config.oauth_token = Some(SecretString::from("test-token".to_string())); + let provider = AnthropicOAuthProvider::new(&config).expect("provider"); + let error = provider + .complete(CompletionRequest::new(vec![ChatMessage::user("hello")])) + .await + .expect_err("scripted provider error"); + server.await.expect("loopback server"); + + assert!(matches!( + error, + LlmError::BadGateway { + provider, + status: 502, + retry_after: None, + } if provider == "anthropic_oauth" + )); +} + +/// One-shot loopback capture server: returns the base URL and a handle +/// resolving to the captured request body. Replies 400 — these tests +/// assert the request wire shape, not response handling. +async fn capture_one_request() -> (String, tokio::sync::oneshot::Receiver) { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::TcpListener; + + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("loopback listener"); + let base_url = format!( + "http://{}", + listener.local_addr().expect("loopback address") + ); + let (tx, rx) = tokio::sync::oneshot::channel(); + tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.expect("accept request"); + let mut request = Vec::new(); + let mut buffer = [0_u8; 4096]; + loop { + let read = socket.read(&mut buffer).await.expect("read request"); + if read == 0 { + return; + } + request.extend_from_slice(&buffer[..read]); + let Some(header_end) = request.windows(4).position(|bytes| bytes == b"\r\n\r\n") else { + continue; + }; + let headers = + std::str::from_utf8(&request[..header_end]).expect("request headers are UTF-8"); + let content_length = headers + .lines() + .find_map(|line| { + let (name, value) = line.split_once(':')?; + name.eq_ignore_ascii_case("content-length") + .then_some(value.trim()) + }) + .expect("content length header") + .parse::() + .expect("content length is numeric"); + let body_start = header_end + 4; + if request.len() < body_start + content_length { + continue; + } + tx.send( + String::from_utf8(request[body_start..body_start + content_length].to_vec()) + .expect("request body is UTF-8"), + ) + .expect("test receives captured request"); + socket + .write_all(b"HTTP/1.1 400 Bad Request\r\nContent-Length: 0\r\n\r\n") + .await + .expect("write test response"); + return; + } + }); + (base_url, rx) +} + +fn provider_with_retention(base_url: &str, retention: CacheRetention) -> AnthropicOAuthProvider { + let mut config = RegistryProviderConfig::generic( + crate::registry::ProviderProtocol::Anthropic, + "anthropic_oauth", + None, + base_url, + "claude-opus-4-6", + ); + config.oauth_token = Some(SecretString::from("test-token".to_string())); + config.cache_retention = retention; + AnthropicOAuthProvider::new(&config).expect("provider") +} + +async fn captured_json(rx: tokio::sync::oneshot::Receiver) -> serde_json::Value { + let body = tokio::time::timeout(std::time::Duration::from_secs(5), rx) + .await + .expect("request capture timed out") + .expect("captured request body"); + serde_json::from_str(&body).expect("captured body is JSON") +} + +/// Wire-level pin for issue #6984: under `Short` retention the OAuth +/// transport emits the three explicit cache breakpoints — system prompt +/// block, last tool definition, and the last content block of the last +/// message (here a tool_result, the common agent-loop tail). +#[tokio::test] +async fn oauth_short_retention_places_explicit_cache_breakpoints() { + let (base_url, captured) = capture_one_request().await; + let provider = provider_with_retention(&base_url, CacheRetention::Short); + + let request = ToolCompletionRequest::new( + vec![ + ChatMessage::system("You are helpful."), + ChatMessage::user("Run the tool."), + ChatMessage::assistant_with_tool_calls( + None, + vec![ToolCall { + id: "call_1".to_string(), + name: "alpha".to_string(), + arguments: serde_json::json!({}), + ..ToolCall::default() + }], + ), + ChatMessage::tool_result("call_1", "alpha", "tool says hi"), + ], + vec![ + ToolDefinition { + name: "alpha".to_string(), + description: "First tool".to_string(), + parameters: serde_json::json!({"type": "object", "properties": {}}), + }, + ToolDefinition { + name: "beta".to_string(), + description: "Second tool".to_string(), + parameters: serde_json::json!({"type": "object", "properties": {}}), + }, + ], + ); + let _ = provider.complete_with_tools(request).await; + let body = captured_json(captured).await; + + let system = body["system"] + .as_array() + .expect("system serialized as blocks when caching is on"); + assert_eq!( + system.last().expect("system block")["cache_control"]["type"], + "ephemeral" + ); + assert!(system.last().unwrap()["cache_control"].get("ttl").is_none()); + + let tools = body["tools"].as_array().expect("tools array"); + assert_eq!(tools.len(), 2); + assert!(tools[0].get("cache_control").is_none()); + assert_eq!(tools[1]["name"], "beta"); + assert_eq!(tools[1]["cache_control"]["type"], "ephemeral"); + + let messages = body["messages"].as_array().expect("messages"); + let last_content = messages.last().expect("last message")["content"] + .as_array() + .expect("last message content blocks"); + let last_block = last_content.last().expect("content block"); + assert_eq!(last_block["type"], "tool_result"); + assert_eq!(last_block["cache_control"]["type"], "ephemeral"); + // Only the final block carries a marker. + for message in &messages[..messages.len() - 1] { + if let Some(blocks) = message["content"].as_array() { + for block in blocks { + assert!(block.get("cache_control").is_none(), "{body}"); + } + } + } +} + +/// Wire-level pin: `Long` retention stamps a 1h TTL on every breakpoint +/// (uniform TTLs satisfy Anthropic's longer-before-shorter ordering rule). +#[tokio::test] +async fn oauth_long_retention_uses_1h_ttl_markers() { + let (base_url, captured) = capture_one_request().await; + let provider = provider_with_retention(&base_url, CacheRetention::Long); + + let _ = provider + .complete(CompletionRequest::new(vec![ + ChatMessage::system("You are helpful."), + ChatMessage::user("Question"), + ])) + .await; + let body = captured_json(captured).await; + + let system = body["system"].as_array().expect("system blocks"); + assert_eq!(system.last().unwrap()["cache_control"]["ttl"], "1h"); + + let messages = body["messages"].as_array().expect("messages"); + let last_content = messages.last().unwrap()["content"] + .as_array() + .expect("last message content blocks"); + assert_eq!(last_content.last().unwrap()["cache_control"]["ttl"], "1h"); +} + +/// Wire-level pin: a model without prompt-cache support (claude-2 era) +/// downgrades to no caching at construction, keeping the legacy shape. +#[tokio::test] +async fn oauth_unsupported_model_downgrades_to_no_caching() { + let (base_url, captured) = capture_one_request().await; + let mut config = RegistryProviderConfig::generic( + crate::registry::ProviderProtocol::Anthropic, + "anthropic_oauth", + None, + &base_url, + "claude-2.1", + ); + config.oauth_token = Some(SecretString::from("test-token".to_string())); + config.cache_retention = CacheRetention::Short; + let provider = AnthropicOAuthProvider::new(&config).expect("provider"); + + let _ = provider + .complete(CompletionRequest::new(vec![ + ChatMessage::system("You are helpful."), + ChatMessage::user("Question"), + ])) + .await; + let body = captured_json(captured).await; + + assert!(body["system"].is_string(), "{body}"); + assert!( + !serde_json::to_string(&body) + .unwrap() + .contains("cache_control") + ); +} + +fn request_with_messages(messages: Vec) -> AnthropicRequest { + AnthropicRequest { + model: "claude-opus-4-6".to_string(), + messages, + system: None, + max_tokens: 128, + temperature: None, + thinking: None, + tools: None, + tool_choice: None, + } +} + +/// Branch coverage for `apply_cache_breakpoints`: a tool_use-tailed +/// assistant message (no system, no tools) gets its last block marked. +#[test] +fn apply_breakpoints_marks_tool_use_tail_without_system_or_tools() { + let mut request = request_with_messages(vec![AnthropicMessage { + role: "assistant".to_string(), + content: AnthropicContent::Blocks(vec![AnthropicContentBlock::ToolUse { + id: "call_1".to_string(), + name: "alpha".to_string(), + input: serde_json::json!({}), + cache_control: None, + }]), + }]); + apply_cache_breakpoints(&mut request, CacheRetention::Short); + + assert!(request.system.is_none()); + let AnthropicContent::Blocks(blocks) = &request.messages[0].content else { + panic!("blocks expected"); + }; + let AnthropicContentBlock::ToolUse { cache_control, .. } = &blocks[0] else { + panic!("tool_use expected"); + }; + assert!(cache_control.is_some()); +} + +/// Branch coverage: an image-tailed message gets its image block marked. +#[test] +fn apply_breakpoints_marks_image_tail() { + let mut request = request_with_messages(vec![AnthropicMessage { + role: "user".to_string(), + content: AnthropicContent::Blocks(vec![AnthropicContentBlock::Image { + source: AnthropicImageSource { + source_type: "base64", + media_type: "image/png".to_string(), + data: "aGk=".to_string(), + }, + cache_control: None, + }]), + }]); + apply_cache_breakpoints(&mut request, CacheRetention::Short); + + let AnthropicContent::Blocks(blocks) = &request.messages[0].content else { + panic!("blocks expected"); + }; + let AnthropicContentBlock::Image { cache_control, .. } = &blocks[0] else { + panic!("image expected"); + }; + assert!(cache_control.is_some()); +} + +/// Branch coverage: an empty trailing text message stays in string form — +/// the API rejects cache_control on empty text blocks — and an empty +/// message list is a no-op. +#[test] +fn apply_breakpoints_skips_empty_text_tail_and_empty_transcript() { + let mut request = request_with_messages(vec![AnthropicMessage { + role: "user".to_string(), + content: AnthropicContent::Text(String::new()), + }]); + apply_cache_breakpoints(&mut request, CacheRetention::Short); + assert!(matches!( + &request.messages[0].content, + AnthropicContent::Text(text) if text.is_empty() + )); + + let mut empty = request_with_messages(Vec::new()); + apply_cache_breakpoints(&mut empty, CacheRetention::Long); + assert!(empty.messages.is_empty()); +} + +/// Wire-level pin: retention `None` keeps the legacy wire shape — system +/// as a plain string, and no cache_control anywhere. +#[tokio::test] +async fn oauth_no_retention_keeps_legacy_wire_shape() { + let (base_url, captured) = capture_one_request().await; + let provider = provider_with_retention(&base_url, CacheRetention::None); + + let _ = provider + .complete(CompletionRequest::new(vec![ + ChatMessage::system("You are helpful."), + ChatMessage::user("Question"), + ])) + .await; + let body = captured_json(captured).await; + + assert!( + body["system"].is_string(), + "system stays a plain string when caching is off: {body}" + ); + assert!( + !serde_json::to_string(&body) + .unwrap() + .contains("cache_control"), + "no cache_control may be emitted when caching is off: {body}" + ); +} + +#[test] +fn context_overflow_413_maps_to_context_length_exceeded() { + // A raw HTTP 413 (payload too large) must become ContextLengthExceeded + // so the loop's context-shrink recovery fires. + match context_length_error_for_status(413, "Request Entity Too Large") { + Some(LlmError::ContextLengthExceeded { .. }) => {} + other => panic!("expected ContextLengthExceeded, got {other:?}"), + } +} + +#[test] +fn context_overflow_400_body_maps_to_context_length_exceeded() { + let body = r#"{"type":"error","error":{"type":"invalid_request_error","message":"prompt is too long: 234872 tokens > 200000 maximum"}}"#; + match context_length_error_for_status(400, body) { + Some(LlmError::ContextLengthExceeded { used, limit }) => { + assert_eq!(used, 234872); + assert_eq!(limit, 200000); + } + other => panic!("expected ContextLengthExceeded, got {other:?}"), + } +} + +#[test] +fn unrelated_400_is_not_context_overflow() { + // A plain bad-request (e.g. invalid request shape) must NOT be + // classified as context overflow — the caller falls through to + // RequestFailed. + assert!( + context_length_error_for_status(400, r#"{"error":{"message":"invalid request body"}}"#) + .is_none() + ); +} + +#[test] +fn unrelated_5xx_is_not_context_overflow() { + assert!(context_length_error_for_status(503, "service unavailable").is_none()); +} + +#[test] +fn test_convert_messages_extracts_system() { + let messages = vec![ + ChatMessage::system("You are helpful."), + ChatMessage::user("Hello"), + ]; + let (system, msgs) = convert_messages(messages); + assert_eq!(system, Some("You are helpful.".to_string())); + assert_eq!(msgs.len(), 1); + assert_eq!(msgs[0].role, "user"); +} + +#[test] +fn test_convert_messages_multiple_systems() { + let messages = vec![ + ChatMessage::system("System 1"), + ChatMessage::system("System 2"), + ChatMessage::user("Hello"), + ]; + let (system, msgs) = convert_messages(messages); + assert_eq!(system, Some("System 1\n\nSystem 2".to_string())); + assert_eq!(msgs.len(), 1); +} + +#[test] +fn test_convert_messages_user_image_becomes_base64_image_block() { + let messages = vec![ChatMessage::user_with_parts( + "what is this?", + vec![ContentPart::ImageUrl { + image_url: crate::provider::ImageUrl { + url: "data:image/png;base64,AQIDBA==".to_string(), + detail: None, + }, + }], + )]; + let (_system, msgs) = convert_messages(messages); + assert_eq!(msgs.len(), 1); + // Text rides as the first block, the image as a base64 `image` block. + let value = serde_json::to_value(&msgs[0]).expect("serialize"); + let blocks = value["content"].as_array().expect("content blocks"); + assert_eq!(blocks.len(), 2); + assert_eq!(blocks[0]["type"], "text"); + assert_eq!(blocks[0]["text"], "what is this?"); + assert_eq!(blocks[1]["type"], "image"); + assert_eq!(blocks[1]["source"]["type"], "base64"); + assert_eq!(blocks[1]["source"]["media_type"], "image/png"); + assert_eq!(blocks[1]["source"]["data"], "AQIDBA=="); +} + +#[test] +fn test_convert_messages_text_only_user_stays_a_string() { + let messages = vec![ChatMessage::user("just text")]; + let (_system, msgs) = convert_messages(messages); + let value = serde_json::to_value(&msgs[0]).expect("serialize"); + // No inline images → compact string content, not a blocks array. + assert_eq!(value["content"], "just text"); +} + +#[test] +fn test_convert_messages_tool_calls() { + let tool_calls = vec![ToolCall { + id: "call_1".to_string(), + name: "search".to_string(), + arguments: serde_json::json!({"q": "test"}), + reasoning: None, + signature: None, + arguments_parse_error: None, + }]; + let messages = vec![ + ChatMessage::user("Search for test"), + ChatMessage::assistant_with_tool_calls(Some("Let me search.".to_string()), tool_calls), + ChatMessage::tool_result("call_1", "search", "found it"), + ]; + let (system, msgs) = convert_messages(messages); + assert!(system.is_none()); + assert_eq!(msgs.len(), 3); + assert_eq!(msgs[0].role, "user"); + assert_eq!(msgs[1].role, "assistant"); + // Tool result should be a user message + assert_eq!(msgs[2].role, "user"); +} + +#[test] +fn test_extract_response_text_only() { + let response = AnthropicResponse { + content: vec![AnthropicResponseBlock::Text { + text: "Hello!".to_string(), + }], + stop_reason: Some("end_turn".to_string()), + usage: AnthropicUsage { + input_tokens: 10, + output_tokens: 5, + cache_creation_input_tokens: 0, + cache_read_input_tokens: 0, + }, + }; + let extracted = extract_response_content(&response); + assert_eq!(extracted.content, Some("Hello!".to_string())); + assert!(extracted.tool_calls.is_empty()); +} + +#[test] +fn test_extract_response_with_tool_use() { + let response = AnthropicResponse { + content: vec![ + AnthropicResponseBlock::Text { + text: "Let me search.".to_string(), + }, + AnthropicResponseBlock::ToolUse { + id: "call_1".to_string(), + name: "search".to_string(), + input: serde_json::json!({"q": "test"}), + }, + ], + stop_reason: Some("tool_use".to_string()), + usage: AnthropicUsage { + input_tokens: 20, + output_tokens: 15, + cache_creation_input_tokens: 0, + cache_read_input_tokens: 0, + }, + }; + let extracted = extract_response_content(&response); + assert_eq!(extracted.content, Some("Let me search.".to_string())); + assert_eq!(extracted.tool_calls.len(), 1); + assert_eq!(extracted.tool_calls[0].name, "search"); +} + +#[test] +fn test_extract_response_preserves_thinking_as_reasoning() { + let response = AnthropicResponse { + content: vec![ + AnthropicResponseBlock::Thinking { + thinking: Some("Raw thinking".to_string()), + summary: Some("Summarized thinking".to_string()), + _signature: Some("sig".to_string()), + }, + AnthropicResponseBlock::Text { + text: "Done.".to_string(), + }, + ], + stop_reason: Some("end_turn".to_string()), + usage: AnthropicUsage { + input_tokens: 20, + output_tokens: 15, + cache_creation_input_tokens: 0, + cache_read_input_tokens: 0, + }, + }; + let extracted = extract_response_content(&response); + assert_eq!(extracted.content, Some("Done.".to_string())); + assert_eq!(extracted.reasoning, Some("Summarized thinking".to_string())); +} + +#[test] +fn test_extract_response_uses_thinking_when_summary_absent() { + let response = AnthropicResponse { + content: vec![ + AnthropicResponseBlock::Thinking { + thinking: Some("Raw thinking fallback".to_string()), + summary: None, + _signature: Some("sig".to_string()), + }, + AnthropicResponseBlock::Text { + text: "Done.".to_string(), + }, + ], + stop_reason: Some("end_turn".to_string()), + usage: AnthropicUsage { + input_tokens: 20, + output_tokens: 15, + cache_creation_input_tokens: 0, + cache_read_input_tokens: 0, + }, + }; + + let extracted = extract_response_content(&response); + + assert_eq!(extracted.content, Some("Done.".to_string())); + assert_eq!( + extracted.reasoning, + Some("Raw thinking fallback".to_string()) + ); +} + +/// Regression test for #1136: token field must be mutable via RwLock +/// so that a refreshed token persists across subsequent requests. +#[test] +fn test_token_update_persists() { + let original = SecretString::from("old_token".to_string()); + let token = std::sync::RwLock::new(original); + + // Read the original + assert_eq!(token.read().unwrap().expose_secret(), "old_token"); + + // Simulate a successful refresh + let refreshed = SecretString::from("new_token".to_string()); + *token.write().unwrap() = refreshed; + + // Subsequent reads see the updated token + assert_eq!(token.read().unwrap().expose_secret(), "new_token"); +} diff --git a/crates/ironclaw_llm/src/lib.rs b/crates/ironclaw_llm/src/lib.rs index e29e1e570e1..88116f076a8 100644 --- a/crates/ironclaw_llm/src/lib.rs +++ b/crates/ironclaw_llm/src/lib.rs @@ -505,13 +505,8 @@ fn create_anthropic_from_registry( // Downgrade retention up front for models without prompt-cache support so // the rig `prompt_caching` flag below agrees with the adapter's own // `with_cache_retention` validation. - let cache_retention = if config.cache_retention != CacheRetention::None - && !rig_adapter::supports_prompt_cache(&config.model) - { - CacheRetention::None - } else { - config.cache_retention - }; + let cache_retention = + rig_adapter::effective_cache_retention(config.cache_retention, &config.model); let mut model = client.completion_model(&config.model); @@ -1863,4 +1858,31 @@ mod tests { forwarding `request_timeout_secs` to `provider_http_client`.", ); } + + /// Construction-path coverage for the Anthropic cache wiring (#6984): + /// every retention mode builds the rig provider, including the + /// unsupported-model downgrade that disables rig's typed breakpoints. + #[test] + fn anthropic_registry_provider_builds_for_every_cache_retention() { + use crate::config::CacheRetention; + + for (model, retention) in [ + ("claude-opus-4-6", CacheRetention::Short), + ("claude-opus-4-6", CacheRetention::Long), + ("claude-opus-4-6", CacheRetention::None), + ("claude-2.1", CacheRetention::Short), + ] { + let mut config = RegistryProviderConfig::generic( + crate::registry::ProviderProtocol::Anthropic, + "anthropic", + Some(secrecy::SecretString::from("sk-test".to_string())), + "http://127.0.0.1:9", + model, + ); + config.cache_retention = retention; + let provider = create_anthropic_from_registry(&config, 5) + .expect("anthropic provider construction"); + assert_eq!(provider.model_name(), model); + } + } } diff --git a/crates/ironclaw_llm/src/rig_adapter.rs b/crates/ironclaw_llm/src/rig_adapter.rs index c8d0430038a..74888c225c3 100644 --- a/crates/ironclaw_llm/src/rig_adapter.rs +++ b/crates/ironclaw_llm/src/rig_adapter.rs @@ -912,11 +912,26 @@ fn saturate_u32(val: u64) -> u32 { val.min(u32::MAX as u64) as u32 } +/// Downgrade a requested cache retention to `None` for models without +/// prompt-cache support. Shared by both Anthropic transports so the +/// construction-time decision (rig `prompt_caching` flag, OAuth breakpoint +/// application) always agrees with `with_cache_retention`'s validation. +pub(crate) fn effective_cache_retention( + retention: CacheRetention, + model_name: &str, +) -> CacheRetention { + if retention != CacheRetention::None && !supports_prompt_cache(model_name) { + CacheRetention::None + } else { + retention + } +} + /// Returns `true` if the model supports Anthropic prompt caching. /// /// Per Anthropic docs, only Claude 3+ models support prompt caching. /// Unsupported: claude-2, claude-2.1, claude-instant-*. -pub(crate) fn supports_prompt_cache(name: &str) -> bool { +fn supports_prompt_cache(name: &str) -> bool { let lower = name.to_lowercase(); // Strip optional provider prefix (e.g. "anthropic/claude-...") let model = lower.strip_prefix("anthropic/").unwrap_or(&lower); @@ -3586,6 +3601,30 @@ mod tests { assert_eq!(moved[0]["cache_control"]["ttl"], "1h"); } + #[test] + fn test_effective_cache_retention_downgrades_unsupported_models() { + assert_eq!( + effective_cache_retention(CacheRetention::Short, "claude-2.1"), + CacheRetention::None + ); + assert_eq!( + effective_cache_retention(CacheRetention::Long, "claude-instant-1.2"), + CacheRetention::None + ); + assert_eq!( + effective_cache_retention(CacheRetention::Short, "claude-opus-4-6"), + CacheRetention::Short + ); + assert_eq!( + effective_cache_retention(CacheRetention::Long, "anthropic/claude-sonnet-4-5"), + CacheRetention::Long + ); + assert_eq!( + effective_cache_retention(CacheRetention::None, "claude-opus-4-6"), + CacheRetention::None + ); + } + #[test] fn test_build_rig_request_keeps_tools_typed_when_none() { let req = build_rig_request( diff --git a/tests/integration/coverage-exemptions.toml b/tests/integration/coverage-exemptions.toml index ca0e55ba28b..3663fea4dde 100644 --- a/tests/integration/coverage-exemptions.toml +++ b/tests/integration/coverage-exemptions.toml @@ -47,6 +47,11 @@ module = "crates/ironclaw_llm/src/rig_adapter/tests/finish_reason_tests.rs" reason = "Test-only module stored under src/ for private adapter access; cargo-llvm-cov omits test harness source from production LCOV while the exercised rig_adapter.rs production lines remain coverage-gated." issue = "https://github.com/nearai/ironclaw/issues/6284" +[[exemption]] +module = "crates/ironclaw_llm/src/anthropic_oauth/tests.rs" +reason = "Test-only module stored under src/ for private transport access; cargo-llvm-cov omits test harness source from production LCOV while the exercised anthropic_oauth.rs production lines remain coverage-gated." +issue = "https://github.com/nearai/ironclaw/issues/6984" + [[exemption]] crate = "ironclaw_gateway" reason = "v1-only: consumed only by root `ironclaw` (src/channels/web/platform/static_files.rs, src/channels/web/handlers/frontend.rs); no crates/* dependents. Covered by \"Tests (Legacy)\"." From 33dd6172754bcab1723456d1982f692491f8a90b Mon Sep 17 00:00:00 2001 From: Illia Polosukhin Date: Sat, 1 Aug 2026 06:10:52 +0000 Subject: [PATCH 3/4] test(llm): cover remaining cache-breakpoint branch arms MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The changed-coverage gate flagged three remnants: the Text arm of set_cache_control, the empty-blocks tail, and the non-Text side of the system take-and-rebuild — which was also a latent drop: a system value already in block form was taken and never restored. Restore it untouched and pin all three paths with direct tests. Co-Authored-By: Claude Fable 5 --- crates/ironclaw_llm/src/anthropic_oauth.rs | 19 +++-- .../ironclaw_llm/src/anthropic_oauth/tests.rs | 71 +++++++++++++++++++ 2 files changed, 83 insertions(+), 7 deletions(-) diff --git a/crates/ironclaw_llm/src/anthropic_oauth.rs b/crates/ironclaw_llm/src/anthropic_oauth.rs index a4d13a82983..99c862adf7e 100644 --- a/crates/ironclaw_llm/src/anthropic_oauth.rs +++ b/crates/ironclaw_llm/src/anthropic_oauth.rs @@ -698,13 +698,18 @@ fn apply_cache_breakpoints( return; }; - if let Some(AnthropicSystem::Text(text)) = request.system.take() { - request.system = Some(AnthropicSystem::Blocks(vec![AnthropicSystemBlock { - block_type: "text", - text, - cache_control: Some(marker.clone()), - }])); - } + // convert_messages only ever produces the string form, but restore any + // other value untouched rather than dropping it on the floor. + request.system = match request.system.take() { + Some(AnthropicSystem::Text(text)) => { + Some(AnthropicSystem::Blocks(vec![AnthropicSystemBlock { + block_type: "text", + text, + cache_control: Some(marker.clone()), + }])) + } + other => other, + }; if let Some(tools) = request.tools.as_mut() && let Some(last) = tools.last_mut() diff --git a/crates/ironclaw_llm/src/anthropic_oauth/tests.rs b/crates/ironclaw_llm/src/anthropic_oauth/tests.rs index edf4b199410..70df6879038 100644 --- a/crates/ironclaw_llm/src/anthropic_oauth/tests.rs +++ b/crates/ironclaw_llm/src/anthropic_oauth/tests.rs @@ -354,6 +354,77 @@ fn apply_breakpoints_skips_empty_text_tail_and_empty_transcript() { assert!(empty.messages.is_empty()); } +/// Branch coverage: a text block at the tail of a multimodal message gets the +/// marker (the `Text` arm of `set_cache_control`), and an empty block list is +/// a no-op. +#[test] +fn apply_breakpoints_marks_text_block_tail_and_skips_empty_blocks() { + let mut request = request_with_messages(vec![AnthropicMessage { + role: "user".to_string(), + content: AnthropicContent::Blocks(vec![ + AnthropicContentBlock::Image { + source: AnthropicImageSource { + source_type: "base64", + media_type: "image/png".to_string(), + data: "aGk=".to_string(), + }, + cache_control: None, + }, + AnthropicContentBlock::Text { + text: "caption".to_string(), + cache_control: None, + }, + ]), + }]); + apply_cache_breakpoints(&mut request, CacheRetention::Short); + + let AnthropicContent::Blocks(blocks) = &request.messages[0].content else { + panic!("blocks expected"); + }; + let AnthropicContentBlock::Image { cache_control, .. } = &blocks[0] else { + panic!("image expected"); + }; + assert!(cache_control.is_none(), "only the tail block is marked"); + let AnthropicContentBlock::Text { cache_control, .. } = &blocks[1] else { + panic!("text expected"); + }; + assert!(cache_control.is_some()); + + let mut empty_blocks = request_with_messages(vec![AnthropicMessage { + role: "user".to_string(), + content: AnthropicContent::Blocks(Vec::new()), + }]); + apply_cache_breakpoints(&mut empty_blocks, CacheRetention::Short); + assert!(matches!( + &empty_blocks.messages[0].content, + AnthropicContent::Blocks(blocks) if blocks.is_empty() + )); +} + +/// Branch coverage: a system value already in block form is restored +/// untouched — never dropped by the take-and-rebuild — pinning the +/// restore-on-non-Text arm. +#[test] +fn apply_breakpoints_preserves_prebuilt_system_blocks() { + let mut request = request_with_messages(vec![AnthropicMessage { + role: "user".to_string(), + content: AnthropicContent::Text("question".to_string()), + }]); + request.system = Some(AnthropicSystem::Blocks(vec![AnthropicSystemBlock { + block_type: "text", + text: "prebuilt".to_string(), + cache_control: None, + }])); + apply_cache_breakpoints(&mut request, CacheRetention::Short); + + let Some(AnthropicSystem::Blocks(blocks)) = &request.system else { + panic!("prebuilt system blocks must survive apply_cache_breakpoints"); + }; + assert_eq!(blocks.len(), 1); + assert_eq!(blocks[0].text, "prebuilt"); + assert!(blocks[0].cache_control.is_none(), "restored untouched"); +} + /// Wire-level pin: retention `None` keeps the legacy wire shape — system /// as a plain string, and no cache_control anywhere. #[tokio::test] From 9ef55d38213a59d8a7e2b3033a38981ff40c09d0 Mon Sep 17 00:00:00 2001 From: Illia Polosukhin Date: Sat, 1 Aug 2026 06:48:18 +0000 Subject: [PATCH 4/4] test(llm): cover the empty-tools arm of apply_cache_breakpoints MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The last uncovered branch on the changed-coverage gate: Some(tools) with an empty vec (only constructible directly — complete_with_tools maps empty to None). Extend the empty-cases test to pin the no-op. Co-Authored-By: Claude Fable 5 --- crates/ironclaw_llm/src/anthropic_oauth/tests.rs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/crates/ironclaw_llm/src/anthropic_oauth/tests.rs b/crates/ironclaw_llm/src/anthropic_oauth/tests.rs index 70df6879038..dd1330be48c 100644 --- a/crates/ironclaw_llm/src/anthropic_oauth/tests.rs +++ b/crates/ironclaw_llm/src/anthropic_oauth/tests.rs @@ -394,11 +394,13 @@ fn apply_breakpoints_marks_text_block_tail_and_skips_empty_blocks() { role: "user".to_string(), content: AnthropicContent::Blocks(Vec::new()), }]); + empty_blocks.tools = Some(Vec::new()); apply_cache_breakpoints(&mut empty_blocks, CacheRetention::Short); assert!(matches!( &empty_blocks.messages[0].content, AnthropicContent::Blocks(blocks) if blocks.is_empty() )); + assert!(matches!(&empty_blocks.tools, Some(tools) if tools.is_empty())); } /// Branch coverage: a system value already in block form is restored