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