diff --git a/crates/domains/ironclaw_llm/CONTRACT.md b/crates/domains/ironclaw_llm/CONTRACT.md index 24bb54bd05c..8d4c40803c2 100644 --- a/crates/domains/ironclaw_llm/CONTRACT.md +++ b/crates/domains/ironclaw_llm/CONTRACT.md @@ -71,7 +71,7 @@ owner, and a deleted one fails until its entry goes. | `model-catalog` | Facts about **models**: what an endpoint lists, and which models see images, generate images, or think natively | Provider identity or routing | `models.rs`, `reasoning_models.rs`, `vision_models.rs`, `image_models.rs` | | `recording` | Trace capture and replay, and binding recorded tool arguments to earlier results | Live provider behavior | `recording.rs`, `trace_binding.rs` | | `transcription` | The `TranscriptionProvider` trait and its implementations — a **different trait** from `LlmProvider`, sharing only transports | Anything implementing `LlmProvider` | `transcription/mod.rs`, `transcription/chat_completions.rs`, `transcription/openai.rs` | -| `test-support` | Fixtures and fault injection, including the published `test-support` feature downstream harnesses consume | Production behavior | `testing/mod.rs`, `testing/fault_injection.rs`, `codex_test_helpers.rs`, `rig_adapter/tests/finish_reason_tests.rs` | +| `test-support` | Fixtures and fault injection, including the published `test-support` feature downstream harnesses consume | Production behavior | `testing/mod.rs`, `testing/fault_injection.rs`, `codex_test_helpers.rs`, `rig_adapter/tests/finish_reason_tests.rs`, `anthropic_oauth/tests.rs` | Four placement calls worth stating, because each is a file whose *shape* suggests one owner and whose *purpose* is another: @@ -337,6 +337,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/domains/ironclaw_llm/src/anthropic_oauth.rs b/crates/domains/ironclaw_llm/src/anthropic_oauth.rs index 19d8d1dd59a..77ca78acd60 100644 --- a/crates/domains/ironclaw_llm/src/anthropic_oauth.rs +++ b/crates/domains/ironclaw_llm/src/anthropic_oauth.rs @@ -104,6 +104,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 { @@ -139,6 +142,9 @@ impl AnthropicOAuthProvider { let unsupported_params: HashSet = config.unsupported_params.iter().cloned().collect(); + let cache_retention = + crate::rig_adapter::effective_cache_retention(config.cache_retention, &config.model); + Ok(Self { client, streaming_client, @@ -148,6 +154,7 @@ impl AnthropicOAuthProvider { base_url, active_model, unsupported_params, + cache_retention, }) } @@ -472,18 +479,20 @@ 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 { stream: false, 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); @@ -521,7 +530,7 @@ impl LlmProvider for AnthropicOAuthProvider { 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, @@ -560,7 +569,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 { stream: false, thinking: if has_tools { None @@ -569,13 +578,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); @@ -628,7 +639,7 @@ impl LlmProvider for AnthropicOAuthProvider { }, model, messages, - system, + system: system.map(AnthropicSystem::Text), max_tokens, temperature: req.temperature, tools: has_tools.then_some(tools), @@ -689,7 +700,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, @@ -707,6 +718,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)] @@ -719,22 +749,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 { @@ -749,6 +803,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)] @@ -1031,6 +1087,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(), @@ -1043,6 +1100,61 @@ 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; + }; + + // 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() + { + 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); + } + } + } + } +} + fn convert_anthropic_tools(tools: Vec) -> Vec { tools .into_iter() @@ -1050,6 +1162,9 @@ fn convert_anthropic_tools(tools: Vec) -> Vec { name: tool.name, description: tool.description, input_schema: tool.parameters, + // Markers are applied per-request by `apply_cache_breakpoints`; + // the legacy shape carries none. + cache_control: None, }) .collect() } @@ -1099,7 +1214,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) @@ -1115,13 +1233,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 { @@ -1144,6 +1266,7 @@ fn convert_messages(messages: Vec) -> (Option, Vec 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::*; - - #[derive(Default)] - struct RecordingSink(std::sync::Mutex>); - - #[async_trait] - impl CompletionStreamSink for RecordingSink { - async fn text_delta(&self, delta: String) { - self.0 - .lock() - .unwrap_or_else(|error| error.into_inner()) - .push(delta); - } - } - - #[tokio::test] - async fn anthropic_stream_emits_text_and_preserves_terminal_tools_and_usage() { - let sink = RecordingSink::default(); - let mut response = AnthropicStreamingResponse::default(); - ingest_anthropic_event( - &mut response, - "message_start", - r#"{"type":"message_start","message":{"usage":{"input_tokens":11,"cache_read_input_tokens":3}}}"#, - &sink, - ) - .await - .expect("message start"); - ingest_anthropic_event( - &mut response, - "content_block_delta", - r#"{"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"hello "}}"#, - &sink, - ) - .await - .expect("text delta"); - assert_eq!( - sink.0 - .lock() - .unwrap_or_else(|error| error.into_inner()) - .as_slice(), - ["hello "] - ); - assert!(!response.terminal, "text must arrive before completion"); - ingest_anthropic_event( - &mut response, - "content_block_start", - r#"{"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call-1","name":"weather","input":{}}}"#, - &sink, - ) - .await - .expect("tool start"); - for partial_json in ["{\"city\":\"", "Istanbul\"}"] { - ingest_anthropic_event( - &mut response, - "content_block_delta", - &format!( - r#"{{"type":"content_block_delta","index":1,"delta":{{"type":"input_json_delta","partial_json":{}}}}}"#, - serde_json::to_string(partial_json).expect("partial JSON string") - ), - &sink, - ) - .await - .expect("tool delta"); - } - ingest_anthropic_event( - &mut response, - "message_delta", - r#"{"type":"message_delta","delta":{"stop_reason":"tool_use"},"usage":{"output_tokens":7}}"#, - &sink, - ) - .await - .expect("terminal delta"); - let response = response.finish().expect("complete stream"); - assert_eq!(response.content, "hello "); - assert_eq!(response.stop_reason.as_deref(), Some("tool_use")); - assert_eq!(response.usage.input_tokens, 11); - assert_eq!(response.usage.output_tokens, 7); - assert_eq!(response.usage.cache_read_input_tokens, 3); - assert_eq!(response.tool_calls.len(), 1); - assert_eq!(response.tool_calls[0].name, "weather"); - assert_eq!( - response.tool_calls[0].arguments, - serde_json::json!({"city":"Istanbul"}) - ); - } - - #[test] - fn anthropic_stream_rejects_tool_state_missing_id_or_name() { - for (id, name) in [("", "weather"), ("call-1", "")] { - let mut response = AnthropicStreamingResponse::default(); - response.tool_call_parts.insert( - 0, - AnthropicStreamingToolCall { - id: id.to_string(), - name: name.to_string(), - input_json: "{}".to_string(), - }, - ); - - assert!(matches!( - response.finish(), - Err(LlmError::InvalidResponse { provider, reason }) - if provider == "anthropic_oauth" - && reason == "streamed tool_use block is missing its id or name" - )); - } - } - - #[test] - fn anthropic_stream_rejects_malformed_accumulated_tool_arguments() { - let mut response = AnthropicStreamingResponse::default(); - response.tool_call_parts.insert( - 0, - AnthropicStreamingToolCall { - id: "call-1".to_string(), - name: "weather".to_string(), - input_json: r#"{"city":"Istanbul""#.to_string(), - }, - ); - - match response.finish() { - Err(LlmError::InvalidResponse { provider, reason }) => { - assert_eq!(provider, "anthropic_oauth"); - assert!(reason.starts_with("streamed tool arguments are invalid JSON: ")); - } - other => panic!("expected invalid streamed tool arguments, got {other:?}"), - } - } - - #[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" - )); - } - - #[tokio::test] - async fn complete_streaming_rejects_eof_without_terminal_event() { - 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 = concat!( - "event: message_start\n", - "data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":1}}}\n\n", - "event: content_block_delta\n", - "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"partial\"}}\n\n" - ); - let response = format!( - "HTTP/1.1 200 OK\r\ncontent-type: text/event-stream\r\n\ - content-length: {}\r\n\r\n{body}", - body.len() - ); - socket - .write_all(response.as_bytes()) - .await - .expect("write streaming 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 sink = Arc::new(RecordingSink::default()); - let error = provider - .complete_streaming( - CompletionRequest::new(vec![ChatMessage::user("hello")]), - sink.clone(), - ) - .await - .expect_err("unterminated stream must fail"); - server.await.expect("loopback server"); - - assert!(matches!( - error, - LlmError::StreamInterrupted { provider, reason } - if provider == "anthropic_oauth" - && reason == "stream ended before message_stop or a stop reason" - )); - assert_eq!( - sink.0 - .lock() - .unwrap_or_else(|error| error.into_inner()) - .as_slice(), - ["partial"] - ); - } - - #[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/domains/ironclaw_llm/src/anthropic_oauth/tests.rs b/crates/domains/ironclaw_llm/src/anthropic_oauth/tests.rs new file mode 100644 index 00000000000..21e31a2c7fa --- /dev/null +++ b/crates/domains/ironclaw_llm/src/anthropic_oauth/tests.rs @@ -0,0 +1,891 @@ +//! 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 { + stream: false, + 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()); +} + +/// 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()), + }]); + 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 +/// 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] +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"); +} + +#[derive(Default)] +struct RecordingSink(std::sync::Mutex>); + +#[async_trait] +impl CompletionStreamSink for RecordingSink { + async fn text_delta(&self, delta: String) { + self.0 + .lock() + .unwrap_or_else(|error| error.into_inner()) + .push(delta); + } +} + +#[tokio::test] +async fn anthropic_stream_emits_text_and_preserves_terminal_tools_and_usage() { + let sink = RecordingSink::default(); + let mut response = AnthropicStreamingResponse::default(); + ingest_anthropic_event( + &mut response, + "message_start", + r#"{"type":"message_start","message":{"usage":{"input_tokens":11,"cache_read_input_tokens":3}}}"#, + &sink, + ) + .await + .expect("message start"); + ingest_anthropic_event( + &mut response, + "content_block_delta", + r#"{"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"hello "}}"#, + &sink, + ) + .await + .expect("text delta"); + assert_eq!( + sink.0 + .lock() + .unwrap_or_else(|error| error.into_inner()) + .as_slice(), + ["hello "] + ); + assert!(!response.terminal, "text must arrive before completion"); + ingest_anthropic_event( + &mut response, + "content_block_start", + r#"{"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call-1","name":"weather","input":{}}}"#, + &sink, + ) + .await + .expect("tool start"); + for partial_json in ["{\"city\":\"", "Istanbul\"}"] { + ingest_anthropic_event( + &mut response, + "content_block_delta", + &format!( + r#"{{"type":"content_block_delta","index":1,"delta":{{"type":"input_json_delta","partial_json":{}}}}}"#, + serde_json::to_string(partial_json).expect("partial JSON string") + ), + &sink, + ) + .await + .expect("tool delta"); + } + ingest_anthropic_event( + &mut response, + "message_delta", + r#"{"type":"message_delta","delta":{"stop_reason":"tool_use"},"usage":{"output_tokens":7}}"#, + &sink, + ) + .await + .expect("terminal delta"); + let response = response.finish().expect("complete stream"); + assert_eq!(response.content, "hello "); + assert_eq!(response.stop_reason.as_deref(), Some("tool_use")); + assert_eq!(response.usage.input_tokens, 11); + assert_eq!(response.usage.output_tokens, 7); + assert_eq!(response.usage.cache_read_input_tokens, 3); + assert_eq!(response.tool_calls.len(), 1); + assert_eq!(response.tool_calls[0].name, "weather"); + assert_eq!( + response.tool_calls[0].arguments, + serde_json::json!({"city":"Istanbul"}) + ); +} + +#[test] +fn anthropic_stream_rejects_tool_state_missing_id_or_name() { + for (id, name) in [("", "weather"), ("call-1", "")] { + let mut response = AnthropicStreamingResponse::default(); + response.tool_call_parts.insert( + 0, + AnthropicStreamingToolCall { + id: id.to_string(), + name: name.to_string(), + input_json: "{}".to_string(), + }, + ); + + assert!(matches!( + response.finish(), + Err(LlmError::InvalidResponse { provider, reason }) + if provider == "anthropic_oauth" + && reason == "streamed tool_use block is missing its id or name" + )); + } +} + +#[test] +fn anthropic_stream_rejects_malformed_accumulated_tool_arguments() { + let mut response = AnthropicStreamingResponse::default(); + response.tool_call_parts.insert( + 0, + AnthropicStreamingToolCall { + id: "call-1".to_string(), + name: "weather".to_string(), + input_json: r#"{"city":"Istanbul""#.to_string(), + }, + ); + + match response.finish() { + Err(LlmError::InvalidResponse { provider, reason }) => { + assert_eq!(provider, "anthropic_oauth"); + assert!(reason.starts_with("streamed tool arguments are invalid JSON: ")); + } + other => panic!("expected invalid streamed tool arguments, got {other:?}"), + } +} + +#[tokio::test] +async fn complete_streaming_rejects_eof_without_terminal_event() { + 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 = concat!( + "event: message_start\n", + "data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":1}}}\n\n", + "event: content_block_delta\n", + "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"partial\"}}\n\n" + ); + let response = format!( + "HTTP/1.1 200 OK\r\ncontent-type: text/event-stream\r\n\ + content-length: {}\r\n\r\n{body}", + body.len() + ); + socket + .write_all(response.as_bytes()) + .await + .expect("write streaming 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 sink = Arc::new(RecordingSink::default()); + let error = provider + .complete_streaming( + CompletionRequest::new(vec![ChatMessage::user("hello")]), + sink.clone(), + ) + .await + .expect_err("unterminated stream must fail"); + server.await.expect("loopback server"); + + assert!(matches!( + error, + LlmError::StreamInterrupted { provider, reason } + if provider == "anthropic_oauth" + && reason == "stream ended before message_stop or a stop reason" + )); + assert_eq!( + sink.0 + .lock() + .unwrap_or_else(|error| error.into_inner()) + .as_slice(), + ["partial"] + ); +} diff --git a/crates/domains/ironclaw_llm/src/config.rs b/crates/domains/ironclaw_llm/src/config.rs index 343d53d50aa..3626de35e30 100644 --- a/crates/domains/ironclaw_llm/src/config.rs +++ b/crates/domains/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/domains/ironclaw_llm/src/lib.rs b/crates/domains/ironclaw_llm/src/lib.rs index 4cc0d108236..648dc7ed013 100644 --- a/crates/domains/ironclaw_llm/src/lib.rs +++ b/crates/domains/ironclaw_llm/src/lib.rs @@ -504,15 +504,27 @@ fn create_anthropic_from_registry( reason: format!("Failed to create Anthropic client: {e}"), })?; - let cache_retention = config.cache_retention; - - let model = client.completion_model(&config.model); + // 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 = + rig_adapter::effective_cache_retention(config.cache_retention, &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)" ); } @@ -2234,4 +2246,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/domains/ironclaw_llm/src/rig_adapter.rs b/crates/domains/ironclaw_llm/src/rig_adapter.rs index 463edfacdf9..f483e62611e 100644 --- a/crates/domains/ironclaw_llm/src/rig_adapter.rs +++ b/crates/domains/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. @@ -297,9 +299,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. @@ -926,6 +931,21 @@ 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. @@ -1076,17 +1096,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, @@ -1102,15 +1137,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 { @@ -3712,6 +3755,311 @@ 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_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( + 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). diff --git a/tests/integration/coverage-exemptions.toml b/tests/integration/coverage-exemptions.toml index 4bdf2ed9bb3..562a5e86e48 100644 --- a/tests/integration/coverage-exemptions.toml +++ b/tests/integration/coverage-exemptions.toml @@ -47,6 +47,11 @@ module = "crates/domains/ironclaw_llm/src/rig_adapter/tests/finish_reason_tests. 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/domains/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)\"."