mod stream;
use self::stream::AnthropicStreamParser;
pub(super) fn parse_captured_sse(
bytes: &[u8],
) -> anyhow::Result<Vec<crate::providers::ProviderEvent>> {
let text = std::str::from_utf8(bytes)?;
let mut parser = AnthropicStreamParser::default();
let events = parser.push_chunk_outcome(text)?.events;
parser.finish()?;
if !matches!(
parser.stop_reason(),
Some("end_turn" | "tool_use" | "max_tokens" | "refusal")
) {
anyhow::bail!("Claude upstream missing supported stop reason");
}
let has_input = events.iter().any(|event| {
matches!(event,
crate::providers::ProviderEvent::UsageObserved(observation) if observation.presence.input)
});
let has_output = events.iter().any(|event| {
matches!(event,
crate::providers::ProviderEvent::UsageObserved(observation) if observation.presence.output)
});
if !(has_input && has_output) {
anyhow::bail!("Claude upstream missing complete token usage");
}
Ok(events)
}
use crate::tools::mvp_tool_definitions_json_with_subagents;
use crate::{
cancellation::AgentCancellation,
config::AnthropicCacheTtl,
model_catalog::types::ModelCatalogEntry,
providers::{
ANTHROPIC_PROVIDER, HttpRequest, HttpTransport, MessageRole, ProviderConversationItem,
ProviderEvent, ProviderRequest, ProviderToolResult,
openai::{
fetch_model_catalog_response_text_cancellable, package_user_agent,
response_instructions,
},
openai_stream::stream_with_transport_parser,
replay_trace::ReplayDropTrace,
},
thinking::ThinkingLevel,
};
use serde_json::{Value, json};
use std::{
collections::{BTreeMap, HashSet},
fmt,
sync::{Arc, OnceLock},
};
pub(crate) const ANTHROPIC_MESSAGES_URL: &str = "https://api.anthropic.com/v1/messages";
pub(crate) const ANTHROPIC_MODELS_URL: &str = "https://api.anthropic.com/v1/models";
pub(crate) const ANTHROPIC_VERSION: &str = "2023-06-01";
pub(crate) const DEFAULT_MAX_TOKENS: u64 = 4096;
#[derive(Clone)]
pub struct AnthropicProvider<T> {
model: String,
api_key: String,
transport: T,
cache_ttl: Option<AnthropicCacheTtl>,
max_output_tokens: Option<u64>,
thinking_level: ThinkingLevel,
}
impl<T: fmt::Debug> fmt::Debug for AnthropicProvider<T> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("AnthropicProvider")
.field("model", &self.model)
.field("api_key", &"<redacted>")
.field("transport", &self.transport)
.field("cache_ttl", &self.cache_ttl)
.field("max_output_tokens", &self.max_output_tokens)
.field("thinking_level", &self.thinking_level)
.finish()
}
}
impl<T> AnthropicProvider<T> {
pub fn new(model: impl Into<String>, api_key: impl Into<String>, transport: T) -> Self {
Self {
model: model.into(),
api_key: api_key.into(),
transport,
cache_ttl: None,
max_output_tokens: None,
thinking_level: ThinkingLevel::Default,
}
}
pub fn with_cache_ttl(mut self, cache_ttl: Option<AnthropicCacheTtl>) -> Self {
self.cache_ttl = cache_ttl;
self
}
pub fn with_max_output_tokens(mut self, max: Option<u64>) -> Self {
self.max_output_tokens = max;
self
}
pub fn with_thinking_level(mut self, level: ThinkingLevel) -> Self {
self.thinking_level = level;
self
}
pub fn build_http_request(&self, request: &ProviderRequest) -> HttpRequest {
HttpRequest {
method: "POST".to_string(),
url: ANTHROPIC_MESSAGES_URL.to_string(),
headers: anthropic_headers(&self.api_key, "text/event-stream"),
body: Arc::new(anthropic_messages_body_with_cache_ttl(
&self.model,
request,
self.cache_ttl,
self.max_output_tokens,
self.thinking_level,
)),
}
}
pub fn build_model_catalog_request(&self) -> HttpRequest {
HttpRequest {
method: "GET".to_string(),
url: format!("{ANTHROPIC_MODELS_URL}?limit=1000"),
headers: anthropic_headers(&self.api_key, "application/json"),
body: Arc::new(Value::Null),
}
}
pub fn discover_model_catalog(&self) -> anyhow::Result<Vec<ModelCatalogEntry>> {
self.discover_model_catalog_cancellable(&AgentCancellation::default())
}
pub fn discover_model_catalog_cancellable(
&self,
cancellation: &AgentCancellation,
) -> anyhow::Result<Vec<ModelCatalogEntry>> {
let text = fetch_model_catalog_response_text_cancellable(
self.build_model_catalog_request(),
ANTHROPIC_PROVIDER,
cancellation,
)?;
parse_anthropic_model_catalog_response(&text)
}
}
impl<T: HttpTransport + Send + Sync> crate::providers::Provider for AnthropicProvider<T> {
fn stream_cancellable(
&self,
request: ProviderRequest,
cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
let semantic_progress_timeout = request.semantic_progress_timeout_or_default();
stream_with_transport_parser(
&self.transport,
self.build_http_request(&request),
cancellation,
semantic_progress_timeout,
AnthropicStreamParser::default,
on_event,
)
}
}
fn anthropic_headers(api_key: &str, accept: &str) -> BTreeMap<String, String> {
BTreeMap::from([
("accept".to_string(), accept.to_string()),
(
"anthropic-version".to_string(),
ANTHROPIC_VERSION.to_string(),
),
("content-type".to_string(), "application/json".to_string()),
("user-agent".to_string(), package_user_agent()),
("x-api-key".to_string(), api_key.to_string()),
])
}
pub(crate) fn anthropic_messages_body_with_cache_ttl(
model: &str,
request: &ProviderRequest,
cache_ttl: Option<AnthropicCacheTtl>,
max_output_tokens: Option<u64>,
thinking_level: ThinkingLevel,
) -> Value {
let mut state = AnthropicBodyBuildState::default();
for item in request.conversation_items_iter() {
state.append_conversation_item(item);
}
for message in &mut state.messages {
if message["role"] == "user"
&& let Some(content) = message["content"].as_array_mut()
{
content.sort_by_key(|block| block["type"] != "tool_result");
}
}
state.trace_dropped_replay_items();
let max_tokens = max_output_tokens.unwrap_or(DEFAULT_MAX_TOKENS);
let mut body = json!({
"model": model,
"max_tokens": max_tokens,
"stream": request.stream,
"messages": state.messages,
});
match anthropic_thinking_mode(model) {
AnthropicThinkingMode::Manual => {
if let Some(budget_tokens) =
anthropic_thinking_budget_tokens(thinking_level, max_tokens)
{
body["thinking"] = json!({"type":"enabled","budget_tokens": budget_tokens});
}
}
AnthropicThinkingMode::Adaptive => {
if let Some(effort) = thinking_level.explicit_effort() {
body["thinking"] = json!({"type":"adaptive"});
body["output_config"] = json!({"effort": effort});
}
}
}
if let Some(cache_ttl) = cache_ttl {
body["cache_control"] = anthropic_cache_control_json(cache_ttl);
}
let system = response_instructions(request);
if !system.is_empty() {
body["system"] = json!(system);
}
if let Some(tool_definitions) = anthropic_tools_for_request(request) {
body["tools"] = tool_definitions;
}
body
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum AnthropicThinkingMode {
Manual,
Adaptive,
}
fn anthropic_thinking_mode(model: &str) -> AnthropicThinkingMode {
let adaptive = model.starts_with("claude-fable-")
|| model.starts_with("claude-mythos-")
|| anthropic_model_version_at_least(model, "claude-sonnet-", 5, 0)
|| anthropic_model_version_at_least(model, "claude-opus-", 4, 7);
if adaptive {
AnthropicThinkingMode::Adaptive
} else {
AnthropicThinkingMode::Manual
}
}
fn anthropic_model_version_at_least(
model: &str,
family_prefix: &str,
minimum_major: u64,
minimum_minor: u64,
) -> bool {
let Some(version) = model.strip_prefix(family_prefix) else {
return false;
};
let mut parts = version.split('-');
let Some(major) = parts.next().and_then(|part| part.parse::<u64>().ok()) else {
return false;
};
let minor = parts
.next()
.and_then(|part| part.parse::<u64>().ok())
.unwrap_or(0);
(major, minor) >= (minimum_major, minimum_minor)
}
fn anthropic_thinking_budget_tokens(level: ThinkingLevel, max_tokens: u64) -> Option<u64> {
let budget = match level {
ThinkingLevel::High => 16_384,
ThinkingLevel::Max => 32_768,
_ => return None,
};
(max_tokens > 1).then_some(budget.min(max_tokens.saturating_sub(1)))
}
fn anthropic_cache_control_json(cache_ttl: AnthropicCacheTtl) -> Value {
match cache_ttl {
AnthropicCacheTtl::FiveMinutes => json!({"type":"ephemeral"}),
AnthropicCacheTtl::OneHour => json!({"type":"ephemeral","ttl":"1h"}),
}
}
struct AnthropicBodyBuildState {
messages: Vec<Value>,
seen_tool_calls: HashSet<String>,
seen_tool_results: HashSet<String>,
replay_drop_trace: ReplayDropTrace,
}
impl Default for AnthropicBodyBuildState {
fn default() -> Self {
Self {
messages: Vec::new(),
seen_tool_calls: HashSet::new(),
seen_tool_results: HashSet::new(),
replay_drop_trace: ReplayDropTrace::new("anthropic"),
}
}
}
impl AnthropicBodyBuildState {
fn append_conversation_item(&mut self, item: &ProviderConversationItem) {
match item {
ProviderConversationItem::ReasoningSelection { .. } => {}
ProviderConversationItem::Message(message) if message.role != MessageRole::System => {
self.push_message(
message.role.as_api_str(),
vec![json!({"type":"text","text": message.content})],
);
}
ProviderConversationItem::Message(_) => {}
ProviderConversationItem::ResponseItem(item) => self.append_response_item(item),
ProviderConversationItem::ToolResult(result) => self.append_tool_result(result),
ProviderConversationItem::LegacyReplayNote {
event_type,
content,
} => {
self.push_message(
"user",
vec![json!({
"type":"text",
"text": ProviderConversationItem::legacy_note_text(event_type, content),
})],
);
}
}
}
fn append_response_item(&mut self, item: &Value) {
if let Some(block) = anthropic_thinking_block(item) {
self.push_message("assistant", vec![block]);
return;
}
if item.get("type").and_then(Value::as_str) == Some("function_call") {
let Some(call_id) = item.get("call_id").and_then(Value::as_str) else {
self.drop_replay_item("function_call_missing_call_id");
return;
};
let Some(name) = item.get("name").and_then(Value::as_str) else {
self.drop_replay_item("function_call_missing_name");
return;
};
if !self.seen_tool_calls.insert(call_id.to_string()) {
return;
}
self.push_message(
"assistant",
vec![json!({
"type":"tool_use",
"id": call_id,
"name": name,
"input": anthropic_tool_input(item.get("arguments").unwrap_or(&Value::Null)),
})],
);
return;
}
if item.get("type").and_then(Value::as_str) == Some("function_call_output") {
let Some(call_id) = item.get("call_id").and_then(Value::as_str) else {
self.drop_replay_item("function_call_output_missing_call_id");
return;
};
if !self.seen_tool_results.insert(call_id.to_string()) {
return;
}
self.push_message(
"user",
vec![json!({
"type":"tool_result",
"tool_use_id": call_id,
"content": text_from_value(item.get("output")),
})],
);
return;
}
let Some(role) = item.get("role").and_then(Value::as_str) else {
self.drop_replay_item("message_missing_role");
return;
};
if !matches!(role, "user" | "assistant") {
self.drop_replay_item("message_unsupported_role");
return;
}
self.push_message(
role,
vec![json!({
"type":"text",
"text": text_from_value(item.get("content")),
})],
);
}
fn append_tool_result(&mut self, result: &ProviderToolResult) {
if !self.seen_tool_results.insert(result.call_id.clone()) {
return;
}
self.push_message(
"user",
vec![json!({
"type":"tool_result",
"tool_use_id": result.call_id,
"content": anthropic_tool_result_content(result),
"is_error": !result.success,
})],
);
}
fn push_message(&mut self, role: &str, mut content: Vec<Value>) {
if role == "tool" {
return;
}
if let Some(last) = self.messages.last_mut()
&& last.get("role").and_then(Value::as_str) == Some(role)
&& let Some(parts) = last.get_mut("content").and_then(Value::as_array_mut)
{
parts.append(&mut content);
return;
}
self.messages
.push(json!({"role": role, "content": content}));
}
fn drop_replay_item(&mut self, reason: &str) {
self.replay_drop_trace.drop_item(reason);
}
fn trace_dropped_replay_items(&self) {
self.replay_drop_trace.trace_summary();
}
}
fn anthropic_tool_result_content(result: &ProviderToolResult) -> String {
if !result.success && result.output.trim().is_empty() {
return "Tool failed with no output.".to_string();
}
result.output.clone()
}
fn anthropic_tool_input(arguments: &Value) -> Value {
match arguments {
Value::String(text) => serde_json::from_str(text).unwrap_or_else(|_| json!({})),
Value::Null => json!({}),
value => value.clone(),
}
}
fn anthropic_thinking_block(item: &Value) -> Option<Value> {
match item.get("type").and_then(Value::as_str) {
Some("thinking") => Some(json!({
"type":"thinking",
"thinking": item.get("thinking").and_then(Value::as_str).unwrap_or_default(),
"signature": item.get("signature").and_then(Value::as_str).unwrap_or_default(),
})),
Some("redacted_thinking") => Some(json!({
"type":"redacted_thinking",
"data": item.get("data").and_then(Value::as_str).unwrap_or_default(),
})),
_ => None,
}
}
fn text_from_value(value: Option<&Value>) -> String {
match value {
Some(Value::String(text)) => text.clone(),
Some(Value::Null) | None => " ".to_string(),
Some(value) => value.to_string(),
}
}
static ANTHROPIC_TOOLS_JSON: OnceLock<Value> = OnceLock::new();
static ANTHROPIC_TOOLS_WITHOUT_SUBAGENTS_JSON: OnceLock<Value> = OnceLock::new();
fn anthropic_tools_for_request(request: &ProviderRequest) -> Option<Value> {
if let Some(include_subagents) = request.static_tool_definitions_variant() {
let cache = if include_subagents {
&ANTHROPIC_TOOLS_JSON
} else {
&ANTHROPIC_TOOLS_WITHOUT_SUBAGENTS_JSON
};
return Some(
cache
.get_or_init(|| {
anthropic_tools_json(mvp_tool_definitions_json_with_subagents(
include_subagents,
))
})
.clone(),
);
}
request
.tool_definitions_json_if_enabled()
.map(anthropic_tools_json)
}
fn anthropic_tools_json(tool_definitions: Value) -> Value {
let Value::Array(tools) = tool_definitions else {
return Value::Array(Vec::new());
};
Value::Array(
tools
.into_iter()
.map(|mut tool| {
let (name, description, input_schema) = match tool.as_object_mut() {
Some(tool) => (
tool.remove("name").unwrap_or(Value::Null),
tool.remove("description").unwrap_or(Value::Null),
tool.remove("parameters").unwrap_or(Value::Null),
),
None => (Value::Null, Value::Null, Value::Null),
};
json!({
"name": name,
"description": description,
"input_schema": input_schema,
})
})
.collect(),
)
}
pub(crate) fn parse_anthropic_model_catalog_response(
text: &str,
) -> anyhow::Result<Vec<ModelCatalogEntry>> {
let value: Value = serde_json::from_str(text)?;
let data = value
.get("data")
.and_then(Value::as_array)
.ok_or_else(|| anyhow::anyhow!("anthropic model discovery response missing data array"))?;
let mut entries = Vec::new();
for item in data {
let Some(model) = item
.get("id")
.and_then(Value::as_str)
.map(str::trim)
.filter(|id| !id.is_empty())
else {
continue;
};
let mut entry = ModelCatalogEntry::new(ANTHROPIC_PROVIDER, model);
entry.reasoning_efforts = Some(ThinkingLevel::HIGH_MAX.to_vec());
entry.supports_reasoning = Some(true);
entry.display_name = item
.get("display_name")
.and_then(Value::as_str)
.filter(|value| !value.is_empty())
.map(str::to_string);
entry.context_window = item.get("max_input_tokens").and_then(Value::as_u64);
entry.max_output_tokens = item.get("max_tokens").and_then(Value::as_u64);
entries.push(entry);
}
if entries.is_empty() {
anyhow::bail!("anthropic model discovery returned no usable models");
}
Ok(entries)
}