266 lines
9.1 KiB
Rust
266 lines
9.1 KiB
Rust
use anyhow::Result;
|
|
use async_stream::try_stream;
|
|
use async_trait::async_trait;
|
|
use futures::TryStreamExt;
|
|
use reqwest::StatusCode;
|
|
use serde_json::Value;
|
|
use std::io;
|
|
use tokio::pin;
|
|
use tokio_util::io::StreamReader;
|
|
|
|
use super::api_client::{ApiClient, ApiResponse, AuthMethod};
|
|
use super::base::{ConfigKey, MessageStream, ModelInfo, Provider, ProviderMetadata, ProviderUsage};
|
|
use super::errors::ProviderError;
|
|
use super::formats::anthropic::{
|
|
create_request, get_usage, response_to_message, response_to_streaming_message,
|
|
};
|
|
use super::utils::{emit_debug_trace, get_model, map_http_error_to_provider_error};
|
|
use crate::conversation::message::Message;
|
|
use crate::impl_provider_default;
|
|
use crate::model::ModelConfig;
|
|
use crate::providers::retry::ProviderRetry;
|
|
use rmcp::model::Tool;
|
|
|
|
const ANTHROPIC_DEFAULT_MODEL: &str = "claude-sonnet-4-0";
|
|
const ANTHROPIC_KNOWN_MODELS: &[&str] = &[
|
|
"claude-sonnet-4-0",
|
|
"claude-sonnet-4-20250514",
|
|
"claude-opus-4-0",
|
|
"claude-opus-4-20250514",
|
|
"claude-3-7-sonnet-latest",
|
|
"claude-3-7-sonnet-20250219",
|
|
"claude-3-5-sonnet-latest",
|
|
"claude-3-5-haiku-latest",
|
|
"claude-3-opus-latest",
|
|
];
|
|
|
|
const ANTHROPIC_DOC_URL: &str = "https://docs.anthropic.com/en/docs/about-claude/models";
|
|
const ANTHROPIC_API_VERSION: &str = "2023-06-01";
|
|
|
|
#[derive(serde::Serialize)]
|
|
pub struct AnthropicProvider {
|
|
#[serde(skip)]
|
|
api_client: ApiClient,
|
|
model: ModelConfig,
|
|
}
|
|
|
|
impl_provider_default!(AnthropicProvider);
|
|
|
|
impl AnthropicProvider {
|
|
pub fn from_env(model: ModelConfig) -> Result<Self> {
|
|
let config = crate::config::Config::global();
|
|
let api_key: String = config.get_secret("ANTHROPIC_API_KEY")?;
|
|
let host: String = config
|
|
.get_param("ANTHROPIC_HOST")
|
|
.unwrap_or_else(|_| "https://api.anthropic.com".to_string());
|
|
|
|
let auth = AuthMethod::ApiKey {
|
|
header_name: "x-api-key".to_string(),
|
|
key: api_key,
|
|
};
|
|
|
|
let api_client =
|
|
ApiClient::new(host, auth)?.with_header("anthropic-version", ANTHROPIC_API_VERSION)?;
|
|
|
|
Ok(Self { api_client, model })
|
|
}
|
|
|
|
fn get_conditional_headers(&self) -> Vec<(&str, &str)> {
|
|
let mut headers = Vec::new();
|
|
|
|
let is_thinking_enabled = std::env::var("CLAUDE_THINKING_ENABLED").is_ok();
|
|
if self.model.model_name.starts_with("claude-3-7-sonnet-") {
|
|
if is_thinking_enabled {
|
|
headers.push(("anthropic-beta", "output-128k-2025-02-19"));
|
|
}
|
|
headers.push(("anthropic-beta", "token-efficient-tools-2025-02-19"));
|
|
}
|
|
|
|
headers
|
|
}
|
|
|
|
async fn post(&self, payload: &Value) -> Result<ApiResponse, ProviderError> {
|
|
let mut request = self.api_client.request("v1/messages");
|
|
|
|
for (key, value) in self.get_conditional_headers() {
|
|
request = request.header(key, value)?;
|
|
}
|
|
|
|
Ok(request.api_post(payload).await?)
|
|
}
|
|
|
|
fn anthropic_api_call_result(response: ApiResponse) -> Result<Value, ProviderError> {
|
|
match response.status {
|
|
StatusCode::OK => response.payload.ok_or_else(|| {
|
|
ProviderError::RequestFailed("Response body is not valid JSON".to_string())
|
|
}),
|
|
_ => {
|
|
if response.status == StatusCode::BAD_REQUEST {
|
|
if let Some(error_msg) = response
|
|
.payload
|
|
.as_ref()
|
|
.and_then(|p| p.get("error"))
|
|
.and_then(|e| e.get("message"))
|
|
.and_then(|m| m.as_str())
|
|
{
|
|
let msg = error_msg.to_string();
|
|
if msg.to_lowercase().contains("too long")
|
|
|| msg.to_lowercase().contains("too many")
|
|
{
|
|
return Err(ProviderError::ContextLengthExceeded(msg));
|
|
}
|
|
}
|
|
}
|
|
Err(map_http_error_to_provider_error(
|
|
response.status,
|
|
response.payload,
|
|
))
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
#[async_trait]
|
|
impl Provider for AnthropicProvider {
|
|
fn metadata() -> ProviderMetadata {
|
|
let models: Vec<ModelInfo> = ANTHROPIC_KNOWN_MODELS
|
|
.iter()
|
|
.map(|&model_name| ModelInfo::new(model_name, 200_000))
|
|
.collect();
|
|
|
|
ProviderMetadata::with_models(
|
|
"anthropic",
|
|
"Anthropic",
|
|
"Claude and other models from Anthropic",
|
|
ANTHROPIC_DEFAULT_MODEL,
|
|
models,
|
|
ANTHROPIC_DOC_URL,
|
|
vec![
|
|
ConfigKey::new("ANTHROPIC_API_KEY", true, true, None),
|
|
ConfigKey::new(
|
|
"ANTHROPIC_HOST",
|
|
true,
|
|
false,
|
|
Some("https://api.anthropic.com"),
|
|
),
|
|
],
|
|
)
|
|
}
|
|
|
|
fn get_model_config(&self) -> ModelConfig {
|
|
self.model.clone()
|
|
}
|
|
|
|
#[tracing::instrument(
|
|
skip(self, system, messages, tools),
|
|
fields(model_config, input, output, input_tokens, output_tokens, total_tokens)
|
|
)]
|
|
async fn complete(
|
|
&self,
|
|
system: &str,
|
|
messages: &[Message],
|
|
tools: &[Tool],
|
|
) -> Result<(Message, ProviderUsage), ProviderError> {
|
|
let payload = create_request(&self.model, system, messages, tools)?;
|
|
|
|
let response = self
|
|
.with_retry(|| async { self.post(&payload).await })
|
|
.await?;
|
|
|
|
let json_response = Self::anthropic_api_call_result(response)?;
|
|
|
|
let message = response_to_message(&json_response)?;
|
|
let usage = get_usage(&json_response)?;
|
|
tracing::debug!("🔍 Anthropic non-streaming parsed usage: input_tokens={:?}, output_tokens={:?}, total_tokens={:?}",
|
|
usage.input_tokens, usage.output_tokens, usage.total_tokens);
|
|
|
|
let model = get_model(&json_response);
|
|
emit_debug_trace(&self.model, &payload, &json_response, &usage);
|
|
let provider_usage = ProviderUsage::new(model, usage);
|
|
tracing::debug!(
|
|
"🔍 Anthropic non-streaming returning ProviderUsage: {:?}",
|
|
provider_usage
|
|
);
|
|
Ok((message, provider_usage))
|
|
}
|
|
|
|
async fn fetch_supported_models(&self) -> Result<Option<Vec<String>>, ProviderError> {
|
|
let response = self.api_client.api_get("v1/models").await?;
|
|
|
|
if response.status != StatusCode::OK {
|
|
return Err(map_http_error_to_provider_error(
|
|
response.status,
|
|
response.payload,
|
|
));
|
|
}
|
|
|
|
let json = response.payload.unwrap_or_default();
|
|
let arr = match json.get("models").and_then(|v| v.as_array()) {
|
|
Some(arr) => arr,
|
|
None => return Ok(None),
|
|
};
|
|
|
|
let mut models: Vec<String> = arr
|
|
.iter()
|
|
.filter_map(|m| {
|
|
if let Some(s) = m.as_str() {
|
|
Some(s.to_string())
|
|
} else if let Some(obj) = m.as_object() {
|
|
obj.get("id").and_then(|v| v.as_str()).map(str::to_string)
|
|
} else {
|
|
None
|
|
}
|
|
})
|
|
.collect();
|
|
models.sort();
|
|
Ok(Some(models))
|
|
}
|
|
|
|
async fn stream(
|
|
&self,
|
|
system: &str,
|
|
messages: &[Message],
|
|
tools: &[Tool],
|
|
) -> Result<MessageStream, ProviderError> {
|
|
let mut payload = create_request(&self.model, system, messages, tools)?;
|
|
payload
|
|
.as_object_mut()
|
|
.unwrap()
|
|
.insert("stream".to_string(), Value::Bool(true));
|
|
|
|
let mut request = self.api_client.request("v1/messages");
|
|
|
|
for (key, value) in self.get_conditional_headers() {
|
|
request = request.header(key, value)?;
|
|
}
|
|
|
|
let response = request.response_post(&payload).await?;
|
|
if !response.status().is_success() {
|
|
let status = response.status();
|
|
let error_text = response.text().await.unwrap_or_default();
|
|
let error_json = serde_json::from_str::<Value>(&error_text).ok();
|
|
return Err(map_http_error_to_provider_error(status, error_json));
|
|
}
|
|
|
|
let stream = response.bytes_stream().map_err(io::Error::other);
|
|
|
|
let model_config = self.model.clone();
|
|
Ok(Box::pin(try_stream! {
|
|
let stream_reader = StreamReader::new(stream);
|
|
let framed = tokio_util::codec::FramedRead::new(stream_reader, tokio_util::codec::LinesCodec::new()).map_err(anyhow::Error::from);
|
|
|
|
let message_stream = response_to_streaming_message(framed);
|
|
pin!(message_stream);
|
|
while let Some(message) = futures::StreamExt::next(&mut message_stream).await {
|
|
let (message, usage) = message.map_err(|e| ProviderError::RequestFailed(format!("Stream decode error: {}", e)))?;
|
|
emit_debug_trace(&model_config, &payload, &message, &usage.as_ref().map(|f| f.usage).unwrap_or_default());
|
|
yield (message, usage);
|
|
}
|
|
}))
|
|
}
|
|
|
|
fn supports_streaming(&self) -> bool {
|
|
true
|
|
}
|
|
}
|