1use 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#[derive(Clone)]
26pub struct AnthropicProvider {
27 client: &'static Client,
28 api_key: Option<String>,
29 base_url: Option<String>,
31 extra_headers: Vec<(String, String)>,
34 native: bool,
39}
40
41impl AnthropicProvider {
42 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 (
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 #[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 #[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 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
133fn 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 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 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 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 let mut body = serde_json::json!({
180 "model": model.id,
181 "messages": messages,
182 "stream": true,
183 });
184
185 if let Some(ref prompt) = context.system_prompt {
187 body["system"] = serde_json::json!(prompt);
188 }
189
190 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 if !context.tools.is_empty() {
207 body["tools"] = build_anthropic_tools(&context.tools)?;
208 }
209
210 if let Some(choice) = build_tool_choice(options.tool_choice.as_ref()) {
212 body["tool_choice"] = choice;
213 }
214
215 if model.reasoning && self.native {
228 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 match opts.thinking_type.as_deref() {
237 Some("adaptive") => {
238 body["thinking"] = serde_json::json!({
240 "type": "adaptive",
241 });
242 if let Some(ref effort) = opts.effort {
243 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 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 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 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 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 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 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 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 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 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 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 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 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 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 let model_name = model.id.clone();
467
468 struct AnthropicScanState {
486 pending_bytes: Vec<u8>,
487 partial: AssistantMessage,
488 usage: Usage,
489 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
535fn 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 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 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
612fn 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 }
651 }
652 }
653
654 Ok(items)
655}
656
657fn 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
670fn 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
693struct AnthropicPendingToolCall {
711 id: String,
713 name: String,
715 partial_json: String,
717}
718
719fn 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 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 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 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 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 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 }
922 }
923 _ => {}
924 }
925 }
926 }
927 Some("content_block_stop") => {
928 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 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 }
980 _ => {}
981 }
982 }
983
984 events
985}
986
987#[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
997fn 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#[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 message: Option<AnthropicMessageStart>,
1017}
1018
1019#[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 id: Option<String>,
1033 name: Option<String>,
1035 #[allow(dead_code)]
1037 thinking: Option<String>,
1038 #[allow(dead_code)]
1040 signature: Option<String>,
1041 #[allow(dead_code)]
1043 text: Option<String>,
1044 tool_use_id: Option<String>,
1046 #[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 #[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 #[serde(rename = "cache_read_input_tokens", alias = "cache_read", default)]
1077 cache_read: usize,
1078 #[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 #[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 #[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 #[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 #[test]
1197 fn parse_content_block_stop_finalizes_tool_call() {
1198 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 assert_eq!(events.len(), 3);
1210
1211 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 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 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 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 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 assert_eq!(tc.arguments["path"].as_str(), Some("test.rs"));
1263 }
1264
1265 #[test]
1266 fn parse_tool_call_in_done_message() {
1267 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 assert!(events.len() >= 6);
1291
1292 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 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 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 #[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 #[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 #[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 #[test]
1398 fn parse_usage_from_message_start() {
1399 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 assert_eq!(events.len(), 3);
1412 match &events[2] {
1413 ProviderEvent::Done { message, .. } => {
1414 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 #[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 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 #[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 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 #[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 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 #[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 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 assert_eq!(events1.len(), 3);
1584
1585 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 assert_eq!(events2.len(), 5);
1596
1597 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 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); assert_eq!(partial.content.len(), 1);
1647
1648 let events2 =
1649 parse_anthropic_events_stateful(chunk2, &mut partial, &mut usage, &mut pending_tc);
1650 assert!(events2.len() >= 2); 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 #[test]
1664 fn url_strips_trailing_v1_for_minimax() {
1665 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 assert_eq!(
1694 anthropic_messages_url("https://example.com/v1beta"),
1695 "https://example.com/v1beta/v1/messages"
1696 );
1697 }
1698}