1pub(crate) mod data;
8
9use crate::attachments::validate_request_attachments;
10use crate::provider::{LlmProvider, thinking_for_forced_tool};
11use crate::streaming::{
12 StreamBox, StreamDelta, StreamErrorKind, reqwest_body_error_delta, reqwest_error_delta,
13};
14use agent_sdk_foundation::llm::{
15 CacheTtl, ChatOutcome, ChatRequest, ChatResponse, ContentBlock, ThinkingConfig, ThinkingMode,
16 Usage,
17};
18use anyhow::Result;
19use async_trait::async_trait;
20use data::{
21 ApiMessagesRequest, ApiOutputConfig, ApiThinkingConfig, ApiToolChoice, build_api_messages,
22 build_api_tools_with_cache, is_message_stop_event, map_content_blocks, map_stop_reason,
23 parse_sse_event, take_next_sse_event,
24};
25use futures::StreamExt;
26use reqwest::StatusCode;
27
28const API_BASE_URL: &str = "https://api.anthropic.com";
29const API_VERSION: &str = "2023-06-01";
30const CLAUDE_CODE_VERSION: &str = "2.1.75";
31const DEFAULT_SAFE_MAX_OUTPUT_TOKENS: u32 = 32_000;
32const MODELS_PAGE_LIMIT: u32 = 1000;
34const MODELS_MAX_PAGES: usize = 100;
38const STREAM_HEADERS_TIMEOUT: std::time::Duration = std::time::Duration::from_mins(1);
51const SSE_BYTE_IDLE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(90);
61const CHAT_REQUEST_TIMEOUT: std::time::Duration = std::time::Duration::from_mins(5);
65const POOL_IDLE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
69
70pub const MODEL_HAIKU_35: &str = "claude-3-5-haiku-20241022";
71pub const MODEL_SONNET_35: &str = "claude-3-5-sonnet-20241022";
72pub const MODEL_SONNET_4: &str = "claude-sonnet-4-20250514";
73pub const MODEL_OPUS_4: &str = "claude-opus-4-20250514";
74
75pub const MODEL_HAIKU_45: &str = "claude-haiku-4-5-20251001";
76pub const MODEL_SONNET_45: &str = "claude-sonnet-4-5-20250929";
77pub const MODEL_SONNET_46: &str = "claude-sonnet-4-6";
78pub const MODEL_SONNET_5: &str = "claude-sonnet-5";
79pub const MODEL_OPUS_46: &str = "claude-opus-4-6";
80pub const MODEL_OPUS_47: &str = "claude-opus-4-7";
81pub const MODEL_OPUS_48: &str = "claude-opus-4-8";
82pub const MODEL_FABLE_5: &str = "claude-fable-5";
83
84const CLAUDE_CODE_TOOLS: &[&str] = &[
91 "Read",
92 "Write",
93 "Edit",
94 "Bash",
95 "Grep",
96 "Glob",
97 "AskUserQuestion",
98 "EnterPlanMode",
99 "ExitPlanMode",
100 "KillShell",
101 "NotebookEdit",
102 "Skill",
103 "Task",
104 "TaskOutput",
105 "TodoWrite",
106 "WebFetch",
107 "WebSearch",
108];
109
110fn to_claude_code_name(name: &str) -> String {
112 let lower = name.to_lowercase();
113 for cc_name in CLAUDE_CODE_TOOLS {
114 if cc_name.to_lowercase() == lower {
115 return (*cc_name).to_string();
116 }
117 }
118 name.to_string()
119}
120
121fn from_claude_code_name(name: &str, original_names: &[String]) -> String {
123 let lower = name.to_lowercase();
124 for original in original_names {
125 if original.to_lowercase() == lower {
126 return original.clone();
127 }
128 }
129 name.to_string()
130}
131
132fn oauth_tool_name_collision(
140 tools: Option<&[agent_sdk_foundation::llm::Tool]>,
141) -> Option<(String, String)> {
142 let tools = tools?;
143 for (index, tool) in tools.iter().enumerate() {
144 for other in &tools[index + 1..] {
145 if tool.name != other.name && tool.name.eq_ignore_ascii_case(&other.name) {
146 return Some((tool.name.clone(), other.name.clone()));
147 }
148 }
149 }
150 None
151}
152
153fn oauth_tool_collision_message(first: &str, second: &str) -> String {
154 format!(
155 "OAuth tool names collide case-insensitively: '{first}' and '{second}' would map to the same Claude Code tool name; rename one to disambiguate"
156 )
157}
158
159#[must_use]
161pub fn is_oauth_token(api_key: &str) -> bool {
162 api_key.starts_with("sk-ant-oat")
163}
164
165struct AnthropicModelsPage {
168 models: Vec<crate::provider::ModelInfo>,
169 has_more: bool,
170 last_id: Option<String>,
171}
172
173fn parse_models_page(body: &str) -> Result<AnthropicModelsPage> {
180 #[derive(serde::Deserialize)]
181 struct ListResponse {
182 #[serde(default)]
183 data: Vec<ModelRow>,
184 #[serde(default)]
185 has_more: bool,
186 #[serde(default)]
187 last_id: Option<String>,
188 }
189 #[derive(serde::Deserialize)]
190 struct ModelRow {
191 id: String,
192 #[serde(default)]
193 display_name: Option<String>,
194 }
195 let parsed: ListResponse = serde_json::from_str(body)
196 .map_err(|e| anyhow::anyhow!("failed to parse Anthropic models list: {e}"))?;
197 let models = parsed
198 .data
199 .into_iter()
200 .map(|row| crate::provider::ModelInfo {
201 id: row.id,
202 display_name: row.display_name,
203 context_window: None,
204 max_output_tokens: None,
205 })
206 .collect();
207 Ok(AnthropicModelsPage {
208 models,
209 has_more: parsed.has_more,
210 last_id: parsed.last_id,
211 })
212}
213
214struct CacheRegions {
218 tools: Option<data::ApiCacheControl>,
219 system: Option<data::ApiCacheControl>,
220 messages: Option<data::ApiCacheControl>,
221}
222
223impl CacheRegions {
224 const DISABLED: Self = Self {
226 tools: None,
227 system: None,
228 messages: None,
229 };
230}
231
232#[derive(Clone, Debug)]
234enum AuthMode {
235 ApiKey,
237 OAuth,
239}
240
241#[derive(Clone)]
243pub struct AnthropicProvider {
244 client: reqwest::Client,
245 api_key: String,
246 model: String,
247 base_url: String,
248 auth_mode: AuthMode,
249 thinking: Option<ThinkingConfig>,
250 extra_headers: Vec<(String, String)>,
252 stream_headers_timeout: std::time::Duration,
254 sse_byte_idle_timeout: std::time::Duration,
256}
257
258impl AnthropicProvider {
259 pub const API_KEY_ENV: &'static str = "ANTHROPIC_API_KEY";
261
262 #[must_use]
267 pub fn new(api_key: impl Into<String>, model: impl Into<String>) -> Self {
268 let api_key = api_key.into();
269 let model = model.into();
270 let auth_mode = if is_oauth_token(&api_key) {
271 AuthMode::OAuth
272 } else {
273 AuthMode::ApiKey
274 };
275
276 let client = reqwest::Client::builder()
283 .connect_timeout(std::time::Duration::from_secs(30))
284 .tcp_keepalive(std::time::Duration::from_secs(30))
285 .pool_idle_timeout(POOL_IDLE_TIMEOUT)
286 .build()
287 .unwrap_or_else(|error| {
288 log::error!(
291 "failed to build Anthropic HTTP client with timeouts, \
292 falling back to reqwest defaults: {error}"
293 );
294 reqwest::Client::default()
295 });
296
297 Self {
298 client,
299 api_key,
300 model,
301 base_url: API_BASE_URL.to_owned(),
302 auth_mode,
303 thinking: None,
304 extra_headers: Vec::new(),
305 stream_headers_timeout: STREAM_HEADERS_TIMEOUT,
306 sse_byte_idle_timeout: SSE_BYTE_IDLE_TIMEOUT,
307 }
308 }
309
310 #[must_use]
312 pub const fn is_oauth(&self) -> bool {
313 matches!(self.auth_mode, AuthMode::OAuth)
314 }
315
316 fn apply_auth(&self, builder: reqwest::RequestBuilder) -> reqwest::RequestBuilder {
322 let builder = if self.api_key.is_empty() {
323 builder.header("anthropic-version", API_VERSION)
324 } else {
325 match self.auth_mode {
326 AuthMode::ApiKey => {
327 let builder = builder
328 .header("x-api-key", &self.api_key)
329 .header("anthropic-version", API_VERSION);
330 if self.is_adaptive_thinking_model() {
338 builder
339 } else {
340 builder.header("anthropic-beta", "interleaved-thinking-2025-05-14")
341 }
342 }
343 AuthMode::OAuth => {
344 let mut beta_features = vec![
348 "claude-code-20250219",
349 "oauth-2025-04-20",
350 "fine-grained-tool-streaming-2025-05-14",
351 ];
352 if !self.is_adaptive_thinking_model() {
353 beta_features.push("interleaved-thinking-2025-05-14");
354 }
355 builder
356 .header("Authorization", format!("Bearer {}", self.api_key))
357 .header("anthropic-version", API_VERSION)
358 .header("anthropic-beta", beta_features.join(","))
359 .header("user-agent", format!("claude-cli/{CLAUDE_CODE_VERSION}"))
360 .header("x-app", "cli")
361 }
362 }
363 };
364 self.extra_headers
365 .iter()
366 .fold(builder, |b, (k, v)| b.header(k.as_str(), v.as_str()))
367 }
368
369 const OAUTH_IDENTITY: &'static str =
370 "You are Claude Code, Anthropic's official CLI for Claude.";
371
372 fn build_system_prompt_for_request<'a>(
378 &self,
379 system: &'a str,
380 cache_control: Option<data::ApiCacheControl>,
381 ) -> Option<data::ApiSystemPrompt<'a>> {
382 match self.auth_mode {
383 AuthMode::ApiKey => data::build_api_system_prompt(system, cache_control),
384 AuthMode::OAuth => {
385 let mut blocks = vec![data::ApiSystemBlock {
386 block_type: "text",
387 text: Self::OAUTH_IDENTITY,
388 cache_control: cache_control.clone(),
389 }];
390 if !system.is_empty() {
391 blocks.push(data::ApiSystemBlock {
392 block_type: "text",
393 text: system,
394 cache_control,
395 });
396 }
397 Some(data::ApiSystemPrompt::Blocks(blocks))
398 }
399 }
400 }
401
402 fn cache_regions(request: &ChatRequest) -> CacheRegions {
411 let (enabled, ttl, max_breakpoints) =
412 request.cache.as_ref().map_or((true, None, None), |cfg| {
413 (cfg.enabled, cfg.ttl, cfg.max_breakpoints)
414 });
415 if !enabled {
416 return CacheRegions::DISABLED;
417 }
418 let control = data::ApiCacheControl::ephemeral_with_ttl(ttl.map(CacheTtl::as_wire_str));
419 let limit = max_breakpoints.unwrap_or(u8::MAX);
420 CacheRegions {
421 tools: (limit >= 1).then(|| control.clone()),
422 system: (limit >= 2).then(|| control.clone()),
423 messages: (limit >= 3).then_some(control),
424 }
425 }
426
427 fn build_cached_api_messages(
428 request: &ChatRequest,
429 cache_control: Option<data::ApiCacheControl>,
430 ) -> Vec<data::ApiMessage> {
431 let mut messages = build_api_messages(request);
432 if let Some(cache_control) = cache_control {
433 data::apply_cache_control_to_last_user_message(&mut messages, cache_control);
434 }
435 messages
436 }
437
438 fn effective_max_tokens(&self, request: &ChatRequest) -> u32 {
439 if request.max_tokens_explicit {
440 request.max_tokens
441 } else {
442 self.default_max_tokens()
443 }
444 }
445
446 #[must_use]
459 pub fn from_env() -> Self {
460 Self::try_from_env().unwrap_or_else(|e| panic!("{e}"))
461 }
462
463 pub fn try_from_env() -> Result<Self> {
471 let api_key = std::env::var(Self::API_KEY_ENV).map_err(|_| {
472 anyhow::anyhow!("environment variable `{}` is not set", Self::API_KEY_ENV)
473 })?;
474 Ok(Self::sonnet(api_key))
475 }
476
477 #[must_use]
479 pub fn haiku(api_key: impl Into<String>) -> Self {
480 Self::new(api_key, MODEL_HAIKU_45)
481 }
482
483 #[must_use]
485 pub fn sonnet(api_key: impl Into<String>) -> Self {
486 Self::new(api_key, MODEL_SONNET_46)
487 }
488
489 #[must_use]
491 pub fn sonnet_45(api_key: impl Into<String>) -> Self {
492 Self::new(api_key, MODEL_SONNET_45)
493 }
494
495 #[must_use]
497 pub fn sonnet_46(api_key: impl Into<String>) -> Self {
498 Self::new(api_key, MODEL_SONNET_46)
499 }
500
501 #[must_use]
503 pub fn opus(api_key: impl Into<String>) -> Self {
504 Self::new(api_key, MODEL_OPUS_46)
505 }
506
507 #[must_use]
515 pub fn opus_47(api_key: impl Into<String>) -> Self {
516 Self::new(api_key, MODEL_OPUS_47)
517 }
518
519 #[must_use]
527 pub fn opus_48(api_key: impl Into<String>) -> Self {
528 Self::new(api_key, MODEL_OPUS_48)
529 }
530
531 #[must_use]
541 pub fn fable(api_key: impl Into<String>) -> Self {
542 Self::new(api_key, MODEL_FABLE_5)
543 }
544
545 #[must_use]
549 pub fn sonnet_5(api_key: impl Into<String>) -> Self {
550 Self::new(api_key, MODEL_SONNET_5)
551 }
552
553 #[must_use]
555 pub const fn with_thinking(mut self, thinking: ThinkingConfig) -> Self {
556 self.thinking = Some(thinking);
557 self
558 }
559
560 #[must_use]
562 pub fn with_base_url(mut self, base_url: impl Into<String>) -> Self {
563 self.base_url = base_url.into();
564 self
565 }
566
567 #[must_use]
573 pub const fn with_stream_stall_timeouts(
574 mut self,
575 headers_timeout: std::time::Duration,
576 byte_idle_timeout: std::time::Duration,
577 ) -> Self {
578 self.stream_headers_timeout = headers_timeout;
579 self.sse_byte_idle_timeout = byte_idle_timeout;
580 self
581 }
582
583 #[must_use]
585 pub fn with_extra_headers(mut self, headers: Vec<(String, String)>) -> Self {
586 self.extra_headers = headers;
587 self
588 }
589
590 fn is_adaptive_thinking_model(&self) -> bool {
591 matches!(
592 self.model.as_str(),
593 MODEL_SONNET_46
594 | MODEL_SONNET_5
595 | MODEL_OPUS_46
596 | MODEL_OPUS_47
597 | MODEL_OPUS_48
598 | MODEL_FABLE_5
599 )
600 }
601}
602
603#[async_trait]
604#[allow(clippy::too_many_lines)]
605impl LlmProvider for AnthropicProvider {
606 async fn chat(&self, request: ChatRequest) -> Result<ChatOutcome> {
607 let thinking_config = match self.resolve_thinking_config(request.thinking.as_ref()) {
608 Ok(thinking) => thinking_for_forced_tool(thinking, request.tool_choice.as_ref()),
612 Err(error) => return Ok(ChatOutcome::InvalidRequest(error.to_string())),
613 };
614 if let Err(error) = validate_request_attachments(self.provider(), self.model(), &request) {
615 return Ok(ChatOutcome::InvalidRequest(error.to_string()));
616 }
617 if self.is_oauth()
618 && let Some((first, second)) = oauth_tool_name_collision(request.tools.as_deref())
619 {
620 return Ok(ChatOutcome::InvalidRequest(oauth_tool_collision_message(
621 &first, &second,
622 )));
623 }
624 let CacheRegions {
625 tools: tools_cache,
626 system: system_cache,
627 messages: messages_cache,
628 } = Self::cache_regions(&request);
629 let messages = Self::build_cached_api_messages(&request, messages_cache);
630 let tools = if self.is_oauth() {
631 build_api_tools_with_cache(&request, tools_cache).map(|tools| {
632 tools
633 .into_iter()
634 .map(|mut t| {
635 t.name = to_claude_code_name(&t.name);
636 t
637 })
638 .collect::<Vec<_>>()
639 })
640 } else {
641 build_api_tools_with_cache(&request, tools_cache)
642 };
643 let thinking = thinking_config
644 .as_ref()
645 .and_then(ApiThinkingConfig::from_thinking_config);
646 let output_config = thinking_config
647 .as_ref()
648 .and_then(|t| t.effort)
649 .map(|effort| ApiOutputConfig { effort });
650
651 let system = self.build_system_prompt_for_request(&request.system, system_cache);
652 let max_tokens = self.effective_max_tokens(&request);
653 let tool_choice = request
654 .tool_choice
655 .as_ref()
656 .map(ApiToolChoice::from_tool_choice);
657
658 let api_request = ApiMessagesRequest {
659 model: Some(&self.model),
660 max_tokens,
661 system,
662 messages: &messages,
663 tools: tools.as_deref(),
664 tool_choice,
665 stream: false,
666 thinking,
667 output_config,
668 anthropic_version: None,
669 };
670
671 log::debug!(
672 "Anthropic LLM request model={} max_tokens={} oauth={}",
673 self.model,
674 max_tokens,
675 self.is_oauth()
676 );
677
678 if log::log_enabled!(log::Level::Debug) {
680 match serde_json::to_string_pretty(&api_request) {
681 Ok(json) => log::debug!("Anthropic API request payload:\n{json}"),
682 Err(e) => log::debug!("Failed to serialize request for logging: {e}"),
683 }
684 }
685
686 let builder = self
691 .client
692 .post(format!("{}/v1/messages", self.base_url))
693 .timeout(CHAT_REQUEST_TIMEOUT)
694 .header("Content-Type", "application/json");
695 let response = self
696 .apply_auth(builder)
697 .json(&api_request)
698 .send()
699 .await
700 .map_err(|e| anyhow::anyhow!("request failed: {e}"))?;
701
702 let status = response.status();
703 let retry_after = if status == StatusCode::TOO_MANY_REQUESTS {
706 crate::http::retry_after_from_headers(response.headers())
707 } else {
708 None
709 };
710 let bytes = response
711 .bytes()
712 .await
713 .map_err(|e| anyhow::anyhow!("failed to read response body: {e}"))?;
714
715 log::debug!(
716 "Anthropic LLM response status={} body_len={}",
717 status,
718 bytes.len()
719 );
720
721 if status == StatusCode::TOO_MANY_REQUESTS {
722 return Ok(ChatOutcome::RateLimited(retry_after));
723 }
724
725 if status.is_server_error() {
726 let body = String::from_utf8_lossy(&bytes);
727 log::error!("Anthropic server error status={status} body={body}");
728 return Ok(ChatOutcome::ServerError(body.into_owned()));
729 }
730
731 if status.is_client_error() {
732 let body = String::from_utf8_lossy(&bytes);
733 log::warn!("Anthropic client error status={status} body={body}");
734 return Ok(ChatOutcome::InvalidRequest(body.into_owned()));
735 }
736
737 let api_response: data::ApiResponse = serde_json::from_slice(&bytes)
738 .map_err(|e| anyhow::anyhow!("failed to parse response: {e}"))?;
739
740 log::debug!(
742 "Anthropic API response: id={} model={} stop_reason={:?} usage={{input_tokens={}, output_tokens={}}} content_blocks={}",
743 api_response.id,
744 api_response.model,
745 api_response.stop_reason,
746 api_response.usage.total_input_tokens(),
747 api_response.usage.output,
748 api_response.content.len()
749 );
750
751 let mut content = map_content_blocks(api_response.content);
752
753 if self.is_oauth() {
755 let original_names: Vec<String> = request
756 .tools
757 .as_ref()
758 .map(|ts| ts.iter().map(|t| t.name.clone()).collect())
759 .unwrap_or_default();
760 for block in &mut content {
761 if let ContentBlock::ToolUse { name, .. } = block {
762 *name = from_claude_code_name(name, &original_names);
763 }
764 }
765 }
766
767 let stop_reason = api_response.stop_reason.as_ref().map(map_stop_reason);
768
769 Ok(ChatOutcome::Success(ChatResponse {
770 id: api_response.id,
771 content,
772 model: api_response.model,
773 stop_reason,
774 usage: Usage {
775 input_tokens: api_response.usage.total_input_tokens(),
776 output_tokens: api_response.usage.output,
777 cached_input_tokens: api_response.usage.cached_input_tokens(),
778 cache_creation_input_tokens: api_response.usage.cache_creation_input_tokens(),
779 },
780 }))
781 }
782
783 fn chat_stream(&self, request: ChatRequest) -> StreamBox<'_> {
784 Box::pin(async_stream::stream! {
785 let is_oauth = self.is_oauth();
786 let original_tool_names: Vec<String> = request
787 .tools
788 .as_ref()
789 .map(|ts| ts.iter().map(|t| t.name.clone()).collect())
790 .unwrap_or_default();
791
792 if let Err(error) = validate_request_attachments(self.provider(), self.model(), &request) {
793 yield Ok(StreamDelta::Error {
794 message: error.to_string(),
795 kind: StreamErrorKind::InvalidRequest,
796 });
797 return;
798 }
799
800 if is_oauth
801 && let Some((first, second)) = oauth_tool_name_collision(request.tools.as_deref())
802 {
803 yield Ok(StreamDelta::Error {
804 message: oauth_tool_collision_message(&first, &second),
805 kind: StreamErrorKind::InvalidRequest,
806 });
807 return;
808 }
809
810 let CacheRegions {
811 tools: tools_cache,
812 system: system_cache,
813 messages: messages_cache,
814 } = Self::cache_regions(&request);
815 let messages = Self::build_cached_api_messages(&request, messages_cache);
816 let tools = if is_oauth {
817 build_api_tools_with_cache(&request, tools_cache).map(|tools| {
818 tools
819 .into_iter()
820 .map(|mut t| {
821 t.name = to_claude_code_name(&t.name);
822 t
823 })
824 .collect::<Vec<_>>()
825 })
826 } else {
827 build_api_tools_with_cache(&request, tools_cache)
828 };
829 let thinking_config = match self.resolve_thinking_config(request.thinking.as_ref()) {
830 Ok(thinking) => thinking_for_forced_tool(thinking, request.tool_choice.as_ref()),
835 Err(error) => {
836 yield Ok(StreamDelta::Error {
837 message: error.to_string(),
838 kind: StreamErrorKind::InvalidRequest,
839 });
840 return;
841 }
842 };
843 let thinking = thinking_config
844 .as_ref()
845 .and_then(ApiThinkingConfig::from_thinking_config);
846 let output_config = thinking_config
847 .as_ref()
848 .and_then(|t| t.effort)
849 .map(|effort| ApiOutputConfig { effort });
850
851 let system = self.build_system_prompt_for_request(&request.system, system_cache);
852 let max_tokens = self.effective_max_tokens(&request);
853 let tool_choice = request
854 .tool_choice
855 .as_ref()
856 .map(ApiToolChoice::from_tool_choice);
857
858 let api_request = ApiMessagesRequest {
859 model: Some(&self.model),
860 max_tokens,
861 system,
862 messages: &messages,
863 tools: tools.as_deref(),
864 tool_choice,
865 stream: true,
866 thinking,
867 output_config,
868 anthropic_version: None,
869 };
870
871 log::debug!("Anthropic streaming LLM request model={} max_tokens={} oauth={}", self.model, max_tokens, is_oauth);
872
873 if log::log_enabled!(log::Level::Debug) {
875 match serde_json::to_string_pretty(&api_request) {
876 Ok(json) => log::debug!("Anthropic streaming API request payload:\n{json}"),
877 Err(e) => log::debug!("Failed to serialize streaming request for logging: {e}"),
878 }
879 }
880
881 let builder = self
882 .client
883 .post(format!("{}/v1/messages", self.base_url))
884 .header("Content-Type", "application/json");
885 let headers_timeout = self.stream_headers_timeout;
894 let send = self.apply_auth(builder).json(&api_request).send();
895 let response = match tokio::time::timeout(headers_timeout, send).await {
896 Ok(Ok(r)) => r,
897 Ok(Err(error)) => {
898 yield Ok(reqwest_error_delta("request failed", &error));
899 return;
900 }
901 Err(_elapsed) => {
902 log::error!(
903 "Anthropic streaming request timed out awaiting response headers after {}s — stalled connection",
904 headers_timeout.as_secs()
905 );
906 yield Ok(StreamDelta::Error {
907 message: format!(
908 "request timed out awaiting response headers after {}s",
909 headers_timeout.as_secs()
910 ),
911 kind: StreamErrorKind::ConnectionLost,
912 });
913 return;
914 }
915 };
916
917 let status = response.status();
918
919 if status == StatusCode::TOO_MANY_REQUESTS {
920 let retry_after = crate::http::retry_after_from_headers(response.headers());
921 yield Ok(StreamDelta::Error {
922 message: "Rate limited".to_string(),
923 kind: StreamErrorKind::RateLimited(retry_after),
924 });
925 return;
926 }
927
928 if status.is_server_error() {
929 let body = response.text().await.unwrap_or_default();
930 log::error!("Anthropic server error status={status} body={body}");
931 yield Ok(StreamDelta::Error {
932 message: body,
933 kind: StreamErrorKind::ServerError,
934 });
935 return;
936 }
937
938 if status.is_client_error() {
939 let body = response.text().await.unwrap_or_default();
940 log::warn!("Anthropic client error status={status} body={body}");
941 yield Ok(StreamDelta::Error {
942 message: body,
943 kind: StreamErrorKind::InvalidRequest,
944 });
945 return;
946 }
947
948 let mut stream = response.bytes_stream();
950 let mut buffer = String::new();
951 let mut input_tokens: u32 = 0;
952 let mut output_tokens: u32 = 0;
953 let mut cached_input_tokens: u32 = 0;
954 let mut cache_creation_input_tokens: u32 = 0;
955 let mut tool_ids: std::collections::HashMap<usize, String> =
957 std::collections::HashMap::new();
958
959 let mut received_message_stop = false;
960 let mut stream_errored = false;
965 let mut pending_stop_reason: Option<agent_sdk_foundation::llm::StopReason> = None;
966 let mut chunk_count: u64 = 0;
967 let mut total_bytes: u64 = 0;
968
969 struct StreamDropGuard {
971 completed: bool,
972 chunk_count: u64,
973 }
974 impl Drop for StreamDropGuard {
975 fn drop(&mut self) {
976 if !self.completed {
977 log::debug!(
981 "SSE stream dropped before completion at chunk_count={} (task was likely cancelled)",
982 self.chunk_count
983 );
984 }
985 }
986 }
987 let mut drop_guard = StreamDropGuard { completed: false, chunk_count: 0 };
988
989 log::debug!("Starting SSE stream processing");
990
991 let byte_idle_timeout = self.sse_byte_idle_timeout;
996 loop {
997 let next = match tokio::time::timeout(byte_idle_timeout, stream.next()).await {
998 Ok(next) => next,
999 Err(_elapsed) => {
1000 log::error!(
1001 "SSE stream timed out: no bytes for {}s chunk_count={chunk_count} total_bytes={total_bytes} — stalled connection",
1002 byte_idle_timeout.as_secs()
1003 );
1004 yield Ok(StreamDelta::Error {
1005 message: format!(
1006 "SSE stream timed out: no bytes for {}s",
1007 byte_idle_timeout.as_secs()
1008 ),
1009 kind: StreamErrorKind::ConnectionLost,
1010 });
1011 return;
1012 }
1013 };
1014 let Some(chunk_result) = next else { break };
1015 let chunk = match chunk_result {
1016 Ok(c) => c,
1017 Err(error) => {
1018 log::error!("Stream error while reading chunk error={error} chunk_count={chunk_count} total_bytes={total_bytes}");
1019 yield Ok(reqwest_body_error_delta("stream error", &error));
1020 return;
1021 }
1022 };
1023
1024 chunk_count += 1;
1025 total_bytes += chunk.len() as u64;
1026 drop_guard.chunk_count = chunk_count;
1027
1028 if chunk_count.is_multiple_of(10) {
1030 log::debug!("SSE chunk progress: chunk_count={chunk_count} total_bytes={total_bytes}");
1031 }
1032 buffer.push_str(&String::from_utf8_lossy(&chunk));
1033
1034 while let Some(event_block) = take_next_sse_event(&mut buffer) {
1036 if is_message_stop_event(&event_block) {
1038 log::debug!("Received message_stop event chunk_count={chunk_count} total_bytes={total_bytes}");
1039 received_message_stop = true;
1040 }
1041
1042 if let Some(mut delta) = parse_sse_event(
1044 &event_block,
1045 &mut input_tokens,
1046 &mut output_tokens,
1047 &mut cached_input_tokens,
1048 &mut cache_creation_input_tokens,
1049 &mut tool_ids,
1050 &mut pending_stop_reason,
1051 ) {
1052 if is_oauth
1054 && let StreamDelta::ToolUseStart { ref mut name, .. } = delta
1055 {
1056 *name = from_claude_code_name(name, &original_tool_names);
1057 }
1058 if matches!(delta, StreamDelta::Error { .. }) {
1061 stream_errored = true;
1062 }
1063 yield Ok(delta);
1064 }
1065 if is_message_stop_event(&event_block) {
1067 yield Ok(StreamDelta::Done {
1068 stop_reason: pending_stop_reason.take(),
1069 });
1070 }
1071 }
1072 }
1073
1074 log::debug!(
1075 "SSE stream ended chunk_count={chunk_count} total_bytes={total_bytes} buffer_remaining={} received_message_stop={received_message_stop}",
1076 buffer.len()
1077 );
1078
1079 let remaining = buffer.trim();
1081 if !remaining.is_empty() {
1082 log::debug!(
1083 "Processing remaining buffer content remaining_len={} remaining_preview={}",
1084 remaining.len(),
1085 remaining.chars().take(100).collect::<String>()
1086 );
1087
1088 if is_message_stop_event(remaining) {
1090 received_message_stop = true;
1091 }
1092
1093 if let Some(mut delta) = parse_sse_event(
1094 remaining,
1095 &mut input_tokens,
1096 &mut output_tokens,
1097 &mut cached_input_tokens,
1098 &mut cache_creation_input_tokens,
1099 &mut tool_ids,
1100 &mut pending_stop_reason,
1101 ) {
1102 if is_oauth
1103 && let StreamDelta::ToolUseStart { ref mut name, .. } = delta
1104 {
1105 *name = from_claude_code_name(name, &original_tool_names);
1106 }
1107 if matches!(delta, StreamDelta::Error { .. }) {
1108 stream_errored = true;
1109 }
1110 yield Ok(delta);
1111 }
1112 if is_message_stop_event(remaining) {
1114 yield Ok(StreamDelta::Done {
1115 stop_reason: pending_stop_reason.take(),
1116 });
1117 }
1118 }
1119
1120 drop_guard.completed = true;
1122
1123 if !received_message_stop && !stream_errored {
1129 log::warn!(
1130 "SSE stream ended without message_stop event - stream may have been interrupted chunk_count={chunk_count} total_bytes={total_bytes}"
1131 );
1132 yield Ok(StreamDelta::Error {
1133 message: "Stream ended unexpectedly without completion".to_string(),
1134 kind: StreamErrorKind::ServerError,
1135 });
1136 }
1137 })
1138 }
1139
1140 fn validate_thinking_config(&self, thinking: Option<&ThinkingConfig>) -> Result<()> {
1141 let Some(thinking) = thinking else {
1142 return Ok(());
1143 };
1144
1145 if self
1146 .capabilities()
1147 .is_some_and(|caps| !caps.supports_thinking)
1148 {
1149 return Err(anyhow::anyhow!(
1150 "thinking is not supported for provider={} model={}",
1151 self.provider(),
1152 self.model()
1153 ));
1154 }
1155
1156 if matches!(thinking.mode, ThinkingMode::Adaptive)
1157 && !self
1158 .capabilities()
1159 .is_some_and(|caps| caps.supports_adaptive_thinking)
1160 {
1161 return Err(anyhow::anyhow!(
1162 "adaptive thinking is not supported for provider={} model={}",
1163 self.provider(),
1164 self.model()
1165 ));
1166 }
1167
1168 if self.is_adaptive_thinking_model()
1169 && matches!(thinking.mode, ThinkingMode::Enabled { .. })
1170 {
1171 return Err(anyhow::anyhow!(
1172 "budget_tokens thinking is rejected for provider={} model={}; use ThinkingConfig::adaptive() or ThinkingConfig::default_with_effort(_) instead",
1173 self.provider(),
1174 self.model()
1175 ));
1176 }
1177
1178 Ok(())
1179 }
1180
1181 async fn list_models(&self) -> Result<Vec<crate::provider::ModelInfo>> {
1182 let mut models = Vec::new();
1186 let mut after_id: Option<String> = None;
1187 for _ in 0..MODELS_MAX_PAGES {
1188 let mut query: Vec<(&str, String)> = vec![("limit", MODELS_PAGE_LIMIT.to_string())];
1189 if let Some(after) = &after_id {
1190 query.push(("after_id", after.clone()));
1191 }
1192 let builder = self
1193 .client
1194 .get(format!("{}/v1/models", self.base_url))
1195 .header("Content-Type", "application/json")
1196 .query(&query);
1197 let builder = self.apply_auth(builder);
1198 let body =
1199 crate::impls::model_listing::fetch_model_list_body(builder, "Anthropic").await?;
1200 let page = parse_models_page(&body)?;
1201 models.extend(page.models);
1202 if !page.has_more {
1203 return Ok(models);
1204 }
1205 match page.last_id {
1206 Some(last) => after_id = Some(last),
1207 None => return Ok(models),
1209 }
1210 }
1211 Ok(models)
1212 }
1213
1214 async fn probe_connectivity(&self) -> bool {
1215 crate::provider::probe_http_reachability(&self.client, &self.base_url).await
1216 }
1217
1218 fn model(&self) -> &str {
1219 &self.model
1220 }
1221
1222 fn provider(&self) -> &'static str {
1223 "anthropic"
1224 }
1225
1226 fn configured_thinking(&self) -> Option<&ThinkingConfig> {
1227 self.thinking.as_ref()
1228 }
1229
1230 fn default_max_tokens(&self) -> u32 {
1231 let model_max = self
1232 .capabilities()
1233 .and_then(|caps| caps.max_output_tokens)
1234 .or_else(|| {
1235 crate::model_capabilities::default_max_output_tokens(self.provider(), self.model())
1236 })
1237 .unwrap_or(4096);
1238 model_max.clamp(4096, DEFAULT_SAFE_MAX_OUTPUT_TOKENS)
1239 }
1240}
1241
1242#[cfg(test)]
1243mod tests {
1244 use super::*;
1245
1246 const ANTHROPIC_MODELS_FIXTURE: &str = r#"{
1247 "data": [
1248 {"type": "model", "id": "claude-opus-4-8", "display_name": "Claude Opus 4.8"},
1249 {"type": "model", "id": "claude-sonnet-4-5", "display_name": "Claude Sonnet 4.5"}
1250 ],
1251 "has_more": false
1252 }"#;
1253
1254 #[test]
1255 fn parse_models_page_reads_id_and_display_name() -> anyhow::Result<()> {
1256 let page = parse_models_page(ANTHROPIC_MODELS_FIXTURE)?;
1257 assert_eq!(page.models.len(), 2);
1258 assert_eq!(page.models[0].id, "claude-opus-4-8");
1259 assert_eq!(
1260 page.models[0].display_name.as_deref(),
1261 Some("Claude Opus 4.8")
1262 );
1263 assert_eq!(page.models[0].context_window, None);
1265 assert_eq!(page.models[0].max_output_tokens, None);
1266 assert!(!page.has_more);
1268 assert_eq!(page.last_id, None);
1269 Ok(())
1270 }
1271
1272 #[tokio::test]
1273 async fn list_models_follows_pagination_across_pages() -> anyhow::Result<()> {
1274 use wiremock::matchers::{method, path, query_param, query_param_is_missing};
1275 use wiremock::{Mock, MockServer, ResponseTemplate};
1276
1277 let server = MockServer::start().await;
1278
1279 Mock::given(method("GET"))
1282 .and(path("/v1/models"))
1283 .and(query_param_is_missing("after_id"))
1284 .respond_with(ResponseTemplate::new(200).set_body_string(
1285 r#"{
1286 "data": [
1287 {"type": "model", "id": "claude-opus-4-8", "display_name": "Opus"},
1288 {"type": "model", "id": "claude-sonnet-4-5", "display_name": "Sonnet"}
1289 ],
1290 "has_more": true,
1291 "last_id": "claude-sonnet-4-5"
1292 }"#,
1293 ))
1294 .mount(&server)
1295 .await;
1296
1297 Mock::given(method("GET"))
1299 .and(path("/v1/models"))
1300 .and(query_param("after_id", "claude-sonnet-4-5"))
1301 .respond_with(ResponseTemplate::new(200).set_body_string(
1302 r#"{
1303 "data": [
1304 {"type": "model", "id": "claude-haiku-4-5", "display_name": "Haiku"}
1305 ],
1306 "has_more": false,
1307 "last_id": "claude-haiku-4-5"
1308 }"#,
1309 ))
1310 .mount(&server)
1311 .await;
1312
1313 let provider = AnthropicProvider::new("test-key-not-a-secret", "claude-test")
1314 .with_base_url(server.uri());
1315 let models = provider.list_models().await?;
1316
1317 let ids: Vec<&str> = models.iter().map(|m| m.id.as_str()).collect();
1319 assert_eq!(
1320 ids,
1321 vec!["claude-opus-4-8", "claude-sonnet-4-5", "claude-haiku-4-5"]
1322 );
1323 Ok(())
1324 }
1325
1326 #[test]
1331 fn test_new_creates_provider_with_custom_model() {
1332 let provider = AnthropicProvider::new("test-api-key", "custom-model");
1333
1334 assert_eq!(provider.model(), "custom-model");
1335 assert_eq!(provider.provider(), "anthropic");
1336 }
1337
1338 #[test]
1339 fn test_haiku_factory_creates_haiku_provider() {
1340 let provider = AnthropicProvider::haiku("test-api-key".to_string());
1341
1342 assert_eq!(provider.model(), MODEL_HAIKU_45);
1343 assert_eq!(provider.provider(), "anthropic");
1344 }
1345
1346 #[test]
1347 fn test_only_anthropic_46_models_accept_adaptive_thinking() {
1348 let sonnet_46 = AnthropicProvider::sonnet_46("test-api-key".to_string());
1349 assert!(
1350 sonnet_46
1351 .validate_thinking_config(Some(&ThinkingConfig::adaptive()))
1352 .is_ok()
1353 );
1354
1355 let sonnet_45 = AnthropicProvider::sonnet_45("test-api-key".to_string());
1356 let error = sonnet_45
1357 .validate_thinking_config(Some(&ThinkingConfig::adaptive()))
1358 .unwrap_err();
1359 assert!(
1360 error
1361 .to_string()
1362 .contains("adaptive thinking is not supported")
1363 );
1364 }
1365
1366 #[test]
1367 fn test_anthropic_46_models_reject_budgeted_thinking() {
1368 let sonnet_46 = AnthropicProvider::sonnet_46("test-api-key".to_string());
1369 let error = sonnet_46
1370 .validate_thinking_config(Some(&ThinkingConfig::new(10_000)))
1371 .unwrap_err();
1372 assert!(error.to_string().contains("ThinkingConfig::adaptive()"));
1373 }
1374
1375 #[test]
1376 fn test_opus_47_rejects_budgeted_thinking() {
1377 let opus_47 = AnthropicProvider::opus_47("test-api-key".to_string());
1382 let error = opus_47
1383 .validate_thinking_config(Some(&ThinkingConfig::new(10_000)))
1384 .unwrap_err();
1385 assert!(
1386 error.to_string().contains("ThinkingConfig::adaptive()"),
1387 "expected migration hint, got: {error}"
1388 );
1389 }
1390
1391 #[test]
1392 fn test_opus_47_accepts_adaptive_thinking() {
1393 let opus_47 = AnthropicProvider::opus_47("test-api-key".to_string());
1394 assert!(
1395 opus_47
1396 .validate_thinking_config(Some(&ThinkingConfig::adaptive()))
1397 .is_ok()
1398 );
1399 assert!(
1400 opus_47
1401 .validate_thinking_config(Some(&ThinkingConfig::adaptive_with_effort(
1402 agent_sdk_foundation::llm::Effort::High
1403 )))
1404 .is_ok()
1405 );
1406 assert!(
1407 opus_47
1408 .validate_thinking_config(Some(&ThinkingConfig::default_with_effort(
1409 agent_sdk_foundation::llm::Effort::High
1410 )))
1411 .is_ok(),
1412 "effort without adaptive must be accepted",
1413 );
1414 }
1415
1416 #[test]
1417 fn test_opus_47_factory_creates_opus_47_provider() {
1418 let provider = AnthropicProvider::opus_47("test-api-key".to_string());
1419 assert_eq!(provider.model(), MODEL_OPUS_47);
1420 assert_eq!(provider.provider(), "anthropic");
1421 }
1422
1423 #[test]
1424 fn test_opus_48_rejects_budgeted_thinking() {
1425 let opus_48 = AnthropicProvider::opus_48("test-api-key".to_string());
1430 let error = opus_48
1431 .validate_thinking_config(Some(&ThinkingConfig::new(10_000)))
1432 .unwrap_err();
1433 assert!(
1434 error.to_string().contains("ThinkingConfig::adaptive()"),
1435 "expected migration hint, got: {error}"
1436 );
1437 }
1438
1439 #[test]
1440 fn test_opus_48_accepts_adaptive_thinking() {
1441 let opus_48 = AnthropicProvider::opus_48("test-api-key".to_string());
1442 assert!(
1443 opus_48
1444 .validate_thinking_config(Some(&ThinkingConfig::adaptive()))
1445 .is_ok()
1446 );
1447 assert!(
1448 opus_48
1449 .validate_thinking_config(Some(&ThinkingConfig::adaptive_with_effort(
1450 agent_sdk_foundation::llm::Effort::High
1451 )))
1452 .is_ok()
1453 );
1454 assert!(
1455 opus_48
1456 .validate_thinking_config(Some(&ThinkingConfig::default_with_effort(
1457 agent_sdk_foundation::llm::Effort::High
1458 )))
1459 .is_ok(),
1460 "effort without adaptive must be accepted",
1461 );
1462 }
1463
1464 #[test]
1465 fn test_opus_48_factory_creates_opus_48_provider() {
1466 let provider = AnthropicProvider::opus_48("test-api-key".to_string());
1467 assert_eq!(provider.model(), MODEL_OPUS_48);
1468 assert_eq!(provider.provider(), "anthropic");
1469 }
1470
1471 #[test]
1472 fn test_sonnet_5_rejects_budgeted_thinking() {
1473 let sonnet_5 = AnthropicProvider::sonnet_5("test-api-key".to_string());
1476 let error = sonnet_5
1477 .validate_thinking_config(Some(&ThinkingConfig::new(10_000)))
1478 .unwrap_err();
1479 assert!(
1480 error.to_string().contains("ThinkingConfig::adaptive()"),
1481 "expected migration hint, got: {error}"
1482 );
1483 }
1484
1485 #[test]
1486 fn test_sonnet_5_accepts_adaptive_thinking() {
1487 let sonnet_5 = AnthropicProvider::sonnet_5("test-api-key".to_string());
1488 assert!(
1489 sonnet_5
1490 .validate_thinking_config(Some(&ThinkingConfig::adaptive()))
1491 .is_ok()
1492 );
1493 assert!(
1494 sonnet_5
1495 .validate_thinking_config(Some(&ThinkingConfig::adaptive_with_effort(
1496 agent_sdk_foundation::llm::Effort::High
1497 )))
1498 .is_ok()
1499 );
1500 assert!(
1501 sonnet_5
1502 .validate_thinking_config(Some(&ThinkingConfig::default_with_effort(
1503 agent_sdk_foundation::llm::Effort::High
1504 )))
1505 .is_ok(),
1506 "effort without adaptive must be accepted",
1507 );
1508 }
1509
1510 #[test]
1511 fn test_fable_5_rejects_budgeted_thinking() {
1512 let fable = AnthropicProvider::fable("test-api-key".to_string());
1516 let error = fable
1517 .validate_thinking_config(Some(&ThinkingConfig::new(10_000)))
1518 .unwrap_err();
1519 assert!(
1520 error.to_string().contains("ThinkingConfig::adaptive()"),
1521 "expected migration hint, got: {error}"
1522 );
1523 }
1524
1525 #[test]
1526 fn test_fable_5_accepts_adaptive_thinking() {
1527 let fable = AnthropicProvider::fable("test-api-key".to_string());
1528 assert!(
1529 fable
1530 .validate_thinking_config(Some(&ThinkingConfig::adaptive()))
1531 .is_ok()
1532 );
1533 assert!(
1534 fable
1535 .validate_thinking_config(Some(&ThinkingConfig::adaptive_with_effort(
1536 agent_sdk_foundation::llm::Effort::High
1537 )))
1538 .is_ok()
1539 );
1540 assert!(
1541 fable
1542 .validate_thinking_config(Some(&ThinkingConfig::default_with_effort(
1543 agent_sdk_foundation::llm::Effort::High
1544 )))
1545 .is_ok(),
1546 "effort without adaptive must be accepted",
1547 );
1548 }
1549
1550 #[test]
1551 fn test_fable_factory_creates_fable_5_provider() {
1552 let provider = AnthropicProvider::fable("test-api-key".to_string());
1553 assert_eq!(provider.model(), MODEL_FABLE_5);
1554 assert_eq!(provider.provider(), "anthropic");
1555 }
1556
1557 #[test]
1558 fn test_sonnet_factory_creates_sonnet_provider() {
1559 let provider = AnthropicProvider::sonnet("test-api-key".to_string());
1560
1561 assert_eq!(provider.model(), MODEL_SONNET_46);
1562 assert_eq!(provider.provider(), "anthropic");
1563 }
1564
1565 #[test]
1566 fn test_sonnet_45_factory_creates_sonnet_provider() {
1567 let provider = AnthropicProvider::sonnet_45("test-api-key".to_string());
1568
1569 assert_eq!(provider.model(), MODEL_SONNET_45);
1570 assert_eq!(provider.provider(), "anthropic");
1571 }
1572
1573 #[test]
1574 fn test_sonnet_46_factory_creates_sonnet_provider() {
1575 let provider = AnthropicProvider::sonnet_46("test-api-key".to_string());
1576
1577 assert_eq!(provider.model(), MODEL_SONNET_46);
1578 assert_eq!(provider.provider(), "anthropic");
1579 }
1580
1581 #[test]
1582 fn test_opus_factory_creates_opus_provider() {
1583 let provider = AnthropicProvider::opus("test-api-key".to_string());
1584
1585 assert_eq!(provider.model(), MODEL_OPUS_46);
1586 assert_eq!(provider.provider(), "anthropic");
1587 }
1588
1589 #[test]
1594 fn test_model_constants_have_expected_values() {
1595 assert!(MODEL_HAIKU_35.contains("haiku"));
1596 assert!(MODEL_SONNET_35.contains("sonnet"));
1597 assert!(MODEL_SONNET_4.contains("sonnet"));
1598 assert!(MODEL_SONNET_46.contains("sonnet"));
1599 assert!(MODEL_OPUS_4.contains("opus"));
1600 }
1601
1602 #[test]
1607 fn test_provider_is_cloneable() {
1608 let provider = AnthropicProvider::new("test-api-key", "test-model");
1609 let cloned = provider.clone();
1610
1611 assert_eq!(provider.model(), cloned.model());
1612 assert_eq!(provider.provider(), cloned.provider());
1613 }
1614
1615 fn tool(name: &str) -> agent_sdk_foundation::llm::Tool {
1620 agent_sdk_foundation::llm::Tool {
1621 name: name.to_string(),
1622 description: "desc".to_string(),
1623 input_schema: serde_json::json!({ "type": "object" }),
1624 display_name: name.to_string(),
1625 tier: agent_sdk_foundation::ToolTier::Observe,
1626 }
1627 }
1628
1629 fn request_with_tools(tools: Vec<agent_sdk_foundation::llm::Tool>) -> ChatRequest {
1630 ChatRequest {
1631 system: String::new(),
1632 messages: vec![agent_sdk_foundation::llm::Message::user("hi")],
1633 tools: Some(tools),
1634 max_tokens: 1024,
1635 max_tokens_explicit: true,
1636 session_id: None,
1637 cached_content: None,
1638 thinking: None,
1639 tool_choice: None,
1640 response_format: None,
1641 cache: None,
1642 }
1643 }
1644
1645 #[test]
1646 fn test_oauth_tool_name_collision_detects_case_variants() {
1647 let tools = vec![tool("task"), tool("Task")];
1648 let collision = oauth_tool_name_collision(Some(&tools));
1649 assert!(collision.is_some());
1650 }
1651
1652 #[test]
1653 fn test_oauth_tool_name_collision_allows_distinct_names() {
1654 let tools = vec![tool("read"), tool("write"), tool("Read_File")];
1655 assert!(oauth_tool_name_collision(Some(&tools)).is_none());
1656 assert!(oauth_tool_name_collision(None).is_none());
1657 }
1658
1659 #[tokio::test]
1660 async fn test_oauth_chat_rejects_case_colliding_tools() -> anyhow::Result<()> {
1661 let provider = AnthropicProvider::new("sk-ant-oat-test", MODEL_SONNET_45);
1663 assert!(provider.is_oauth());
1664 let request = request_with_tools(vec![tool("task"), tool("Task")]);
1665 let outcome = provider.chat(request).await?;
1666 match outcome {
1667 ChatOutcome::InvalidRequest(msg) => {
1668 assert!(msg.contains("collide case-insensitively"), "got: {msg}");
1669 }
1670 other => panic!("expected InvalidRequest, got {other:?}"),
1671 }
1672 Ok(())
1673 }
1674
1675 #[tokio::test]
1676 async fn test_api_key_chat_does_not_apply_oauth_collision_gate() -> anyhow::Result<()> {
1677 let provider = AnthropicProvider::new("sk-ant-api-test", MODEL_SONNET_45);
1682 assert!(!provider.is_oauth());
1683 let tools = vec![tool("task"), tool("Task")];
1684 assert!(oauth_tool_name_collision(Some(&tools)).is_some());
1686 Ok(())
1687 }
1688
1689 fn apply_auth_beta_header(provider: &AnthropicProvider) -> anyhow::Result<Option<String>> {
1694 let builder = reqwest::Client::new().post("http://localhost/v1/messages");
1695 let request = provider.apply_auth(builder).build()?;
1696 Ok(request
1697 .headers()
1698 .get("anthropic-beta")
1699 .and_then(|value| value.to_str().ok())
1700 .map(str::to_owned))
1701 }
1702
1703 #[test]
1704 fn api_key_auth_sends_interleaved_beta_for_budget_thinking_models() -> anyhow::Result<()> {
1705 for provider in [
1708 AnthropicProvider::sonnet_45("test-key-not-a-secret"),
1709 AnthropicProvider::haiku("test-key-not-a-secret"),
1710 ] {
1711 assert!(!provider.is_oauth());
1712 assert_eq!(
1713 apply_auth_beta_header(&provider)?.as_deref(),
1714 Some("interleaved-thinking-2025-05-14"),
1715 "expected interleaved beta for {}",
1716 provider.model()
1717 );
1718 }
1719 Ok(())
1720 }
1721
1722 #[test]
1723 fn api_key_auth_omits_interleaved_beta_for_adaptive_models() -> anyhow::Result<()> {
1724 for provider in [
1727 AnthropicProvider::opus_48("test-key-not-a-secret"),
1728 AnthropicProvider::sonnet_5("test-key-not-a-secret"),
1729 AnthropicProvider::fable("test-key-not-a-secret"),
1730 ] {
1731 assert!(!provider.is_oauth());
1732 assert_eq!(
1733 apply_auth_beta_header(&provider)?,
1734 None,
1735 "expected no beta header for adaptive model {}",
1736 provider.model()
1737 );
1738 }
1739 Ok(())
1740 }
1741
1742 async fn captured_request_body(
1751 provider: &AnthropicProvider,
1752 request: ChatRequest,
1753 server: &wiremock::MockServer,
1754 ) -> serde_json::Value {
1755 use wiremock::matchers::{method, path};
1756 use wiremock::{Mock, ResponseTemplate};
1757
1758 Mock::given(method("POST"))
1759 .and(path("/v1/messages"))
1760 .respond_with(ResponseTemplate::new(200).set_body_string("{}"))
1761 .mount(server)
1762 .await;
1763
1764 let _ = provider.chat(request).await;
1765
1766 let received = server
1767 .received_requests()
1768 .await
1769 .expect("mock server records requests");
1770 assert_eq!(received.len(), 1, "expected exactly one request");
1771 serde_json::from_slice(&received[0].body).expect("request body is JSON")
1772 }
1773
1774 #[tokio::test]
1775 async fn forced_tool_drops_configured_thinking_on_the_wire() {
1776 let server = wiremock::MockServer::start().await;
1783 let provider = AnthropicProvider::sonnet_45("sk-ant-api-test")
1784 .with_thinking(ThinkingConfig::new(10_000))
1785 .with_base_url(server.uri());
1786
1787 let mut request = request_with_tools(vec![tool("respond")]);
1788 request.tool_choice = Some(agent_sdk_foundation::llm::ToolChoice::Tool(
1789 "respond".to_owned(),
1790 ));
1791
1792 let body = captured_request_body(&provider, request, &server).await;
1793
1794 assert!(
1795 body.get("thinking").is_none(),
1796 "thinking must be absent when a tool is forced, got: {body}"
1797 );
1798 assert_eq!(
1799 body["tool_choice"]["type"], "tool",
1800 "the forced tool_choice must survive, got: {body}"
1801 );
1802 }
1803
1804 #[tokio::test]
1805 async fn configured_thinking_survives_without_forced_tool() {
1806 let server = wiremock::MockServer::start().await;
1811 let provider = AnthropicProvider::sonnet_45("sk-ant-api-test")
1812 .with_thinking(ThinkingConfig::new(10_000))
1813 .with_base_url(server.uri());
1814
1815 let mut request = request_with_tools(vec![tool("read")]);
1816 request.tool_choice = Some(agent_sdk_foundation::llm::ToolChoice::Auto);
1817
1818 let body = captured_request_body(&provider, request, &server).await;
1819
1820 assert_eq!(
1821 body["thinking"]["type"], "enabled",
1822 "configured thinking must survive when no tool is forced, got: {body}"
1823 );
1824 }
1825
1826 #[tokio::test]
1827 async fn thinking_display_defaults_to_omitted_and_honors_the_override() {
1828 let server = wiremock::MockServer::start().await;
1829 let provider = AnthropicProvider::fable("sk-ant-api-test")
1830 .with_thinking(ThinkingConfig::adaptive())
1831 .with_base_url(server.uri());
1832 let body =
1833 captured_request_body(&provider, request_with_tools(vec![tool("read")]), &server).await;
1834 assert_eq!(body["thinking"]["display"], "omitted");
1835
1836 let server = wiremock::MockServer::start().await;
1837 let provider = AnthropicProvider::fable("sk-ant-api-test")
1838 .with_thinking(
1839 ThinkingConfig::adaptive()
1840 .with_display(agent_sdk_foundation::llm::ThinkingDisplay::Summarized),
1841 )
1842 .with_base_url(server.uri());
1843 let body =
1844 captured_request_body(&provider, request_with_tools(vec![tool("read")]), &server).await;
1845 assert_eq!(
1846 body["thinking"]["display"], "summarized",
1847 "the configured display must reach the wire, got: {body}"
1848 );
1849 }
1850
1851 #[tokio::test]
1852 async fn default_mode_effort_sends_output_config_without_thinking_block() {
1853 let server = wiremock::MockServer::start().await;
1854 let provider = AnthropicProvider::fable("sk-ant-api-test")
1855 .with_thinking(ThinkingConfig::default_with_effort(
1856 agent_sdk_foundation::llm::Effort::High,
1857 ))
1858 .with_base_url(server.uri());
1859
1860 let request = request_with_tools(vec![tool("read")]);
1861 let body = captured_request_body(&provider, request, &server).await;
1862
1863 assert!(
1864 body.get("thinking").is_none(),
1865 "default mode must not send a thinking block, got: {body}"
1866 );
1867 assert_eq!(
1868 body["output_config"]["effort"], "high",
1869 "the effort must reach output_config, got: {body}"
1870 );
1871 }
1872
1873 #[tokio::test]
1878 async fn streaming_headers_stall_yields_connection_lost_error() {
1879 use wiremock::matchers::{method, path};
1880 use wiremock::{Mock, MockServer, ResponseTemplate};
1881
1882 let server = MockServer::start().await;
1883 Mock::given(method("POST"))
1884 .and(path("/v1/messages"))
1885 .respond_with(ResponseTemplate::new(200).set_delay(std::time::Duration::from_secs(30)))
1886 .mount(&server)
1887 .await;
1888
1889 let provider = AnthropicProvider::sonnet_45("sk-ant-api-test")
1890 .with_base_url(server.uri())
1891 .with_stream_stall_timeouts(
1892 std::time::Duration::from_millis(100),
1893 std::time::Duration::from_secs(5),
1894 );
1895
1896 let items: Vec<_> = provider
1897 .chat_stream(request_with_tools(vec![]))
1898 .collect()
1899 .await;
1900
1901 assert_eq!(items.len(), 1, "a stalled send yields exactly one item");
1902 let Some(Ok(StreamDelta::Error { message, kind })) = items.first() else {
1903 panic!("headers stall must surface as a classified error: {items:?}")
1904 };
1905 assert_eq!(*kind, StreamErrorKind::ConnectionLost);
1906 assert!(
1907 message.contains("timed out awaiting response headers"),
1908 "message must name the headers stall: {message}"
1909 );
1910 }
1911
1912 #[tokio::test]
1917 async fn sse_byte_stall_yields_connection_lost_error() {
1918 use tokio::io::{AsyncReadExt, AsyncWriteExt};
1919
1920 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
1921 .await
1922 .expect("bind test listener");
1923 let addr = listener.local_addr().expect("listener addr");
1924 let server = tokio::spawn(async move {
1925 let (mut socket, _) = listener.accept().await.expect("accept");
1926 let mut buf = [0u8; 4096];
1928 let _ = socket.read(&mut buf).await;
1929 let headers = "HTTP/1.1 200 OK\r\ncontent-type: text/event-stream\r\ntransfer-encoding: chunked\r\n\r\n";
1930 socket.write_all(headers.as_bytes()).await.expect("headers");
1931 let ping = "event: ping\ndata: {\"type\": \"ping\"}\n\n";
1934 let chunk = format!("{:x}\r\n{ping}\r\n", ping.len());
1935 socket
1936 .write_all(chunk.as_bytes())
1937 .await
1938 .expect("ping chunk");
1939 socket.flush().await.expect("flush");
1940 std::future::pending::<()>().await;
1941 });
1942
1943 let provider = AnthropicProvider::sonnet_45("sk-ant-api-test")
1944 .with_base_url(format!("http://{addr}"))
1945 .with_stream_stall_timeouts(
1946 std::time::Duration::from_secs(5),
1947 std::time::Duration::from_millis(200),
1948 );
1949
1950 let items: Vec<_> = provider
1951 .chat_stream(request_with_tools(vec![]))
1952 .collect()
1953 .await;
1954
1955 let Some(Ok(StreamDelta::Error { message, kind })) = items.last() else {
1956 panic!("byte stall must surface as a classified error: {items:?}")
1957 };
1958 assert_eq!(*kind, StreamErrorKind::ConnectionLost);
1959 assert!(
1960 message.contains("no bytes for"),
1961 "message must name the byte stall: {message}"
1962 );
1963 server.abort();
1964 }
1965}