1use futures::{Stream, StreamExt};
4use reqwest::Client;
5use serde::Deserialize;
6use serde_json::Value as JsonValue;
7use std::future::Future;
8use std::pin::Pin;
9use std::sync::Arc;
10
11use super::openai_responses_shared::parse_streaming_json;
12use super::shared_client;
13use super::sse::split_complete_lines;
14use crate::{
15 Api, AssistantMessage, ContentBlock, Context, Model, Provider, ProviderEvent, StopReason,
16 StreamOptions, StreamResult, TextContent, ThinkingContent, ToolCall, Usage,
17 error::ProviderError,
18};
19
20#[derive(Clone)]
25pub struct AnthropicProvider {
26 client: &'static Client,
27 api_key: Option<String>,
28 base_url: Option<String>,
30 extra_headers: Vec<(String, String)>,
33 native: bool,
38}
39
40impl AnthropicProvider {
41 pub fn new() -> Self {
46 Self {
47 client: shared_client(),
48 api_key: None,
49 base_url: None,
50 extra_headers: vec![
51 ("anthropic-version".to_string(), "2023-06-01".to_string()),
52 (
55 "anthropic-beta".to_string(),
56 "interleaved-thinking-2025-05-14,fine-grained-tool-streaming-2025-05-14"
57 .to_string(),
58 ),
59 ],
60 native: true,
61 }
62 }
63
64 #[allow(dead_code)]
66 pub fn with_api_key(api_key: impl Into<String>) -> Self {
67 Self {
68 client: shared_client(),
69 api_key: Some(api_key.into()),
70 base_url: None,
71 extra_headers: vec![
72 ("anthropic-version".to_string(), "2023-06-01".to_string()),
73 (
74 "anthropic-beta".to_string(),
75 "interleaved-thinking-2025-05-14,fine-grained-tool-streaming-2025-05-14"
76 .to_string(),
77 ),
78 ],
79 native: true,
80 }
81 }
82
83 #[allow(dead_code)]
89 pub fn with_base_url(base_url: &str) -> Self {
90 Self {
91 client: shared_client(),
92 api_key: None,
93 base_url: Some(base_url.to_string()),
94 extra_headers: vec![("anthropic-version".to_string(), "2023-06-01".to_string())],
95 native: false,
96 }
97 }
98
99 pub fn with_config(
110 base_url: &str,
111 api_key: Option<String>,
112 extra_headers: Vec<(String, String)>,
113 ) -> Self {
114 let mut headers = vec![("anthropic-version".to_string(), "2023-06-01".to_string())];
115 headers.extend(extra_headers);
116 Self {
117 client: shared_client(),
118 api_key,
119 base_url: Some(base_url.to_string()),
120 extra_headers: headers,
121 native: false,
122 }
123 }
124}
125
126impl Default for AnthropicProvider {
127 fn default() -> Self {
128 Self::new()
129 }
130}
131
132fn anthropic_messages_url(base_url: &str) -> String {
142 let trimmed = base_url.trim_end_matches('/');
143 let stripped = trimmed.strip_suffix("/v1").unwrap_or(trimmed);
144 format!("{}/v1/messages", stripped)
145}
146impl Provider for AnthropicProvider {
147 fn stream<'a>(
148 &'a self,
149 model: &'a Model,
150 context: &'a Context,
151 options: Option<StreamOptions>,
152 ) -> Pin<Box<dyn Future<Output = StreamResult> + Send + 'a>> {
153 Box::pin(async move {
154 let options = options.unwrap_or_default();
155
156 let effective_base_url = self.base_url.as_deref().unwrap_or(&model.base_url);
158 let url = anthropic_messages_url(effective_base_url);
159
160 let api_key = options
162 .api_key
163 .as_ref()
164 .or(self.api_key.as_ref())
165 .ok_or_else(|| ProviderError::MissingApiKey)?;
166
167 let normalized = crate::providers::openai::normalize_messages(
170 &context.messages,
171 &model.provider,
172 &model.id,
173 );
174 let messages =
175 build_anthropic_messages_from_normalized(&context.system_prompt, &normalized)?;
176
177 let mut body = serde_json::json!({
179 "model": model.id,
180 "messages": messages,
181 "stream": true,
182 });
183
184 if let Some(ref prompt) = context.system_prompt {
186 body["system"] = serde_json::json!(prompt);
187 }
188
189 let thinking_active = body.get("thinking").is_some();
194 if let Some(temp) = options.temperature
195 && !thinking_active
196 {
197 body["temperature"] = serde_json::json!(temp);
198 }
199
200 if let Some(max) = options.max_tokens {
201 body["max_tokens"] = serde_json::json!(max);
202 }
203
204 if !context.tools.is_empty() {
206 body["tools"] = build_anthropic_tools(&context.tools)?;
207 }
208
209 if model.reasoning && self.native {
222 let anthropic_opts = options
224 .provider_options
225 .as_ref()
226 .and_then(|po| po.anthropic.as_ref());
227
228 if let Some(opts) = anthropic_opts {
229 match opts.thinking_type.as_deref() {
231 Some("adaptive") => {
232 body["thinking"] = serde_json::json!({
234 "type": "adaptive",
235 });
236 if let Some(ref effort) = opts.effort {
237 let budget = match effort.as_str() {
240 "max" => model.max_tokens.min(31999),
241 "xhigh" => (model.max_tokens * 4 / 5).min(31999),
242 "high" => (model.max_tokens / 2).min(31999),
243 "medium" => (model.max_tokens / 4).min(16000),
244 "low" => (model.max_tokens / 8).min(8000),
245 _ => (model.max_tokens / 4).min(16000),
246 };
247 if body.get("max_tokens").is_none() {
248 body["max_tokens"] =
249 serde_json::json!((budget + 1024).min(model.max_tokens));
250 }
251 }
252 }
253 Some("enabled") => {
254 let budget = opts.thinking_budget.unwrap_or_else(|| {
256 compute_thinking_budget(&options.thinking_level, model.max_tokens)
257 });
258 if budget > 0 {
259 if body.get("max_tokens").is_none() {
260 body["max_tokens"] =
261 serde_json::json!((budget + 1024).min(model.max_tokens));
262 }
263 body["thinking"] = serde_json::json!({
264 "type": "enabled",
265 "budget_tokens": budget,
266 });
267 }
268 }
269 _ => {
270 let budget =
272 compute_thinking_budget(&options.thinking_level, model.max_tokens);
273 if budget > 0 {
274 if body.get("max_tokens").is_none() {
275 body["max_tokens"] =
276 serde_json::json!((budget + 1024).min(model.max_tokens));
277 }
278 body["thinking"] = serde_json::json!({
279 "type": "enabled",
280 "budget_tokens": budget,
281 });
282 }
283 }
284 }
285 } else if let Some(ref level) = options.thinking_level {
286 let budget = compute_thinking_budget(&Some(*level), model.max_tokens);
288 if budget > 0 {
289 if body.get("max_tokens").is_none() {
290 body["max_tokens"] =
291 serde_json::json!((budget + 1024).min(model.max_tokens));
292 }
293 body["thinking"] = serde_json::json!({
294 "type": "enabled",
295 "budget_tokens": budget,
296 });
297 }
298 }
299 }
300
301 if model.reasoning
312 && let Some(current) = body.get("max_tokens").and_then(|v| v.as_u64())
313 && current < 16_384
314 {
315 body["max_tokens"] = serde_json::json!(model.max_tokens.min(32_768));
316 }
317 if body.get("max_tokens").is_none() {
318 body["max_tokens"] = serde_json::json!(model.max_tokens.min(16384));
319 }
320
321 const ANTHROPIC_BREAKPOINT_CAP: usize = 4;
330 let want_cache = options.cache_retention == Some(crate::CacheRetention::Short)
331 || options.cache_retention == Some(crate::CacheRetention::Long);
332
333 if want_cache {
334 let mut remaining = ANTHROPIC_BREAKPOINT_CAP;
335 let cache_marker = serde_json::json!({ "type": "ephemeral" });
336
337 if remaining > 0
339 && let Some(system) = body.get_mut("system")
340 && system.is_string()
341 {
342 *system = serde_json::json!([{
343 "type": "text",
344 "text": system,
345 "cache_control": cache_marker.clone(),
346 }]);
347 remaining -= 1;
348 }
349
350 if remaining > 0
352 && let Some(messages) = body.get_mut("messages").and_then(|m| m.as_array_mut())
353 {
354 if let Some(last_msg) = messages.last_mut()
355 && let Some(content) = last_msg.get_mut("content")
356 {
357 if let Some(parts) = content.as_array_mut() {
358 if let Some(last_part) = parts.last_mut() {
359 last_part["cache_control"] = cache_marker.clone();
360 remaining -= 1;
361 }
362 } else if content.is_string() {
363 let text = content.take();
364 *content = serde_json::json!([{
365 "type": "text",
366 "text": text,
367 "cache_control": cache_marker.clone(),
368 }]);
369 remaining -= 1;
370 }
371 }
372
373 if remaining > 0 {
375 let msg_count = messages.len();
376 if msg_count >= 3
377 && let Some(msg) = messages.get_mut(msg_count - 3)
378 && let Some(content) = msg.get_mut("content")
379 && let Some(parts) = content.as_array_mut()
380 && let Some(last_part) = parts.last_mut()
381 {
382 last_part["cache_control"] = cache_marker;
383 remaining -= 1;
384 }
385 }
386 }
387
388 if remaining < ANTHROPIC_BREAKPOINT_CAP {
389 tracing::debug!(
390 used = ANTHROPIC_BREAKPOINT_CAP - remaining,
391 cap = ANTHROPIC_BREAKPOINT_CAP,
392 "Anthropic cache breakpoints applied"
393 );
394 }
395 }
396
397 let mut headers = reqwest::header::HeaderMap::new();
399 headers.insert(
400 "x-api-key",
401 api_key.parse().map_err(|e| {
402 ProviderError::InvalidResponse(format!("invalid header value: {e}"))
403 })?,
404 );
405 headers.insert(
406 "content-type",
407 "application/json".parse().map_err(|e| {
408 ProviderError::InvalidResponse(format!("invalid header value: {e}"))
409 })?,
410 );
411
412 for (k, v) in &self.extra_headers {
414 if let (Ok(name), Ok(value)) = (
415 k.parse::<reqwest::header::HeaderName>(),
416 v.parse::<reqwest::header::HeaderValue>(),
417 ) {
418 headers.insert(name, value);
419 }
420 }
421
422 for (k, v) in &options.headers {
424 if let (Ok(name), Ok(value)) = (
425 k.parse::<reqwest::header::HeaderName>(),
426 v.parse::<reqwest::header::HeaderValue>(),
427 ) {
428 headers.insert(name, value);
429 }
430 }
431
432 let response = self
434 .client
435 .post(&url)
436 .headers(headers)
437 .json(&body)
438 .send()
439 .await
440 .map_err(ProviderError::RequestFailed)?;
441
442 if !response.status().is_success() {
443 let status = response.status();
444 let request_id = response
447 .headers()
448 .get("request-id")
449 .and_then(|v| v.to_str().ok())
450 .map(str::to_string);
451 let body: String = response.text().await.unwrap_or_default();
452 return Err(ProviderError::HttpError(
453 crate::error::HttpErrorDetail::new(status.as_u16(), body)
454 .with_provider("anthropic")
455 .with_request_id(request_id),
456 ));
457 }
458
459 let model_name = model.id.clone();
461
462 struct AnthropicScanState {
480 pending_bytes: Vec<u8>,
481 partial: AssistantMessage,
482 usage: Usage,
483 pending_tool_calls: std::collections::HashMap<usize, AnthropicPendingToolCall>,
485 }
486
487 let initial_state = AnthropicScanState {
488 pending_bytes: Vec::new(),
489 partial: AssistantMessage::new(Api::AnthropicMessages, "anthropic", &model_name),
490 usage: Usage::default(),
491 pending_tool_calls: std::collections::HashMap::new(),
492 };
493
494 let stream = response
495 .bytes_stream()
496 .scan(
497 initial_state,
498 move |state, chunk: Result<bytes::Bytes, reqwest::Error>| {
499 let events = match chunk {
500 Ok(bytes) => {
501 let mut combined =
502 Vec::with_capacity(state.pending_bytes.len() + bytes.len());
503 combined.extend_from_slice(&state.pending_bytes);
504 combined.extend_from_slice(&bytes);
505 let (text, trailing) = split_complete_lines(&combined);
506 state.pending_bytes = trailing;
507 parse_anthropic_events_stateful(
508 &text,
509 &mut state.partial,
510 &mut state.usage,
511 &mut state.pending_tool_calls,
512 )
513 }
514 Err(e) => vec![ProviderEvent::Error {
515 reason: StopReason::Error,
516 error: create_error_message(&e.to_string()),
517 }],
518 };
519 async move { Some(futures::stream::iter(events)) }
520 },
521 )
522 .flatten();
523
524 Ok(Box::pin(stream) as Pin<Box<dyn Stream<Item = ProviderEvent> + Send>>)
525 })
526 }
527}
528
529fn build_anthropic_messages_from_normalized(
536 _system_prompt: &Option<String>,
537 messages_in: &[crate::Message],
538) -> Result<Vec<JsonValue>, ProviderError> {
539 let mut messages: Vec<JsonValue> = Vec::new();
540 let mut i = 0;
541
542 while i < messages_in.len() {
543 let msg = &messages_in[i];
544 match msg {
545 crate::Message::User(u) => {
546 let content = match &u.content {
547 crate::MessageContent::Text(s) => vec![serde_json::json!({
548 "type": "text",
549 "text": s,
550 })],
551 crate::MessageContent::Blocks(blocks) => blocks_to_anthropic_content(blocks)?,
552 };
553 messages.push(serde_json::json!({
554 "role": "user",
555 "content": content,
556 }));
557 i += 1;
558 }
559 crate::Message::Assistant(a) => {
560 let content = blocks_to_anthropic_content(&a.content)?;
561 messages.push(serde_json::json!({
562 "role": "assistant",
563 "content": content,
564 }));
565 i += 1;
566 }
567 crate::Message::ToolResult(t) => {
568 let mut tool_results: Vec<JsonValue> = Vec::new();
571 let content = blocks_to_anthropic_content(&t.content)?;
572 tool_results.push(serde_json::json!({
573 "type": "tool_result",
574 "tool_use_id": t.tool_call_id,
575 "content": content,
576 }));
577
578 let mut j = i + 1;
580 while j < messages_in.len() {
581 if let crate::Message::ToolResult(next_t) = &messages_in[j] {
582 let next_content = blocks_to_anthropic_content(&next_t.content)?;
583 tool_results.push(serde_json::json!({
584 "type": "tool_result",
585 "tool_use_id": next_t.tool_call_id,
586 "content": next_content,
587 }));
588 j += 1;
589 } else {
590 break;
591 }
592 }
593
594 messages.push(serde_json::json!({
595 "role": "user",
596 "content": tool_results,
597 }));
598 i = j;
599 }
600 }
601 }
602
603 Ok(messages)
604}
605
606fn blocks_to_anthropic_content(blocks: &[ContentBlock]) -> Result<Vec<JsonValue>, ProviderError> {
608 let mut items = Vec::new();
609
610 for block in blocks {
611 match block {
612 ContentBlock::Text(t) => {
613 items.push(serde_json::json!({
614 "type": "text",
615 "text": t.text,
616 }));
617 }
618 ContentBlock::ToolCall(tc) => {
619 items.push(serde_json::json!({
620 "type": "tool_use",
621 "id": tc.id,
622 "name": tc.name,
623 "input": tc.arguments,
624 }));
625 }
626 ContentBlock::Thinking(th) => {
627 items.push(serde_json::json!({
628 "type": "thinking",
629 "thinking": th.thinking,
630 }));
631 }
632 ContentBlock::Image(img) => {
633 items.push(serde_json::json!({
634 "type": "image",
635 "source": {
636 "type": "base64",
637 "media_type": img.mime_type,
638 "data": img.data,
639 },
640 }));
641 }
642 ContentBlock::Unknown(_) => {
643 }
645 }
646 }
647
648 Ok(items)
649}
650
651fn compute_thinking_budget(level: &Option<crate::ThinkingLevel>, max_tokens: usize) -> usize {
653 match level {
654 Some(crate::ThinkingLevel::High) | Some(crate::ThinkingLevel::XHigh) => {
655 (max_tokens / 2).min(31999)
656 }
657 Some(crate::ThinkingLevel::Medium) => (max_tokens / 4).min(16000),
658 Some(crate::ThinkingLevel::Low) => (max_tokens / 8).min(8000),
659 Some(crate::ThinkingLevel::Minimal) => (max_tokens / 16).min(4000),
660 _ => 0,
661 }
662}
663
664fn build_anthropic_tools(tools: &[crate::Tool]) -> Result<JsonValue, ProviderError> {
665 let items: Vec<_> = tools
666 .iter()
667 .map(|tool| {
668 serde_json::json!({
669 "name": tool.name,
670 "description": tool.description,
671 "input_schema": tool.parameters,
672 })
673 })
674 .collect();
675
676 Ok(serde_json::json!(items))
677}
678
679struct AnthropicPendingToolCall {
697 id: String,
699 name: String,
701 partial_json: String,
703}
704
705fn parse_anthropic_events_stateful(
718 text: &str,
719 partial_message: &mut AssistantMessage,
720 accumulated_usage: &mut Usage,
721 pending_tool_calls: &mut std::collections::HashMap<usize, AnthropicPendingToolCall>,
722) -> Vec<ProviderEvent> {
723 let mut events = Vec::with_capacity(text.len() / 80);
727
728 for line in text.split('\n') {
729 let line = line.trim_end_matches('\r');
730 if line.is_empty() {
731 continue;
732 }
733
734 if !line.starts_with("data: ") {
735 continue;
736 }
737
738 let data = &line[6..];
739
740 if data == "[DONE]" || data.is_empty() {
741 continue;
742 }
743
744 let event = match serde_json::from_str::<AnthropicEvent>(data) {
745 Ok(e) => e,
746 Err(_) => continue,
747 };
748
749 let event_type = event.type_.as_deref();
750
751 if let Some(usage) = &event.usage {
755 accumulated_usage.input = usage.input_tokens.max(accumulated_usage.input);
756 accumulated_usage.output = usage.output_tokens.max(accumulated_usage.output);
757 accumulated_usage.cache_read = usage.cache_read.max(accumulated_usage.cache_read);
758 accumulated_usage.cache_write = usage.cache_creation.max(accumulated_usage.cache_write);
759 accumulated_usage.total_tokens = accumulated_usage.input + accumulated_usage.output;
760 } else if let Some(msg) = &event.message
761 && let Some(usage) = &msg.usage
762 {
763 accumulated_usage.input = usage.input_tokens.max(accumulated_usage.input);
764 accumulated_usage.output = usage.output_tokens.max(accumulated_usage.output);
765 accumulated_usage.cache_read = usage.cache_read.max(accumulated_usage.cache_read);
766 accumulated_usage.cache_write = usage.cache_creation.max(accumulated_usage.cache_write);
767 accumulated_usage.total_tokens = accumulated_usage.input + accumulated_usage.output;
768 }
769
770 match event_type {
771 Some("message_start") => {
772 events.push(ProviderEvent::Start {
773 partial: Arc::new(partial_message.clone()),
774 });
775 }
776 Some("content_block_start") => {
777 if let Some(block) = &event.content_block {
778 let idx = block.index.or(event.index).unwrap_or(0);
779 match block.type_.as_deref() {
780 Some("text") => {
781 events.push(ProviderEvent::TextStart {
782 content_index: idx,
783 partial: Arc::new(partial_message.clone()),
784 });
785 }
786 Some("thinking") => {
787 events.push(ProviderEvent::ThinkingStart {
788 content_index: idx,
789 partial: Arc::new(partial_message.clone()),
790 });
791 }
792 Some("tool_use") | Some("server_tool_use") => {
793 let tc_id = block.id.clone().unwrap_or_default();
797 let tc_name = block.name.clone().unwrap_or_default();
798 pending_tool_calls.insert(
799 idx,
800 AnthropicPendingToolCall {
801 id: tc_id.clone(),
802 name: tc_name.clone(),
803 partial_json: String::new(),
804 },
805 );
806 events.push(ProviderEvent::ToolCallStart {
807 content_index: idx,
808 tool_call_id: Some(tc_id),
809 tool_name: Some(tc_name),
810 partial: Arc::new(partial_message.clone()),
811 });
812 }
813 Some(t) if t.ends_with("_tool_result") => {
814 let name = match t {
815 "web_search_tool_result" => Some("web_search".to_string()),
816 "code_execution_tool_result" => Some("code_execution".to_string()),
817 "web_fetch_tool_result" => Some("web_fetch".to_string()),
818 _ => None,
819 };
820 if let Some(tool_name) = name {
821 let tc = ToolCall::new(
822 block.tool_use_id.clone().unwrap_or_default(),
823 tool_name,
824 serde_json::json!({}),
825 );
826 partial_message
827 .content
828 .push(ContentBlock::ToolCall(tc.clone()));
829 events.push(ProviderEvent::ToolCallEnd {
830 content_index: idx,
831 tool_call: tc,
832 partial: Arc::new(partial_message.clone()),
833 });
834 }
835 }
836 _ => {}
837 }
838 }
839 }
840 Some("content_block_delta") => {
841 if let Some(delta) = &event.delta {
842 match delta.type_.as_deref() {
843 Some("text_delta") => {
844 if let Some(text) = &delta.text {
845 let last_text_idx = partial_message
848 .content
849 .iter()
850 .rposition(|b| matches!(b, ContentBlock::Text(_)));
851 if let Some(idx) = last_text_idx
852 && let ContentBlock::Text(t) = &mut partial_message.content[idx]
853 {
854 t.text.push_str(text);
855 } else {
856 partial_message
857 .content
858 .push(ContentBlock::Text(TextContent::new(text.clone())));
859 }
860 events.push(ProviderEvent::TextDelta {
861 content_index: event.index.unwrap_or(0),
862 delta: text.clone(),
863 partial: Arc::new(partial_message.clone()),
864 });
865 }
866 }
867 Some("thinking_delta") => {
868 if let Some(text) = &delta.thinking {
869 let last_think_idx = partial_message
870 .content
871 .iter()
872 .rposition(|b| matches!(b, ContentBlock::Thinking(_)));
873 if let Some(idx) = last_think_idx
874 && let ContentBlock::Thinking(t) =
875 &mut partial_message.content[idx]
876 {
877 t.thinking.push_str(text);
878 } else {
879 partial_message.content.push(ContentBlock::Thinking(
880 ThinkingContent::new(text.clone()),
881 ));
882 }
883 events.push(ProviderEvent::ThinkingDelta {
884 content_index: event.index.unwrap_or(0),
885 delta: text.clone(),
886 partial: Arc::new(partial_message.clone()),
887 });
888 }
889 }
890 Some("input_json_delta") => {
891 if let Some(args) = &delta.partial_json {
892 let block_idx = event.index.unwrap_or(0);
894 if let Some(ptc) = pending_tool_calls.get_mut(&block_idx) {
895 ptc.partial_json.push_str(args);
896 }
897 events.push(ProviderEvent::ToolCallDelta {
898 content_index: block_idx,
899 delta: args.clone(),
900 partial: Arc::new(partial_message.clone()),
901 });
902 }
903 }
904 Some("signature_delta") => {
905 if let Some(_sig) = &delta.signature {
906 }
908 }
909 _ => {}
910 }
911 }
912 }
913 Some("content_block_stop") => {
914 let block_idx = event.index.unwrap_or(0);
918 if let Some(ptc) = pending_tool_calls.remove(&block_idx) {
919 let args_value = parse_streaming_json(&ptc.partial_json);
920 let tc = ToolCall::new(ptc.id, ptc.name, args_value);
921
922 partial_message
926 .content
927 .push(ContentBlock::ToolCall(tc.clone()));
928
929 tracing::debug!(
930 block_idx,
931 tool_id = %tc.id,
932 tool_name = %tc.name,
933 "content_block_stop: finalized tool call"
934 );
935
936 events.push(ProviderEvent::ToolCallEnd {
937 content_index: block_idx,
938 tool_call: tc,
939 partial: Arc::new(partial_message.clone()),
940 });
941 }
942 }
943 Some("message_delta") => {
944 if let Some(delta) = &event.delta {
945 let reason = match delta.stop_reason.as_deref() {
946 Some("end_turn") | Some("stop_sequence") | Some("pause_turn") => {
947 StopReason::Stop
948 }
949 Some("max_tokens") => StopReason::Length,
950 Some("tool_use") => StopReason::ToolUse,
951 Some("refusal") | Some("sensitive") => StopReason::Error,
952 _ => StopReason::Stop,
953 };
954
955 let mut done_msg = partial_message.clone();
956 done_msg.usage = accumulated_usage.clone();
957 events.push(ProviderEvent::Done {
958 reason,
959 message: done_msg,
960 });
961 }
962 }
963 Some("message_stop") => {
964 }
966 _ => {}
967 }
968 }
969
970 events
971}
972
973#[cfg(test)]
976fn parse_anthropic_events(text: &str, model_id: &str) -> Vec<ProviderEvent> {
977 let mut partial = AssistantMessage::new(Api::AnthropicMessages, "anthropic", model_id);
978 let mut usage = Usage::default();
979 let mut pending_tool_calls = std::collections::HashMap::new();
980 parse_anthropic_events_stateful(text, &mut partial, &mut usage, &mut pending_tool_calls)
981}
982
983fn create_error_message(msg: &str) -> AssistantMessage {
985 let mut message = AssistantMessage::new(Api::AnthropicMessages, "anthropic", "unknown");
986 message.stop_reason = StopReason::Error;
987 message.error_message = Some(msg.to_string());
988 message
989}
990
991#[derive(Debug, Deserialize)]
993struct AnthropicEvent {
994 #[serde(rename = "type")]
995 type_: Option<String>,
996 #[serde(rename = "index")]
997 index: Option<usize>,
998 content_block: Option<ContentBlockStart>,
999 delta: Option<Delta>,
1000 usage: Option<AnthropicUsage>,
1001 message: Option<AnthropicMessageStart>,
1003}
1004
1005#[derive(Debug, Deserialize)]
1008struct AnthropicMessageStart {
1009 usage: Option<AnthropicUsage>,
1010}
1011
1012#[derive(Debug, Deserialize)]
1013struct ContentBlockStart {
1014 #[serde(rename = "type")]
1015 type_: Option<String>,
1016 index: Option<usize>,
1017 id: Option<String>,
1019 name: Option<String>,
1021 #[allow(dead_code)]
1023 thinking: Option<String>,
1024 #[allow(dead_code)]
1026 signature: Option<String>,
1027 #[allow(dead_code)]
1029 text: Option<String>,
1030 tool_use_id: Option<String>,
1032 #[allow(dead_code)]
1034 content: Option<JsonValue>,
1035}
1036
1037#[derive(Debug, Deserialize)]
1038struct Delta {
1039 #[serde(rename = "type")]
1040 type_: Option<String>,
1041 text: Option<String>,
1042 thinking: Option<String>,
1043 partial_json: Option<String>,
1044 #[allow(dead_code)]
1047 signature: Option<String>,
1048 #[serde(rename = "stop_reason")]
1049 stop_reason: Option<String>,
1050 #[serde(rename = "stop_sequence")]
1051 #[allow(dead_code)]
1052 stop_sequence: Option<String>,
1053}
1054
1055#[derive(Debug, Deserialize)]
1056struct AnthropicUsage {
1057 #[serde(rename = "input_tokens", default)]
1058 input_tokens: usize,
1059 #[serde(rename = "output_tokens", default)]
1060 output_tokens: usize,
1061 #[serde(rename = "cache_read_input_tokens", alias = "cache_read", default)]
1063 cache_read: usize,
1064 #[serde(
1066 rename = "cache_creation_input_tokens",
1067 alias = "cache_creation",
1068 default
1069 )]
1070 cache_creation: usize,
1071}
1072
1073#[cfg(test)]
1074mod tests {
1075 use super::*;
1076
1077 const MODEL: &str = "claude-3-5-sonnet-20241022";
1078
1079 #[test]
1082 fn parse_message_start() {
1083 let sse = "data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\"}}\n";
1084 let events = parse_anthropic_events(sse, MODEL);
1085 assert_eq!(events.len(), 1);
1086 assert!(matches!(&events[0], ProviderEvent::Start { .. }));
1087 }
1088
1089 #[test]
1092 fn parse_text_block_start() {
1093 let sse = "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n";
1094 let events = parse_anthropic_events(sse, MODEL);
1095 assert_eq!(events.len(), 1);
1096 match &events[0] {
1097 ProviderEvent::TextStart { content_index, .. } => assert_eq!(*content_index, 0),
1098 other => panic!("expected TextStart, got {other:?}"),
1099 }
1100 }
1101
1102 #[test]
1103 fn parse_thinking_block_start() {
1104 let sse = "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\",\"thinking\":\"\"}}\n";
1105 let events = parse_anthropic_events(sse, MODEL);
1106 assert_eq!(events.len(), 1);
1107 assert!(matches!(&events[0], ProviderEvent::ThinkingStart { .. }));
1108 }
1109
1110 #[test]
1111 fn parse_tool_use_block_start() {
1112 let sse = "data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tool_1\",\"name\":\"search\"}}\n";
1113 let events = parse_anthropic_events(sse, MODEL);
1114 assert_eq!(events.len(), 1);
1115 match &events[0] {
1116 ProviderEvent::ToolCallStart { content_index, .. } => assert_eq!(*content_index, 1),
1117 other => panic!("expected ToolCallStart, got {other:?}"),
1118 }
1119 }
1120
1121 #[test]
1124 fn parse_text_delta() {
1125 let sse = "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"Hello\"}}\n";
1126 let events = parse_anthropic_events(sse, MODEL);
1127 assert_eq!(events.len(), 1);
1128 match &events[0] {
1129 ProviderEvent::TextDelta {
1130 delta,
1131 content_index,
1132 ..
1133 } => {
1134 assert_eq!(delta, "Hello");
1135 assert_eq!(*content_index, 0);
1136 }
1137 other => panic!("expected TextDelta, got {other:?}"),
1138 }
1139 }
1140
1141 #[test]
1142 fn parse_thinking_delta() {
1143 let sse = "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"Let me reason...\"}}\n";
1144 let events = parse_anthropic_events(sse, MODEL);
1145 assert_eq!(events.len(), 1);
1146 match &events[0] {
1147 ProviderEvent::ThinkingDelta { delta, .. } => assert_eq!(delta, "Let me reason..."),
1148 other => panic!("expected ThinkingDelta, got {other:?}"),
1149 }
1150 }
1151
1152 #[test]
1153 fn parse_input_json_delta() {
1154 let sse = "data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"city\\\":\\\"SF\\\"}\"}}\n";
1155 let events = parse_anthropic_events(sse, MODEL);
1156 assert_eq!(events.len(), 1);
1157 match &events[0] {
1158 ProviderEvent::ToolCallDelta {
1159 delta,
1160 content_index,
1161 ..
1162 } => {
1163 assert_eq!(delta, "{\"city\":\"SF\"}");
1164 assert_eq!(*content_index, 1);
1165 }
1166 other => panic!("expected ToolCallDelta, got {other:?}"),
1167 }
1168 }
1169
1170 #[test]
1173 fn parse_content_block_stop_finalizes_tool_call() {
1174 let sse = concat!(
1176 "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tool_abc\",\"name\":\"bash\"}}\n",
1177 "\n",
1178 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"command\\\":\\\"ls\\\"}\"}}\n",
1179 "\n",
1180 "data: {\"type\":\"content_block_stop\",\"index\":0}\n",
1181 "\n"
1182 );
1183 let events = parse_anthropic_events(sse, MODEL);
1184 assert_eq!(events.len(), 3);
1186
1187 let tc_end = events.iter().find_map(|e| match e {
1189 ProviderEvent::ToolCallEnd { tool_call, .. } => Some(tool_call.clone()),
1190 _ => None,
1191 });
1192 let tc = tc_end.expect("Should have ToolCallEnd");
1193 assert_eq!(tc.id, "tool_abc");
1194 assert_eq!(tc.name, "bash");
1195 assert_eq!(tc.arguments, serde_json::json!({"command": "ls"}));
1196 }
1197
1198 #[test]
1199 fn parse_content_block_stop_ignores_non_tool() {
1200 let sse = concat!(
1202 "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n",
1203 "\n",
1204 "data: {\"type\":\"content_block_stop\",\"index\":0}\n",
1205 "\n"
1206 );
1207 let events = parse_anthropic_events(sse, MODEL);
1208 assert_eq!(events.len(), 1);
1210 assert!(matches!(&events[0], ProviderEvent::TextStart { .. }));
1211 }
1212
1213 #[test]
1214 fn parse_tool_call_accumulates_across_deltas() {
1215 let sse = concat!(
1217 "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tool_1\",\"name\":\"edit\"}}\n",
1218 "\n",
1219 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"path\\\":\\\"tes\"}}\n",
1220 "\n",
1221 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"t.rs\"}}\n",
1222 "\n",
1223 "data: {\"type\":\"content_block_stop\",\"index\":0}\n",
1224 "\n"
1225 );
1226 let events = parse_anthropic_events(sse, MODEL);
1227 assert_eq!(events.len(), 4);
1229
1230 let tc_end = events.iter().find_map(|e| match e {
1231 ProviderEvent::ToolCallEnd { tool_call, .. } => Some(tool_call.clone()),
1232 _ => None,
1233 });
1234 let tc = tc_end.expect("Should have ToolCallEnd");
1235 assert_eq!(tc.id, "tool_1");
1236 assert_eq!(tc.name, "edit");
1237 assert_eq!(tc.arguments["path"].as_str(), Some("test.rs"));
1239 }
1240
1241 #[test]
1242 fn parse_tool_call_in_done_message() {
1243 let sse = concat!(
1245 "data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\"}}\n",
1246 "\n",
1247 "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\",\"thinking\":\"\"}}\n",
1248 "\n",
1249 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"I need to search.\"}}\n",
1250 "\n",
1251 "data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tool_search\",\"name\":\"web_search\"}}\n",
1252 "\n",
1253 "data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"query\\\":\\\"rust async\\\"}\"}}\n",
1254 "\n",
1255 "data: {\"type\":\"content_block_stop\",\"index\":0}\n",
1256 "\n",
1257 "data: {\"type\":\"content_block_stop\",\"index\":1}\n",
1258 "\n",
1259 "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"tool_use\"}}\n"
1260 );
1261 let events = parse_anthropic_events(sse, MODEL);
1262
1263 assert!(events.len() >= 6);
1267
1268 let tc_end = events.iter().find_map(|e| match e {
1270 ProviderEvent::ToolCallEnd { tool_call, .. } => Some(tool_call.clone()),
1271 _ => None,
1272 });
1273 let tc = tc_end.expect("Should have ToolCallEnd");
1274 assert_eq!(tc.id, "tool_search");
1275 assert_eq!(tc.name, "web_search");
1276 assert_eq!(tc.arguments, serde_json::json!({"query": "rust async"}));
1277
1278 let done = events.iter().find_map(|e| match e {
1280 ProviderEvent::Done { reason, .. } => Some(*reason),
1281 _ => None,
1282 });
1283 assert_eq!(done, Some(StopReason::ToolUse));
1284
1285 let done_msg = events.iter().find_map(|e| match e {
1287 ProviderEvent::Done { message, .. } => Some(message.clone()),
1288 _ => None,
1289 });
1290 let msg = done_msg.expect("Should have Done event");
1291 let tool_calls: Vec<_> = msg
1292 .content
1293 .iter()
1294 .filter(|b| matches!(b, ContentBlock::ToolCall(_)))
1295 .collect();
1296 assert_eq!(
1297 tool_calls.len(),
1298 1,
1299 "Done message should contain exactly 1 tool call"
1300 );
1301 }
1302
1303 #[test]
1306 fn parse_message_delta_end_turn() {
1307 let sse = "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"}}\n";
1308 let events = parse_anthropic_events(sse, MODEL);
1309 assert_eq!(events.len(), 1);
1310 match &events[0] {
1311 ProviderEvent::Done { reason, .. } => assert!(matches!(reason, StopReason::Stop)),
1312 other => panic!("expected Done, got {other:?}"),
1313 }
1314 }
1315
1316 #[test]
1317 fn parse_message_delta_max_tokens() {
1318 let sse = "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"max_tokens\"}}\n";
1319 let events = parse_anthropic_events(sse, MODEL);
1320 match &events[0] {
1321 ProviderEvent::Done { reason, .. } => assert!(matches!(reason, StopReason::Length)),
1322 other => panic!("expected Done with Length, got {other:?}"),
1323 }
1324 }
1325
1326 #[test]
1327 fn parse_message_delta_stop_sequence() {
1328 let sse =
1329 "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"stop_sequence\"}}\n";
1330 let events = parse_anthropic_events(sse, MODEL);
1331 match &events[0] {
1332 ProviderEvent::Done { reason, .. } => assert!(matches!(reason, StopReason::Stop)),
1333 other => panic!("expected Done with Stop, got {other:?}"),
1334 }
1335 }
1336
1337 #[test]
1340 fn parse_message_stop_no_event_emitted() {
1341 let sse = "data: {\"type\":\"message_stop\"}\n";
1342 let events = parse_anthropic_events(sse, MODEL);
1343 assert!(events.is_empty());
1344 }
1345
1346 #[test]
1349 fn parse_thinking_block_flow() {
1350 let sse = concat!(
1351 "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\",\"thinking\":\"\"}}\n",
1352 "\n",
1353 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"I should\"}}\n",
1354 "\n",
1355 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\" check this.\"}}\n",
1356 "\n"
1357 );
1358 let events = parse_anthropic_events(sse, MODEL);
1359 assert_eq!(events.len(), 3);
1360 assert!(matches!(&events[0], ProviderEvent::ThinkingStart { .. }));
1361 let thinking: Vec<&str> = events[1..]
1362 .iter()
1363 .filter_map(|e| match e {
1364 ProviderEvent::ThinkingDelta { delta, .. } => Some(delta.as_str()),
1365 _ => None,
1366 })
1367 .collect();
1368 assert_eq!(thinking, vec!["I should", " check this."]);
1369 }
1370
1371 #[test]
1374 fn parse_usage_from_message_start() {
1375 let sse = concat!(
1379 "data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\"},\"usage\":{\"input_tokens\":100,\"output_tokens\":0,\"cache_read\":80,\"cache_creation\":20}}\n",
1380 "\n",
1381 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"hi\"}}\n",
1382 "\n",
1383 "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"}}\n"
1384 );
1385 let events = parse_anthropic_events(sse, MODEL);
1386 assert_eq!(events.len(), 3);
1388 match &events[2] {
1389 ProviderEvent::Done { message, .. } => {
1390 assert_eq!(message.usage.input, 100);
1392 assert_eq!(message.usage.output, 0);
1393 assert_eq!(message.usage.total_tokens, 100);
1394 assert_eq!(message.usage.cache_read, 80);
1395 assert_eq!(message.usage.cache_write, 20);
1396 }
1397 other => panic!("expected Done, got {other:?}"),
1398 }
1399 }
1400
1401 #[test]
1404 fn parse_cache_metrics() {
1405 let sse = concat!(
1406 "data: {\"type\":\"message_start\",\"usage\":{\"input_tokens\":50,\"output_tokens\":0,\"cache_read\":40,\"cache_creation\":10}}\n",
1407 "\n",
1408 "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"},\"usage\":{\"input_tokens\":50,\"output_tokens\":20,\"cache_read\":40,\"cache_creation\":10}}\n"
1409 );
1410 let events = parse_anthropic_events(sse, MODEL);
1411 assert_eq!(events.len(), 2);
1413 match &events[1] {
1414 ProviderEvent::Done { message, .. } => {
1415 assert_eq!(message.usage.cache_read, 40);
1416 assert_eq!(message.usage.cache_write, 10);
1417 }
1418 other => panic!("expected Done, got {other:?}"),
1419 }
1420 }
1421
1422 #[test]
1425 fn parse_empty_input() {
1426 let events = parse_anthropic_events("", MODEL);
1427 assert!(events.is_empty());
1428 }
1429
1430 #[test]
1431 fn parse_done_marker_is_ignored() {
1432 let sse = "data: [DONE]\n";
1434 let events = parse_anthropic_events(sse, MODEL);
1435 assert!(events.is_empty());
1436 }
1437
1438 #[test]
1439 fn parse_malformed_json_is_skipped() {
1440 let sse = "data: {broken\ndata: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"ok\"}}\n";
1441 let events = parse_anthropic_events(sse, MODEL);
1442 assert_eq!(events.len(), 1);
1443 match &events[0] {
1444 ProviderEvent::TextDelta { delta, .. } => assert_eq!(delta, "ok"),
1445 other => panic!("expected TextDelta, got {other:?}"),
1446 }
1447 }
1448
1449 #[test]
1450 fn parse_non_data_lines_ignored() {
1451 let sse = "event: ping\nid: 42\ndata: {\"type\":\"message_start\"}\n";
1452 let events = parse_anthropic_events(sse, MODEL);
1453 assert_eq!(events.len(), 1);
1454 }
1455
1456 #[test]
1457 fn parse_empty_data_line_skipped() {
1458 let sse = "data: \ndata: {\"type\":\"message_start\"}\n";
1459 let events = parse_anthropic_events(sse, MODEL);
1460 assert_eq!(events.len(), 1);
1461 }
1462
1463 #[test]
1464 fn parse_unknown_event_type_ignored() {
1465 let sse = "data: {\"type\":\"ping\"}\ndata: {\"type\":\"message_start\"}\n";
1466 let events = parse_anthropic_events(sse, MODEL);
1467 assert_eq!(events.len(), 1);
1468 }
1469
1470 #[test]
1471 fn parse_carriage_return_line_endings() {
1472 let sse = "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"CR\"}}\r\n\r\n";
1473 let events = parse_anthropic_events(sse, MODEL);
1474 assert_eq!(events.len(), 1);
1475 match &events[0] {
1476 ProviderEvent::TextDelta { delta, .. } => assert_eq!(delta, "CR"),
1477 other => panic!("expected TextDelta, got {other:?}"),
1478 }
1479 }
1480
1481 #[test]
1484 fn parse_full_anthropic_stream() {
1485 let sse = concat!(
1486 "data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\"}}\n",
1487 "\n",
1488 "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n",
1489 "\n",
1490 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"Hello\"}}\n",
1491 "\n",
1492 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\" world\"}}\n",
1493 "\n",
1494 "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"}}\n",
1495 "\n",
1496 "data: {\"type\":\"message_stop\"}\n"
1497 );
1498 let events = parse_anthropic_events(sse, MODEL);
1499 assert_eq!(events.len(), 5);
1501
1502 assert!(matches!(&events[0], ProviderEvent::Start { .. }));
1503 assert!(matches!(&events[1], ProviderEvent::TextStart { .. }));
1504
1505 let texts: Vec<&str> = events[2..4]
1506 .iter()
1507 .filter_map(|e| match e {
1508 ProviderEvent::TextDelta { delta, .. } => Some(delta.as_str()),
1509 _ => None,
1510 })
1511 .collect();
1512 assert_eq!(texts, vec!["Hello", " world"]);
1513
1514 assert!(matches!(
1515 &events[4],
1516 ProviderEvent::Done {
1517 reason: StopReason::Stop,
1518 ..
1519 }
1520 ));
1521 }
1522
1523 #[test]
1528 fn parse_stateful_across_two_chunks() {
1529 let chunk1 = concat!(
1530 "data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\"}}\n",
1531 "\n",
1532 "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\",\"thinking\":\"\"}}\n",
1533 "\n",
1534 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"Let me\"}}\n",
1535 "\n"
1536 );
1537 let chunk2 = concat!(
1538 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\" think.\"}}\n",
1539 "\n",
1540 "data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n",
1541 "\n",
1542 "data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"text_delta\",\"text\":\"Hello\"}}\n",
1543 "\n",
1544 "data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"text_delta\",\"text\":\" world\"}}\n",
1545 "\n",
1546 "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"}}\n",
1547 "\n"
1548 );
1549
1550 let mut partial = AssistantMessage::new(Api::AnthropicMessages, "anthropic", MODEL);
1552 let mut usage = Usage::default();
1553 let mut pending_tc: std::collections::HashMap<usize, AnthropicPendingToolCall> =
1554 std::collections::HashMap::new();
1555
1556 let events1 =
1557 parse_anthropic_events_stateful(chunk1, &mut partial, &mut usage, &mut pending_tc);
1558 assert_eq!(events1.len(), 3);
1560
1561 assert_eq!(partial.content.len(), 1);
1563 match &partial.content[0] {
1564 ContentBlock::Thinking(t) => assert_eq!(t.thinking, "Let me"),
1565 other => panic!("Expected Thinking block, got {:?}", other),
1566 }
1567
1568 let events2 =
1569 parse_anthropic_events_stateful(chunk2, &mut partial, &mut usage, &mut pending_tc);
1570 assert_eq!(events2.len(), 5);
1572
1573 assert_eq!(partial.content.len(), 2);
1575 match &partial.content[0] {
1576 ContentBlock::Thinking(t) => assert_eq!(t.thinking, "Let me think."),
1577 other => panic!("Expected Thinking block, got {:?}", other),
1578 }
1579 match &partial.content[1] {
1580 ContentBlock::Text(t) => assert_eq!(t.text, "Hello world"),
1581 other => panic!("Expected Text block, got {:?}", other),
1582 }
1583
1584 let done = events2.iter().find_map(|e| match e {
1586 ProviderEvent::Done { message, .. } => Some(message.clone()),
1587 _ => None,
1588 });
1589 let done_msg = done.expect("Should have Done event");
1590 assert_eq!(done_msg.content.len(), 2);
1591 }
1592
1593 #[test]
1594 fn parse_stateful_tool_use_across_chunks() {
1595 let chunk1 = concat!(
1596 "data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\"}}\n",
1597 "\n",
1598 "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\",\"thinking\":\"\"}}\n",
1599 "\n",
1600 "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"I should search.\"}}\n",
1601 "\n"
1602 );
1603 let chunk2 = concat!(
1604 "data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tool_1\",\"name\":\"search\"}}\n",
1605 "\n",
1606 "data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"q\\\":\\\"rust\\\"}\"}}\n",
1607 "\n",
1608 "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"tool_use\"}}\n",
1609 "\n"
1610 );
1611
1612 let mut partial = AssistantMessage::new(Api::AnthropicMessages, "anthropic", MODEL);
1613 let mut usage = Usage::default();
1614 let mut pending_tc: std::collections::HashMap<usize, AnthropicPendingToolCall> =
1615 std::collections::HashMap::new();
1616
1617 let events1 =
1618 parse_anthropic_events_stateful(chunk1, &mut partial, &mut usage, &mut pending_tc);
1619 assert_eq!(events1.len(), 3); assert_eq!(partial.content.len(), 1);
1623
1624 let events2 =
1625 parse_anthropic_events_stateful(chunk2, &mut partial, &mut usage, &mut pending_tc);
1626 assert!(events2.len() >= 2); let done = events2.iter().find_map(|e| match e {
1632 ProviderEvent::Done { reason, .. } => Some(*reason),
1633 _ => None,
1634 });
1635 assert_eq!(done, Some(StopReason::ToolUse));
1636 }
1637 #[test]
1640 fn url_strips_trailing_v1_for_minimax() {
1641 assert_eq!(
1645 anthropic_messages_url("https://api.minimax.io/anthropic/v1"),
1646 "https://api.minimax.io/anthropic/v1/messages"
1647 );
1648 }
1649
1650 #[test]
1651 fn url_plain_anthropic_base() {
1652 assert_eq!(
1653 anthropic_messages_url("https://api.anthropic.com"),
1654 "https://api.anthropic.com/v1/messages"
1655 );
1656 }
1657
1658 #[test]
1659 fn url_handles_trailing_slash_then_v1() {
1660 assert_eq!(
1661 anthropic_messages_url("https://api.minimax.io/anthropic/v1/"),
1662 "https://api.minimax.io/anthropic/v1/messages"
1663 );
1664 }
1665
1666 #[test]
1667 fn url_preserves_non_v1_path_suffix() {
1668 assert_eq!(
1670 anthropic_messages_url("https://example.com/v1beta"),
1671 "https://example.com/v1beta/v1/messages"
1672 );
1673 }
1674}