返回 CodeWhale
anthropic.rs
根目录 / crates / tui / src / client / anthropic.rs
1 //! Native Anthropic Messages API adapter (#3014).
2 //!
3 //! CodeWhale's internal wire types are already Anthropic-shaped (the harness
4 //! speaks Messages internally and translates *out* to OpenAI dialects), so
5 //! this adapter is mostly native serialization plus an SSE pass-through:
6 //! `StreamEvent` deserializes Anthropic's `message_start` /
7 //! `content_block_*` / `message_delta` / `message_stop` / `ping` events
8 //! directly. What the adapter adds on top:
9 //!
10 //! - request shaping: adaptive thinking + `output_config.effort` from
11 //! CodeWhale's `reasoning_effort` tiers, sampling-parameter rules for
12 //! models that reject them, and `cache_control` breakpoint placement
13 //! aligned with the prefix-zone model in `prefix_cache.rs`;
14 //! - usage normalization (#2961 / #4318): `prompt_cache_hit_tokens` comes from
15 //! `cache_read_input_tokens`, `prompt_cache_write_tokens` from
16 //! `cache_creation_input_tokens`, `prompt_cache_miss_tokens` is the raw
17 //! non-cached `input_tokens`, and the normalized `input_tokens` is the sum
18 //! of all three (total prompt, the DeepSeek convention);
19 //! - signed-thinking handling: `signature_delta` is captured into
20 //! [`codewhale_models::Delta::SignatureDelta`] and assistant thinking blocks
21 //! replay verbatim (signature included); unsigned thinking blocks are
22 //! dropped from replay because the API rejects them.
23 //!
24 //! Modeled on `client/responses.rs` (separate file per dialect, no protocol
25 //! hacks in the shared paths).
26
27 use anyhow::{Context, Result};
28 use serde::Deserialize;
29 use serde_json::{Value, json};
30
31 use crate::config::{ProviderKind, wire_model_for_provider_route};
32 use crate::llm_client::StreamEventBox;
33 use crate::logging;
34 use crate::tools::schema_sanitize;
35 use codewhale_models::{ContentBlock, MessageRequest, MessageResponse, StreamEvent, Usage};
36
37 use super::CodewhaleClient;
38 use super::prepared::WireDialect;
39 use super::role_placement::{RolePlacement, role_placement};
40 use super::wire::{extract_sse_data_value, next_sse_line, push_sse_event_data};
41
42 /// Maximum `cache_control` breakpoints Anthropic accepts per request.
43 const MAX_CACHE_BREAKPOINTS: usize = 4;
44
45 impl CodewhaleClient {
46 /// Build the native Messages API request body from a [`MessageRequest`].
47 pub(super) fn build_anthropic_body(&self, request: &MessageRequest, stream: bool) -> Value {
48 let model =
49 wire_model_for_provider_route(self.api_provider, &self.base_url, &request.model);
50 let mut body = json!({
51 "model": model,
52 "max_tokens": request.max_tokens,
53 "stream": stream,
54 });
55
56 if let Some(system) = request.system.as_ref() {
57 body["system"] = match system {
58 codewhale_models::SystemPrompt::Text(text) => json!(text),
59 codewhale_models::SystemPrompt::Blocks(blocks) => json!(
60 blocks
61 .iter()
62 .map(|block| {
63 let mut value = json!({
64 "type": "text",
65 "text": block.text,
66 });
67 if let Some(cache) = block.cache_control.as_ref() {
68 value["cache_control"] = json!({ "type": cache.cache_type });
69 }
70 value
71 })
72 .collect::<Vec<_>>()
73 ),
74 };
75 }
76
77 let mut messages: Vec<Value> = request
78 .messages
79 .iter()
80 .filter_map(message_to_anthropic)
81 .collect();
82 merge_split_tool_results(&mut messages);
83 repair_dangling_tool_uses(&mut messages);
84 body["messages"] = Value::Array(messages);
85
86 if let Some(tools) = request.tools.as_ref()
87 && !tools.is_empty()
88 {
89 body["tools"] = json!(
90 tools
91 .iter()
92 .map(|tool| {
93 // Sanitize the tool's input_schema the same way the
94 // OpenAI Responses adapter does: strip top-level
95 // oneOf/anyOf/allOf (which Anthropic rejects), merge
96 // alternative properties into the root, and surface
97 // the dropped constraint as a description note so the
98 // model still knows which parameters are expected.
99 let mut schema = tool.input_schema.clone();
100 let constraint_note = schema_sanitize::sanitize_for_responses(&mut schema);
101 let description = match constraint_note {
102 Some(note) if tool.description.trim().is_empty() => note,
103 Some(note) => format!("{}\n\n{}", tool.description.trim(), note),
104 None => tool.description.clone(),
105 };
106 let mut value = json!({
107 "name": tool.name,
108 "description": description,
109 "input_schema": schema,
110 });
111 if let Some(strict) = tool.strict {
112 value["strict"] = json!(strict);
113 }
114 if let Some(cache) = tool.cache_control.as_ref() {
115 value["cache_control"] = json!({ "type": cache.cache_type });
116 }
117 value
118 })
119 .collect::<Vec<_>>()
120 );
121 }
122
123 if let Some(tool_choice) = request.tool_choice.as_ref() {
124 body["tool_choice"] = anthropic_tool_choice(tool_choice);
125 }
126
127 // Thinking + effort shaping. MiniMax supports adaptive/disabled but
128 // not Anthropic's output_config effort field; native Anthropic routes
129 // keep the existing effort mapping. Other Messages-compatible
130 // gateways (#4978, e.g. Sensenova) only accept the documented
131 // enabled/disabled/auto thinking types, so non-native routes get the
132 // portable `{"type":"enabled","budget_tokens":N}` shape instead.
133 let thinking_capable = codewhale_models::model_supports_reasoning(&model);
134 let is_minimax_provider = self.api_provider == ProviderKind::MinimaxAnthropic;
135 let is_minimax = crate::config::is_exact_minimax_anthropic_m3_route(
136 self.api_provider,
137 &self.base_url,
138 &model,
139 );
140 let is_deepseek = self.api_provider == ProviderKind::DeepseekAnthropic;
141 // Model Studio's Anthropic-compatible endpoint documents the portable
142 // `{"type":"enabled","budget_tokens":N}` shape AND `{"type":"disabled"}`
143 // (alibabacloud.com/help/en/model-studio/anthropic-api-messages), so
144 // an explicit "off" can be honored on the wire instead of silently
145 // falling through to the server default (which is thinking-ON for the
146 // qwen3.x families).
147 let is_modelstudio = matches!(
148 self.api_provider,
149 ProviderKind::ModelstudioTokenPlan
150 | ProviderKind::ModelstudioTokenPlanAnthropic
151 | ProviderKind::ModelstudioCodingPlan
152 | ProviderKind::ModelstudioCodingPlanAnthropic
153 );
154 // MiniMax's exact M3 route and DeepSeek's Messages dialect both
155 // document adaptive support; everything else needs the native host.
156 let supports_adaptive =
157 is_native_anthropic_base_url(&self.base_url) || is_minimax || is_deepseek;
158 let effort = request
159 .reasoning_effort
160 .as_deref()
161 .map(|raw| raw.trim().to_ascii_lowercase());
162 match effort.as_deref() {
163 _ if is_minimax_provider && !is_minimax => {}
164 Some("off" | "disabled" | "none" | "false")
165 if (is_minimax || is_deepseek || is_modelstudio) && thinking_capable =>
166 {
167 // Deliberately includes thinking-only Model Studio models
168 // (qwen3.8-max family): unlike the chat dialect's
169 // enable_thinking switch, the Messages endpoint documents the
170 // portable {"type":"disabled"} shape for them
171 // (alibabacloud.com/help/en/model-studio/anthropic-api-messages)
172 // — pinned by modelstudio_messages_body_requests_thinking_
173 // with_budget. Re-checked 2026-08-04.
174 body["thinking"] = json!({ "type": "disabled" });
175 }
176 Some("off" | "disabled" | "none" | "false") => {}
177 Some(level) if thinking_capable && supports_adaptive => {
178 body["thinking"] = json!({ "type": "adaptive" });
179 if !is_minimax {
180 let mapped = match level {
181 "low" | "minimal" => "low",
182 "medium" | "mid" => "medium",
183 "max" | "xhigh" | "highest" => "max",
184 _ => "high",
185 };
186 body["output_config"] = json!({ "effort": mapped });
187 }
188 }
189 None if thinking_capable && supports_adaptive => {
190 body["thinking"] = json!({ "type": "adaptive" });
191 }
192 _ if thinking_capable => {
193 if let Some(budget) = compat_thinking_budget(effort.as_deref(), request.max_tokens)
194 {
195 body["thinking"] = json!({ "type": "enabled", "budget_tokens": budget });
196 }
197 }
198 _ => {}
199 }
200
201 // Sampling parameters: Claude 4.7+ rejects temperature/top_p
202 // entirely; earlier models reject the two together. Send at most one
203 // (temperature wins), or neither for models that forbid them.
204 if !anthropic_model_rejects_sampling(&request.model) {
205 if let Some(temperature) = request.temperature {
206 body["temperature"] = json!(temperature);
207 } else if let Some(top_p) = request.top_p {
208 body["top_p"] = json!(top_p);
209 }
210 }
211
212 apply_anthropic_cache_breakpoints(&mut body);
213 body
214 }
215
216 /// Non-streaming send through the shared typed retry path
217 /// (`send_with_retry`): 429/5xx/transport failures retry with backoff and
218 /// honor `Retry-After`, and a final failure stays a downcastable
219 /// `LlmError` so the engine can classify auth, rate-limit, context and
220 /// invalid-request failures like every other wire.
221 async fn send_anthropic_request(&self, url: &str, body: &Value) -> Result<reqwest::Response> {
222 let url = self.messages_transport_url(url);
223 let request_body =
224 serde_json::to_vec(body).context("Failed to serialize Anthropic Messages request")?;
225 self.send_with_retry(|| {
226 self.http_client
227 .post(&url)
228 .header(reqwest::header::CONTENT_TYPE, "application/json")
229 .header("Accept", "text/event-stream")
230 .body(request_body.clone())
231 })
232 .await
233 .context("Anthropic Messages API request failed")
234 }
235
236 /// Open the streaming Messages request through the shared stream-entry
237 /// transport policy: bounded header wait, dual-client selection, and at
238 /// most one HTTP/1.1 fallback retry on a classified H2 header stall.
239 /// Inside each open attempt the provider retry loop
240 /// (`send_stream_open_with_retry`) handles rate limits and transient
241 /// upstream failures before any stream body exists — with no total
242 /// deadline, which would ride on the returned body — as the Chat and
243 /// Responses adapters do. Wire-specific request construction (headers,
244 /// endpoint, body) stays here at the adapter edge.
245 async fn open_anthropic_stream_response(
246 &self,
247 url: &str,
248 body: &Value,
249 ) -> Result<reqwest::Response> {
250 let url = self.messages_transport_url(url);
251 let request_body =
252 serde_json::to_vec(body).context("Failed to serialize Anthropic Messages request")?;
253 let open_req = self.stream_open_request();
254 super::stream_entry::open_sse_response(&open_req, |policy| {
255 let url = url.clone();
256 let request_body = request_body.clone();
257 async move {
258 let client = super::stream_entry::client_for_policy(
259 &self.http_client,
260 self.http1_fallback_client(),
261 policy,
262 );
263 self.send_stream_open_with_retry(|| {
264 client
265 .post(&url)
266 .header(reqwest::header::CONTENT_TYPE, "application/json")
267 .header("Accept", "text/event-stream")
268 .body(request_body.clone())
269 })
270 .await
271 .context("Anthropic Messages API request failed")
272 }
273 })
274 .await
275 }
276
277 /// Handle a streaming Messages API request.
278 pub(super) async fn handle_anthropic_stream(
279 &self,
280 prepared: &super::PreparedOutboundRequest,
281 ) -> Result<StreamEventBox> {
282 // Body and endpoint come from the shared prepared-request seam
283 // (`prepare_outbound_request`), never from a second builder.
284 let body = &prepared.body;
285 let response = self
286 .open_anthropic_stream_response(&prepared.endpoint.url, body)
287 .await?;
288
289 let stream_idle_timeout = self.stream_idle_timeout;
290 let first_byte = super::stream_entry::first_byte_timeout(stream_idle_timeout);
291 let provider_label = self.api_provider.provider().display_name();
292 let byte_stream = response.bytes_stream();
293
294 let stream = async_stream::stream! {
295 use futures_util::StreamExt;
296
297 // Raw byte buffer: decode only COMPLETE lines (or the stream-end
298 // tail) via the shared take_sse_line / flush_sse_line helpers so a
299 // multi-byte UTF-8 char (CJK/emoji) split across HTTP/2 DATA is
300 // never corrupted to U+FFFD. Genuine invalid bytes fail closed.
301 let mut buffer: Vec<u8> = Vec::new();
302 // `data:` fields of the event being assembled, dispatched at the
303 // blank line that ends it (or at stream end).
304 let mut event_data = String::new();
305 let stream_start = std::time::Instant::now();
306 let mut last_chunk_at = std::time::Instant::now();
307 let mut bytes_received: usize = 0;
308 let mut ended = false;
309 tokio::pin!(byte_stream);
310
311 loop {
312 if !ended {
313 let wait = super::stream_entry::next_chunk_timeout(
314 stream_idle_timeout,
315 first_byte,
316 bytes_received,
317 );
318 match tokio::time::timeout(wait, byte_stream.next()).await {
319 Ok(Some(Ok(chunk))) => {
320 bytes_received += chunk.len();
321 last_chunk_at = std::time::Instant::now();
322 buffer.extend_from_slice(&chunk);
323 }
324 Ok(Some(Err(e))) => {
325 yield Err(anyhow::anyhow!("Stream read error: {e}"));
326 return;
327 }
328 Ok(None) => ended = true,
329 Err(_) => {
330 yield Err(anyhow::anyhow!(super::stream_entry::body_timeout_message(
331 wait,
332 bytes_received,
333 stream_start.elapsed(),
334 last_chunk_at.elapsed(),
335 provider_label,
336 )));
337 return;
338 }
339 }
340 }
341
342 loop {
343 let data = match next_sse_line(&mut buffer, ended) {
344 // A blank line ends the event.
345 Ok(Some(line)) if line.is_empty() => std::mem::take(&mut event_data),
346 Ok(Some(line)) => {
347 // `event:` lines are redundant (the data payload
348 // carries `type`) and comment/heartbeat lines are
349 // ignorable.
350 if let Some(value) = extract_sse_data_value(&line)
351 && let Err(err) = push_sse_event_data(&mut event_data, value)
352 {
353 yield Err(anyhow::anyhow!("{err}"));
354 return;
355 }
356 continue;
357 }
358 // The final event may arrive without its blank line.
359 Ok(None) if ended && !event_data.is_empty() => {
360 std::mem::take(&mut event_data)
361 }
362 Ok(None) => break,
363 Err(err) => {
364 yield Err(anyhow::anyhow!("{err}"));
365 return;
366 }
367 };
368 if data.is_empty() {
369 continue;
370 }
371
372 match convert_anthropic_sse_data(&data) {
373 Some(Ok(StreamEvent::Error { error })) => {
374 let (error_type, message) = anthropic_error_fields(&error);
375 yield Err(anyhow::anyhow!(
376 "Anthropic stream error ({error_type}): {message}"
377 ));
378 return;
379 }
380 Some(Ok(event)) => {
381 let is_stop = matches!(event, StreamEvent::MessageStop);
382 yield Ok(event);
383 if is_stop {
384 return;
385 }
386 }
387 Some(Err(e)) => {
388 logging::warn(format!("Failed to parse Anthropic SSE event: {e}"));
389 }
390 None => {}
391 }
392 }
393
394 if ended {
395 break;
396 }
397 }
398 // Only `message_stop` (returned above) or a provider error proves
399 // the response is whole. A bare HTTP EOF is truncation, and ending
400 // the stream quietly would hand the turn a partial answer as done.
401 yield Err(anyhow::anyhow!(
402 "Anthropic Messages stream closed before message_stop"
403 ));
404 };
405
406 Ok(Box::pin(stream))
407 }
408
409 /// Handle a non-streaming Messages API request.
410 pub(super) async fn handle_anthropic_message(
411 &self,
412 prepared: &super::PreparedOutboundRequest,
413 ) -> Result<MessageResponse> {
414 let response = self
415 .send_anthropic_request(&prepared.endpoint.url, &prepared.body)
416 .await?;
417 let value: Value = response
418 .json()
419 .await
420 .context("Failed to parse Anthropic Messages response")?;
421 decode_anthropic_message(value)
422 }
423 }
424
425 fn discard_provider_execution_ids(value: &mut Value) {
426 // Shared history retains execution ids on disk. Incoming provider content
427 // cannot choose one, including a value that would not deserialize as the
428 // host field. Live execution must mint its own correlation instead.
429 if let Some(blocks) = value.get_mut("content").and_then(Value::as_array_mut) {
430 for block in blocks {
431 if let Some(block) = block.as_object_mut() {
432 block.remove("execution_id");
433 }
434 }
435 }
436 }
437
438 fn decode_anthropic_message(mut value: Value) -> Result<MessageResponse> {
439 discard_provider_execution_ids(&mut value);
440 if let Some(usage) = value.get_mut("usage") {
441 *usage = json!(parse_anthropic_usage(usage));
442 }
443 serde_json::from_value(value).context("Failed to decode Anthropic Messages response")
444 }
445
446 /// Build the `/v1/messages` endpoint URL, tolerating base URLs that already
447 /// carry a `/v1` suffix.
448 pub(super) fn anthropic_messages_url(base_url: &str) -> String {
449 let trimmed = base_url.trim_end_matches('/');
450 if trimmed.ends_with("/v1") {
451 format!("{trimmed}/messages")
452 } else {
453 format!("{trimmed}/v1/messages")
454 }
455 }
456
457 /// Whether the route targets first-party Anthropic (`api.anthropic.com`),
458 /// where the `{"type":"adaptive"}` thinking control is valid. Strict
459 /// Anthropic-compatible gateways reject it (#4978).
460 fn is_native_anthropic_base_url(base_url: &str) -> bool {
461 let rest = base_url
462 .trim()
463 .trim_start_matches("https://")
464 .trim_start_matches("http://");
465 let host = rest
466 .split(['/', ':', '?', '#'])
467 .next()
468 .unwrap_or("")
469 .to_ascii_lowercase();
470 host == "api.anthropic.com" || host.ends_with(".anthropic.com")
471 }
472
473 /// Minimum `budget_tokens` the Messages API accepts for extended thinking.
474 const MIN_THINKING_BUDGET_TOKENS: u32 = 1024;
475
476 /// Effort-tier `budget_tokens` for gateways that only accept the documented
477 /// `{"type":"enabled","budget_tokens":N}` thinking shape (#4978). The wire
478 /// contract requires `budget_tokens >= 1024` and `< max_tokens`, so requests
479 /// too small to fit the minimum budget send no thinking block at all.
480 fn compat_thinking_budget(effort: Option<&str>, max_tokens: u32) -> Option<u32> {
481 let tier: u32 = match effort {
482 Some("low" | "minimal") => 4_096,
483 Some("medium" | "mid") => 8_192,
484 Some("max" | "xhigh" | "highest") => 32_768,
485 // "high" and unspecified effort share the adaptive default tier.
486 _ => 16_384,
487 };
488 let budget = tier.min(max_tokens.checked_sub(1)?);
489 (budget >= MIN_THINKING_BUDGET_TOKENS).then_some(budget)
490 }
491
492 /// Fold a user turn that carries `tool_result`s into the user turn before it
493 /// (#6378).
494 ///
495 /// The engine records each tool result as its own user message, so a parallel
496 /// tool-call batch arrives here as `assistant{use_a, use_b}`, `user{result_a}`,
497 /// `user{result_b}`. Anthropic wants every result in the user turn right after
498 /// the `tool_use`s, and the repair below reads only that turn: it would answer
499 /// `use_b` with an error placeholder while the real result sits in the next
500 /// message. Only turns carrying a `tool_result` are folded, so any other
501 /// consecutive user turns keep their shape; inside the merged turn the results
502 /// stay ahead of other content so they still lead it.
503 fn merge_split_tool_results(messages: &mut Vec<Value>) {
504 let is_user = |message: &Value| message.get("role").and_then(Value::as_str) == Some("user");
505 let is_tool_result =
506 |block: &Value| block.get("type").and_then(Value::as_str) == Some("tool_result");
507 let carries_tool_result = |message: &Value| {
508 message
509 .get("content")
510 .and_then(Value::as_array)
511 .is_some_and(|blocks| blocks.iter().any(is_tool_result))
512 };
513 let mut merged: Vec<Value> = Vec::with_capacity(messages.len());
514 for mut message in messages.drain(..) {
515 if is_user(&message)
516 && carries_tool_result(&message)
517 && let Some(previous) = merged.last_mut()
518 && is_user(previous)
519 {
520 let mut blocks = match previous["content"].take() {
521 Value::Array(blocks) => blocks,
522 other => vec![other],
523 };
524 match message["content"].take() {
525 Value::Array(incoming) => blocks.extend(incoming),
526 other => blocks.push(other),
527 }
528 // Stable sort: results keep their order and lead the turn.
529 blocks.sort_by_key(|block| !is_tool_result(block));
530 previous["content"] = Value::Array(blocks);
531 } else {
532 merged.push(message);
533 }
534 }
535 *messages = merged;
536 }
537
538 /// Placeholder body for a `tool_use` that never produced a `tool_result`.
539 const UNEXECUTED_TOOL_RESULT: &str = "tool call was not executed";
540
541 /// Defensive wire repair (#5002): every assistant `tool_use` must be answered
542 /// by a `tool_result` in the immediately following user message, or the API
543 /// rejects the whole conversation with a 400 on every retry. Pre-dispatch
544 /// failure paths (e.g. the model calling an unavailable tool) can strand an
545 /// orphaned `tool_use` in history, so missing results get an explicit
546 /// error placeholder instead of poisoning the session.
547 fn repair_dangling_tool_uses(messages: &mut Vec<Value>) {
548 let mut index = 0;
549 while index < messages.len() {
550 let ids = assistant_tool_use_ids(&messages[index]);
551 if ids.is_empty() {
552 index += 1;
553 continue;
554 }
555 let next_is_user = messages
556 .get(index + 1)
557 .and_then(|message| message.get("role"))
558 .and_then(Value::as_str)
559 == Some("user");
560 if !next_is_user {
561 messages.insert(index + 1, json!({ "role": "user", "content": [] }));
562 }
563 if let Some(blocks) = messages[index + 1]
564 .get_mut("content")
565 .and_then(Value::as_array_mut)
566 {
567 let answered: std::collections::HashSet<String> = blocks
568 .iter()
569 .filter(|block| block.get("type").and_then(Value::as_str) == Some("tool_result"))
570 .filter_map(|block| block.get("tool_use_id").and_then(Value::as_str))
571 .map(str::to_string)
572 .collect();
573 // tool_result blocks must lead the user turn, so placeholders are
574 // prepended in tool_use order.
575 for (offset, id) in ids
576 .iter()
577 .filter(|id| !answered.contains(id.as_str()))
578 .enumerate()
579 {
580 blocks.insert(
581 offset,
582 json!({
583 "type": "tool_result",
584 "tool_use_id": id,
585 "content": UNEXECUTED_TOOL_RESULT,
586 "is_error": true,
587 }),
588 );
589 }
590 }
591 index += 1;
592 }
593 }
594
595 fn assistant_tool_use_ids(message: &Value) -> Vec<String> {
596 if message.get("role").and_then(Value::as_str) != Some("assistant") {
597 return Vec::new();
598 }
599 message
600 .get("content")
601 .and_then(Value::as_array)
602 .map(|blocks| {
603 blocks
604 .iter()
605 .filter(|block| block.get("type").and_then(Value::as_str) == Some("tool_use"))
606 .filter_map(|block| block.get("id").and_then(Value::as_str))
607 .map(str::to_string)
608 .collect()
609 })
610 .unwrap_or_default()
611 }
612
613 /// Models that reject `temperature` / `top_p` outright (Claude 4.7+).
614 fn anthropic_model_rejects_sampling(model: &str) -> bool {
615 let lower = model.to_ascii_lowercase();
616 lower.contains("opus-4-7")
617 || lower.contains("opus-4-8")
618 || lower.contains("fable")
619 || lower.contains("mythos")
620 }
621
622 /// Convert the engine's `tool_choice` value (OpenAI-style string or object)
623 /// to the Anthropic object form.
624 fn anthropic_tool_choice(tool_choice: &Value) -> Value {
625 match tool_choice.as_str() {
626 Some("auto") => json!({ "type": "auto" }),
627 Some("none") => json!({ "type": "none" }),
628 Some("any" | "required") => json!({ "type": "any" }),
629 Some(name) => json!({ "type": "tool", "name": name }),
630 None => tool_choice.clone(),
631 }
632 }
633
634 /// Convert one internal message to the Anthropic wire shape. Returns `None`
635 /// when no blocks survive conversion (Anthropic rejects empty content) or
636 /// when the role has no Anthropic channel.
637 ///
638 /// The wire role used to be `message.role` forwarded verbatim, which is how a
639 /// `system` message ended up on the wire for the provider to 400 on. It now
640 /// comes from the shared placement table, and pairs the table rejects are
641 /// refused at the outbound seam before this function ever runs.
642 pub(super) fn message_to_anthropic(message: &codewhale_models::Message) -> Option<Value> {
643 let placement = role_placement(&message.role, WireDialect::AnthropicMessages);
644 let wire_role = match placement {
645 RolePlacement::User | RolePlacement::Developer => "user",
646 RolePlacement::Assistant | RolePlacement::InterruptedAssistant => "assistant",
647 // Unreachable in production: `reject_unsupported_roles` refuses these
648 // pairs at the seam. Failing closed here keeps a future caller that
649 // skips the seam from putting an unrepresentable role on the wire.
650 RolePlacement::System | RolePlacement::Omitted | RolePlacement::Rejected => return None,
651 };
652 let mut blocks: Vec<Value> = message
653 .content
654 .iter()
655 .filter_map(content_block_to_anthropic)
656 .collect();
657 if blocks.is_empty() {
658 return None;
659 }
660 if placement == RolePlacement::InterruptedAssistant
661 && let Some(text) = blocks
662 .iter_mut()
663 .find(|block| block.get("type").and_then(Value::as_str) == Some("text"))
664 {
665 let existing = text
666 .get("text")
667 .and_then(Value::as_str)
668 .unwrap_or_default()
669 .to_string();
670 text["text"] = json!(format!(
671 "{}{}",
672 codewhale_models::INTERRUPTED_ASSISTANT_CONTEXT_PREFIX,
673 existing
674 ));
675 }
676 Some(json!({ "role": wire_role, "content": blocks }))
677 }
678
679 /// Project the shared `ImageUrl` block onto Anthropic's tagged image source.
680 ///
681 /// The OpenAI dialects carry an image as a single URL string, so that is what
682 /// [`ContentBlock::ImageUrl`] stores. Anthropic instead models the source as a
683 /// tagged union, and — this is the part that used to be wrong here — it does
684 /// **not** accept a `data:` URL under `{"type":"url"}`. Sending a local
685 /// screenshot that way earns an opaque provider-side 400, which is exactly the
686 /// confusing failure this whole path exists to avoid, so the data URL is taken
687 /// back apart into `{"type":"base64", media_type, data}`.
688 fn anthropic_image_block(url: &str) -> Value {
689 if let Some((media_type, data)) = crate::image_attach::parse_data_url(url) {
690 return json!({
691 "type": "image",
692 "source": { "type": "base64", "media_type": media_type, "data": data },
693 });
694 }
695 if crate::image_attach::is_remote_image_url(url) {
696 return json!({
697 "type": "image",
698 "source": { "type": "url", "url": url },
699 });
700 }
701 // Anything else (a bare path, a `file://`, a truncated data URL) has no
702 // Anthropic representation. Degrade to visible text rather than emitting a
703 // source the API will reject: the turn survives and the model can see that
704 // something was meant to be here.
705 json!({
706 "type": "text",
707 "text": format!("[unsupported image reference: {url}]"),
708 })
709 }
710
711 pub(super) fn anthropic_tool_result_content(
712 content: &str,
713 content_blocks: Option<&[Value]>,
714 ) -> Value {
715 let (image, omitted) = crate::image_attach::provider_tool_result_image_refs(content_blocks);
716 let content = crate::image_attach::tool_result_text_with_omission(content, omitted);
717 let Some((mime_type, data)) = image else {
718 return json!(content);
719 };
720 let mut blocks = Vec::with_capacity(2);
721 if !content.is_empty() {
722 blocks.push(json!({ "type": "text", "text": content }));
723 }
724 blocks.push(json!({
725 "type": "image",
726 "source": { "type": "base64", "media_type": mime_type, "data": data },
727 }));
728 json!(blocks)
729 }
730
731 fn content_block_to_anthropic(block: &ContentBlock) -> Option<Value> {
732 match block {
733 ContentBlock::Text {
734 text,
735 cache_control,
736 } => {
737 let mut value = json!({ "type": "text", "text": text });
738 if let Some(cache) = cache_control {
739 value["cache_control"] = json!({ "type": cache.cache_type });
740 }
741 Some(value)
742 }
743 ContentBlock::Thinking {
744 thinking,
745 signature,
746 ..
747 } => {
748 // Anthropic rejects unsigned thinking blocks on replay (and the
749 // DeepSeek-era "(reasoning omitted)" placeholders mean nothing to
750 // it), so only signed blocks are replayed — verbatim, signature
751 // included.
752 signature.as_ref().map(|signature| {
753 json!({
754 "type": "thinking",
755 "thinking": thinking,
756 "signature": signature,
757 })
758 })
759 }
760 ContentBlock::ToolUse {
761 id, name, input, ..
762 } => Some(json!({
763 "type": "tool_use",
764 "id": id,
765 "name": name,
766 "input": input,
767 })),
768 ContentBlock::ToolResult {
769 tool_use_id,
770 content,
771 is_error,
772 content_blocks,
773 ..
774 } => {
775 let mut value = json!({
776 "type": "tool_result",
777 "tool_use_id": tool_use_id,
778 "content": anthropic_tool_result_content(content, content_blocks.as_deref()),
779 });
780 if let Some(is_error) = is_error {
781 value["is_error"] = json!(is_error);
782 }
783 Some(value)
784 }
785 ContentBlock::ImageUrl { image_url } => Some(anthropic_image_block(&image_url.url)),
786 // Server-tool block types are DeepSeek/internal concepts with no
787 // Anthropic client-side wire equivalent.
788 ContentBlock::ServerToolUse { .. }
789 | ContentBlock::ToolSearchToolResult { .. }
790 | ContentBlock::CodeExecutionToolResult { .. } => None,
791 }
792 }
793
794 /// Enforce the prefix-zone breakpoint policy (#3014):
795 /// 1. the last tool in the catalog (or, with no tools, the last system
796 /// block) — caches the immutable prefix;
797 /// 2. the last content block of the most recent user turn — caches the
798 /// append-only history.
799 ///
800 /// Caller-provided breakpoints are preserved, but the total is capped at
801 /// [`MAX_CACHE_BREAKPOINTS`] by dropping the earliest markers first (the
802 /// latest markers cover the longest prefixes).
803 fn apply_anthropic_cache_breakpoints(body: &mut Value) {
804 // Place breakpoint 1: prefer the last tool; otherwise last system block.
805 let mut placed_prefix = false;
806 if let Some(tools) = body.get_mut("tools").and_then(Value::as_array_mut)
807 && let Some(last) = tools.last_mut()
808 {
809 last["cache_control"] = json!({ "type": "ephemeral" });
810 placed_prefix = true;
811 }
812 if !placed_prefix
813 && let Some(system) = body.get_mut("system").and_then(Value::as_array_mut)
814 && let Some(last) = system.last_mut()
815 {
816 last["cache_control"] = json!({ "type": "ephemeral" });
817 }
818
819 // Place breakpoint 2: last content block of the latest user message.
820 if let Some(messages) = body.get_mut("messages").and_then(Value::as_array_mut)
821 && let Some(last_user) = messages
822 .iter_mut()
823 .rev()
824 .find(|message| message.get("role").and_then(Value::as_str) == Some("user"))
825 && let Some(last_block) = last_user
826 .get_mut("content")
827 .and_then(Value::as_array_mut)
828 .and_then(|blocks| blocks.last_mut())
829 {
830 last_block["cache_control"] = json!({ "type": "ephemeral" });
831 }
832
833 // Cap at MAX_CACHE_BREAKPOINTS in render order (tools → system →
834 // messages), dropping the earliest extras.
835 let mut marked: Vec<*mut Value> = Vec::new();
836 let collect = |value: Option<&mut Value>| {
837 let Some(array) = value.and_then(Value::as_array_mut) else {
838 return Vec::new();
839 };
840 array
841 .iter_mut()
842 .filter(|item| item.get("cache_control").is_some())
843 .map(|item| item as *mut Value)
844 .collect::<Vec<_>>()
845 };
846 marked.extend(collect(body.get_mut("tools")));
847 marked.extend(collect(body.get_mut("system")));
848 if let Some(messages) = body.get_mut("messages").and_then(Value::as_array_mut) {
849 for message in messages.iter_mut() {
850 if let Some(blocks) = message.get_mut("content").and_then(Value::as_array_mut) {
851 marked.extend(
852 blocks
853 .iter_mut()
854 .filter(|block| block.get("cache_control").is_some())
855 .map(|block| block as *mut Value),
856 );
857 }
858 }
859 }
860 if marked.len() > MAX_CACHE_BREAKPOINTS {
861 let excess = marked.len() - MAX_CACHE_BREAKPOINTS;
862 for pointer in marked.into_iter().take(excess) {
863 // SAFETY: the pointers were collected from `body`, which is
864 // exclusively borrowed for the duration of this function, and
865 // each pointer targets a distinct JSON node.
866 unsafe {
867 if let Some(map) = (*pointer).as_object_mut() {
868 map.remove("cache_control");
869 }
870 }
871 }
872 }
873 }
874
875 /// Provider event types [`convert_anthropic_sse_data`] accepts. Anything else
876 /// with a string `type` is tolerated as `None` (future additions); note
877 /// `tool_projection_warning` is deliberately absent — it is local-only and
878 /// must never decode from provider SSE.
879 fn is_known_sse_type(event_type: &str) -> bool {
880 matches!(
881 event_type,
882 "message_start"
883 | "content_block_start"
884 | "content_block_delta"
885 | "content_block_stop"
886 | "message_delta"
887 | "message_stop"
888 | "ping"
889 | "error"
890 )
891 }
892
893 /// Peek at an SSE payload's `type` without building a DOM.
894 #[derive(Deserialize)]
895 struct SseTagPeek<'a> {
896 #[serde(borrow)]
897 r#type: Option<&'a str>,
898 }
899
900 /// Convert one SSE `data:` payload into a [`StreamEvent`], normalizing usage
901 /// objects to the #2961 convention. Returns `None` for ignorable payloads.
902 ///
903 /// #6213 T7: the per-token path deserializes directly into the tagged
904 /// [`StreamEvent`] instead of building a `Value` DOM and converting it.
905 /// Usage-bearing events (two per stream) keep the exact legacy path — the
906 /// usage rewrite reads wire fields the normalized [`Usage`] cannot
907 /// represent — and decode failures keep their exact legacy outcomes.
908 fn convert_anthropic_sse_data(data: &str) -> Option<Result<StreamEvent>> {
909 let trimmed = data.trim();
910 if trimmed.is_empty() {
911 return None;
912 }
913 let usage_event = matches!(
914 serde_json::from_str::<SseTagPeek>(trimmed).map(|peek| peek.r#type),
915 Ok(Some("message_start" | "message_delta"))
916 );
917 if usage_event {
918 return convert_anthropic_sse_usage_event(trimmed);
919 }
920 match serde_json::from_str::<StreamEvent>(trimmed) {
921 // Local-only receipt: the legacy path ignored it (not a provider
922 // type), so it stays ignored rather than decoding.
923 Ok(StreamEvent::ToolProjectionWarning { .. }) => None,
924 Ok(event) => Some(Ok(event)),
925 Err(error) => {
926 // Cold path, reached only when direct decode fails: invalid JSON
927 // and unknown types keep their exact legacy outcomes.
928 let value: Value = match serde_json::from_str(trimmed) {
929 Ok(value) => value,
930 Err(e) => return Some(Err(anyhow::anyhow!("invalid SSE JSON: {e}"))),
931 };
932 match value.get("type").and_then(Value::as_str) {
933 // Tolerate unknown event types (e.g. future additions) silently.
934 Some(known) if !is_known_sse_type(known) => None,
935 _ => Some(Err(anyhow::anyhow!("unrecognized SSE event: {error}"))),
936 }
937 }
938 }
939 }
940
941 /// Legacy `Value` path for `message_start`/`message_delta`: the usage
942 /// rewrite reads wire fields the normalized [`Usage`] cannot represent, so
943 /// these two events normalize before decoding, exactly as before.
944 fn convert_anthropic_sse_usage_event(trimmed: &str) -> Option<Result<StreamEvent>> {
945 let mut value: Value = match serde_json::from_str(trimmed) {
946 Ok(value) => value,
947 Err(e) => return Some(Err(anyhow::anyhow!("invalid SSE JSON: {e}"))),
948 };
949
950 match value.get("type").and_then(Value::as_str) {
951 Some("message_start") => {
952 if let Some(message) = value.get_mut("message") {
953 discard_provider_execution_ids(message);
954 if let Some(usage) = message.get_mut("usage") {
955 *usage = json!(parse_anthropic_usage(usage));
956 }
957 }
958 }
959 Some("message_delta") => {
960 if let Some(usage) = value.get_mut("usage") {
961 *usage = json!(parse_anthropic_usage(usage));
962 }
963 }
964 // Tolerate unknown event types (e.g. future additions) silently.
965 Some(known) if !is_known_sse_type(known) => {
966 return None;
967 }
968 _ => {}
969 }
970
971 Some(serde_json::from_value(value).map_err(|e| anyhow::anyhow!("unrecognized SSE event: {e}")))
972 }
973
974 /// Map Anthropic's usage payload onto the normalized [`Usage`] convention
975 /// (#2961 / #4318): hit = cache reads, write = cache creation, miss = raw
976 /// uncached input, `input_tokens` = the total prompt across all three.
977 fn parse_anthropic_usage(usage: &Value) -> Usage {
978 let field = |name: &str| {
979 usage
980 .get(name)
981 .and_then(Value::as_u64)
982 .and_then(|value| u32::try_from(value).ok())
983 .unwrap_or(0)
984 };
985 let input_raw = field("input_tokens");
986 let cache_creation = field("cache_creation_input_tokens");
987 let cache_read = field("cache_read_input_tokens");
988 let output = field("output_tokens");
989
990 Usage {
991 input_tokens: input_raw
992 .saturating_add(cache_creation)
993 .saturating_add(cache_read),
994 output_tokens: output,
995 prompt_cache_hit_tokens: Some(cache_read),
996 prompt_cache_miss_tokens: Some(input_raw),
997 prompt_cache_write_tokens: Some(cache_creation),
998 reasoning_tokens: None,
999 reasoning_replay_tokens: None,
1000 server_tool_use: None,
1001 }
1002 }
1003
1004 /// Extract `error.type` / `error.message` from an Anthropic error envelope
1005 /// (`{"type":"error","error":{"type":...,"message":...}}`), falling back to
1006 /// the raw body so nothing is swallowed.
1007 fn anthropic_error_fields(error: &Value) -> (String, String) {
1008 let error_type = error
1009 .get("type")
1010 .and_then(Value::as_str)
1011 .unwrap_or("unknown")
1012 .to_string();
1013 let message = error
1014 .get("message")
1015 .and_then(Value::as_str)
1016 .map(str::to_string)
1017 .unwrap_or_else(|| error.to_string());
1018 (error_type, message)
1019 }
1020
1021 #[cfg(test)]
1022 mod tests {
1023 use super::*;
1024 use codewhale_models::Role;
1025 use codewhale_models::{CacheControl, Message, SystemBlock, SystemPrompt, Tool};
1026
1027 #[test]
1028 fn provider_messages_cannot_supply_host_execution_identity() {
1029 for supplied in [json!("forged-local"), json!({"invalid": "host id"})] {
1030 let value = json!({
1031 "id": "provider-message", "type": "message", "role": "assistant",
1032 "model": "claude-sonnet-4-6", "stop_reason": "tool_use",
1033 "content": [
1034 {"type": "tool_use", "id": "wire-call", "name": "read",
1035 "input": {"execution_id": "ordinary argument"},
1036 "execution_id": supplied,
1037 "caller": {"type": "code_execution", "tool_id": "parent-wire"},
1038 "thought_signature": "provider-signature"},
1039 {"type": "tool_result", "tool_use_id": "wire-call", "content": "result",
1040 "execution_id": supplied},
1041 ],
1042 "usage": {"input_tokens": 3, "output_tokens": 2},
1043 });
1044 let decoded = decode_anthropic_message(value.clone()).unwrap();
1045 let event = convert_anthropic_sse_data(
1046 &json!({
1047 "type": "message_start", "message": value,
1048 })
1049 .to_string(),
1050 )
1051 .unwrap()
1052 .unwrap();
1053 let StreamEvent::MessageStart { message: streamed } = event else {
1054 panic!("expected message start")
1055 };
1056 for message in [decoded, streamed] {
1057 let ContentBlock::ToolUse {
1058 id,
1059 input,
1060 execution_id,
1061 caller,
1062 thought_signature,
1063 ..
1064 } = &message.content[0]
1065 else {
1066 panic!("expected tool use")
1067 };
1068 assert!(execution_id.is_none());
1069 assert_eq!(id, "wire-call");
1070 assert_eq!(input["execution_id"], "ordinary argument");
1071 assert_eq!(
1072 caller.as_ref().unwrap().tool_id.as_deref(),
1073 Some("parent-wire")
1074 );
1075 assert_eq!(thought_signature.as_deref(), Some("provider-signature"));
1076 assert!(matches!(
1077 &message.content[1],
1078 ContentBlock::ToolResult {
1079 execution_id: None,
1080 ..
1081 }
1082 ));
1083 assert_eq!(message.usage.input_tokens, 3);
1084 }
1085 }
1086 }
1087
1088 fn request_with(
1089 model: &str,
1090 reasoning_effort: Option<&str>,
1091 temperature: Option<f32>,
1092 top_p: Option<f32>,
1093 ) -> MessageRequest {
1094 MessageRequest {
1095 model: model.to_string(),
1096 messages: vec![Message {
1097 role: Role::User,
1098 content: vec![ContentBlock::Text {
1099 text: "hello".to_string(),
1100 cache_control: None,
1101 }],
1102 }],
1103 max_tokens: 1024,
1104 system: Some(SystemPrompt::Blocks(vec![SystemBlock {
1105 block_type: "text".to_string(),
1106 text: "be helpful".to_string(),
1107 cache_control: Some(CacheControl {
1108 cache_type: "ephemeral".to_string(),
1109 }),
1110 }])),
1111 tools: None,
1112 tool_choice: None,
1113 metadata: None,
1114 thinking: None,
1115 reasoning_effort: reasoning_effort.map(str::to_string),
1116 stream: Some(true),
1117 temperature,
1118 top_p,
1119 }
1120 }
1121
1122 fn test_client() -> CodewhaleClient {
1123 anthropic_test_client(None)
1124 }
1125
1126 /// #6378: the engine stores each tool result as its own user message, so
1127 /// a parallel batch reaches the wire as `assistant{a, b}`, `user{a}`,
1128 /// `user{b}`. Both results must land in the one user turn after the batch,
1129 /// and the dangling-use repair must not answer `b` a second time.
1130 #[test]
1131 fn parallel_tool_results_split_across_user_turns_are_answered_once() {
1132 let client = test_client();
1133 let mut request = request_with("claude-sonnet-4-6", None, None, None);
1134 let tool_use = |id: &str, path: &str| ContentBlock::ToolUse {
1135 execution_id: None,
1136 id: id.to_string(),
1137 name: "read".to_string(),
1138 input: json!({ "path": path }),
1139 caller: None,
1140 thought_signature: None,
1141 };
1142 let tool_result = |id: &str, content: &str| ContentBlock::ToolResult {
1143 execution_id: None,
1144 tool_use_id: id.to_string(),
1145 content: content.to_string(),
1146 is_error: None,
1147 content_blocks: None,
1148 };
1149 request.messages = vec![
1150 Message {
1151 role: Role::User,
1152 content: vec![ContentBlock::Text {
1153 text: "Read a.txt and b.txt".to_string(),
1154 cache_control: None,
1155 }],
1156 },
1157 Message {
1158 role: Role::Assistant,
1159 content: vec![
1160 ContentBlock::Text {
1161 text: "I will read both files in parallel.".to_string(),
1162 cache_control: None,
1163 },
1164 tool_use("toolu_a", "a.txt"),
1165 tool_use("toolu_b", "b.txt"),
1166 ],
1167 },
1168 Message {
1169 role: Role::User,
1170 content: vec![tool_result("toolu_a", "content of file A")],
1171 },
1172 Message {
1173 role: Role::User,
1174 content: vec![tool_result("toolu_b", "content of file B")],
1175 },
1176 ];
1177
1178 let body = client.build_anthropic_body(&request, true);
1179 let messages = body["messages"].as_array().expect("messages array");
1180 assert_eq!(
1181 messages.len(),
1182 3,
1183 "both results share one user turn: {body}"
1184 );
1185 let results = messages[2]["content"].as_array().expect("user content");
1186 assert_eq!(
1187 results
1188 .iter()
1189 .map(|block| (block["tool_use_id"].as_str(), block["content"].as_str()))
1190 .collect::<Vec<_>>(),
1191 vec![
1192 (Some("toolu_a"), Some("content of file A")),
1193 (Some("toolu_b"), Some("content of file B")),
1194 ],
1195 "{body}"
1196 );
1197 assert!(
1198 results.iter().all(|block| block.get("is_error").is_none()),
1199 "{body}"
1200 );
1201 assert!(!body.to_string().contains(UNEXECUTED_TOOL_RESULT), "{body}");
1202 }
1203
1204 fn anthropic_test_client(base_url: Option<&str>) -> CodewhaleClient {
1205 let _ = rustls::crypto::ring::default_provider().install_default();
1206 let config = crate::config::Config {
1207 provider: Some("anthropic".to_string()),
1208 providers: Some(crate::config::ProvidersConfig {
1209 anthropic: crate::config::ProviderConfig {
1210 api_key: Some("test-key".to_string()),
1211 base_url: base_url.map(str::to_string),
1212 ..Default::default()
1213 },
1214 ..Default::default()
1215 }),
1216 ..Default::default()
1217 };
1218 CodewhaleClient::new(&config).expect("anthropic client constructs")
1219 }
1220
1221 fn minimax_test_client() -> CodewhaleClient {
1222 minimax_test_client_for(crate::config::DEFAULT_MINIMAX_ANTHROPIC_BASE_URL)
1223 }
1224
1225 fn minimax_test_client_for(base_url: &str) -> CodewhaleClient {
1226 let _ = rustls::crypto::ring::default_provider().install_default();
1227 let config = crate::config::Config {
1228 provider: Some("minimax-anthropic".to_string()),
1229 providers: Some(crate::config::ProvidersConfig {
1230 minimax_anthropic: crate::config::ProviderConfig {
1231 api_key: Some("test-key".to_string()),
1232 base_url: Some(base_url.to_string()),
1233 ..Default::default()
1234 },
1235 ..Default::default()
1236 }),
1237 ..Default::default()
1238 };
1239 CodewhaleClient::new(&config).expect("MiniMax Messages client constructs")
1240 }
1241
1242 fn deepseek_test_client(base_url: &str) -> CodewhaleClient {
1243 let _ = rustls::crypto::ring::default_provider().install_default();
1244 let config = crate::config::Config {
1245 provider: Some("deepseek-anthropic".to_string()),
1246 providers: Some(crate::config::ProvidersConfig {
1247 deepseek_anthropic: crate::config::ProviderConfig {
1248 api_key: Some("test-key".to_string()),
1249 base_url: Some(base_url.to_string()),
1250 ..Default::default()
1251 },
1252 ..Default::default()
1253 }),
1254 ..Default::default()
1255 };
1256 CodewhaleClient::new(&config).expect("DeepSeek Messages client constructs")
1257 }
1258
1259 fn modelstudio_test_client(base_url: &str) -> CodewhaleClient {
1260 let _ = rustls::crypto::ring::default_provider().install_default();
1261 let config = crate::config::Config {
1262 provider: Some("modelstudio-token-plan-anthropic".to_string()),
1263 providers: Some(crate::config::ProvidersConfig {
1264 // Durable secret-store keys share a family slot, but a literal
1265 // config key belongs to the selected route's own table.
1266 modelstudio_token_plan_anthropic: crate::config::ProviderConfig {
1267 api_key: Some("test-key".to_string()),
1268 base_url: Some(base_url.to_string()),
1269 ..Default::default()
1270 },
1271 ..Default::default()
1272 }),
1273 ..Default::default()
1274 };
1275 CodewhaleClient::new(&config).expect("Model Studio Messages client constructs")
1276 }
1277
1278 #[test]
1279 fn body_keeps_native_cache_control_on_system_and_tools() {
1280 let client = test_client();
1281 let mut request = request_with("claude-sonnet-4-6", Some("high"), None, None);
1282 request.tools = Some(vec![Tool {
1283 tool_type: None,
1284 name: "read_file".to_string(),
1285 description: "Read a file".to_string(),
1286 input_schema: json!({"type": "object", "additionalProperties": false}),
1287 allowed_callers: None,
1288 defer_loading: None,
1289 input_examples: None,
1290 strict: Some(true),
1291 cache_control: None,
1292 }]);
1293
1294 let body = client.build_anthropic_body(&request, true);
1295
1296 assert_eq!(
1297 body.pointer("/system/0/cache_control/type")
1298 .and_then(Value::as_str),
1299 Some("ephemeral"),
1300 "system cache_control must survive natively: {body}"
1301 );
1302 assert_eq!(
1303 body.pointer("/tools/0/strict").and_then(Value::as_bool),
1304 Some(true)
1305 );
1306 assert_eq!(
1307 body.pointer("/tools/0/cache_control/type")
1308 .and_then(Value::as_str),
1309 Some("ephemeral"),
1310 "breakpoint 1 lands on the last tool: {body}"
1311 );
1312 // Breakpoint 2 lands on the latest user turn's last block.
1313 assert_eq!(
1314 body.pointer("/messages/0/content/0/cache_control/type")
1315 .and_then(Value::as_str),
1316 Some("ephemeral")
1317 );
1318 }
1319
1320 #[test]
1321 fn body_maps_reasoning_effort_to_adaptive_thinking_and_effort() {
1322 let client = test_client();
1323
1324 let body = client.build_anthropic_body(
1325 &request_with("claude-sonnet-4-6", Some("high"), None, None),
1326 true,
1327 );
1328 assert_eq!(
1329 body.pointer("/thinking/type").and_then(Value::as_str),
1330 Some("adaptive")
1331 );
1332 assert_eq!(
1333 body.pointer("/output_config/effort")
1334 .and_then(Value::as_str),
1335 Some("high")
1336 );
1337
1338 let body = client.build_anthropic_body(
1339 &request_with("claude-opus-4-8", Some("xhigh"), None, None),
1340 true,
1341 );
1342 assert_eq!(
1343 body.pointer("/output_config/effort")
1344 .and_then(Value::as_str),
1345 Some("max")
1346 );
1347
1348 let body = client.build_anthropic_body(
1349 &request_with("claude-sonnet-4-6", Some("off"), None, None),
1350 true,
1351 );
1352 assert!(body.get("thinking").is_none(), "off omits thinking: {body}");
1353 assert!(body.get("output_config").is_none());
1354
1355 // Haiku is not thinking-capable: no thinking, no effort.
1356 let body = client.build_anthropic_body(
1357 &request_with("claude-haiku-4-5", Some("high"), None, None),
1358 true,
1359 );
1360 assert!(body.get("thinking").is_none(), "{body}");
1361 assert!(body.get("output_config").is_none(), "{body}");
1362 }
1363
1364 #[test]
1365 fn compat_gateway_sends_enabled_budget_thinking_instead_of_adaptive() {
1366 // #4978: strict Anthropic-compatible gateways (e.g. Sensenova) reject
1367 // {"type":"adaptive"} with a 400; non-native routes must send the
1368 // documented enabled+budget shape and no output_config.
1369 let client = anthropic_test_client(Some("https://api.sensenova.example/v1"));
1370
1371 let mut request = request_with("claude-sonnet-4-6", Some("high"), None, None);
1372 request.max_tokens = 64_000;
1373 let body = client.build_anthropic_body(&request, true);
1374 assert_eq!(
1375 body.pointer("/thinking/type").and_then(Value::as_str),
1376 Some("enabled"),
1377 "{body}"
1378 );
1379 assert_eq!(
1380 body.pointer("/thinking/budget_tokens")
1381 .and_then(Value::as_u64),
1382 Some(16_384)
1383 );
1384 assert!(body.get("output_config").is_none(), "{body}");
1385
1386 // Effort tiers map onto budgets, capped below max_tokens.
1387 let mut request = request_with("claude-sonnet-4-6", Some("max"), None, None);
1388 request.max_tokens = 64_000;
1389 let body = client.build_anthropic_body(&request, true);
1390 assert_eq!(
1391 body.pointer("/thinking/budget_tokens")
1392 .and_then(Value::as_u64),
1393 Some(32_768)
1394 );
1395 let mut request = request_with("claude-sonnet-4-6", Some("max"), None, None);
1396 request.max_tokens = 8_000;
1397 let body = client.build_anthropic_body(&request, true);
1398 assert_eq!(
1399 body.pointer("/thinking/budget_tokens")
1400 .and_then(Value::as_u64),
1401 Some(7_999),
1402 "budget stays below max_tokens: {body}"
1403 );
1404
1405 // Unspecified effort defaults to the "high" tier.
1406 let mut request = request_with("claude-sonnet-4-6", None, None, None);
1407 request.max_tokens = 64_000;
1408 let body = client.build_anthropic_body(&request, true);
1409 assert_eq!(
1410 body.pointer("/thinking/type").and_then(Value::as_str),
1411 Some("enabled")
1412 );
1413 assert_eq!(
1414 body.pointer("/thinking/budget_tokens")
1415 .and_then(Value::as_u64),
1416 Some(16_384)
1417 );
1418
1419 // "off" and requests too small for the 1024-token minimum budget
1420 // omit thinking entirely.
1421 let mut request = request_with("claude-sonnet-4-6", Some("off"), None, None);
1422 request.max_tokens = 64_000;
1423 let body = client.build_anthropic_body(&request, true);
1424 assert!(body.get("thinking").is_none(), "{body}");
1425 let body = client.build_anthropic_body(
1426 &request_with("claude-sonnet-4-6", Some("high"), None, None),
1427 true,
1428 );
1429 assert!(
1430 body.get("thinking").is_none(),
1431 "max_tokens=1024 cannot fit the minimum budget: {body}"
1432 );
1433
1434 // The native route keeps adaptive; the compat shape is only for
1435 // non-anthropic.com hosts.
1436 let native = test_client().build_anthropic_body(
1437 &request_with("claude-sonnet-4-6", Some("high"), None, None),
1438 true,
1439 );
1440 assert_eq!(
1441 native.pointer("/thinking/type").and_then(Value::as_str),
1442 Some("adaptive")
1443 );
1444 }
1445
1446 #[test]
1447 fn dangling_tool_use_gets_placeholder_tool_result() {
1448 // #5002: an orphaned tool_use with no matching tool_result poisons
1449 // the conversation with repeated 400s; request preparation must
1450 // repair it with an explicit placeholder result.
1451 let client = test_client();
1452 let mut request = request_with("claude-sonnet-4-6", None, None, None);
1453 request.messages = vec![
1454 Message {
1455 role: Role::User,
1456 content: vec![ContentBlock::Text {
1457 text: "run both tools".to_string(),
1458 cache_control: None,
1459 }],
1460 },
1461 Message {
1462 role: Role::Assistant,
1463 content: vec![
1464 ContentBlock::ToolUse {
1465 execution_id: None,
1466 id: "toolu_ok".to_string(),
1467 name: "read_file".to_string(),
1468 input: json!({"path": "a.txt"}),
1469 caller: None,
1470 thought_signature: None,
1471 },
1472 ContentBlock::ToolUse {
1473 execution_id: None,
1474 id: "toolu_orphan".to_string(),
1475 name: "task".to_string(),
1476 input: json!({}),
1477 caller: None,
1478 thought_signature: None,
1479 },
1480 ],
1481 },
1482 // Pre-dispatch failure left only one tool_result behind.
1483 Message {
1484 role: Role::User,
1485 content: vec![ContentBlock::ToolResult {
1486 execution_id: None,
1487 tool_use_id: "toolu_ok".to_string(),
1488 content: "contents".to_string(),
1489 is_error: None,
1490 content_blocks: None,
1491 }],
1492 },
1493 // Trailing assistant tool_use with no user turn at all.
1494 Message {
1495 role: Role::Assistant,
1496 content: vec![ContentBlock::ToolUse {
1497 execution_id: None,
1498 id: "toolu_tail".to_string(),
1499 name: "task".to_string(),
1500 input: json!({}),
1501 caller: None,
1502 thought_signature: None,
1503 }],
1504 },
1505 ];
1506
1507 let body = client.build_anthropic_body(&request, true);
1508 let messages = body["messages"].as_array().expect("messages array");
1509 assert_eq!(messages.len(), 5, "a repair turn is appended: {body}");
1510
1511 // The orphaned id gets a leading placeholder; the answered one is
1512 // untouched (no duplicate result).
1513 let repaired = messages[2]["content"].as_array().expect("user content");
1514 assert_eq!(repaired.len(), 2, "{body}");
1515 assert_eq!(repaired[0]["type"].as_str(), Some("tool_result"));
1516 assert_eq!(repaired[0]["tool_use_id"].as_str(), Some("toolu_orphan"));
1517 assert_eq!(
1518 repaired[0]["content"].as_str(),
1519 Some(UNEXECUTED_TOOL_RESULT)
1520 );
1521 assert_eq!(repaired[0]["is_error"].as_bool(), Some(true));
1522 assert_eq!(repaired[1]["tool_use_id"].as_str(), Some("toolu_ok"));
1523 assert_eq!(repaired[1]["content"].as_str(), Some("contents"));
1524
1525 // The trailing tool_use gains a synthesized user turn.
1526 assert_eq!(messages[4]["role"].as_str(), Some("user"));
1527 let tail = messages[4]["content"].as_array().expect("tail content");
1528 assert_eq!(tail.len(), 1, "{body}");
1529 assert_eq!(tail[0]["type"].as_str(), Some("tool_result"));
1530 assert_eq!(tail[0]["tool_use_id"].as_str(), Some("toolu_tail"));
1531 assert_eq!(tail[0]["content"].as_str(), Some(UNEXECUTED_TOOL_RESULT));
1532
1533 // A fully answered history is left alone.
1534 request.messages.truncate(3);
1535 request.messages[1].content.retain(
1536 |block| !matches!(block, ContentBlock::ToolUse { id, ..} if id == "toolu_orphan"),
1537 );
1538 let body = client.build_anthropic_body(&request, true);
1539 let messages = body["messages"].as_array().expect("messages array");
1540 assert_eq!(messages.len(), 3, "no repair turn appended: {body}");
1541 let untouched = messages[2]["content"].as_array().expect("user content");
1542 assert_eq!(untouched.len(), 1, "{body}");
1543 assert_eq!(untouched[0]["tool_use_id"].as_str(), Some("toolu_ok"));
1544 }
1545
1546 #[test]
1547 fn modelstudio_messages_body_requests_thinking_with_budget() {
1548 // Model Studio's Anthropic-compatible endpoint documents the portable
1549 // {"type":"enabled","budget_tokens":N} shape plus {"type":"disabled"}
1550 // (alibabacloud.com/help/en/model-studio/anthropic-api-messages).
1551 let client = modelstudio_test_client(
1552 "https://token-plan.ap-southeast-1.maas.aliyuncs.com/apps/anthropic",
1553 );
1554
1555 let mut request = request_with("qwen3.8-max", Some("high"), None, None);
1556 request.max_tokens = 64_000;
1557 let body = client.build_anthropic_body(&request, true);
1558 assert_eq!(
1559 body.pointer("/thinking/type").and_then(Value::as_str),
1560 Some("enabled"),
1561 "{body}"
1562 );
1563 assert!(
1564 body.pointer("/thinking/budget_tokens")
1565 .and_then(Value::as_u64)
1566 .is_some(),
1567 "{body}"
1568 );
1569 assert!(body.get("output_config").is_none(), "{body}");
1570 assert_eq!(
1571 body.get("model").and_then(Value::as_str),
1572 Some("qwen3.8-max"),
1573 "{body}"
1574 );
1575
1576 // An explicit "off" is honored on the wire instead of silently
1577 // falling through to the server default (thinking-ON for qwen3.x).
1578 let mut request = request_with("qwen3.8-max", Some("off"), None, None);
1579 request.max_tokens = 64_000;
1580 let body = client.build_anthropic_body(&request, true);
1581 assert_eq!(
1582 body.pointer("/thinking/type").and_then(Value::as_str),
1583 Some("disabled"),
1584 "{body}"
1585 );
1586 }
1587
1588 #[test]
1589 fn deepseek_messages_body_retires_aliases_and_keeps_thinking_control() {
1590 let client = deepseek_test_client(crate::config::DEFAULT_DEEPSEEK_ANTHROPIC_BASE_URL);
1591
1592 let chat = client.build_anthropic_body(
1593 &request_with("deepseek-chat", Some("off"), None, None),
1594 true,
1595 );
1596 assert_eq!(
1597 chat.get("model").and_then(Value::as_str),
1598 Some(crate::config::DEEPSEEK_ALIAS_REPLACEMENT)
1599 );
1600 assert_eq!(
1601 chat.pointer("/thinking/type").and_then(Value::as_str),
1602 Some("disabled")
1603 );
1604
1605 let reasoner = client.build_anthropic_body(
1606 &request_with("deepseek-reasoner", Some("high"), None, None),
1607 true,
1608 );
1609 assert_eq!(
1610 reasoner.get("model").and_then(Value::as_str),
1611 Some(crate::config::DEEPSEEK_ALIAS_REPLACEMENT)
1612 );
1613 assert_eq!(
1614 reasoner.pointer("/thinking/type").and_then(Value::as_str),
1615 Some("adaptive")
1616 );
1617 assert_eq!(
1618 reasoner
1619 .pointer("/output_config/effort")
1620 .and_then(Value::as_str),
1621 Some("high")
1622 );
1623
1624 let custom = deepseek_test_client("https://messages.example/v1");
1625 let custom_body = custom.build_anthropic_body(
1626 &request_with("deepseek-reasoner", Some("high"), None, None),
1627 true,
1628 );
1629 assert_eq!(
1630 custom_body.get("model").and_then(Value::as_str),
1631 Some("deepseek-reasoner")
1632 );
1633 }
1634
1635 #[test]
1636 fn omitted_alias_effort_is_migrated_into_deepseek_messages_body() {
1637 for (alias, expected_effort, expected_thinking) in [
1638 ("deepseek-chat", "off", "disabled"),
1639 ("deepseek-reasoner", "high", "adaptive"),
1640 ] {
1641 let mut config = crate::config::Config {
1642 provider: Some("deepseek-anthropic".to_string()),
1643 providers: Some(crate::config::ProvidersConfig {
1644 deepseek_anthropic: crate::config::ProviderConfig {
1645 api_key: Some("test-key".to_string()),
1646 model: Some(alias.to_string()),
1647 ..Default::default()
1648 },
1649 ..Default::default()
1650 }),
1651 ..Default::default()
1652 };
1653 assert!(
1654 config.reasoning_effort().is_none(),
1655 "fixture must omit effort"
1656 );
1657
1658 crate::config::normalize_model_config_for_test(&mut config);
1659 let client = CodewhaleClient::new(&config).expect("DeepSeek Messages client");
1660 let model = config.default_model();
1661 let body = client.build_anthropic_body(
1662 &request_with(&model, config.reasoning_effort(), None, None),
1663 true,
1664 );
1665
1666 assert_eq!(
1667 body.get("model").and_then(Value::as_str),
1668 Some(crate::config::DEEPSEEK_ALIAS_REPLACEMENT),
1669 "{alias}: {body}"
1670 );
1671 assert_eq!(config.reasoning_effort(), Some(expected_effort));
1672 assert_eq!(
1673 body.pointer("/thinking/type").and_then(Value::as_str),
1674 Some(expected_thinking),
1675 "{alias}: {body}"
1676 );
1677 if alias == "deepseek-reasoner" {
1678 assert_eq!(
1679 body.pointer("/output_config/effort")
1680 .and_then(Value::as_str),
1681 Some("high"),
1682 "{body}"
1683 );
1684 } else {
1685 assert!(body.get("output_config").is_none(), "{body}");
1686 }
1687 }
1688 }
1689
1690 #[test]
1691 fn minimax_body_uses_supported_thinking_controls() {
1692 let client = minimax_test_client();
1693 let body =
1694 client.build_anthropic_body(&request_with("MiniMax-M3", Some("off"), None, None), true);
1695 assert_eq!(
1696 body.pointer("/thinking/type").and_then(Value::as_str),
1697 Some("disabled")
1698 );
1699 assert!(body.get("output_config").is_none(), "{body}");
1700
1701 let mut enabled_bodies = Vec::new();
1702 for effort in ["high", "max"] {
1703 let body = client
1704 .build_anthropic_body(&request_with("MiniMax-M3", Some(effort), None, None), true);
1705 assert_eq!(
1706 body.pointer("/thinking/type").and_then(Value::as_str),
1707 Some("adaptive"),
1708 "{effort}: {body}"
1709 );
1710 assert!(body.get("output_config").is_none(), "{effort}: {body}");
1711 enabled_bodies.push(body);
1712 }
1713 assert_eq!(
1714 enabled_bodies[0].get("thinking"),
1715 enabled_bodies[1].get("thinking"),
1716 "MiniMax high/max select the same untiered adaptive wire control"
1717 );
1718 }
1719
1720 #[test]
1721 fn minimax_messages_reasoning_controls_require_exact_first_party_m3_route() {
1722 for (base_url, model) in [
1723 (
1724 "https://gateway.example/anthropic",
1725 crate::config::DEFAULT_MINIMAX_MODEL,
1726 ),
1727 (
1728 crate::config::DEFAULT_MINIMAX_ANTHROPIC_BASE_URL,
1729 "MiniMax-M2",
1730 ),
1731 ] {
1732 let client = minimax_test_client_for(base_url);
1733 for effort in ["off", "high", "max"] {
1734 let body = client
1735 .build_anthropic_body(&request_with(model, Some(effort), None, None), true);
1736 assert!(
1737 body.get("thinking").is_none(),
1738 "{base_url} {model} {effort}: {body}"
1739 );
1740 assert!(
1741 body.get("output_config").is_none(),
1742 "{base_url} {model} {effort}: {body}"
1743 );
1744 }
1745 }
1746 }
1747
1748 #[test]
1749 fn body_drops_sampling_params_for_models_that_reject_them() {
1750 let client = test_client();
1751
1752 let body = client.build_anthropic_body(
1753 &request_with("claude-opus-4-8", None, Some(0.7), Some(0.9)),
1754 true,
1755 );
1756 assert!(body.get("temperature").is_none(), "{body}");
1757 assert!(body.get("top_p").is_none(), "{body}");
1758
1759 // Older models accept ONE of temperature / top_p (temperature wins).
1760 let body = client.build_anthropic_body(
1761 &request_with("claude-sonnet-4-6", None, Some(0.7), Some(0.9)),
1762 true,
1763 );
1764 assert_eq!(
1765 body.get("temperature").and_then(Value::as_f64),
1766 Some(f64::from(0.7f32))
1767 );
1768 assert!(body.get("top_p").is_none(), "never send both: {body}");
1769 }
1770
1771 #[test]
1772 fn body_replays_signed_thinking_and_drops_unsigned_placeholders() {
1773 let client = test_client();
1774 let mut request = request_with("claude-sonnet-4-6", None, None, None);
1775 request.messages = vec![
1776 Message {
1777 role: Role::User,
1778 content: vec![ContentBlock::Text {
1779 text: "do the thing".to_string(),
1780 cache_control: None,
1781 }],
1782 },
1783 Message {
1784 role: Role::Assistant,
1785 content: vec![
1786 ContentBlock::Thinking {
1787 thinking: "signed reasoning".to_string(),
1788 signature: Some("sig-abc".to_string()),
1789 state: None,
1790 },
1791 ContentBlock::Thinking {
1792 thinking: "(reasoning omitted)".to_string(),
1793 signature: None,
1794 state: None,
1795 },
1796 ContentBlock::ToolUse {
1797 execution_id: None,
1798 id: "toolu_1".to_string(),
1799 name: "read_file".to_string(),
1800 input: json!({"path": "a.txt"}),
1801 caller: None,
1802 thought_signature: None,
1803 },
1804 ],
1805 },
1806 Message {
1807 role: Role::User,
1808 content: vec![ContentBlock::ToolResult {
1809 execution_id: None,
1810 tool_use_id: "toolu_1".to_string(),
1811 content: "contents".to_string(),
1812 is_error: None,
1813 content_blocks: None,
1814 }],
1815 },
1816 ];
1817
1818 let body = client.build_anthropic_body(&request, true);
1819 let assistant = &body["messages"][1]["content"];
1820 assert_eq!(assistant.as_array().map(Vec::len), Some(2));
1821 assert_eq!(
1822 assistant[0]["signature"].as_str(),
1823 Some("sig-abc"),
1824 "signed thinking replays verbatim: {assistant}"
1825 );
1826 assert_eq!(assistant[1]["type"].as_str(), Some("tool_use"));
1827 assert!(
1828 assistant[1].get("caller").is_none(),
1829 "internal caller metadata must not reach the wire"
1830 );
1831 assert_eq!(
1832 body["messages"][2]["content"][0]["type"].as_str(),
1833 Some("tool_result")
1834 );
1835 }
1836
1837 #[test]
1838 fn breakpoints_are_capped_at_four_dropping_earliest() {
1839 let client = test_client();
1840 let mut request = request_with("claude-sonnet-4-6", None, None, None);
1841 // Five caller-marked user turns + the two placed breakpoints.
1842 request.messages = (0..5)
1843 .map(|i| Message {
1844 role: Role::User,
1845 content: vec![ContentBlock::Text {
1846 text: format!("turn {i}"),
1847 cache_control: Some(CacheControl {
1848 cache_type: "ephemeral".to_string(),
1849 }),
1850 }],
1851 })
1852 .collect();
1853
1854 let body = client.build_anthropic_body(&request, true);
1855 let mut count = 0;
1856 if body.pointer("/system/0/cache_control").is_some() {
1857 count += 1;
1858 }
1859 for message in body["messages"].as_array().unwrap() {
1860 for block in message["content"].as_array().unwrap() {
1861 if block.get("cache_control").is_some() {
1862 count += 1;
1863 }
1864 }
1865 }
1866 assert!(
1867 count <= MAX_CACHE_BREAKPOINTS,
1868 "breakpoints must be capped at {MAX_CACHE_BREAKPOINTS}, got {count}: {body}"
1869 );
1870 // The latest user turn keeps its marker (longest prefix coverage).
1871 assert!(
1872 body.pointer("/messages/4/content/0/cache_control")
1873 .is_some(),
1874 "{body}"
1875 );
1876 }
1877
1878 #[test]
1879 fn sse_fixture_decodes_text_thinking_signature_and_tool_use() {
1880 use codewhale_models::{ContentBlockStart, Delta};
1881
1882 let events = [
1883 r#"{"type":"message_start","message":{"id":"msg_01","type":"message","role":"assistant","content":[],"model":"claude-sonnet-4-6","stop_reason":null,"stop_sequence":null,"usage":{"input_tokens":3,"cache_creation_input_tokens":2045,"cache_read_input_tokens":18000,"output_tokens":1}}}"#,
1884 r#"{"type":"content_block_start","index":0,"content_block":{"type":"thinking","thinking":""}}"#,
1885 r#"{"type":"content_block_delta","index":0,"delta":{"type":"thinking_delta","thinking":"Let me check"}}"#,
1886 r#"{"type":"content_block_delta","index":0,"delta":{"type":"signature_delta","signature":"sig-xyz"}}"#,
1887 r#"{"type":"content_block_stop","index":0}"#,
1888 r#"{"type":"content_block_start","index":1,"content_block":{"type":"text","text":""}}"#,
1889 r#"{"type":"content_block_delta","index":1,"delta":{"type":"text_delta","text":"Reading the file."}}"#,
1890 r#"{"type":"content_block_stop","index":1}"#,
1891 r#"{"type":"content_block_start","index":2,"content_block":{"type":"tool_use","id":"toolu_9","name":"read_file","input":{}}}"#,
1892 r#"{"type":"content_block_delta","index":2,"delta":{"type":"input_json_delta","partial_json":"{\"path\":"}}"#,
1893 r#"{"type":"content_block_delta","index":2,"delta":{"type":"input_json_delta","partial_json":"\"a.txt\"}"}}"#,
1894 r#"{"type":"content_block_stop","index":2}"#,
1895 r#"{"type":"ping"}"#,
1896 r#"{"type":"message_delta","delta":{"stop_reason":"tool_use","stop_sequence":null},"usage":{"output_tokens":42}}"#,
1897 r#"{"type":"message_stop"}"#,
1898 ];
1899
1900 let decoded: Vec<StreamEvent> = events
1901 .iter()
1902 .map(|data| {
1903 convert_anthropic_sse_data(data)
1904 .expect("known event")
1905 .expect("decodes")
1906 })
1907 .collect();
1908
1909 // message_start usage normalized to the #2961 convention.
1910 let StreamEvent::MessageStart { message } = &decoded[0] else {
1911 panic!("expected MessageStart, got {:?}", decoded[0]);
1912 };
1913 assert_eq!(message.usage.input_tokens, 3 + 2045 + 18000);
1914 assert_eq!(message.usage.prompt_cache_hit_tokens, Some(18000));
1915 assert_eq!(message.usage.prompt_cache_miss_tokens, Some(3));
1916 assert_eq!(message.usage.prompt_cache_write_tokens, Some(2045));
1917
1918 assert!(matches!(
1919 &decoded[1],
1920 StreamEvent::ContentBlockStart {
1921 content_block: ContentBlockStart::Thinking { .. },
1922 ..
1923 }
1924 ));
1925 assert!(matches!(
1926 &decoded[3],
1927 StreamEvent::ContentBlockDelta {
1928 delta: Delta::SignatureDelta { signature },
1929 ..
1930 } if signature == "sig-xyz"
1931 ));
1932 assert!(matches!(
1933 &decoded[6],
1934 StreamEvent::ContentBlockDelta {
1935 delta: Delta::TextDelta { text },
1936 ..
1937 } if text == "Reading the file."
1938 ));
1939 let mut tool_json = String::new();
1940 for event in &decoded {
1941 if let StreamEvent::ContentBlockDelta {
1942 delta: Delta::InputJsonDelta { partial_json },
1943 ..
1944 } = event
1945 {
1946 tool_json.push_str(partial_json);
1947 }
1948 }
1949 assert_eq!(
1950 serde_json::from_str::<Value>(&tool_json).expect("accumulated tool args parse"),
1951 json!({"path": "a.txt"})
1952 );
1953 assert!(matches!(&decoded[12], StreamEvent::Ping));
1954 let StreamEvent::MessageDelta { delta, usage } = &decoded[13] else {
1955 panic!("expected MessageDelta");
1956 };
1957 assert_eq!(delta.stop_reason.as_deref(), Some("tool_use"));
1958 assert_eq!(usage.as_ref().map(|u| u.output_tokens), Some(42));
1959 assert!(matches!(&decoded[14], StreamEvent::MessageStop));
1960 }
1961
1962 #[test]
1963 fn sse_error_event_and_unknown_events_are_handled() {
1964 let error = convert_anthropic_sse_data(
1965 r#"{"type":"error","error":{"type":"overloaded_error","message":"Overloaded"}}"#,
1966 )
1967 .expect("error event decodes")
1968 .expect("error event is a StreamEvent");
1969 let StreamEvent::Error { error } = error else {
1970 panic!("expected StreamEvent::Error");
1971 };
1972 let (error_type, message) = anthropic_error_fields(&error);
1973 assert_eq!(error_type, "overloaded_error");
1974 assert_eq!(message, "Overloaded");
1975
1976 assert!(
1977 convert_anthropic_sse_data(r#"{"type":"content_block_started_v2","index":0}"#)
1978 .is_none(),
1979 "unknown event types are tolerated"
1980 );
1981 assert!(convert_anthropic_sse_data(" ").is_none());
1982 }
1983
1984 #[test]
1985 fn sse_decode_failures_keep_legacy_outcomes_on_the_direct_path() {
1986 // Malformed JSON: the invalid-input error, not the unrecognized one.
1987 let error = convert_anthropic_sse_data("{oops")
1988 .expect("malformed is Some")
1989 .expect_err("malformed is Err");
1990 assert!(error.to_string().contains("invalid SSE JSON"), "{error:?}");
1991 // Structurally invalid known event: unrecognized, not tolerated.
1992 let error = convert_anthropic_sse_data(r#"{"type":"content_block_stop"}"#)
1993 .expect("known type is Some")
1994 .expect_err("missing index is Err");
1995 assert!(
1996 error.to_string().contains("unrecognized SSE event"),
1997 "{error:?}"
1998 );
1999 // Local-only receipt: never provider SSE, stays ignored.
2000 assert!(
2001 convert_anthropic_sse_data(
2002 r#"{"type":"tool_projection_warning","provider":"x","omitted_tool_names":[],"omitted_tool_count":0}"#
2003 )
2004 .is_none()
2005 );
2006 }
2007
2008 #[test]
2009 fn usage_mapping_handles_missing_cache_fields() {
2010 let usage = parse_anthropic_usage(&json!({"input_tokens": 10, "output_tokens": 5}));
2011 assert_eq!(usage.input_tokens, 10);
2012 assert_eq!(usage.output_tokens, 5);
2013 assert_eq!(usage.prompt_cache_hit_tokens, Some(0));
2014 assert_eq!(usage.prompt_cache_miss_tokens, Some(10));
2015 assert_eq!(usage.prompt_cache_write_tokens, Some(0));
2016 }
2017
2018 #[test]
2019 fn usage_mapping_keeps_cache_write_separate_from_miss() {
2020 let usage = parse_anthropic_usage(&json!({
2021 "input_tokens": 3,
2022 "cache_creation_input_tokens": 2045,
2023 "cache_read_input_tokens": 18000,
2024 "output_tokens": 1,
2025 }));
2026 assert_eq!(usage.input_tokens, 3 + 2045 + 18000);
2027 assert_eq!(usage.prompt_cache_hit_tokens, Some(18000));
2028 assert_eq!(usage.prompt_cache_miss_tokens, Some(3));
2029 assert_eq!(usage.prompt_cache_write_tokens, Some(2045));
2030 }
2031
2032 #[test]
2033 fn data_url_image_becomes_a_base64_source_not_a_url_source() {
2034 // Anthropic rejects a `data:` URL under `{"type":"url"}`. This is the
2035 // whole reason the projection exists; if it regresses, every locally
2036 // attached screenshot 400s on the native route.
2037 let block = content_block_to_anthropic(&ContentBlock::ImageUrl {
2038 image_url: codewhale_models::ImageUrlContent {
2039 url: "data:image/png;base64,QUJD".to_string(),
2040 },
2041 })
2042 .expect("image block");
2043
2044 assert_eq!(block["type"], "image");
2045 assert_eq!(block["source"]["type"], "base64");
2046 assert_eq!(block["source"]["media_type"], "image/png");
2047 assert_eq!(block["source"]["data"], "QUJD");
2048 assert!(
2049 block["source"].get("url").is_none(),
2050 "base64 sources must not carry a url field: {block}"
2051 );
2052 }
2053
2054 #[test]
2055 fn tool_result_image_stays_inside_the_native_tool_result_block() {
2056 let content = anthropic_tool_result_content(
2057 "screenshot captured",
2058 Some(&[json!({
2059 "type": "image",
2060 "mime_type": "image/png",
2061 "data": "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR4nGP4z8DwHwAFAAH/iZk9HQAAAABJRU5ErkJggg==",
2062 })]),
2063 );
2064 let blocks = content.as_array().expect("rich tool_result content");
2065
2066 assert_eq!(
2067 blocks[0],
2068 json!({"type": "text", "text": "screenshot captured"})
2069 );
2070 assert_eq!(blocks[1]["type"], "image");
2071 assert_eq!(blocks[1]["source"]["type"], "base64");
2072 assert_eq!(blocks[1]["source"]["media_type"], "image/png");
2073 assert_eq!(
2074 blocks[1]["source"]["data"],
2075 "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR4nGP4z8DwHwAFAAH/iZk9HQAAAABJRU5ErkJggg=="
2076 );
2077 }
2078
2079 #[test]
2080 fn remote_image_url_stays_a_url_source() {
2081 let block = content_block_to_anthropic(&ContentBlock::ImageUrl {
2082 image_url: codewhale_models::ImageUrlContent {
2083 url: "https://example.com/shot.png".to_string(),
2084 },
2085 })
2086 .expect("image block");
2087
2088 assert_eq!(block["type"], "image");
2089 assert_eq!(block["source"]["type"], "url");
2090 assert_eq!(block["source"]["url"], "https://example.com/shot.png");
2091 }
2092
2093 #[test]
2094 fn unrepresentable_image_reference_degrades_to_visible_text() {
2095 for url in [
2096 "file:///tmp/shot.png",
2097 "/tmp/shot.png",
2098 "data:image/png,QUJD",
2099 ] {
2100 let block = content_block_to_anthropic(&ContentBlock::ImageUrl {
2101 image_url: codewhale_models::ImageUrlContent {
2102 url: url.to_string(),
2103 },
2104 })
2105 .expect("block");
2106
2107 assert_eq!(block["type"], "text", "{url} should degrade: {block}");
2108 assert!(
2109 block["text"].as_str().expect("text").contains(url),
2110 "the degraded text should name the reference: {block}"
2111 );
2112 }
2113 }
2114
2115 #[test]
2116 fn messages_url_tolerates_v1_suffix() {
2117 assert_eq!(
2118 anthropic_messages_url("https://api.anthropic.com"),
2119 "https://api.anthropic.com/v1/messages"
2120 );
2121 assert_eq!(
2122 anthropic_messages_url("https://api.anthropic.com/"),
2123 "https://api.anthropic.com/v1/messages"
2124 );
2125 assert_eq!(
2126 anthropic_messages_url("https://gateway.example/v1"),
2127 "https://gateway.example/v1/messages"
2128 );
2129 assert_eq!(
2130 anthropic_messages_url("https://api.deepseek.com/anthropic"),
2131 "https://api.deepseek.com/anthropic/v1/messages"
2132 );
2133 assert_eq!(
2134 anthropic_messages_url("https://api.minimax.io/anthropic"),
2135 "https://api.minimax.io/anthropic/v1/messages"
2136 );
2137 assert_eq!(
2138 anthropic_messages_url("https://api.minimaxi.com/anthropic"),
2139 "https://api.minimaxi.com/anthropic/v1/messages"
2140 );
2141 }
2142
2143 #[test]
2144 fn anthropic_body_serializes_the_child_catalog_without_duplication() {
2145 // The real child catalog fixture (not a hand-built tool list) must
2146 // survive Messages serialization with exactly one canonical `read`
2147 // entry — no dedup, filter, or sanitizer may drop or duplicate it.
2148 // `load_skill` is eager in DEFAULT_ACTIVE_NATIVE_TOOLS, and children
2149 // resolve the same catalog authority the parent does, so it appears
2150 // here exactly once like any other default tool.
2151 let tools = crate::tools::subagent::kimi_general_child_request_tools_fixture();
2152 assert_eq!(
2153 tools.iter().filter(|tool| tool.name == "read").count(),
2154 1,
2155 "catalog fixture carries one canonical read"
2156 );
2157 assert_eq!(
2158 tools
2159 .iter()
2160 .filter(|tool| tool.name == "load_skill")
2161 .count(),
2162 1,
2163 "child wire catalog carries one canonical load_skill"
2164 );
2165 let client = test_client();
2166 let mut request = request_with("claude-sonnet-4-6", None, None, None);
2167 request.tools = Some(tools);
2168 let body = client.build_anthropic_body(&request, true);
2169 let serialized = body["tools"]
2170 .as_array()
2171 .expect("tools serialize as an array");
2172 let reads: Vec<_> = serialized
2173 .iter()
2174 .filter(|tool| tool["name"] == "read")
2175 .collect();
2176 assert_eq!(
2177 reads.len(),
2178 1,
2179 "exactly one canonical read definition reaches the Messages wire"
2180 );
2181 assert!(
2182 reads[0]["input_schema"]["properties"].is_object(),
2183 "read keeps a valid object schema: {}",
2184 reads[0]
2185 );
2186 assert_eq!(
2187 serialized
2188 .iter()
2189 .filter(|tool| tool["name"] == "load_skill")
2190 .count(),
2191 1,
2192 "exactly one canonical load_skill definition reaches the Messages wire"
2193 );
2194 }
2195
2196 #[tokio::test]
2197 async fn anthropic_stream_opens_through_shared_seam_preserving_headers() {
2198 use futures_util::StreamExt;
2199 use wiremock::matchers::{header, method, path};
2200 use wiremock::{Mock, MockServer, ResponseTemplate};
2201
2202 let server = MockServer::start().await;
2203 // The wire-specific Accept header must survive the shared stream-entry
2204 // open path; the mock only answers when it is present.
2205 Mock::given(method("POST"))
2206 .and(path("/v1/messages"))
2207 .and(header("Accept", "text/event-stream"))
2208 .respond_with(
2209 ResponseTemplate::new(200)
2210 .insert_header("Content-Type", "text/event-stream")
2211 .set_body_string("data: {\"type\":\"message_stop\"}\n\n"),
2212 )
2213 .expect(1)
2214 .mount(&server)
2215 .await;
2216
2217 let client = deepseek_test_client(&server.uri());
2218 let mut stream = client
2219 .handle_anthropic_stream(
2220 &client
2221 .prepare_outbound_request(request_with("deepseek-v4", None, None, None), true)
2222 .expect("anthropic request prepares"),
2223 )
2224 .await
2225 .expect("stream opens through the shared seam");
2226
2227 let mut saw_stop = false;
2228 tokio::time::timeout(std::time::Duration::from_secs(5), async {
2229 while let Some(event) = stream.next().await {
2230 if matches!(event.expect("stream event"), StreamEvent::MessageStop) {
2231 saw_stop = true;
2232 }
2233 }
2234 })
2235 .await
2236 .expect("stream finishes after message_stop");
2237 assert!(saw_stop, "message_stop should arrive through the seam");
2238 }
2239
2240 async fn collect_anthropic_stream(body: &'static str) -> Vec<Result<StreamEvent>> {
2241 use futures_util::StreamExt;
2242 use wiremock::matchers::{method, path};
2243 use wiremock::{Mock, MockServer, ResponseTemplate};
2244
2245 let server = MockServer::start().await;
2246 Mock::given(method("POST"))
2247 .and(path("/v1/messages"))
2248 .respond_with(
2249 ResponseTemplate::new(200)
2250 .insert_header("Content-Type", "text/event-stream")
2251 .set_body_string(body),
2252 )
2253 .mount(&server)
2254 .await;
2255 let client = deepseek_test_client(&server.uri());
2256 let stream = client
2257 .handle_anthropic_stream(
2258 &client
2259 .prepare_outbound_request(request_with("deepseek-v4", None, None, None), true)
2260 .expect("anthropic request prepares"),
2261 )
2262 .await
2263 .expect("stream opens");
2264 tokio::time::timeout(std::time::Duration::from_secs(5), stream.collect())
2265 .await
2266 .expect("stream ends")
2267 }
2268
2269 #[tokio::test]
2270 async fn anthropic_stream_joins_multiline_data_fields_into_one_event() {
2271 use codewhale_models::Delta;
2272 // One JSON payload split across two `data:` fields of a single event.
2273 let events = collect_anthropic_stream(concat!(
2274 "event: content_block_delta\n",
2275 "data: {\"type\":\"content_block_delta\",\"index\":0,\n",
2276 "data: \"delta\":{\"type\":\"text_delta\",\"text\":\"joined\"}}\n\n",
2277 "data: {\"type\":\"message_stop\"}\n\n",
2278 ))
2279 .await;
2280 assert!(
2281 events.iter().any(|event| matches!(
2282 event,
2283 Ok(StreamEvent::ContentBlockDelta {
2284 delta: Delta::TextDelta { text },
2285 ..
2286 }) if text == "joined"
2287 )),
2288 "the split event was lost: {events:?}"
2289 );
2290 assert!(matches!(events.last(), Some(Ok(StreamEvent::MessageStop))));
2291 }
2292
2293 #[tokio::test]
2294 async fn anthropic_stream_eof_without_message_stop_is_an_error() {
2295 let events = collect_anthropic_stream(
2296 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"partial\"}}\n\n",
2297 )
2298 .await;
2299 let last = events.last().expect("the delta and a terminal item");
2300 assert!(
2301 last.as_ref()
2302 .is_err_and(|error| error.to_string().contains("closed before message_stop")),
2303 "a truncated stream must not end quietly: {events:?}"
2304 );
2305 }
2306
2307 /// Fault injection (#6184): a provider that answers the headers and then
2308 /// sends nothing fails the stream at the first-byte bound with a
2309 /// distinct error and a `crashes/` stall record, instead of holding the
2310 /// turn for the full idle budget.
2311 #[tokio::test]
2312 async fn stall_first_byte_timeout_fails_stream_and_records_stall() {
2313 use futures_util::StreamExt;
2314 use tokio::io::{AsyncReadExt, AsyncWriteExt};
2315
2316 let dir = tempfile::tempdir().expect("tempdir");
2317 crate::core::engine::turn_heartbeat::set_test_stall_record_dir(Some(
2318 dir.path().to_path_buf(),
2319 ));
2320 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
2321 .await
2322 .expect("bind");
2323 let base_url = format!("http://{}", listener.local_addr().expect("addr"));
2324 let server = tokio::spawn(async move {
2325 let (mut socket, _) = listener.accept().await.expect("accept");
2326 let mut buf = vec![0u8; 64 * 1024];
2327 let _ = socket.read(&mut buf).await;
2328 socket
2329 .write_all(
2330 b"HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nTransfer-Encoding: chunked\r\n\r\n",
2331 )
2332 .await
2333 .expect("headers");
2334 // Hold the connection open with no body bytes.
2335 tokio::time::sleep(std::time::Duration::from_secs(60)).await;
2336 drop(socket);
2337 });
2338
2339 let mut client = deepseek_test_client(&base_url);
2340 client.stream_idle_timeout = std::time::Duration::from_secs(1);
2341 let started = std::time::Instant::now();
2342 let mut stream = client
2343 .handle_anthropic_stream(
2344 &client
2345 .prepare_outbound_request(request_with("deepseek-v4", None, None, None), true)
2346 .expect("anthropic request prepares"),
2347 )
2348 .await
2349 .expect("headers arrive");
2350 let error = tokio::time::timeout(std::time::Duration::from_secs(10), async {
2351 loop {
2352 match stream.next().await {
2353 Some(Err(error)) => break error,
2354 Some(Ok(_)) => continue,
2355 None => panic!("stream ended without the first-byte error"),
2356 }
2357 }
2358 })
2359 .await
2360 .expect("first-byte bound fires");
2361 assert!(error.to_string().contains("first-byte timeout"), "{error}");
2362 assert!(started.elapsed() < std::time::Duration::from_secs(10));
2363 let records: Vec<String> = std::fs::read_dir(dir.path())
2364 .expect("record dir")
2365 .flatten()
2366 .filter_map(|entry| std::fs::read_to_string(entry.path()).ok())
2367 .collect();
2368 assert_eq!(records.len(), 1, "{records:?}");
2369 assert!(records[0].contains("first byte"), "{}", records[0]);
2370 server.abort();
2371 crate::core::engine::turn_heartbeat::set_test_stall_record_dir(None);
2372 }
2373
2374 #[tokio::test]
2375 async fn anthropic_stream_open_error_is_not_retried() {
2376 use wiremock::matchers::{method, path};
2377 use wiremock::{Mock, MockServer, ResponseTemplate};
2378
2379 let server = MockServer::start().await;
2380 // A definitive provider error before any stream body must fail fast:
2381 // exactly one request, no H1 fallback, envelope preserved.
2382 Mock::given(method("POST"))
2383 .and(path("/v1/messages"))
2384 .respond_with(ResponseTemplate::new(401).set_body_string(
2385 "{\"error\":{\"type\":\"authentication_error\",\"message\":\"bad key\"}}",
2386 ))
2387 .expect(1)
2388 .mount(&server)
2389 .await;
2390
2391 let client = deepseek_test_client(&server.uri());
2392 let err = match client
2393 .handle_anthropic_stream(
2394 &client
2395 .prepare_outbound_request(request_with("deepseek-v4", None, None, None), true)
2396 .expect("anthropic request prepares"),
2397 )
2398 .await
2399 {
2400 Ok(_) => panic!("auth errors must fail fast"),
2401 Err(err) => err,
2402 };
2403 // Typed through the shared classifier, provider message kept.
2404 assert!(
2405 matches!(
2406 err.chain()
2407 .find_map(|cause| cause.downcast_ref::<crate::llm_client::LlmError>()),
2408 Some(crate::llm_client::LlmError::AuthenticationError(_))
2409 ),
2410 "a 401 must stay a typed authentication error: {err:#}"
2411 );
2412 assert!(format!("{err:#}").contains("bad key"), "{err:#}");
2413 }
2414
2415 #[tokio::test]
2416 async fn anthropic_rate_limit_is_retried_honoring_retry_after() {
2417 use wiremock::matchers::{method, path};
2418 use wiremock::{Mock, MockServer, ResponseTemplate};
2419
2420 let server = MockServer::start().await;
2421 Mock::given(method("POST"))
2422 .and(path("/v1/messages"))
2423 .respond_with(
2424 ResponseTemplate::new(429)
2425 .insert_header("retry-after", "1")
2426 .set_body_string(
2427 "{\"type\":\"error\",\"error\":{\"type\":\"rate_limit_error\",\"message\":\"slow down\"}}",
2428 ),
2429 )
2430 .up_to_n_times(1)
2431 .expect(1)
2432 .mount(&server)
2433 .await;
2434 Mock::given(method("POST"))
2435 .and(path("/v1/messages"))
2436 .respond_with(
2437 ResponseTemplate::new(200)
2438 .insert_header("Content-Type", "text/event-stream")
2439 .set_body_string("data: {\"type\":\"message_stop\"}\n\n"),
2440 )
2441 .expect(1)
2442 .mount(&server)
2443 .await;
2444
2445 let mut client = deepseek_test_client(&server.uri());
2446 client.retry.enabled = true;
2447 client.retry.max_retries = 2;
2448 client.retry.initial_delay = 0.0;
2449 client.retry.max_delay = 0.0;
2450 let started = std::time::Instant::now();
2451 let stream = client
2452 .handle_anthropic_stream(
2453 &client
2454 .prepare_outbound_request(request_with("deepseek-v4", None, None, None), true)
2455 .expect("anthropic request prepares"),
2456 )
2457 .await;
2458 assert!(
2459 stream.is_ok(),
2460 "a 429 before the stream body must be retried, not surfaced"
2461 );
2462 // The backoff is configured to zero, so only the provider's
2463 // `Retry-After: 1` can account for the wait before the retry.
2464 assert!(
2465 started.elapsed() >= std::time::Duration::from_millis(900),
2466 "the retry must wait out Retry-After, waited {:?}",
2467 started.elapsed()
2468 );
2469 }
2470
2471 /// A `system`-role history message — what a compaction summary, a branch
2472 /// summary, or an imported journal `system` entry becomes once it reaches
2473 /// `MessageRequest::messages` — must not be emitted verbatim: the Messages
2474 /// API accepts only `user` and `assistant` in `messages[].role` and 400s
2475 /// the whole conversation otherwise, on every retry.
2476 #[test]
2477 fn system_role_history_message_is_not_emitted_verbatim_on_the_messages_wire() {
2478 let mut request = request_with("claude-opus-4-6", None, None, None);
2479 request.messages.insert(
2480 0,
2481 Message {
2482 role: Role::System,
2483 content: vec![ContentBlock::Text {
2484 text: "[compaction summary] the user is porting the parser".to_string(),
2485 cache_control: None,
2486 }],
2487 },
2488 );
2489
2490 let body = test_client().build_anthropic_body(&request, false);
2491 let messages = body["messages"].as_array().expect("messages array");
2492
2493 for message in messages {
2494 let role = message["role"].as_str().expect("role is a string");
2495 assert!(
2496 role == "user" || role == "assistant",
2497 "Messages API rejects role {role:?}"
2498 );
2499 }
2500 let carried = messages.iter().any(|message| {
2501 message["content"].as_array().is_some_and(|blocks| {
2502 blocks.iter().any(|block| {
2503 block["text"].as_str()
2504 == Some("[compaction summary] the user is porting the parser")
2505 })
2506 })
2507 });
2508 assert!(carried, "the summary text must survive: {messages:?}");
2509 }
2510 }
2511
2511 lines RUST