Skip to main content

oxicode_ai/providers/
anthropic.rs

1//! Anthropic provider implementation
2
3use futures::{Stream, StreamExt};
4use reqwest::Client;
5use serde::Deserialize;
6use serde_json::Value as JsonValue;
7use std::future::Future;
8use std::pin::Pin;
9use std::sync::Arc;
10
11use super::openai_responses_shared::parse_streaming_json;
12use super::shared_client;
13use super::sse::split_complete_lines;
14use crate::{
15    Api, AssistantMessage, ContentBlock, Context, Model, Provider, ProviderEvent, StopReason,
16    StreamOptions, StreamResult, TextContent, ThinkingContent, ToolCall, Usage,
17    error::ProviderError,
18};
19
20/// Anthropic provider
21///
22/// Supports the Anthropic Messages API (v1/messages) and Anthropic-compatible
23/// providers (MiniMax, etc.) via custom base URLs and extra headers.
24#[derive(Clone)]
25pub struct AnthropicProvider {
26    client: &'static Client,
27    api_key: Option<String>,
28    /// Override base URL. When `None`, falls back to `model.base_url`.
29    base_url: Option<String>,
30    /// Extra HTTP headers to include in every request (e.g. anthropic-version,
31    /// anthropic-beta).
32    extra_headers: Vec<(String, String)>,
33    /// Whether this is the native Anthropic API endpoint.
34    /// When `false` (compatible providers like MiniMax), Anthropic-specific
35    /// features such as the `thinking` parameter and beta headers are
36    /// suppressed because they may not be supported.
37    native: bool,
38}
39
40impl AnthropicProvider {
41    /// Create a new Anthropic provider without an API key.
42    ///
43    /// API keys are resolved at request time via auth.json or StreamOptions.
44    /// Use `with_api_key()` for explicit key injection.
45    pub fn new() -> Self {
46        Self {
47            client: shared_client(),
48            api_key: None,
49            base_url: None,
50            extra_headers: vec![
51                ("anthropic-version".to_string(), "2023-06-01".to_string()),
52                // Enable interleaved thinking and fine-grained tool streaming
53                // (matches opencode's anthropic custom loader defaults)
54                (
55                    "anthropic-beta".to_string(),
56                    "interleaved-thinking-2025-05-14,fine-grained-tool-streaming-2025-05-14"
57                        .to_string(),
58                ),
59            ],
60            native: true,
61        }
62    }
63
64    /// Create with an explicit API key.
65    #[allow(dead_code)]
66    pub fn with_api_key(api_key: impl Into<String>) -> Self {
67        Self {
68            client: shared_client(),
69            api_key: Some(api_key.into()),
70            base_url: None,
71            extra_headers: vec![
72                ("anthropic-version".to_string(), "2023-06-01".to_string()),
73                (
74                    "anthropic-beta".to_string(),
75                    "interleaved-thinking-2025-05-14,fine-grained-tool-streaming-2025-05-14"
76                        .to_string(),
77                ),
78            ],
79            native: true,
80        }
81    }
82
83    /// Create with a custom base URL (for Anthropic-compatible providers like
84    /// MiniMax that expose the Messages API at a different host).
85    ///
86    /// Note: Anthropic-specific beta headers are NOT included because
87    /// third-party compatible providers may not support them.
88    #[allow(dead_code)]
89    pub fn with_base_url(base_url: &str) -> Self {
90        Self {
91            client: shared_client(),
92            api_key: None,
93            base_url: Some(base_url.to_string()),
94            extra_headers: vec![("anthropic-version".to_string(), "2023-06-01".to_string())],
95            native: false,
96        }
97    }
98
99    /// Create with a custom base URL, API key, and extra headers.
100    ///
101    /// Used for registering Anthropic-compatible providers (MiniMax, etc.).
102    ///
103    /// Note: Anthropic-specific beta headers (`interleaved-thinking`,
104    /// `fine-grained-tool-streaming`) are NOT included here because
105    /// third-party compatible providers may not support them, leading
106    /// to stream hangs or protocol errors. Use [`Self::new`] or
107    /// [`Self::with_api_key`] for the native Anthropic endpoint which includes
108    /// beta headers.
109    pub fn with_config(
110        base_url: &str,
111        api_key: Option<String>,
112        extra_headers: Vec<(String, String)>,
113    ) -> Self {
114        let mut headers = vec![("anthropic-version".to_string(), "2023-06-01".to_string())];
115        headers.extend(extra_headers);
116        Self {
117            client: shared_client(),
118            api_key,
119            base_url: Some(base_url.to_string()),
120            extra_headers: headers,
121            native: false,
122        }
123    }
124}
125
126impl Default for AnthropicProvider {
127    fn default() -> Self {
128        Self::new()
129    }
130}
131
132/// Build the Anthropic Messages API URL from a base URL.
133///
134/// Normalizes the base URL by trimming any trailing slash and stripping a
135/// trailing `/v1` segment before appending `/v1/messages`. This prevents the
136/// double-`/v1` path (`.../v1/v1/messages`) for Anthropic-compatible providers
137/// (e.g. MiniMax) whose registered base URL already includes `/v1`
138/// (e.g. `https://api.minimax.io/anthropic/v1`). For a base URL without a
139/// trailing `/v1` (e.g. the default `https://api.anthropic.com`), this is a
140/// no-op and yields the standard `.../v1/messages`.
141fn anthropic_messages_url(base_url: &str) -> String {
142    let trimmed = base_url.trim_end_matches('/');
143    let stripped = trimmed.strip_suffix("/v1").unwrap_or(trimmed);
144    format!("{}/v1/messages", stripped)
145}
146impl Provider for AnthropicProvider {
147    fn stream<'a>(
148        &'a self,
149        model: &'a Model,
150        context: &'a Context,
151        options: Option<StreamOptions>,
152    ) -> Pin<Box<dyn Future<Output = StreamResult> + Send + 'a>> {
153        Box::pin(async move {
154            let options = options.unwrap_or_default();
155
156            // Build the request – use provider base_url override, fall back to model.base_url
157            let effective_base_url = self.base_url.as_deref().unwrap_or(&model.base_url);
158            let url = anthropic_messages_url(effective_base_url);
159
160            // Get API key
161            let api_key = options
162                .api_key
163                .as_ref()
164                .or(self.api_key.as_ref())
165                .ok_or_else(|| ProviderError::MissingApiKey)?;
166
167            // Build messages (apply provider-specific normalization for
168            // Anthropic: filter empty content, reorder tool_use blocks)
169            let normalized = crate::providers::openai::normalize_messages(
170                &context.messages,
171                &model.provider,
172                &model.id,
173            );
174            let messages =
175                build_anthropic_messages_from_normalized(&context.system_prompt, &normalized)?;
176
177            // Build request body
178            let mut body = serde_json::json!({
179                "model": model.id,
180                "messages": messages,
181                "stream": true,
182            });
183
184            // Add system prompt
185            if let Some(ref prompt) = context.system_prompt {
186                body["system"] = serde_json::json!(prompt);
187            }
188
189            // Add optional parameters
190            // Temperature is incompatible with extended thinking (adaptive or
191            // budget-based). Only send when thinking is not active.
192            // Matches pi: `if (options?.temperature !== undefined && !options?.thinkingEnabled)`
193            let thinking_active = body.get("thinking").is_some();
194            if let Some(temp) = options.temperature
195                && !thinking_active
196            {
197                body["temperature"] = serde_json::json!(temp);
198            }
199
200            if let Some(max) = options.max_tokens {
201                body["max_tokens"] = serde_json::json!(max);
202            }
203
204            // Add tools if present
205            if !context.tools.is_empty() {
206                body["tools"] = build_anthropic_tools(&context.tools)?;
207            }
208
209            // ── Thinking / Extended Reasoning ──────────────────────────────
210            // Supports two modes via provider_options.anthropic:
211            //   1. "enabled" with explicit budget_tokens (from thinking_level or custom)
212            //   2. "adaptive" with effort level (Anthropic dynamically allocates budget)
213            //
214            // Falls back to thinking_level-based budget when provider_options is absent.
215            //
216            // IMPORTANT: Only send the `thinking` parameter when connected to the
217            // native Anthropic API (`self.native == true`). Third-party compatible
218            // providers (MiniMax, etc.) may generate thinking content natively but
219            // don't support the explicit `thinking` request parameter — sending it
220            // can cause incomplete responses or prevent tool calls.
221            if model.reasoning && self.native {
222                // Check provider_options first for fine-grained control
223                let anthropic_opts = options
224                    .provider_options
225                    .as_ref()
226                    .and_then(|po| po.anthropic.as_ref());
227
228                if let Some(opts) = anthropic_opts {
229                    // Provider-level override
230                    match opts.thinking_type.as_deref() {
231                        Some("adaptive") => {
232                            // Adaptive thinking — Anthropic chooses budget
233                            body["thinking"] = serde_json::json!({
234                                "type": "adaptive",
235                            });
236                            if let Some(ref effort) = opts.effort {
237                                // effort is not a body param but could influence
238                                // max_tokens allocation
239                                let budget = match effort.as_str() {
240                                    "max" => model.max_tokens.min(31999),
241                                    "xhigh" => (model.max_tokens * 4 / 5).min(31999),
242                                    "high" => (model.max_tokens / 2).min(31999),
243                                    "medium" => (model.max_tokens / 4).min(16000),
244                                    "low" => (model.max_tokens / 8).min(8000),
245                                    _ => (model.max_tokens / 4).min(16000),
246                                };
247                                if body.get("max_tokens").is_none() {
248                                    body["max_tokens"] =
249                                        serde_json::json!((budget + 1024).min(model.max_tokens));
250                                }
251                            }
252                        }
253                        Some("enabled") => {
254                            // Explicit budget from provider_options or thinking_level
255                            let budget = opts.thinking_budget.unwrap_or_else(|| {
256                                compute_thinking_budget(&options.thinking_level, model.max_tokens)
257                            });
258                            if budget > 0 {
259                                if body.get("max_tokens").is_none() {
260                                    body["max_tokens"] =
261                                        serde_json::json!((budget + 1024).min(model.max_tokens));
262                                }
263                                body["thinking"] = serde_json::json!({
264                                    "type": "enabled",
265                                    "budget_tokens": budget,
266                                });
267                            }
268                        }
269                        _ => {
270                            // No explicit thinking_type — use thinking_level fallback
271                            let budget =
272                                compute_thinking_budget(&options.thinking_level, model.max_tokens);
273                            if budget > 0 {
274                                if body.get("max_tokens").is_none() {
275                                    body["max_tokens"] =
276                                        serde_json::json!((budget + 1024).min(model.max_tokens));
277                                }
278                                body["thinking"] = serde_json::json!({
279                                    "type": "enabled",
280                                    "budget_tokens": budget,
281                                });
282                            }
283                        }
284                    }
285                } else if let Some(ref level) = options.thinking_level {
286                    // No provider_options — use thinking_level directly
287                    let budget = compute_thinking_budget(&Some(*level), model.max_tokens);
288                    if budget > 0 {
289                        if body.get("max_tokens").is_none() {
290                            body["max_tokens"] =
291                                serde_json::json!((budget + 1024).min(model.max_tokens));
292                        }
293                        body["thinking"] = serde_json::json!({
294                            "type": "enabled",
295                            "budget_tokens": budget,
296                        });
297                    }
298                }
299            }
300
301            // Ensure max_tokens is always set (Anthropic requires it)
302            // For reasoning models, ensure max_tokens is large enough for
303            // the model to think AND respond. opencode uses OUTPUT_TOKEN_MAX = 32_000
304            // for MiniMax and other reasoning models.
305            //
306            // When the caller sets a small max_tokens (e.g., 4096 default),
307            // reasoning models like MiniMax may think for thousands of tokens
308            // and then stop before generating tool calls because they calculate
309            // they won't have enough room. Bumping to a minimum of 16_384 for
310            // reasoning models ensures the model has space to think + call tools.
311            if model.reasoning
312                && let Some(current) = body.get("max_tokens").and_then(|v| v.as_u64())
313                && current < 16_384
314            {
315                body["max_tokens"] = serde_json::json!(model.max_tokens.min(32_768));
316            }
317            if body.get("max_tokens").is_none() {
318                body["max_tokens"] = serde_json::json!(model.max_tokens.min(16384));
319            }
320
321            // ── Cache Control ─────────────────────────────────────────────
322            // When cache_retention is set, add cache_control breakpoints to
323            // the system prompt and last few messages.
324            //
325            // Anthropic allows at most 4 cache_control breakpoints per request.
326            // We use a counter to enforce this limit, allocating in priority
327            // order: system prompt → last message → second-to-last message.
328            // Mirrors opencode's Cache.Breakpoints with ANTHROPIC_BREAKPOINT_CAP.
329            const ANTHROPIC_BREAKPOINT_CAP: usize = 4;
330            let want_cache = options.cache_retention == Some(crate::CacheRetention::Short)
331                || options.cache_retention == Some(crate::CacheRetention::Long);
332
333            if want_cache {
334                let mut remaining = ANTHROPIC_BREAKPOINT_CAP;
335                let cache_marker = serde_json::json!({ "type": "ephemeral" });
336
337                // 1. System prompt (highest priority)
338                if remaining > 0
339                    && let Some(system) = body.get_mut("system")
340                    && system.is_string()
341                {
342                    *system = serde_json::json!([{
343                        "type": "text",
344                        "text": system,
345                        "cache_control": cache_marker.clone(),
346                    }]);
347                    remaining -= 1;
348                }
349
350                // 2. Last message (high priority)
351                if remaining > 0
352                    && let Some(messages) = body.get_mut("messages").and_then(|m| m.as_array_mut())
353                {
354                    if let Some(last_msg) = messages.last_mut()
355                        && let Some(content) = last_msg.get_mut("content")
356                    {
357                        if let Some(parts) = content.as_array_mut() {
358                            if let Some(last_part) = parts.last_mut() {
359                                last_part["cache_control"] = cache_marker.clone();
360                                remaining -= 1;
361                            }
362                        } else if content.is_string() {
363                            let text = content.take();
364                            *content = serde_json::json!([{
365                                "type": "text",
366                                "text": text,
367                                "cache_control": cache_marker.clone(),
368                            }]);
369                            remaining -= 1;
370                        }
371                    }
372
373                    // 3. Second-to-last message (tool results)
374                    if remaining > 0 {
375                        let msg_count = messages.len();
376                        if msg_count >= 3
377                            && let Some(msg) = messages.get_mut(msg_count - 3)
378                            && let Some(content) = msg.get_mut("content")
379                            && let Some(parts) = content.as_array_mut()
380                            && let Some(last_part) = parts.last_mut()
381                        {
382                            last_part["cache_control"] = cache_marker;
383                            remaining -= 1;
384                        }
385                    }
386                }
387
388                if remaining < ANTHROPIC_BREAKPOINT_CAP {
389                    tracing::debug!(
390                        used = ANTHROPIC_BREAKPOINT_CAP - remaining,
391                        cap = ANTHROPIC_BREAKPOINT_CAP,
392                        "Anthropic cache breakpoints applied"
393                    );
394                }
395            }
396
397            // Build headers
398            let mut headers = reqwest::header::HeaderMap::new();
399            headers.insert(
400                "x-api-key",
401                api_key.parse().map_err(|e| {
402                    ProviderError::InvalidResponse(format!("invalid header value: {e}"))
403                })?,
404            );
405            headers.insert(
406                "content-type",
407                "application/json".parse().map_err(|e| {
408                    ProviderError::InvalidResponse(format!("invalid header value: {e}"))
409                })?,
410            );
411
412            // Provider-level default headers (e.g. anthropic-version, anthropic-beta)
413            for (k, v) in &self.extra_headers {
414                if let (Ok(name), Ok(value)) = (
415                    k.parse::<reqwest::header::HeaderName>(),
416                    v.parse::<reqwest::header::HeaderValue>(),
417                ) {
418                    headers.insert(name, value);
419                }
420            }
421
422            // Per-request headers (from StreamOptions)
423            for (k, v) in &options.headers {
424                if let (Ok(name), Ok(value)) = (
425                    k.parse::<reqwest::header::HeaderName>(),
426                    v.parse::<reqwest::header::HeaderValue>(),
427                ) {
428                    headers.insert(name, value);
429                }
430            }
431
432            // Make request
433            let response = self
434                .client
435                .post(&url)
436                .headers(headers)
437                .json(&body)
438                .send()
439                .await
440                .map_err(ProviderError::RequestFailed)?;
441
442            if !response.status().is_success() {
443                let status = response.status();
444                // Anthropic returns a request-id header — capture it so errors
445                // carry the traceable id (omp AnthropicApiError align).
446                let request_id = response
447                    .headers()
448                    .get("request-id")
449                    .and_then(|v| v.to_str().ok())
450                    .map(str::to_string);
451                let body: String = response.text().await.unwrap_or_default();
452                return Err(ProviderError::HttpError(
453                    crate::error::HttpErrorDetail::new(status.as_u16(), body)
454                        .with_provider("anthropic")
455                        .with_request_id(request_id),
456                ));
457            }
458
459            // Create event stream
460            let model_name = model.id.clone();
461
462            // Stateful scan: persists partial_message and usage ACROSS chunks.
463            //
464            // Previous implementation created a fresh partial_message per chunk,
465            // which caused content loss at chunk boundaries. When an HTTP response
466            // is split into multiple chunks, content_block_delta events from
467            // earlier chunks would be lost because the new chunk's partial_message
468            // started empty. This particularly affected compatible providers
469            // (MiniMax, etc.) that split responses across many small chunks.
470            //
471            // pi (TypeScript) avoids this by keeping a single `output` object
472            // across the entire stream. We replicate that pattern here by
473            // including partial_message in the scan state.
474            //
475            // Tool calls are tracked in `pending_tool_calls` so that partial JSON
476            // arguments are accumulated across `input_json_delta` events and
477            // finalized when `content_block_stop` fires.
478
479            struct AnthropicScanState {
480                pending_bytes: Vec<u8>,
481                partial: AssistantMessage,
482                usage: Usage,
483                /// In-flight tool calls keyed by content block index.
484                pending_tool_calls: std::collections::HashMap<usize, AnthropicPendingToolCall>,
485            }
486
487            let initial_state = AnthropicScanState {
488                pending_bytes: Vec::new(),
489                partial: AssistantMessage::new(Api::AnthropicMessages, "anthropic", &model_name),
490                usage: Usage::default(),
491                pending_tool_calls: std::collections::HashMap::new(),
492            };
493
494            let stream = response
495                .bytes_stream()
496                .scan(
497                    initial_state,
498                    move |state, chunk: Result<bytes::Bytes, reqwest::Error>| {
499                        let events = match chunk {
500                            Ok(bytes) => {
501                                let mut combined =
502                                    Vec::with_capacity(state.pending_bytes.len() + bytes.len());
503                                combined.extend_from_slice(&state.pending_bytes);
504                                combined.extend_from_slice(&bytes);
505                                let (text, trailing) = split_complete_lines(&combined);
506                                state.pending_bytes = trailing;
507                                parse_anthropic_events_stateful(
508                                    &text,
509                                    &mut state.partial,
510                                    &mut state.usage,
511                                    &mut state.pending_tool_calls,
512                                )
513                            }
514                            Err(e) => vec![ProviderEvent::Error {
515                                reason: StopReason::Error,
516                                error: create_error_message(&e.to_string()),
517                            }],
518                        };
519                        async move { Some(futures::stream::iter(events)) }
520                    },
521                )
522                .flatten();
523
524            Ok(Box::pin(stream) as Pin<Box<dyn Stream<Item = ProviderEvent> + Send>>)
525        })
526    }
527}
528
529/// Build messages in Anthropic format from normalized Message structs.
530///
531/// Key behavior for Anthropic compatibility:
532/// - Consecutive `ToolResult` messages are merged into a single `user` message
533///   with multiple `tool_result` content blocks (Anthropic API requirement).
534///   Matches pi's look-ahead pattern.
535fn build_anthropic_messages_from_normalized(
536    _system_prompt: &Option<String>,
537    messages_in: &[crate::Message],
538) -> Result<Vec<JsonValue>, ProviderError> {
539    let mut messages: Vec<JsonValue> = Vec::new();
540    let mut i = 0;
541
542    while i < messages_in.len() {
543        let msg = &messages_in[i];
544        match msg {
545            crate::Message::User(u) => {
546                let content = match &u.content {
547                    crate::MessageContent::Text(s) => vec![serde_json::json!({
548                        "type": "text",
549                        "text": s,
550                    })],
551                    crate::MessageContent::Blocks(blocks) => blocks_to_anthropic_content(blocks)?,
552                };
553                messages.push(serde_json::json!({
554                    "role": "user",
555                    "content": content,
556                }));
557                i += 1;
558            }
559            crate::Message::Assistant(a) => {
560                let content = blocks_to_anthropic_content(&a.content)?;
561                messages.push(serde_json::json!({
562                    "role": "assistant",
563                    "content": content,
564                }));
565                i += 1;
566            }
567            crate::Message::ToolResult(t) => {
568                // Anthropic requires consecutive tool_results to be grouped
569                // into a single user message (pi's look-ahead pattern).
570                let mut tool_results: Vec<JsonValue> = Vec::new();
571                let content = blocks_to_anthropic_content(&t.content)?;
572                tool_results.push(serde_json::json!({
573                    "type": "tool_result",
574                    "tool_use_id": t.tool_call_id,
575                    "content": content,
576                }));
577
578                // Look ahead for consecutive ToolResult messages
579                let mut j = i + 1;
580                while j < messages_in.len() {
581                    if let crate::Message::ToolResult(next_t) = &messages_in[j] {
582                        let next_content = blocks_to_anthropic_content(&next_t.content)?;
583                        tool_results.push(serde_json::json!({
584                            "type": "tool_result",
585                            "tool_use_id": next_t.tool_call_id,
586                            "content": next_content,
587                        }));
588                        j += 1;
589                    } else {
590                        break;
591                    }
592                }
593
594                messages.push(serde_json::json!({
595                    "role": "user",
596                    "content": tool_results,
597                }));
598                i = j;
599            }
600        }
601    }
602
603    Ok(messages)
604}
605
606/// Convert content blocks to Anthropic format
607fn blocks_to_anthropic_content(blocks: &[ContentBlock]) -> Result<Vec<JsonValue>, ProviderError> {
608    let mut items = Vec::new();
609
610    for block in blocks {
611        match block {
612            ContentBlock::Text(t) => {
613                items.push(serde_json::json!({
614                    "type": "text",
615                    "text": t.text,
616                }));
617            }
618            ContentBlock::ToolCall(tc) => {
619                items.push(serde_json::json!({
620                    "type": "tool_use",
621                    "id": tc.id,
622                    "name": tc.name,
623                    "input": tc.arguments,
624                }));
625            }
626            ContentBlock::Thinking(th) => {
627                items.push(serde_json::json!({
628                    "type": "thinking",
629                    "thinking": th.thinking,
630                }));
631            }
632            ContentBlock::Image(img) => {
633                items.push(serde_json::json!({
634                    "type": "image",
635                    "source": {
636                        "type": "base64",
637                        "media_type": img.mime_type,
638                        "data": img.data,
639                    },
640                }));
641            }
642            ContentBlock::Unknown(_) => {
643                // Skip unknown blocks
644            }
645        }
646    }
647
648    Ok(items)
649}
650
651/// Compute thinking budget from ThinkingLevel enum + model max_tokens.
652fn compute_thinking_budget(level: &Option<crate::ThinkingLevel>, max_tokens: usize) -> usize {
653    match level {
654        Some(crate::ThinkingLevel::High) | Some(crate::ThinkingLevel::XHigh) => {
655            (max_tokens / 2).min(31999)
656        }
657        Some(crate::ThinkingLevel::Medium) => (max_tokens / 4).min(16000),
658        Some(crate::ThinkingLevel::Low) => (max_tokens / 8).min(8000),
659        Some(crate::ThinkingLevel::Minimal) => (max_tokens / 16).min(4000),
660        _ => 0,
661    }
662}
663
664fn build_anthropic_tools(tools: &[crate::Tool]) -> Result<JsonValue, ProviderError> {
665    let items: Vec<_> = tools
666        .iter()
667        .map(|tool| {
668            serde_json::json!({
669                "name": tool.name,
670                "description": tool.description,
671                "input_schema": tool.parameters,
672            })
673        })
674        .collect();
675
676    Ok(serde_json::json!(items))
677}
678
679/// Parse Anthropic SSE event stream (stateful across chunks).
680///
681/// Unlike the old `parse_anthropic_events` which created a fresh `partial_message`
682/// per invocation, this version takes `&mut AssistantMessage` and `&mut Usage`
683/// so content accumulates correctly across HTTP chunk boundaries.
684///
685/// Tool calls are tracked in `pending_tool_calls` so that partial JSON arguments
686/// are accumulated across `input_json_delta` events and finalized when
687/// Tracks a pending tool call being assembled from Anthropic streaming events.
688///
689/// Anthropic streams tool calls as three SSE events:
690///   1. `content_block_start` → id, name (arguments empty)
691///   2. `content_block_delta` (input_json_delta) → partial JSON fragments
692///   3. `content_block_stop` → complete
693///
694/// We accumulate the partial JSON in `partial_json` and build the final
695/// `ToolCall` when the block completes.
696struct AnthropicPendingToolCall {
697    /// Tool call ID from `content_block_start`.
698    id: String,
699    /// Tool name from `content_block_start`.
700    name: String,
701    /// Accumulated partial JSON arguments from `input_json_delta` events.
702    partial_json: String,
703}
704
705/// Parse Anthropic SSE event stream (stateful across chunks).
706///
707/// Unlike the old `parse_anthropic_events` which created a fresh `partial_message`
708/// per invocation, this version takes `&mut AssistantMessage` and `&mut Usage`
709/// so content accumulates correctly across HTTP chunk boundaries.
710///
711/// Tool calls are tracked in `pending_tool_calls` so that partial JSON arguments
712/// are accumulated across `input_json_delta` events and finalized when
713/// `content_block_stop` fires (matching pi's `content_block_stop → toolcall_end`
714/// pattern). On `content_block_start` (tool_use), a placeholder `ToolCall` is
715/// added to `partial_message.content` so that downstream consumers (streaming.rs)
716/// can see tool calls in the partial snapshot immediately.
717fn parse_anthropic_events_stateful(
718    text: &str,
719    partial_message: &mut AssistantMessage,
720    accumulated_usage: &mut Usage,
721    pending_tool_calls: &mut std::collections::HashMap<usize, AnthropicPendingToolCall>,
722) -> Vec<ProviderEvent> {
723    // F-6 (audit 2026-06-21): length-based estimate replaces 2-pass
724    // count-then-parse scan. See oxicode-ai/src/providers/openai.rs for the
725    // rationale (same heuristic applied uniformly across providers).
726    let mut events = Vec::with_capacity(text.len() / 80);
727
728    for line in text.split('\n') {
729        let line = line.trim_end_matches('\r');
730        if line.is_empty() {
731            continue;
732        }
733
734        if !line.starts_with("data: ") {
735            continue;
736        }
737
738        let data = &line[6..];
739
740        if data == "[DONE]" || data.is_empty() {
741            continue;
742        }
743
744        let event = match serde_json::from_str::<AnthropicEvent>(data) {
745            Ok(e) => e,
746            Err(_) => continue,
747        };
748
749        let event_type = event.type_.as_deref();
750
751        // ── Accumulate usage BEFORE the match statement ──────────────────
752        // Use max() to avoid overwriting values from earlier events with 0
753        // when a later event only includes a subset of fields.
754        if let Some(usage) = &event.usage {
755            accumulated_usage.input = usage.input_tokens.max(accumulated_usage.input);
756            accumulated_usage.output = usage.output_tokens.max(accumulated_usage.output);
757            accumulated_usage.cache_read = usage.cache_read.max(accumulated_usage.cache_read);
758            accumulated_usage.cache_write = usage.cache_creation.max(accumulated_usage.cache_write);
759            accumulated_usage.total_tokens = accumulated_usage.input + accumulated_usage.output;
760        } else if let Some(msg) = &event.message
761            && let Some(usage) = &msg.usage
762        {
763            accumulated_usage.input = usage.input_tokens.max(accumulated_usage.input);
764            accumulated_usage.output = usage.output_tokens.max(accumulated_usage.output);
765            accumulated_usage.cache_read = usage.cache_read.max(accumulated_usage.cache_read);
766            accumulated_usage.cache_write = usage.cache_creation.max(accumulated_usage.cache_write);
767            accumulated_usage.total_tokens = accumulated_usage.input + accumulated_usage.output;
768        }
769
770        match event_type {
771            Some("message_start") => {
772                events.push(ProviderEvent::Start {
773                    partial: Arc::new(partial_message.clone()),
774                });
775            }
776            Some("content_block_start") => {
777                if let Some(block) = &event.content_block {
778                    let idx = block.index.or(event.index).unwrap_or(0);
779                    match block.type_.as_deref() {
780                        Some("text") => {
781                            events.push(ProviderEvent::TextStart {
782                                content_index: idx,
783                                partial: Arc::new(partial_message.clone()),
784                            });
785                        }
786                        Some("thinking") => {
787                            events.push(ProviderEvent::ThinkingStart {
788                                content_index: idx,
789                                partial: Arc::new(partial_message.clone()),
790                            });
791                        }
792                        Some("tool_use") | Some("server_tool_use") => {
793                            // Register the tool call in pending_tool_calls
794                            // so that input_json_delta events can accumulate
795                            // the partial JSON arguments.
796                            let tc_id = block.id.clone().unwrap_or_default();
797                            let tc_name = block.name.clone().unwrap_or_default();
798                            pending_tool_calls.insert(
799                                idx,
800                                AnthropicPendingToolCall {
801                                    id: tc_id.clone(),
802                                    name: tc_name.clone(),
803                                    partial_json: String::new(),
804                                },
805                            );
806                            events.push(ProviderEvent::ToolCallStart {
807                                content_index: idx,
808                                tool_call_id: Some(tc_id),
809                                tool_name: Some(tc_name),
810                                partial: Arc::new(partial_message.clone()),
811                            });
812                        }
813                        Some(t) if t.ends_with("_tool_result") => {
814                            let name = match t {
815                                "web_search_tool_result" => Some("web_search".to_string()),
816                                "code_execution_tool_result" => Some("code_execution".to_string()),
817                                "web_fetch_tool_result" => Some("web_fetch".to_string()),
818                                _ => None,
819                            };
820                            if let Some(tool_name) = name {
821                                let tc = ToolCall::new(
822                                    block.tool_use_id.clone().unwrap_or_default(),
823                                    tool_name,
824                                    serde_json::json!({}),
825                                );
826                                partial_message
827                                    .content
828                                    .push(ContentBlock::ToolCall(tc.clone()));
829                                events.push(ProviderEvent::ToolCallEnd {
830                                    content_index: idx,
831                                    tool_call: tc,
832                                    partial: Arc::new(partial_message.clone()),
833                                });
834                            }
835                        }
836                        _ => {}
837                    }
838                }
839            }
840            Some("content_block_delta") => {
841                if let Some(delta) = &event.delta {
842                    match delta.type_.as_deref() {
843                        Some("text_delta") => {
844                            if let Some(text) = &delta.text {
845                                // Accumulate into partial_message so the TUI
846                                // can diff against its snapshot tracker.
847                                let last_text_idx = partial_message
848                                    .content
849                                    .iter()
850                                    .rposition(|b| matches!(b, ContentBlock::Text(_)));
851                                if let Some(idx) = last_text_idx
852                                    && let ContentBlock::Text(t) = &mut partial_message.content[idx]
853                                {
854                                    t.text.push_str(text);
855                                } else {
856                                    partial_message
857                                        .content
858                                        .push(ContentBlock::Text(TextContent::new(text.clone())));
859                                }
860                                events.push(ProviderEvent::TextDelta {
861                                    content_index: event.index.unwrap_or(0),
862                                    delta: text.clone(),
863                                    partial: Arc::new(partial_message.clone()),
864                                });
865                            }
866                        }
867                        Some("thinking_delta") => {
868                            if let Some(text) = &delta.thinking {
869                                let last_think_idx = partial_message
870                                    .content
871                                    .iter()
872                                    .rposition(|b| matches!(b, ContentBlock::Thinking(_)));
873                                if let Some(idx) = last_think_idx
874                                    && let ContentBlock::Thinking(t) =
875                                        &mut partial_message.content[idx]
876                                {
877                                    t.thinking.push_str(text);
878                                } else {
879                                    partial_message.content.push(ContentBlock::Thinking(
880                                        ThinkingContent::new(text.clone()),
881                                    ));
882                                }
883                                events.push(ProviderEvent::ThinkingDelta {
884                                    content_index: event.index.unwrap_or(0),
885                                    delta: text.clone(),
886                                    partial: Arc::new(partial_message.clone()),
887                                });
888                            }
889                        }
890                        Some("input_json_delta") => {
891                            if let Some(args) = &delta.partial_json {
892                                // Accumulate partial JSON into the pending tool call.
893                                let block_idx = event.index.unwrap_or(0);
894                                if let Some(ptc) = pending_tool_calls.get_mut(&block_idx) {
895                                    ptc.partial_json.push_str(args);
896                                }
897                                events.push(ProviderEvent::ToolCallDelta {
898                                    content_index: block_idx,
899                                    delta: args.clone(),
900                                    partial: Arc::new(partial_message.clone()),
901                                });
902                            }
903                        }
904                        Some("signature_delta") => {
905                            if let Some(_sig) = &delta.signature {
906                                // Signature for session continuity
907                            }
908                        }
909                        _ => {}
910                    }
911                }
912            }
913            Some("content_block_stop") => {
914                // When a tool_use block completes, finalize the accumulated
915                // partial JSON into a proper ToolCall and emit ToolCallEnd.
916                // This matches pi's `content_block_stop → toolcall_end` pattern.
917                let block_idx = event.index.unwrap_or(0);
918                if let Some(ptc) = pending_tool_calls.remove(&block_idx) {
919                    let args_value = parse_streaming_json(&ptc.partial_json);
920                    let tc = ToolCall::new(ptc.id, ptc.name, args_value);
921
922                    // Add the finalized tool call to partial_message so
923                    // downstream consumers (streaming.rs Done handler) can
924                    // see it in the accumulated message.
925                    partial_message
926                        .content
927                        .push(ContentBlock::ToolCall(tc.clone()));
928
929                    tracing::debug!(
930                        block_idx,
931                        tool_id = %tc.id,
932                        tool_name = %tc.name,
933                        "content_block_stop: finalized tool call"
934                    );
935
936                    events.push(ProviderEvent::ToolCallEnd {
937                        content_index: block_idx,
938                        tool_call: tc,
939                        partial: Arc::new(partial_message.clone()),
940                    });
941                }
942            }
943            Some("message_delta") => {
944                if let Some(delta) = &event.delta {
945                    let reason = match delta.stop_reason.as_deref() {
946                        Some("end_turn") | Some("stop_sequence") | Some("pause_turn") => {
947                            StopReason::Stop
948                        }
949                        Some("max_tokens") => StopReason::Length,
950                        Some("tool_use") => StopReason::ToolUse,
951                        Some("refusal") | Some("sensitive") => StopReason::Error,
952                        _ => StopReason::Stop,
953                    };
954
955                    let mut done_msg = partial_message.clone();
956                    done_msg.usage = accumulated_usage.clone();
957                    events.push(ProviderEvent::Done {
958                        reason,
959                        message: done_msg,
960                    });
961                }
962            }
963            Some("message_stop") => {
964                // Message complete – no event needed
965            }
966            _ => {}
967        }
968    }
969
970    events
971}
972
973/// Test-friendly wrapper that creates fresh state for a single-chunk parse.
974/// Used by unit tests that pass complete SSE streams in one go.
975#[cfg(test)]
976fn parse_anthropic_events(text: &str, model_id: &str) -> Vec<ProviderEvent> {
977    let mut partial = AssistantMessage::new(Api::AnthropicMessages, "anthropic", model_id);
978    let mut usage = Usage::default();
979    let mut pending_tool_calls = std::collections::HashMap::new();
980    parse_anthropic_events_stateful(text, &mut partial, &mut usage, &mut pending_tool_calls)
981}
982
983/// Create error assistant message
984fn create_error_message(msg: &str) -> AssistantMessage {
985    let mut message = AssistantMessage::new(Api::AnthropicMessages, "anthropic", "unknown");
986    message.stop_reason = StopReason::Error;
987    message.error_message = Some(msg.to_string());
988    message
989}
990
991// Anthropic event structure
992#[derive(Debug, Deserialize)]
993struct AnthropicEvent {
994    #[serde(rename = "type")]
995    type_: Option<String>,
996    #[serde(rename = "index")]
997    index: Option<usize>,
998    content_block: Option<ContentBlockStart>,
999    delta: Option<Delta>,
1000    usage: Option<AnthropicUsage>,
1001    /// Nested message from `message_start` event — contains initial usage.
1002    message: Option<AnthropicMessageStart>,
1003}
1004
1005/// Nested message object from Anthropic `message_start` events.
1006/// Carries initial token usage (input_tokens, output_tokens) before streaming begins.
1007#[derive(Debug, Deserialize)]
1008struct AnthropicMessageStart {
1009    usage: Option<AnthropicUsage>,
1010}
1011
1012#[derive(Debug, Deserialize)]
1013struct ContentBlockStart {
1014    #[serde(rename = "type")]
1015    type_: Option<String>,
1016    index: Option<usize>,
1017    /// Tool call ID present for tool_use blocks (Anthropic sends this in content_block_start)
1018    id: Option<String>,
1019    /// Tool name for tool_use blocks.
1020    name: Option<String>,
1021    /// Initial thinking content (may be empty string).
1022    #[allow(dead_code)]
1023    thinking: Option<String>,
1024    /// Signature for extended thinking session continuity.
1025    #[allow(dead_code)]
1026    signature: Option<String>,
1027    /// Initial text content (may be non-empty for pre-filled blocks).
1028    #[allow(dead_code)]
1029    text: Option<String>,
1030    /// Tool use ID for server tool result blocks.
1031    tool_use_id: Option<String>,
1032    /// Structured content for server tool result blocks.
1033    #[allow(dead_code)]
1034    content: Option<JsonValue>,
1035}
1036
1037#[derive(Debug, Deserialize)]
1038struct Delta {
1039    #[serde(rename = "type")]
1040    type_: Option<String>,
1041    text: Option<String>,
1042    thinking: Option<String>,
1043    partial_json: Option<String>,
1044    /// Signature for extended thinking session continuity.
1045    /// Sent in `signature_delta` events after a thinking block.
1046    #[allow(dead_code)]
1047    signature: Option<String>,
1048    #[serde(rename = "stop_reason")]
1049    stop_reason: Option<String>,
1050    #[serde(rename = "stop_sequence")]
1051    #[allow(dead_code)]
1052    stop_sequence: Option<String>,
1053}
1054
1055#[derive(Debug, Deserialize)]
1056struct AnthropicUsage {
1057    #[serde(rename = "input_tokens", default)]
1058    input_tokens: usize,
1059    #[serde(rename = "output_tokens", default)]
1060    output_tokens: usize,
1061    /// Cache read tokens (Anthropic field: `cache_read_input_tokens`)
1062    #[serde(rename = "cache_read_input_tokens", alias = "cache_read", default)]
1063    cache_read: usize,
1064    /// Cache write tokens (Anthropic field: `cache_creation_input_tokens`)
1065    #[serde(
1066        rename = "cache_creation_input_tokens",
1067        alias = "cache_creation",
1068        default
1069    )]
1070    cache_creation: usize,
1071}
1072
1073#[cfg(test)]
1074mod tests {
1075    use super::*;
1076
1077    const MODEL: &str = "claude-3-5-sonnet-20241022";
1078
1079    // ── message_start ──────────────────────────────────────────────────
1080
1081    #[test]
1082    fn parse_message_start() {
1083        let sse = "data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\"}}\n";
1084        let events = parse_anthropic_events(sse, MODEL);
1085        assert_eq!(events.len(), 1);
1086        assert!(matches!(&events[0], ProviderEvent::Start { .. }));
1087    }
1088
1089    // ── content_block_start ────────────────────────────────────────────
1090
1091    #[test]
1092    fn parse_text_block_start() {
1093        let sse = "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n";
1094        let events = parse_anthropic_events(sse, MODEL);
1095        assert_eq!(events.len(), 1);
1096        match &events[0] {
1097            ProviderEvent::TextStart { content_index, .. } => assert_eq!(*content_index, 0),
1098            other => panic!("expected TextStart, got {other:?}"),
1099        }
1100    }
1101
1102    #[test]
1103    fn parse_thinking_block_start() {
1104        let sse = "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\",\"thinking\":\"\"}}\n";
1105        let events = parse_anthropic_events(sse, MODEL);
1106        assert_eq!(events.len(), 1);
1107        assert!(matches!(&events[0], ProviderEvent::ThinkingStart { .. }));
1108    }
1109
1110    #[test]
1111    fn parse_tool_use_block_start() {
1112        let sse = "data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tool_1\",\"name\":\"search\"}}\n";
1113        let events = parse_anthropic_events(sse, MODEL);
1114        assert_eq!(events.len(), 1);
1115        match &events[0] {
1116            ProviderEvent::ToolCallStart { content_index, .. } => assert_eq!(*content_index, 1),
1117            other => panic!("expected ToolCallStart, got {other:?}"),
1118        }
1119    }
1120
1121    // ── content_block_delta ────────────────────────────────────────────
1122
1123    #[test]
1124    fn parse_text_delta() {
1125        let sse = "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"Hello\"}}\n";
1126        let events = parse_anthropic_events(sse, MODEL);
1127        assert_eq!(events.len(), 1);
1128        match &events[0] {
1129            ProviderEvent::TextDelta {
1130                delta,
1131                content_index,
1132                ..
1133            } => {
1134                assert_eq!(delta, "Hello");
1135                assert_eq!(*content_index, 0);
1136            }
1137            other => panic!("expected TextDelta, got {other:?}"),
1138        }
1139    }
1140
1141    #[test]
1142    fn parse_thinking_delta() {
1143        let sse = "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"Let me reason...\"}}\n";
1144        let events = parse_anthropic_events(sse, MODEL);
1145        assert_eq!(events.len(), 1);
1146        match &events[0] {
1147            ProviderEvent::ThinkingDelta { delta, .. } => assert_eq!(delta, "Let me reason..."),
1148            other => panic!("expected ThinkingDelta, got {other:?}"),
1149        }
1150    }
1151
1152    #[test]
1153    fn parse_input_json_delta() {
1154        let sse = "data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"city\\\":\\\"SF\\\"}\"}}\n";
1155        let events = parse_anthropic_events(sse, MODEL);
1156        assert_eq!(events.len(), 1);
1157        match &events[0] {
1158            ProviderEvent::ToolCallDelta {
1159                delta,
1160                content_index,
1161                ..
1162            } => {
1163                assert_eq!(delta, "{\"city\":\"SF\"}");
1164                assert_eq!(*content_index, 1);
1165            }
1166            other => panic!("expected ToolCallDelta, got {other:?}"),
1167        }
1168    }
1169
1170    // ── content_block_stop ────────────────────────────────────────────
1171
1172    #[test]
1173    fn parse_content_block_stop_finalizes_tool_call() {
1174        // Full tool call flow: start → delta → stop
1175        let sse = concat!(
1176            "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tool_abc\",\"name\":\"bash\"}}\n",
1177            "\n",
1178            "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"command\\\":\\\"ls\\\"}\"}}\n",
1179            "\n",
1180            "data: {\"type\":\"content_block_stop\",\"index\":0}\n",
1181            "\n"
1182        );
1183        let events = parse_anthropic_events(sse, MODEL);
1184        // ToolCallStart + ToolCallDelta + ToolCallEnd
1185        assert_eq!(events.len(), 3);
1186
1187        // Verify ToolCallEnd has the finalized tool call
1188        let tc_end = events.iter().find_map(|e| match e {
1189            ProviderEvent::ToolCallEnd { tool_call, .. } => Some(tool_call.clone()),
1190            _ => None,
1191        });
1192        let tc = tc_end.expect("Should have ToolCallEnd");
1193        assert_eq!(tc.id, "tool_abc");
1194        assert_eq!(tc.name, "bash");
1195        assert_eq!(tc.arguments, serde_json::json!({"command": "ls"}));
1196    }
1197
1198    #[test]
1199    fn parse_content_block_stop_ignores_non_tool() {
1200        // content_block_stop for a text block should not emit anything
1201        let sse = concat!(
1202            "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n",
1203            "\n",
1204            "data: {\"type\":\"content_block_stop\",\"index\":0}\n",
1205            "\n"
1206        );
1207        let events = parse_anthropic_events(sse, MODEL);
1208        // Only TextStart — content_block_stop for text doesn't emit
1209        assert_eq!(events.len(), 1);
1210        assert!(matches!(&events[0], ProviderEvent::TextStart { .. }));
1211    }
1212
1213    #[test]
1214    fn parse_tool_call_accumulates_across_deltas() {
1215        // Tool arguments split across multiple deltas
1216        let sse = concat!(
1217            "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tool_1\",\"name\":\"edit\"}}\n",
1218            "\n",
1219            "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"path\\\":\\\"tes\"}}\n",
1220            "\n",
1221            "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"t.rs\"}}\n",
1222            "\n",
1223            "data: {\"type\":\"content_block_stop\",\"index\":0}\n",
1224            "\n"
1225        );
1226        let events = parse_anthropic_events(sse, MODEL);
1227        // ToolCallStart + 2×ToolCallDelta + ToolCallEnd
1228        assert_eq!(events.len(), 4);
1229
1230        let tc_end = events.iter().find_map(|e| match e {
1231            ProviderEvent::ToolCallEnd { tool_call, .. } => Some(tool_call.clone()),
1232            _ => None,
1233        });
1234        let tc = tc_end.expect("Should have ToolCallEnd");
1235        assert_eq!(tc.id, "tool_1");
1236        assert_eq!(tc.name, "edit");
1237        // Streaming JSON parser should handle partial "test.rs"
1238        assert_eq!(tc.arguments["path"].as_str(), Some("test.rs"));
1239    }
1240
1241    #[test]
1242    fn parse_tool_call_in_done_message() {
1243        // Full flow with thinking + tool_use + Done
1244        let sse = concat!(
1245            "data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\"}}\n",
1246            "\n",
1247            "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\",\"thinking\":\"\"}}\n",
1248            "\n",
1249            "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"I need to search.\"}}\n",
1250            "\n",
1251            "data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tool_search\",\"name\":\"web_search\"}}\n",
1252            "\n",
1253            "data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"query\\\":\\\"rust async\\\"}\"}}\n",
1254            "\n",
1255            "data: {\"type\":\"content_block_stop\",\"index\":0}\n",
1256            "\n",
1257            "data: {\"type\":\"content_block_stop\",\"index\":1}\n",
1258            "\n",
1259            "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"tool_use\"}}\n"
1260        );
1261        let events = parse_anthropic_events(sse, MODEL);
1262
1263        // Start + ThinkingStart + ThinkingDelta + ToolCallStart + ToolCallDelta
1264        // + content_block_stop(thinking, no emit) + ToolCallEnd + Done
1265        // content_block_stop for thinking index 0 has no pending tool call → no emit
1266        assert!(events.len() >= 6);
1267
1268        // Verify ToolCallEnd is present
1269        let tc_end = events.iter().find_map(|e| match e {
1270            ProviderEvent::ToolCallEnd { tool_call, .. } => Some(tool_call.clone()),
1271            _ => None,
1272        });
1273        let tc = tc_end.expect("Should have ToolCallEnd");
1274        assert_eq!(tc.id, "tool_search");
1275        assert_eq!(tc.name, "web_search");
1276        assert_eq!(tc.arguments, serde_json::json!({"query": "rust async"}));
1277
1278        // Verify Done has ToolUse stop reason
1279        let done = events.iter().find_map(|e| match e {
1280            ProviderEvent::Done { reason, .. } => Some(*reason),
1281            _ => None,
1282        });
1283        assert_eq!(done, Some(StopReason::ToolUse));
1284
1285        // Verify Done message contains tool call
1286        let done_msg = events.iter().find_map(|e| match e {
1287            ProviderEvent::Done { message, .. } => Some(message.clone()),
1288            _ => None,
1289        });
1290        let msg = done_msg.expect("Should have Done event");
1291        let tool_calls: Vec<_> = msg
1292            .content
1293            .iter()
1294            .filter(|b| matches!(b, ContentBlock::ToolCall(_)))
1295            .collect();
1296        assert_eq!(
1297            tool_calls.len(),
1298            1,
1299            "Done message should contain exactly 1 tool call"
1300        );
1301    }
1302
1303    // ── message_delta (completion) ─────────────────────────────────────
1304
1305    #[test]
1306    fn parse_message_delta_end_turn() {
1307        let sse = "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"}}\n";
1308        let events = parse_anthropic_events(sse, MODEL);
1309        assert_eq!(events.len(), 1);
1310        match &events[0] {
1311            ProviderEvent::Done { reason, .. } => assert!(matches!(reason, StopReason::Stop)),
1312            other => panic!("expected Done, got {other:?}"),
1313        }
1314    }
1315
1316    #[test]
1317    fn parse_message_delta_max_tokens() {
1318        let sse = "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"max_tokens\"}}\n";
1319        let events = parse_anthropic_events(sse, MODEL);
1320        match &events[0] {
1321            ProviderEvent::Done { reason, .. } => assert!(matches!(reason, StopReason::Length)),
1322            other => panic!("expected Done with Length, got {other:?}"),
1323        }
1324    }
1325
1326    #[test]
1327    fn parse_message_delta_stop_sequence() {
1328        let sse =
1329            "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"stop_sequence\"}}\n";
1330        let events = parse_anthropic_events(sse, MODEL);
1331        match &events[0] {
1332            ProviderEvent::Done { reason, .. } => assert!(matches!(reason, StopReason::Stop)),
1333            other => panic!("expected Done with Stop, got {other:?}"),
1334        }
1335    }
1336
1337    // ── message_stop ───────────────────────────────────────────────────
1338
1339    #[test]
1340    fn parse_message_stop_no_event_emitted() {
1341        let sse = "data: {\"type\":\"message_stop\"}\n";
1342        let events = parse_anthropic_events(sse, MODEL);
1343        assert!(events.is_empty());
1344    }
1345
1346    // ── Thinking block full flow ───────────────────────────────────────
1347
1348    #[test]
1349    fn parse_thinking_block_flow() {
1350        let sse = concat!(
1351            "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\",\"thinking\":\"\"}}\n",
1352            "\n",
1353            "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"I should\"}}\n",
1354            "\n",
1355            "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\" check this.\"}}\n",
1356            "\n"
1357        );
1358        let events = parse_anthropic_events(sse, MODEL);
1359        assert_eq!(events.len(), 3);
1360        assert!(matches!(&events[0], ProviderEvent::ThinkingStart { .. }));
1361        let thinking: Vec<&str> = events[1..]
1362            .iter()
1363            .filter_map(|e| match e {
1364                ProviderEvent::ThinkingDelta { delta, .. } => Some(delta.as_str()),
1365                _ => None,
1366            })
1367            .collect();
1368        assert_eq!(thinking, vec!["I should", " check this."]);
1369    }
1370
1371    // ── Usage accumulation ─────────────────────────────────────────────
1372
1373    #[test]
1374    fn parse_usage_from_message_start() {
1375        // Usage accumulates from earlier events; Done captures what was accumulated
1376        // *before* the message_delta chunk (since usage updates happen after event emission).
1377        // The message_start carries initial usage, which gets captured in the Done event.
1378        let sse = concat!(
1379            "data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\"},\"usage\":{\"input_tokens\":100,\"output_tokens\":0,\"cache_read\":80,\"cache_creation\":20}}\n",
1380            "\n",
1381            "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"hi\"}}\n",
1382            "\n",
1383            "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"}}\n"
1384        );
1385        let events = parse_anthropic_events(sse, MODEL);
1386        // Start + TextDelta + Done
1387        assert_eq!(events.len(), 3);
1388        match &events[2] {
1389            ProviderEvent::Done { message, .. } => {
1390                // Captures usage from message_start (output_tokens was 0 there)
1391                assert_eq!(message.usage.input, 100);
1392                assert_eq!(message.usage.output, 0);
1393                assert_eq!(message.usage.total_tokens, 100);
1394                assert_eq!(message.usage.cache_read, 80);
1395                assert_eq!(message.usage.cache_write, 20);
1396            }
1397            other => panic!("expected Done, got {other:?}"),
1398        }
1399    }
1400
1401    // ── Cache metrics ──────────────────────────────────────────────────
1402
1403    #[test]
1404    fn parse_cache_metrics() {
1405        let sse = concat!(
1406            "data: {\"type\":\"message_start\",\"usage\":{\"input_tokens\":50,\"output_tokens\":0,\"cache_read\":40,\"cache_creation\":10}}\n",
1407            "\n",
1408            "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"},\"usage\":{\"input_tokens\":50,\"output_tokens\":20,\"cache_read\":40,\"cache_creation\":10}}\n"
1409        );
1410        let events = parse_anthropic_events(sse, MODEL);
1411        // Start + Done
1412        assert_eq!(events.len(), 2);
1413        match &events[1] {
1414            ProviderEvent::Done { message, .. } => {
1415                assert_eq!(message.usage.cache_read, 40);
1416                assert_eq!(message.usage.cache_write, 10);
1417            }
1418            other => panic!("expected Done, got {other:?}"),
1419        }
1420    }
1421
1422    // ── Empty / malformed handling ─────────────────────────────────────
1423
1424    #[test]
1425    fn parse_empty_input() {
1426        let events = parse_anthropic_events("", MODEL);
1427        assert!(events.is_empty());
1428    }
1429
1430    #[test]
1431    fn parse_done_marker_is_ignored() {
1432        // Anthropic uses event-type based termination, but [DONE] should be silently skipped
1433        let sse = "data: [DONE]\n";
1434        let events = parse_anthropic_events(sse, MODEL);
1435        assert!(events.is_empty());
1436    }
1437
1438    #[test]
1439    fn parse_malformed_json_is_skipped() {
1440        let sse = "data: {broken\ndata: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"ok\"}}\n";
1441        let events = parse_anthropic_events(sse, MODEL);
1442        assert_eq!(events.len(), 1);
1443        match &events[0] {
1444            ProviderEvent::TextDelta { delta, .. } => assert_eq!(delta, "ok"),
1445            other => panic!("expected TextDelta, got {other:?}"),
1446        }
1447    }
1448
1449    #[test]
1450    fn parse_non_data_lines_ignored() {
1451        let sse = "event: ping\nid: 42\ndata: {\"type\":\"message_start\"}\n";
1452        let events = parse_anthropic_events(sse, MODEL);
1453        assert_eq!(events.len(), 1);
1454    }
1455
1456    #[test]
1457    fn parse_empty_data_line_skipped() {
1458        let sse = "data: \ndata: {\"type\":\"message_start\"}\n";
1459        let events = parse_anthropic_events(sse, MODEL);
1460        assert_eq!(events.len(), 1);
1461    }
1462
1463    #[test]
1464    fn parse_unknown_event_type_ignored() {
1465        let sse = "data: {\"type\":\"ping\"}\ndata: {\"type\":\"message_start\"}\n";
1466        let events = parse_anthropic_events(sse, MODEL);
1467        assert_eq!(events.len(), 1);
1468    }
1469
1470    #[test]
1471    fn parse_carriage_return_line_endings() {
1472        let sse = "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"CR\"}}\r\n\r\n";
1473        let events = parse_anthropic_events(sse, MODEL);
1474        assert_eq!(events.len(), 1);
1475        match &events[0] {
1476            ProviderEvent::TextDelta { delta, .. } => assert_eq!(delta, "CR"),
1477            other => panic!("expected TextDelta, got {other:?}"),
1478        }
1479    }
1480
1481    // ── Full stream ────────────────────────────────────────────────────
1482
1483    #[test]
1484    fn parse_full_anthropic_stream() {
1485        let sse = concat!(
1486            "data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\"}}\n",
1487            "\n",
1488            "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n",
1489            "\n",
1490            "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"Hello\"}}\n",
1491            "\n",
1492            "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\" world\"}}\n",
1493            "\n",
1494            "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"}}\n",
1495            "\n",
1496            "data: {\"type\":\"message_stop\"}\n"
1497        );
1498        let events = parse_anthropic_events(sse, MODEL);
1499        // Start + TextStart + 2×TextDelta + Done
1500        assert_eq!(events.len(), 5);
1501
1502        assert!(matches!(&events[0], ProviderEvent::Start { .. }));
1503        assert!(matches!(&events[1], ProviderEvent::TextStart { .. }));
1504
1505        let texts: Vec<&str> = events[2..4]
1506            .iter()
1507            .filter_map(|e| match e {
1508                ProviderEvent::TextDelta { delta, .. } => Some(delta.as_str()),
1509                _ => None,
1510            })
1511            .collect();
1512        assert_eq!(texts, vec!["Hello", " world"]);
1513
1514        assert!(matches!(
1515            &events[4],
1516            ProviderEvent::Done {
1517                reason: StopReason::Stop,
1518                ..
1519            }
1520        ));
1521    }
1522
1523    // ── Stateful multi-chunk parsing ────────────────────────────────────
1524
1525    /// Simulate how bytes_stream + scan works: parse two chunks with a
1526    /// shared partial_message and verify content survives across chunks.
1527    #[test]
1528    fn parse_stateful_across_two_chunks() {
1529        let chunk1 = concat!(
1530            "data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\"}}\n",
1531            "\n",
1532            "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\",\"thinking\":\"\"}}\n",
1533            "\n",
1534            "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"Let me\"}}\n",
1535            "\n"
1536        );
1537        let chunk2 = concat!(
1538            "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\" think.\"}}\n",
1539            "\n",
1540            "data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n",
1541            "\n",
1542            "data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"text_delta\",\"text\":\"Hello\"}}\n",
1543            "\n",
1544            "data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"text_delta\",\"text\":\" world\"}}\n",
1545            "\n",
1546            "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"}}\n",
1547            "\n"
1548        );
1549
1550        // State persists across chunks (mimics scan() state)
1551        let mut partial = AssistantMessage::new(Api::AnthropicMessages, "anthropic", MODEL);
1552        let mut usage = Usage::default();
1553        let mut pending_tc: std::collections::HashMap<usize, AnthropicPendingToolCall> =
1554            std::collections::HashMap::new();
1555
1556        let events1 =
1557            parse_anthropic_events_stateful(chunk1, &mut partial, &mut usage, &mut pending_tc);
1558        // Start + ThinkingStart + ThinkingDelta("Let me")
1559        assert_eq!(events1.len(), 3);
1560
1561        // CRITICAL: partial_message should have accumulated thinking content
1562        assert_eq!(partial.content.len(), 1);
1563        match &partial.content[0] {
1564            ContentBlock::Thinking(t) => assert_eq!(t.thinking, "Let me"),
1565            other => panic!("Expected Thinking block, got {:?}", other),
1566        }
1567
1568        let events2 =
1569            parse_anthropic_events_stateful(chunk2, &mut partial, &mut usage, &mut pending_tc);
1570        // ThinkingDelta(" think.") + TextStart + 2×TextDelta + Done
1571        assert_eq!(events2.len(), 5);
1572
1573        // After both chunks, partial_message has both thinking AND text
1574        assert_eq!(partial.content.len(), 2);
1575        match &partial.content[0] {
1576            ContentBlock::Thinking(t) => assert_eq!(t.thinking, "Let me think."),
1577            other => panic!("Expected Thinking block, got {:?}", other),
1578        }
1579        match &partial.content[1] {
1580            ContentBlock::Text(t) => assert_eq!(t.text, "Hello world"),
1581            other => panic!("Expected Text block, got {:?}", other),
1582        }
1583
1584        // Done event should carry the accumulated content
1585        let done = events2.iter().find_map(|e| match e {
1586            ProviderEvent::Done { message, .. } => Some(message.clone()),
1587            _ => None,
1588        });
1589        let done_msg = done.expect("Should have Done event");
1590        assert_eq!(done_msg.content.len(), 2);
1591    }
1592
1593    #[test]
1594    fn parse_stateful_tool_use_across_chunks() {
1595        let chunk1 = concat!(
1596            "data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\"}}\n",
1597            "\n",
1598            "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\",\"thinking\":\"\"}}\n",
1599            "\n",
1600            "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"I should search.\"}}\n",
1601            "\n"
1602        );
1603        let chunk2 = concat!(
1604            "data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tool_1\",\"name\":\"search\"}}\n",
1605            "\n",
1606            "data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"q\\\":\\\"rust\\\"}\"}}\n",
1607            "\n",
1608            "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"tool_use\"}}\n",
1609            "\n"
1610        );
1611
1612        let mut partial = AssistantMessage::new(Api::AnthropicMessages, "anthropic", MODEL);
1613        let mut usage = Usage::default();
1614        let mut pending_tc: std::collections::HashMap<usize, AnthropicPendingToolCall> =
1615            std::collections::HashMap::new();
1616
1617        let events1 =
1618            parse_anthropic_events_stateful(chunk1, &mut partial, &mut usage, &mut pending_tc);
1619        assert_eq!(events1.len(), 3); // Start + ThinkingStart + ThinkingDelta
1620
1621        // Thinking persists into chunk2
1622        assert_eq!(partial.content.len(), 1);
1623
1624        let events2 =
1625            parse_anthropic_events_stateful(chunk2, &mut partial, &mut usage, &mut pending_tc);
1626        // ToolCallStart + ToolCallDelta + Done
1627        // Note: if content_block_stop were included, we'd also get ToolCallEnd
1628        assert!(events2.len() >= 2); // At least ToolCallStart + Done
1629
1630        // Done should have ToolUse stop reason
1631        let done = events2.iter().find_map(|e| match e {
1632            ProviderEvent::Done { reason, .. } => Some(*reason),
1633            _ => None,
1634        });
1635        assert_eq!(done, Some(StopReason::ToolUse));
1636    }
1637    // ── URL construction (double-/v1 prevention) ─────────────────────
1638
1639    #[test]
1640    fn url_strips_trailing_v1_for_minimax() {
1641        // MiniMax registers its base URL with a trailing `/v1`
1642        // (https://api.minimax.io/anthropic/v1). The adapter must not
1643        // double it into `.../v1/v1/messages`.
1644        assert_eq!(
1645            anthropic_messages_url("https://api.minimax.io/anthropic/v1"),
1646            "https://api.minimax.io/anthropic/v1/messages"
1647        );
1648    }
1649
1650    #[test]
1651    fn url_plain_anthropic_base() {
1652        assert_eq!(
1653            anthropic_messages_url("https://api.anthropic.com"),
1654            "https://api.anthropic.com/v1/messages"
1655        );
1656    }
1657
1658    #[test]
1659    fn url_handles_trailing_slash_then_v1() {
1660        assert_eq!(
1661            anthropic_messages_url("https://api.minimax.io/anthropic/v1/"),
1662            "https://api.minimax.io/anthropic/v1/messages"
1663        );
1664    }
1665
1666    #[test]
1667    fn url_preserves_non_v1_path_suffix() {
1668        // `/v1beta` must not be confused with `/v1`.
1669        assert_eq!(
1670            anthropic_messages_url("https://example.com/v1beta"),
1671            "https://example.com/v1beta/v1/messages"
1672        );
1673    }
1674}