locode_provider/anthropic/
mod.rs1pub mod build;
13mod client;
14pub mod config;
15pub mod error;
16pub mod parse;
17pub mod retry;
18pub mod wire;
19
20use std::sync::{Arc, Mutex, PoisonError};
21
22use async_trait::async_trait;
23
24pub use build::{build_request, count_cache_controls, normalize_input_schema};
25pub use config::{ApiBackend, AuthScheme, DeveloperRendering, ModelConfig, ReasoningEncoding};
26pub use error::{HttpFailure, classify, parse_retry_after};
27pub use parse::response_to_completion;
28pub use retry::{RetryPolicy, backoff, run_with_retry};
29
30use crate::completion::{Completion, CompletionDelta};
31use crate::provider::{Provider, ProviderError};
32use crate::repair::repair_pairing;
33use crate::request::ConversationRequest;
34
35pub trait AuthRefresh: Send + Sync {
42 fn refresh(&self) -> Option<AuthScheme>;
44}
45
46pub struct AnthropicProvider {
48 http: reqwest::Client,
49 config: ModelConfig,
50 retry: RetryPolicy,
51 auth: Mutex<AuthScheme>,
53 auth_refresh: Option<Arc<dyn AuthRefresh>>,
54}
55
56impl AnthropicProvider {
57 pub fn new(config: ModelConfig) -> Result<Self, ProviderError> {
62 Ok(Self {
63 http: crate::http::build_http_client()?,
64 auth: Mutex::new(config.auth.clone()),
65 config,
66 retry: RetryPolicy::default(),
67 auth_refresh: None,
68 })
69 }
70
71 pub fn from_env() -> Result<Self, ProviderError> {
78 Self::new(ModelConfig::from_env()?)
79 }
80
81 #[must_use]
83 pub fn with_retry_policy(mut self, retry: RetryPolicy) -> Self {
84 self.retry = retry;
85 self
86 }
87
88 #[must_use]
90 pub fn with_auth_refresh(mut self, refresh: Arc<dyn AuthRefresh>) -> Self {
91 self.auth_refresh = Some(refresh);
92 self
93 }
94
95 #[must_use]
97 pub fn config(&self) -> &ModelConfig {
98 &self.config
99 }
100
101 pub fn config_mut(&mut self) -> &mut ModelConfig {
104 &mut self.config
105 }
106
107 fn current_auth(&self) -> AuthScheme {
108 self.auth
109 .lock()
110 .unwrap_or_else(PoisonError::into_inner)
111 .clone()
112 }
113
114 fn try_refresh(&self) -> bool {
117 let Some(refresher) = &self.auth_refresh else {
118 return false;
119 };
120 let Some(new_auth) = refresher.refresh() else {
121 return false;
122 };
123 let mut current = self.auth.lock().unwrap_or_else(PoisonError::into_inner);
124 if *current == new_auth {
125 return false;
126 }
127 *current = new_auth;
128 true
129 }
130
131 async fn exchange(&self, request: &wire::MessagesRequest) -> Result<Completion, ProviderError> {
133 let auth = self.current_auth();
134 run_with_retry(&self.retry, |_attempt| {
135 client::send_once(&self.http, &self.config, &auth, request)
136 })
137 .await
138 }
139}
140
141#[async_trait]
142impl Provider for AnthropicProvider {
143 #[allow(clippy::unnecessary_literal_bound)] fn api_schema(&self) -> &str {
145 "anthropic"
146 }
147
148 async fn complete(&self, request: &ConversationRequest) -> Result<Completion, ProviderError> {
149 if let Some(crate::request::ReasoningEffort::Other(tier)) =
153 &request.sampling_args.reasoning_effort
154 && self.config.reasoning_encoding == config::ReasoningEncoding::Budget
155 {
156 return Err(ProviderError::Config(format!(
157 "unknown reasoning-effort tier {tier:?} has no budget mapping on the \
158 anthropic wire (Budget encoding)"
159 )));
160 }
161
162 let mut repaired = request.clone();
167 let _ = repair_pairing(&mut repaired.messages);
168
169 let wire_request = build_request(&repaired, &self.config);
171
172 match self.exchange(&wire_request).await {
175 Err(ProviderError::Auth(message)) => {
176 if self.try_refresh() {
177 self.exchange(&wire_request).await
178 } else {
179 Err(ProviderError::Auth(message))
180 }
181 }
182 other => other,
183 }
184 }
185
186 async fn stream(
187 &self,
188 request: &ConversationRequest,
189 on_delta: &mut (dyn FnMut(CompletionDelta) + Send),
190 ) -> Result<Completion, ProviderError> {
191 if let Some(crate::request::ReasoningEffort::Other(tier)) =
193 &request.sampling_args.reasoning_effort
194 && self.config.reasoning_encoding == config::ReasoningEncoding::Budget
195 {
196 return Err(ProviderError::Config(format!(
197 "unknown reasoning-effort tier {tier:?} has no budget mapping on the \
198 anthropic wire (Budget encoding)"
199 )));
200 }
201 let mut repaired = request.clone();
202 let _ = repair_pairing(&mut repaired.messages);
203 let mut wire_request = build_request(&repaired, &self.config);
204 wire_request.stream = Some(true);
205
206 let first = client::send_once_streaming(
212 &self.http,
213 &self.config,
214 &self.current_auth(),
215 &wire_request,
216 on_delta,
217 )
218 .await;
219 match first {
220 Ok(completion) => Ok(completion),
221 Err(failure)
222 if matches!(failure.error, ProviderError::Auth(_)) && self.try_refresh() =>
223 {
224 client::send_once_streaming(
225 &self.http,
226 &self.config,
227 &self.current_auth(),
228 &wire_request,
229 on_delta,
230 )
231 .await
232 .map_err(|f| f.error)
233 }
234 Err(failure) => Err(failure.error),
235 }
236 }
237}