pub mod build;
mod client;
pub mod config;
pub mod error;
pub mod parse;
pub mod retry;
pub mod wire;
use std::sync::{Arc, Mutex, PoisonError};
use async_trait::async_trait;
pub use build::{build_request, count_cache_controls, normalize_input_schema};
pub use config::{ApiBackend, AuthScheme, DeveloperRendering, ModelConfig, ReasoningEncoding};
pub use error::{HttpFailure, classify, parse_retry_after};
pub use parse::response_to_completion;
pub use retry::{RetryPolicy, backoff, run_with_retry};
use crate::completion::Completion;
use crate::provider::{Provider, ProviderError};
use crate::repair::repair_pairing;
use crate::request::ConversationRequest;
pub trait AuthRefresh: Send + Sync {
fn refresh(&self) -> Option<AuthScheme>;
}
pub struct AnthropicProvider {
http: reqwest::Client,
config: ModelConfig,
retry: RetryPolicy,
auth: Mutex<AuthScheme>,
auth_refresh: Option<Arc<dyn AuthRefresh>>,
}
impl AnthropicProvider {
pub fn new(config: ModelConfig) -> Result<Self, ProviderError> {
Ok(Self {
http: crate::http::build_http_client()?,
auth: Mutex::new(config.auth.clone()),
config,
retry: RetryPolicy::default(),
auth_refresh: None,
})
}
pub fn from_env() -> Result<Self, ProviderError> {
Self::new(ModelConfig::from_env()?)
}
#[must_use]
pub fn with_retry_policy(mut self, retry: RetryPolicy) -> Self {
self.retry = retry;
self
}
#[must_use]
pub fn with_auth_refresh(mut self, refresh: Arc<dyn AuthRefresh>) -> Self {
self.auth_refresh = Some(refresh);
self
}
#[must_use]
pub fn config(&self) -> &ModelConfig {
&self.config
}
fn current_auth(&self) -> AuthScheme {
self.auth
.lock()
.unwrap_or_else(PoisonError::into_inner)
.clone()
}
fn try_refresh(&self) -> bool {
let Some(refresher) = &self.auth_refresh else {
return false;
};
let Some(new_auth) = refresher.refresh() else {
return false;
};
let mut current = self.auth.lock().unwrap_or_else(PoisonError::into_inner);
if *current == new_auth {
return false;
}
*current = new_auth;
true
}
async fn exchange(&self, request: &wire::MessagesRequest) -> Result<Completion, ProviderError> {
let auth = self.current_auth();
run_with_retry(&self.retry, |_attempt| {
client::send_once(&self.http, &self.config, &auth, request)
})
.await
}
}
#[async_trait]
impl Provider for AnthropicProvider {
#[allow(clippy::unnecessary_literal_bound)] fn api_schema(&self) -> &str {
"anthropic"
}
async fn complete(&self, request: &ConversationRequest) -> Result<Completion, ProviderError> {
if let Some(crate::request::ReasoningEffort::Other(tier)) =
&request.sampling_args.reasoning_effort
&& self.config.reasoning_encoding == config::ReasoningEncoding::Budget
{
return Err(ProviderError::Config(format!(
"unknown reasoning-effort tier {tier:?} has no budget mapping on the \
anthropic wire (Budget encoding)"
)));
}
let mut repaired = request.clone();
let _ = repair_pairing(&mut repaired.messages);
let wire_request = build_request(&repaired, &self.config);
match self.exchange(&wire_request).await {
Err(ProviderError::Auth(message)) => {
if self.try_refresh() {
self.exchange(&wire_request).await
} else {
Err(ProviderError::Auth(message))
}
}
other => other,
}
}
}