use crate::{
agent::cancellation::AgentCancellation,
config::AnthropicCacheTtl,
model_catalog::ModelCatalogEntry,
providers::{
ANTHROPIC_PROVIDER, HttpRequest, HttpTransport, MessageRole, ProviderConversationItem,
ProviderEvent, ProviderRequest, ProviderToolResult, ToolCall, Usage,
error::{
ProviderError, ProviderStreamTrace, ProviderStreamTraceEvent,
ProviderStreamTracePendingTool, ProviderStreamTraceUsage,
},
openai::{
fetch_model_catalog_response_text_cancellable, package_user_agent,
response_instructions,
},
openai_stream::stream_with_transport_parser,
replay_trace::ReplayDropTrace,
sse::{diagnostic_snippet, next_sse_event_boundary, sse_data},
},
thinking::ThinkingLevel,
};
use serde_json::{Value, json};
use sha2::{Digest, Sha256};
use std::{
collections::{BTreeMap, HashSet, VecDeque},
fmt,
};
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;
const MAX_SSE_EVENT_BUFFER_BYTES: usize = 1024 * 1024;
const MAX_TOOL_ARGUMENT_BYTES: usize = 1024 * 1024;
const ANTHROPIC_STREAM_RECENT_EVENT_LIMIT: usize = 5;
const ANTHROPIC_STREAM_TRACE_EVENT_LIMIT: usize = 16;
const ANTHROPIC_STREAM_TRACE_PENDING_TOOL_LIMIT: usize = 16;
#[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: 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: 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(
model: &str,
request: &ProviderRequest,
max_output_tokens: Option<u64>,
) -> Value {
anthropic_messages_body_with_cache_ttl(
model,
request,
None,
max_output_tokens,
ThinkingLevel::Default,
)
}
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 {
anthropic_messages_body_with_system_prefix_inner(
model,
request,
None,
cache_ttl,
max_output_tokens,
thinking_level,
)
}
pub(crate) fn anthropic_messages_body_with_system_prefix(
model: &str,
request: &ProviderRequest,
system_prefix: Option<&str>,
cache_ttl: AnthropicCacheTtl,
max_output_tokens: Option<u64>,
thinking_level: ThinkingLevel,
) -> Value {
anthropic_messages_body_with_system_prefix_inner(
model,
request,
system_prefix,
Some(cache_ttl),
max_output_tokens,
thinking_level,
)
}
pub(crate) fn anthropic_messages_body_with_system_prefix_without_cache_control(
model: &str,
request: &ProviderRequest,
system_prefix: Option<&str>,
max_output_tokens: Option<u64>,
) -> Value {
anthropic_messages_body_with_system_prefix_inner(
model,
request,
system_prefix,
None,
max_output_tokens,
ThinkingLevel::Default,
)
}
fn anthropic_messages_body_with_system_prefix_inner(
model: &str,
request: &ProviderRequest,
system_prefix: Option<&str>,
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);
}
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,
});
if let Some(budget_tokens) = anthropic_thinking_budget_tokens(thinking_level, max_tokens) {
body["thinking"] = json!({"type":"enabled","budget_tokens": budget_tokens});
}
if let Some(cache_ttl) = cache_ttl {
body["cache_control"] = anthropic_cache_control_json(cache_ttl);
}
let system = response_instructions_with_prefix(request, system_prefix);
if !system.is_empty() {
body["system"] = json!(system);
}
if let Some(tool_definitions) = request.tool_definitions_json_if_enabled() {
body["tools"] = anthropic_tools_json(tool_definitions);
}
body
}
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"}),
}
}
fn response_instructions_with_prefix(
request: &ProviderRequest,
system_prefix: Option<&str>,
) -> String {
let base = response_instructions(request);
match (system_prefix, base.is_empty()) {
(Some(prefix), false) => format!("{prefix}\n\n{base}"),
(Some(prefix), true) => prefix.to_string(),
(None, _) => base,
}
}
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::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(),
}
}
fn anthropic_tools_json(tool_definitions: Value) -> Value {
let tools = tool_definitions
.as_array()
.cloned()
.unwrap_or_default()
.into_iter()
.map(|tool| {
json!({
"name": tool.get("name").cloned().unwrap_or(Value::Null),
"description": tool.get("description").cloned().unwrap_or(Value::Null),
"input_schema": tool.get("parameters").cloned().unwrap_or(Value::Null),
})
})
.collect::<Vec<_>>();
Value::Array(tools)
}
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.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)
}
#[derive(Default)]
pub(crate) struct AnthropicStreamParser {
event_buffer: String,
recent_event_types: VecDeque<String>,
trace_sequence: u64,
trace_events: VecDeque<ProviderStreamTraceEvent>,
message_delta_stop_reason: Option<String>,
pending_tools: BTreeMap<u64, PendingAnthropicTool>,
pending_thinking: BTreeMap<u64, PendingAnthropicThinking>,
usage: Usage,
saw_terminal_completion: bool,
emitted_done: bool,
}
#[derive(Default)]
struct PendingAnthropicTool {
id: String,
name: String,
arguments_text: String,
}
#[derive(Default)]
struct PendingAnthropicThinking {
block_type: AnthropicThinkingBlockType,
thinking: String,
signature: Option<String>,
data: Option<String>,
}
#[derive(Default)]
enum AnthropicThinkingBlockType {
#[default]
Thinking,
RedactedThinking,
}
impl AnthropicStreamParser {
pub(crate) fn push_chunk_outcome(
&mut self,
chunk: &str,
) -> anyhow::Result<crate::providers::stream::StreamParseOutcome> {
self.event_buffer.push_str(chunk);
let mut buffer = std::mem::take(&mut self.event_buffer);
let mut events = Vec::new();
let mut semantic_progress = false;
let mut unsafe_recovery_progress = false;
while let Some((boundary, boundary_len)) = next_sse_event_boundary(&buffer) {
let parsed = sse_data(&buffer[..boundary]).map(|data| self.parse_data_event(&data));
buffer.drain(..boundary + boundary_len);
if let Some(parsed) = parsed {
let parsed = parsed?;
semantic_progress |= parsed.semantic_progress || !parsed.events.is_empty();
unsafe_recovery_progress |= parsed.unsafe_recovery_progress;
events.extend(parsed.events);
}
}
self.event_buffer = buffer;
if self.event_buffer.len() > MAX_SSE_EVENT_BUFFER_BYTES {
anyhow::bail!(
"provider SSE event exceeded maximum buffered size of {MAX_SSE_EVENT_BUFFER_BYTES} bytes before a frame boundary"
);
}
Ok(crate::providers::stream::StreamParseOutcome {
events,
semantic_progress,
unsafe_recovery_progress,
})
}
pub(crate) fn finish(&mut self) -> anyhow::Result<Vec<ProviderEvent>> {
if !self.event_buffer.trim().is_empty() {
let raw_event = std::mem::take(&mut self.event_buffer);
anyhow::bail!(
"provider SSE stream ended with incomplete event buffer: {}",
diagnostic_snippet(&raw_event)
);
}
self.ensure_no_pending_tools("stream finish")?;
self.ensure_no_pending_thinking("stream finish")?;
if !self.saw_terminal_completion {
return Err(crate::providers::error::ProviderError::stream_terminal(
"missing provider stream completion before EOF",
)
.into());
}
Ok(Vec::new())
}
fn parse_data_event(
&mut self,
data: &str,
) -> anyhow::Result<crate::providers::stream::StreamParseOutcome> {
let value = serde_json::from_str::<Value>(data).map_err(|error| {
anyhow::anyhow!(
"malformed provider SSE data JSON: {error}: {}",
diagnostic_snippet(data)
)
})?;
let event_type = value
.get("type")
.and_then(Value::as_str)
.unwrap_or_default();
if !event_type.is_empty() {
self.recent_event_types.push_back(event_type.to_string());
while self.recent_event_types.len() > ANTHROPIC_STREAM_RECENT_EVENT_LIMIT {
self.recent_event_types.pop_front();
}
self.record_trace_event(&value, event_type);
}
if event_type == "ping" {
return Ok(Default::default());
}
if event_type == "error" {
anyhow::bail!(
"anthropic provider stream error: {}",
diagnostic_snippet(data)
);
}
let mut events = Vec::new();
let mut semantic_progress = false;
let mut unsafe_recovery_progress = false;
match event_type {
"message_start" | "message_delta" => {
if let Some(event) = self.parse_usage_event(&value) {
events.push(event);
}
}
"content_block_start" => {
match value.pointer("/content_block/type").and_then(Value::as_str) {
Some("tool_use") => {
let index = event_index(&value, "content_block_start")?;
if self.pending_tools.contains_key(&index) {
anyhow::bail!(
"anthropic provider stream duplicate active tool content block index {index}"
);
}
let id = value
.pointer("/content_block/id")
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| {
anyhow::anyhow!(
"anthropic provider tool_use content block missing non-empty id"
)
})?
.to_string();
let name = value
.pointer("/content_block/name")
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| anyhow::anyhow!("anthropic provider tool_use content block missing non-empty name"))?
.to_string();
self.pending_tools.insert(
index,
PendingAnthropicTool {
id,
name,
arguments_text: String::new(),
},
);
semantic_progress = true;
unsafe_recovery_progress = true;
}
Some("thinking") => {
let index = event_index(&value, "content_block_start thinking")?;
self.pending_thinking.insert(
index,
PendingAnthropicThinking {
block_type: AnthropicThinkingBlockType::Thinking,
thinking: value
.pointer("/content_block/thinking")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string(),
signature: value
.pointer("/content_block/signature")
.and_then(Value::as_str)
.map(str::to_string),
data: None,
},
);
semantic_progress = true;
}
Some("redacted_thinking") => {
let index = event_index(&value, "content_block_start redacted_thinking")?;
self.pending_thinking.insert(
index,
PendingAnthropicThinking {
block_type: AnthropicThinkingBlockType::RedactedThinking,
thinking: String::new(),
signature: None,
data: value
.pointer("/content_block/data")
.and_then(Value::as_str)
.map(str::to_string),
},
);
semantic_progress = true;
}
_ => {}
}
}
"content_block_delta" => match value.pointer("/delta/type").and_then(Value::as_str) {
Some("text_delta") => {
if let Some(text) = value.pointer("/delta/text").and_then(Value::as_str) {
events.push(ProviderEvent::TextDelta(text.to_string()));
semantic_progress = true;
}
}
Some("thinking_delta") => {
let index = event_index(&value, "content_block_delta thinking_delta")?;
if let Some(delta) = value.pointer("/delta/thinking").and_then(Value::as_str)
&& let Some(pending) = self.pending_thinking.get_mut(&index)
{
pending.thinking.push_str(delta);
semantic_progress = true;
}
}
Some("signature_delta") => {
let index = event_index(&value, "content_block_delta signature_delta")?;
if let Some(signature) =
value.pointer("/delta/signature").and_then(Value::as_str)
&& let Some(pending) = self.pending_thinking.get_mut(&index)
{
pending.signature = Some(signature.to_string());
semantic_progress = true;
}
}
Some("redacted_thinking_delta") => {
let index = event_index(&value, "content_block_delta redacted_thinking_delta")?;
if let Some(data) = value.pointer("/delta/data").and_then(Value::as_str)
&& let Some(pending) = self.pending_thinking.get_mut(&index)
{
pending.data.get_or_insert_with(String::new).push_str(data);
semantic_progress = true;
}
}
Some("input_json_delta") => {
let index = event_index(&value, "content_block_delta input_json_delta")?;
if value
.pointer("/delta/partial_json")
.and_then(Value::as_str)
.is_some()
{
unsafe_recovery_progress = true;
}
if let Some(delta) =
value.pointer("/delta/partial_json").and_then(Value::as_str)
&& let Some(pending) = self.pending_tools.get_mut(&index)
{
let next_len = pending.arguments_text.len().saturating_add(delta.len());
if next_len > MAX_TOOL_ARGUMENT_BYTES {
anyhow::bail!(
"provider tool call arguments exceeded maximum size of {MAX_TOOL_ARGUMENT_BYTES} bytes"
);
}
pending.arguments_text.push_str(delta);
semantic_progress = true;
}
}
_ => {}
},
"content_block_stop" => {
let index = event_index(&value, "content_block_stop")?;
if let Some(pending) = self.pending_thinking.remove(&index) {
let response_item = match pending.block_type {
AnthropicThinkingBlockType::Thinking => json!({
"type":"thinking",
"thinking": pending.thinking,
"signature": pending.signature.unwrap_or_default(),
}),
AnthropicThinkingBlockType::RedactedThinking => json!({
"type":"redacted_thinking",
"data": pending.data.unwrap_or_default(),
}),
};
events.push(ProviderEvent::ResponseItem(response_item));
semantic_progress = true;
} else if let Some(pending) = self.pending_tools.remove(&index) {
let arguments = parse_arguments_text(&pending.arguments_text)?;
let response_item = json!({
"type":"function_call",
"call_id": pending.id,
"name": pending.name,
"arguments": pending.arguments_text,
"status":"completed",
});
events.push(ProviderEvent::ResponseItem(response_item));
events.push(ProviderEvent::ToolCall(ToolCall {
id: pending.id,
name: pending.name,
arguments,
}));
semantic_progress = true;
}
}
"message_stop" => {
self.ensure_no_pending_tools("message_stop")?;
self.ensure_no_pending_thinking("message_stop")?;
self.saw_terminal_completion = true;
semantic_progress = true;
if !self.emitted_done {
self.emitted_done = true;
events.push(ProviderEvent::Done);
}
}
_ => {}
}
semantic_progress |= !events.is_empty();
Ok(crate::providers::stream::StreamParseOutcome {
events,
semantic_progress,
unsafe_recovery_progress,
})
}
fn parse_usage_event(&mut self, value: &Value) -> Option<ProviderEvent> {
let usage = value
.get("usage")
.or_else(|| value.pointer("/message/usage"))?;
let input = usage.get("input_tokens").and_then(Value::as_u64);
let output = usage.get("output_tokens").and_then(Value::as_u64);
let cache_read = usage.get("cache_read_input_tokens").and_then(Value::as_u64);
let cache_write = usage
.get("cache_creation_input_tokens")
.and_then(Value::as_u64);
if let Some(input) = input {
self.usage.input = input;
}
if let Some(output) = output {
self.usage.output = output;
}
if let Some(cache_read) = cache_read {
self.usage.cache_read = cache_read;
}
if let Some(cache_write) = cache_write {
self.usage.cache_write = cache_write;
}
self.usage.total = self
.usage
.input
.saturating_add(self.usage.output)
.saturating_add(self.usage.cache_read)
.saturating_add(self.usage.cache_write);
Some(if input.is_some() {
ProviderEvent::Usage(self.usage.clone())
} else {
ProviderEvent::UsagePartial(self.usage.clone())
})
}
fn ensure_no_pending_thinking(&self, context: &str) -> anyhow::Result<()> {
if self.pending_thinking.is_empty() {
return Ok(());
}
let indexes = self
.pending_thinking
.keys()
.map(u64::to_string)
.collect::<Vec<_>>()
.join(", ");
Err(ProviderError::stream_terminal(format!(
"anthropic provider stream {context} with incomplete thinking content block index(es): {indexes}"
))
.into())
}
fn ensure_no_pending_tools(&self, context: &str) -> anyhow::Result<()> {
if self.pending_tools.is_empty() {
return Ok(());
}
let indexes = self
.pending_tools
.keys()
.map(u64::to_string)
.collect::<Vec<_>>()
.join(", ");
let recent_events = self
.recent_event_types
.iter()
.map(String::as_str)
.collect::<Vec<_>>()
.join(", ");
let pending_tools = self
.pending_tools
.iter()
.map(|(index, tool)| {
format!(
"index {index} id {} name {} argument_bytes {}",
tool.id,
tool.name,
tool.arguments_text.len()
)
})
.collect::<Vec<_>>()
.join("; ");
Err(ProviderError::stream_terminal(format!(
"anthropic provider stream {context} with incomplete tool_use content block index(es): {indexes}; recent events: {recent_events}; pending tools: {pending_tools}"
))
.with_stream_trace(self.provider_stream_trace(context))
.into())
}
fn record_trace_event(&mut self, value: &Value, event_type: &str) {
self.trace_sequence = self.trace_sequence.saturating_add(1);
let usage = value
.get("usage")
.or_else(|| value.pointer("/message/usage"));
let message_delta_stop_reason = value
.pointer("/delta/stop_reason")
.and_then(Value::as_str)
.map(str::to_string);
if message_delta_stop_reason.is_some() {
self.message_delta_stop_reason = message_delta_stop_reason.clone();
}
let partial_json = value.pointer("/delta/partial_json").and_then(Value::as_str);
self.trace_events.push_back(ProviderStreamTraceEvent {
seq: self.trace_sequence,
event_type: event_type.to_string(),
index: value.get("index").and_then(Value::as_u64),
content_block_type: value
.pointer("/content_block/type")
.and_then(Value::as_str)
.map(str::to_string),
delta_type: value
.pointer("/delta/type")
.and_then(Value::as_str)
.map(str::to_string),
message_delta_stop_reason,
usage: usage.map(|usage| ProviderStreamTraceUsage {
input_tokens: usage.get("input_tokens").and_then(Value::as_u64),
output_tokens: usage.get("output_tokens").and_then(Value::as_u64),
cache_read_input_tokens: usage
.get("cache_read_input_tokens")
.and_then(Value::as_u64),
cache_creation_input_tokens: usage
.get("cache_creation_input_tokens")
.and_then(Value::as_u64),
}),
partial_json_bytes: partial_json.map(str::len),
partial_json_sha256: partial_json.map(sha256_hex),
});
while self.trace_events.len() > ANTHROPIC_STREAM_TRACE_EVENT_LIMIT {
self.trace_events.pop_front();
}
}
fn provider_stream_trace(&self, context: &str) -> ProviderStreamTrace {
let pending_tools = self
.pending_tools
.iter()
.take(ANTHROPIC_STREAM_TRACE_PENDING_TOOL_LIMIT)
.map(|(index, tool)| ProviderStreamTracePendingTool {
index: *index,
id: tool.id.clone(),
name: tool.name.clone(),
argument_bytes: tool.arguments_text.len(),
argument_sha256: sha256_hex(&tool.arguments_text),
})
.collect();
ProviderStreamTrace {
schema_version: 1,
provider: "anthropic".to_string(),
failure_context: context.to_string(),
message_delta_stop_reason: self.message_delta_stop_reason.clone(),
recent_events: self.trace_events.iter().cloned().collect(),
pending_tool_count: self.pending_tools.len(),
pending_tools_truncated: self.pending_tools.len()
> ANTHROPIC_STREAM_TRACE_PENDING_TOOL_LIMIT,
pending_tools,
}
}
}
fn sha256_hex(text: &str) -> String {
crate::hex::lower_hex(Sha256::digest(text.as_bytes()))
}
fn event_index(value: &Value, event_context: &str) -> anyhow::Result<u64> {
value.get("index").and_then(Value::as_u64).ok_or_else(|| {
anyhow::anyhow!("anthropic provider stream {event_context} missing required index")
})
}
impl crate::providers::openai_stream::ProviderStreamParser for AnthropicStreamParser {
fn push_chunk_outcome(
&mut self,
chunk: &str,
) -> anyhow::Result<crate::providers::stream::StreamParseOutcome> {
self.push_chunk_outcome(chunk)
}
fn finish(&mut self) -> anyhow::Result<Vec<ProviderEvent>> {
self.finish()
}
}
fn parse_arguments_text(text: &str) -> anyhow::Result<Value> {
if text.trim().is_empty() {
return Ok(json!({}));
}
serde_json::from_str(text).map_err(|error| {
anyhow::anyhow!(
"malformed non-empty provider tool call arguments: {error}: {}",
diagnostic_snippet(text)
)
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::providers::{
ChatMessage, Provider, ProviderToolResult, error::provider_stream_trace_from_error,
};
use std::{
collections::VecDeque,
sync::{Arc, Mutex},
};
enum ScriptStep {
Chunks(Vec<&'static str>),
Error(anyhow::Error),
}
struct ScriptedTransport {
steps: Mutex<VecDeque<ScriptStep>>,
attempts: Arc<Mutex<usize>>,
}
impl ScriptedTransport {
fn new(steps: Vec<ScriptStep>) -> Self {
Self {
steps: Mutex::new(steps.into()),
attempts: Arc::new(Mutex::new(0)),
}
}
fn attempts_handle(&self) -> Arc<Mutex<usize>> {
Arc::clone(&self.attempts)
}
}
impl HttpTransport for ScriptedTransport {
fn stream_json(
&self,
request: HttpRequest,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
self.stream_json_cancellable(request, &AgentCancellation::default(), on_chunk)
}
fn stream_json_cancellable(
&self,
_request: HttpRequest,
cancellation: &AgentCancellation,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
cancellation.check()?;
*self.attempts.lock().unwrap() += 1;
match self.steps.lock().unwrap().pop_front().unwrap() {
ScriptStep::Chunks(chunks) => {
for chunk in chunks {
cancellation.check()?;
on_chunk(chunk)?;
}
Ok(())
}
ScriptStep::Error(error) => Err(error),
}
}
fn stream_json_cancellable_with_semantic_deadline(
&self,
request: HttpRequest,
cancellation: &AgentCancellation,
_semantic_deadline: &std::sync::atomic::AtomicU64,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
self.stream_json_cancellable(request, cancellation, on_chunk)
}
}
fn retryable_503_error() -> anyhow::Error {
crate::providers::error::ProviderError::http_status(
503,
"provider request failed for https://api.anthropic.com/v1/messages with status 503 Service Unavailable: upstream overloaded",
)
.into()
}
#[derive(Default)]
struct CapturingTransport {
chunks: Vec<String>,
}
impl HttpTransport for CapturingTransport {
fn stream_json(
&self,
request: HttpRequest,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
self.stream_json_cancellable(request, &AgentCancellation::default(), on_chunk)
}
fn stream_json_cancellable(
&self,
_request: HttpRequest,
cancellation: &AgentCancellation,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
for chunk in &self.chunks {
cancellation.check()?;
on_chunk(chunk)?;
}
Ok(())
}
fn stream_json_cancellable_with_semantic_deadline(
&self,
request: HttpRequest,
cancellation: &AgentCancellation,
_semantic_deadline: &std::sync::atomic::AtomicU64,
on_chunk: &mut dyn FnMut(&str) -> anyhow::Result<()>,
) -> anyhow::Result<()> {
self.stream_json_cancellable(request, cancellation, on_chunk)
}
}
#[test]
fn anthropic_stream_trace_snapshot_on_pending_tool_message_stop() {
let mut parser = AnthropicStreamParser::default();
parser.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"toolu_trace_full_id\",\"name\":\"read\"}}\n\n").unwrap();
parser.push_chunk_outcome("data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"path\\\":\\\"/tmp\"}}\n\n").unwrap();
parser.push_chunk_outcome("data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"tool_use\"},\"usage\":{\"output_tokens\":7}}\n\n").unwrap();
let error = parser
.push_chunk_outcome("data: {\"type\":\"message_stop\"}\n\n")
.unwrap_err();
let trace = provider_stream_trace_from_error(&error).unwrap();
assert_eq!(trace.provider, "anthropic");
assert_eq!(trace.failure_context, "message_stop");
assert_eq!(trace.message_delta_stop_reason.as_deref(), Some("tool_use"));
assert_eq!(
trace
.recent_events
.iter()
.map(|event| event.event_type.as_str())
.collect::<Vec<_>>(),
vec![
"content_block_start",
"content_block_delta",
"message_delta",
"message_stop"
]
);
let pending = trace.pending_tools.first().unwrap();
assert_eq!(pending.index, 1);
assert_eq!(pending.id, "toolu_trace_full_id");
assert_eq!(pending.name, "read");
assert_eq!(pending.argument_bytes, 13);
assert_eq!(pending.argument_sha256.len(), 64);
assert!(
pending
.argument_sha256
.chars()
.all(|ch| ch.is_ascii_hexdigit() && !ch.is_ascii_uppercase())
);
let serialized = serde_json::to_string(&trace).unwrap();
assert!(!serialized.contains("/tmp"));
assert!(!serialized.contains(r#"{\"path\":\"/tmp"#));
}
#[test]
fn anthropic_stream_trace_bounds_recent_events_and_omits_raw_partial_json() {
let mut parser = AnthropicStreamParser::default();
for index in 0..20 {
parser
.push_chunk_outcome(&format!(
"data: {{\"type\":\"content_block_start\",\"index\":{index},\"content_block\":{{\"type\":\"tool_use\",\"id\":\"toolu_{index}\",\"name\":\"read\"}}}}\n\n"
))
.unwrap();
parser
.push_chunk_outcome(&format!(
"data: {{\"type\":\"content_block_delta\",\"index\":{index},\"delta\":{{\"type\":\"input_json_delta\",\"partial_json\":\"raw-/tmp-marker-{index}\"}}}}\n\n"
))
.unwrap();
}
let error = parser.finish().unwrap_err();
let trace = provider_stream_trace_from_error(&error).unwrap();
let serialized = serde_json::to_string(&trace).unwrap();
assert!(trace.recent_events.len() <= ANTHROPIC_STREAM_TRACE_EVENT_LIMIT);
assert_eq!(trace.pending_tool_count, 20);
assert_eq!(
trace.pending_tools.len(),
ANTHROPIC_STREAM_TRACE_PENDING_TOOL_LIMIT
);
assert!(trace.pending_tools_truncated);
assert!(!serialized.contains("raw-/tmp-marker"));
assert!(!serialized.contains("/tmp"));
assert!(serialized.contains("partial_json_sha256"));
}
#[test]
fn anthropic_request_encodes_high_and_max_thinking_with_budget_below_max_tokens() {
let high = AnthropicProvider::new("claude-test", "secret", CapturingTransport::default())
.with_max_output_tokens(Some(20_000))
.with_thinking_level(ThinkingLevel::High)
.build_http_request(&ProviderRequest::new(
"ignored",
vec![ChatMessage::user("hello")],
));
assert_eq!(
high.body["thinking"],
json!({"type":"enabled","budget_tokens":16_384})
);
let max = AnthropicProvider::new("claude-test", "secret", CapturingTransport::default())
.with_max_output_tokens(Some(16_000))
.with_thinking_level(ThinkingLevel::Max)
.build_http_request(&ProviderRequest::new(
"ignored",
vec![ChatMessage::user("hello")],
));
assert_eq!(max.body["thinking"]["budget_tokens"], 15_999);
}
#[test]
fn anthropic_request_omits_thinking_for_default_level() {
let provider =
AnthropicProvider::new("claude-test", "secret", CapturingTransport::default())
.with_thinking_level(ThinkingLevel::Default);
let http = provider.build_http_request(&ProviderRequest::new(
"ignored",
vec![ChatMessage::user("hello")],
));
assert!(http.body.get("thinking").is_none());
}
#[test]
fn anthropic_request_construction_targets_messages_endpoint() {
let provider =
AnthropicProvider::new("claude-test", "secret", CapturingTransport::default());
let http = provider.build_http_request(&ProviderRequest::new(
"ignored",
vec![ChatMessage::user("hello")],
));
assert_eq!(http.method, "POST");
assert_eq!(http.url, ANTHROPIC_MESSAGES_URL);
assert_eq!(http.body["model"], "claude-test");
assert_eq!(http.body["max_tokens"], DEFAULT_MAX_TOKENS);
assert_eq!(http.body["stream"], true);
}
#[test]
fn anthropic_request_uses_configured_max_output_tokens_when_provided() {
let provider =
AnthropicProvider::new("claude-test", "secret", CapturingTransport::default())
.with_max_output_tokens(Some(8192));
let http = provider.build_http_request(&ProviderRequest::new(
"ignored",
vec![ChatMessage::user("hello")],
));
assert_eq!(http.body["max_tokens"], 8192);
}
#[test]
fn anthropic_request_falls_back_to_default_max_tokens_when_none() {
let provider =
AnthropicProvider::new("claude-test", "secret", CapturingTransport::default());
let http = provider.build_http_request(&ProviderRequest::new(
"ignored",
vec![ChatMessage::user("hello")],
));
assert_eq!(http.body["max_tokens"], DEFAULT_MAX_TOKENS);
}
#[test]
fn anthropic_request_default_omits_cache_control() {
let body = anthropic_messages_body(
"claude-test",
&ProviderRequest::new("ignored", vec![ChatMessage::user("hello")]),
None,
);
assert!(body.get("cache_control").is_none());
}
#[test]
fn anthropic_request_with_system_prefix_includes_cache_control() {
let body = anthropic_messages_body_with_system_prefix(
"claude-test",
&ProviderRequest::new(
"ignored",
vec![ChatMessage::system("base sys"), ChatMessage::user("hello")],
),
Some("prefix sys"),
AnthropicCacheTtl::FiveMinutes,
None,
ThinkingLevel::Default,
);
assert_eq!(body["cache_control"], json!({"type":"ephemeral"}));
assert_eq!(body["system"], "prefix sys\n\nbase sys");
}
#[test]
fn anthropic_request_one_hour_cache_control_emits_ttl() {
let provider =
AnthropicProvider::new("claude-test", "secret", CapturingTransport::default())
.with_cache_ttl(Some(AnthropicCacheTtl::OneHour));
let http = provider.build_http_request(&ProviderRequest::new(
"ignored",
vec![ChatMessage::user("hello")],
));
assert_eq!(
http.body["cache_control"],
json!({"type":"ephemeral","ttl":"1h"})
);
assert!(!http.headers.contains_key("anthropic-beta"));
}
#[test]
fn anthropic_request_uses_x_api_key_and_anthropic_version_headers() {
let provider =
AnthropicProvider::new("claude-test", "secret", CapturingTransport::default());
let http = provider.build_http_request(&ProviderRequest::new(
"ignored",
vec![ChatMessage::user("hello")],
));
assert_eq!(http.headers["x-api-key"], "secret");
assert_eq!(http.headers["anthropic-version"], ANTHROPIC_VERSION);
assert_eq!(http.headers["accept"], "text/event-stream");
assert!(!http.headers.contains_key("authorization"));
}
#[test]
fn anthropic_request_moves_system_to_top_level_system() {
let body = anthropic_messages_body(
"claude-test",
&ProviderRequest::new(
"ignored",
vec![ChatMessage::system("sys"), ChatMessage::user("hello")],
),
None,
);
assert_eq!(body["system"], "sys");
assert_eq!(body["messages"][0]["role"], "user");
}
#[test]
fn anthropic_request_converts_tools_to_input_schema() {
let body = anthropic_messages_body(
"claude-test",
&ProviderRequest::new("ignored", vec![ChatMessage::user("hello")]),
None,
);
assert_eq!(body["tools"][0]["name"], "read");
assert!(body["tools"][0].get("input_schema").is_some());
assert!(body["tools"][0].get("parameters").is_none());
}
#[test]
fn anthropic_request_appends_dynamic_mcp_tools() {
let request = ProviderRequest::new("ignored", vec![ChatMessage::user("hello")])
.with_dynamic_tool_definitions(vec![json!({
"type":"function",
"name":"mcp__mock__echo",
"description":"Echo",
"parameters":{"type":"object","properties":{"text":{"type":"string"}}}
})]);
let body = anthropic_messages_body("claude-test", &request, None);
let tool = body["tools"].as_array().unwrap().last().unwrap();
assert_eq!(tool["name"], "mcp__mock__echo");
assert_eq!(tool["input_schema"]["properties"]["text"]["type"], "string");
assert!(tool.get("parameters").is_none());
}
#[test]
fn anthropic_request_converts_canonical_function_call_and_tool_result() {
let request = ProviderRequest::new("ignored", vec![ChatMessage::user("hello")])
.with_response_items(vec![json!({"type":"function_call","call_id":"call_1","name":"read","arguments":"{\"path\":\"a.txt\"}"})])
.with_tool_results(vec![ProviderToolResult { call_id: "call_1".into(), tool_name: "read".into(), success: false, output: "boom".into() }]);
let body = anthropic_messages_body("claude-test", &request, None);
assert_eq!(body["messages"][1]["content"][0]["type"], "tool_use");
assert_eq!(
body["messages"][1]["content"][0]["input"],
json!({"path":"a.txt"})
);
assert_eq!(body["messages"][2]["content"][0]["type"], "tool_result");
assert_eq!(body["messages"][2]["content"][0]["content"], "boom");
assert_eq!(body["messages"][2]["content"][0]["is_error"], true);
}
#[test]
fn anthropic_request_preserves_thinking_blocks_before_tool_use_for_continuation() {
let request = ProviderRequest::new("ignored", vec![ChatMessage::user("hello")])
.with_response_items(vec![
json!({"type":"thinking","thinking":"kept private","signature":"sig_1"}),
json!({"type":"redacted_thinking","data":"encrypted_1"}),
json!({"type":"function_call","call_id":"call_1","name":"read","arguments":"{\"path\":\"a.txt\"}"}),
])
.with_tool_results(vec![ProviderToolResult {
call_id: "call_1".into(),
tool_name: "read".into(),
success: true,
output: "file text".into(),
}]);
let body = anthropic_messages_body("claude-test", &request, None);
let assistant_content = body["messages"][1]["content"].as_array().unwrap();
assert_eq!(
assistant_content[0],
json!({"type":"thinking","thinking":"kept private","signature":"sig_1"})
);
assert_eq!(
assistant_content[1],
json!({"type":"redacted_thinking","data":"encrypted_1"})
);
assert_eq!(assistant_content[2]["type"], "tool_use");
assert_eq!(assistant_content[2]["id"], "call_1");
assert_eq!(body["messages"][2]["content"][0]["type"], "tool_result");
}
#[test]
fn anthropic_request_failed_tool_result_empty_output_uses_fallback() {
let request = ProviderRequest::new("ignored", vec![ChatMessage::user("hello")])
.with_response_items(vec![
json!({"type":"function_call","call_id":"call_1","name":"grep","arguments":"{}"}),
])
.with_tool_results(vec![ProviderToolResult {
call_id: "call_1".into(),
tool_name: "grep".into(),
success: false,
output: "".into(),
}]);
let body = anthropic_messages_body("claude-test", &request, None);
assert_eq!(body["messages"][2]["content"][0]["type"], "tool_result");
assert_eq!(
body["messages"][2]["content"][0]["content"],
"Tool failed with no output."
);
assert_eq!(body["messages"][2]["content"][0]["is_error"], true);
}
#[test]
fn anthropic_request_successful_tool_result_empty_output_stays_empty() {
let request = ProviderRequest::new("ignored", vec![ChatMessage::user("hello")])
.with_response_items(vec![
json!({"type":"function_call","call_id":"call_1","name":"grep","arguments":"{}"}),
])
.with_tool_results(vec![ProviderToolResult {
call_id: "call_1".into(),
tool_name: "grep".into(),
success: true,
output: "".into(),
}]);
let body = anthropic_messages_body("claude-test", &request, None);
assert_eq!(body["messages"][2]["content"][0]["type"], "tool_result");
assert_eq!(body["messages"][2]["content"][0]["content"], "");
assert_eq!(body["messages"][2]["content"][0]["is_error"], false);
}
#[test]
fn anthropic_request_omits_tools_when_disabled() {
let body = anthropic_messages_body(
"claude-test",
&ProviderRequest::new_without_tools("ignored", vec![ChatMessage::user("hello")]),
None,
);
assert!(body.get("tools").is_none());
}
#[test]
fn anthropic_request_omits_tools_when_all_tools_filtered() {
let disabled = crate::tools::MVP_TOOL_CAPABILITIES
.iter()
.map(|tool| tool.canonical_name().to_string())
.collect::<Vec<_>>();
let request = ProviderRequest::new("ignored", vec![ChatMessage::user("hello")])
.with_disabled_tool_names(disabled);
let body = anthropic_messages_body("claude-test", &request, None);
assert!(body.get("tools").is_none());
}
#[test]
fn anthropic_stream_parser_emits_text_usage_and_done() {
let transport = CapturingTransport { chunks: vec![concat!(
"data: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\",\"type\":\"message\",\"role\":\"assistant\",\"content\":[],\"model\":\"claude-test\",\"stop_reason\":null,\"stop_sequence\":null,\"usage\":{\"input_tokens\":2,\"cache_read_input_tokens\":5}}}\n\n",
"data: {\"type\":\"content_block_delta\",\"delta\":{\"type\":\"text_delta\",\"text\":\"hi\"}}\n\n",
"data: {\"type\":\"message_delta\",\"usage\":{\"output_tokens\":3,\"cache_creation_input_tokens\":7}}\n\n",
"data: {\"type\":\"message_stop\"}\n\n"
).to_string()] };
let provider = AnthropicProvider::new("claude-test", "secret", transport);
let mut events = Vec::new();
provider
.stream(
ProviderRequest::new("ignored", vec![ChatMessage::user("hello")]),
&mut |event| {
events.push(event);
Ok(())
},
)
.unwrap();
assert!(events.contains(&ProviderEvent::TextDelta("hi".to_string())));
assert!(events.contains(&ProviderEvent::Usage(Usage {
input: 2,
output: 0,
total: 7,
cache_read: 5,
..Usage::default()
})));
assert!(events.contains(&ProviderEvent::UsagePartial(Usage {
input: 2,
output: 3,
total: 17,
cache_read: 5,
cache_write: 7,
..Usage::default()
})));
assert!(events.contains(&ProviderEvent::Done));
}
#[test]
fn anthropic_stream_parser_emits_tool_call_after_input_json_complete() {
let transport = CapturingTransport { chunks: vec![concat!(
"data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"call_1\",\"name\":\"read\"}}\n\n",
"data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"path\\\":\"}}\n\n",
"data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"\\\"a.txt\\\"}\"}}\n\n",
"data: {\"type\":\"content_block_stop\",\"index\":1}\n\n",
"data: {\"type\":\"message_stop\"}\n\n"
).to_string()] };
let provider = AnthropicProvider::new("claude-test", "secret", transport);
let mut events = Vec::new();
provider
.stream(
ProviderRequest::new("ignored", vec![ChatMessage::user("hello")]),
&mut |event| {
events.push(event);
Ok(())
},
)
.unwrap();
assert!(events.contains(&ProviderEvent::ResponseItem(json!({"type":"function_call","call_id":"call_1","name":"read","arguments":"{\"path\":\"a.txt\"}","status":"completed"}))));
assert!(events.contains(&ProviderEvent::ToolCall(ToolCall {
id: "call_1".into(),
name: "read".into(),
arguments: json!({"path":"a.txt"})
})));
}
#[test]
fn anthropic_stream_parser_emits_thinking_response_items_and_marks_progress() {
let mut parser = AnthropicStreamParser::default();
let start = parser.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\"}}\n\n").unwrap();
assert!(start.semantic_progress);
assert!(start.events.is_empty());
let delta = parser.push_chunk_outcome("data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"private\"}}\n\n").unwrap();
assert!(delta.semantic_progress);
assert!(delta.events.is_empty());
let signature = parser.push_chunk_outcome("data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"signature_delta\",\"signature\":\"sig_1\"}}\n\n").unwrap();
assert!(signature.semantic_progress);
assert!(signature.events.is_empty());
let stop = parser
.push_chunk_outcome("data: {\"type\":\"content_block_stop\",\"index\":0}\n\n")
.unwrap();
assert!(stop.semantic_progress);
assert_eq!(
stop.events,
vec![ProviderEvent::ResponseItem(json!({
"type":"thinking",
"thinking":"private",
"signature":"sig_1"
}))]
);
}
#[test]
fn anthropic_stream_parser_emits_redacted_thinking_response_item() {
let mut parser = AnthropicStreamParser::default();
parser.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"redacted_thinking\",\"data\":\"encrypted_1\"}}\n\n").unwrap();
let stop = parser
.push_chunk_outcome("data: {\"type\":\"content_block_stop\",\"index\":0}\n\n")
.unwrap();
assert_eq!(
stop.events,
vec![ProviderEvent::ResponseItem(json!({
"type":"redacted_thinking",
"data":"encrypted_1"
}))]
);
}
#[test]
fn anthropic_provider_retries_retryable_http_status_before_events() {
let transport = ScriptedTransport::new(vec![
ScriptStep::Error(retryable_503_error()),
ScriptStep::Chunks(vec!["data: {\"type\":\"message_stop\"}\n\n"]),
]);
let attempts = transport.attempts_handle();
let provider = AnthropicProvider::new("claude-test", "secret", transport);
let mut events = Vec::new();
provider
.stream(
ProviderRequest::new("ignored", vec![ChatMessage::user("hello")]),
&mut |event| {
events.push(event);
Ok(())
},
)
.unwrap();
assert_eq!(*attempts.lock().unwrap(), 2);
assert_eq!(events, vec![ProviderEvent::Done]);
}
#[test]
fn anthropic_stream_parser_marks_tool_progress_unsafe_for_recovery() {
let mut parser = AnthropicStreamParser::default();
let start = parser.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"call_1\",\"name\":\"read\"}}\n\n").unwrap();
assert!(start.semantic_progress);
assert!(start.unsafe_recovery_progress);
let delta = parser.push_chunk_outcome("data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"path\\\":\"}}\n\n").unwrap();
assert!(delta.semantic_progress);
assert!(delta.unsafe_recovery_progress);
}
#[test]
fn anthropic_stream_parser_marks_text_progress_safe_for_recovery() {
let mut parser = AnthropicStreamParser::default();
let outcome = parser.push_chunk_outcome("data: {\"type\":\"content_block_delta\",\"delta\":{\"type\":\"text_delta\",\"text\":\"hi\"}}\n\n").unwrap();
assert!(outcome.semantic_progress);
assert!(!outcome.unsafe_recovery_progress);
}
#[test]
fn anthropic_stream_parser_ignores_ping() {
let mut parser = AnthropicStreamParser::default();
assert!(
parser
.push_chunk_outcome("data: {\"type\":\"ping\"}\n\n")
.unwrap()
.events
.is_empty()
);
}
#[test]
fn anthropic_stream_parser_errors_on_error_event() {
let mut parser = AnthropicStreamParser::default();
let error = parser
.push_chunk_outcome("data: {\"type\":\"error\",\"error\":{\"message\":\"bad\"}}\n\n")
.unwrap_err()
.to_string();
assert!(error.contains("anthropic provider stream error"));
}
#[test]
fn anthropic_stream_parser_errors_on_clean_eof_before_message_stop() {
let mut parser = AnthropicStreamParser::default();
parser.push_chunk_outcome("data: {\"type\":\"content_block_delta\",\"delta\":{\"type\":\"text_delta\",\"text\":\"hi\"}}\n\n").unwrap();
assert!(
parser
.finish()
.unwrap_err()
.to_string()
.contains("missing provider stream completion")
);
}
#[test]
fn anthropic_stream_parser_rejects_oversized_tool_arguments() {
let mut parser = AnthropicStreamParser::default();
parser.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"call_1\",\"name\":\"read\"}}\n\n").unwrap();
let delta = "x".repeat(MAX_TOOL_ARGUMENT_BYTES + 1);
let event = json!({"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json": delta}});
assert!(
parser
.push_chunk_outcome(&format!("data: {event}\n\n"))
.unwrap_err()
.to_string()
.contains("maximum size")
);
}
#[test]
fn anthropic_stream_parser_errors_on_pending_tool_at_message_stop() {
let mut parser = AnthropicStreamParser::default();
parser.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"call_1\",\"name\":\"read\"}}\n\n").unwrap();
let error = parser
.push_chunk_outcome("data: {\"type\":\"message_stop\"}\n\n")
.unwrap_err()
.to_string();
assert!(error.contains("incomplete tool_use"));
assert!(error.contains("1"));
}
#[test]
fn anthropic_stream_parser_errors_on_pending_tool_with_partial_input_at_message_stop() {
let mut parser = AnthropicStreamParser::default();
let start = parser.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"call_1\",\"name\":\"read\"}}\n\n").unwrap();
let delta = parser.push_chunk_outcome("data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"path\\\":\\\"/tmp\"}}\n\n").unwrap();
assert!(start.unsafe_recovery_progress);
assert!(delta.unsafe_recovery_progress);
assert!(start.events.iter().chain(&delta.events).all(|event| {
!matches!(
event,
ProviderEvent::ToolCall(_) | ProviderEvent::ResponseItem(_)
)
}));
let error = parser
.push_chunk_outcome("data: {\"type\":\"message_stop\"}\n\n")
.unwrap_err()
.to_string();
assert!(error.contains("incomplete tool_use"), "{error}");
assert!(error.contains("index 1"), "{error}");
assert!(error.contains("recent events"), "{error}");
assert!(error.contains("content_block_delta"), "{error}");
assert!(error.contains("message_stop"), "{error}");
assert!(error.contains("call_1"), "{error}");
assert!(error.contains("read"), "{error}");
assert!(error.contains("argument_bytes 13"), "{error}");
assert!(!error.contains("/tmp"), "{error}");
assert!(!error.contains(r#"{"path":"/tmp"#), "{error}");
}
#[test]
fn anthropic_stream_parser_errors_on_pending_tool_at_finish() {
let mut parser = AnthropicStreamParser::default();
parser.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"call_1\",\"name\":\"read\"}}\n\n").unwrap();
let error = parser.finish().unwrap_err().to_string();
assert!(error.contains("incomplete tool_use"));
assert!(error.contains("1"));
}
#[test]
fn anthropic_stream_parser_requires_tool_event_indexes() {
let mut parser = AnthropicStreamParser::default();
let start_error = parser
.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"content_block\":{\"type\":\"tool_use\",\"id\":\"call_1\",\"name\":\"read\"}}\n\n")
.unwrap_err()
.to_string();
assert!(start_error.contains("missing required index"));
let mut parser = AnthropicStreamParser::default();
parser.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"call_1\",\"name\":\"read\"}}\n\n").unwrap();
let delta_error = parser
.push_chunk_outcome("data: {\"type\":\"content_block_delta\",\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{}\"}}\n\n")
.unwrap_err()
.to_string();
assert!(delta_error.contains("missing required index"));
let mut parser = AnthropicStreamParser::default();
parser.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"call_1\",\"name\":\"read\"}}\n\n").unwrap();
let stop_error = parser
.push_chunk_outcome("data: {\"type\":\"content_block_stop\"}\n\n")
.unwrap_err()
.to_string();
assert!(stop_error.contains("missing required index"));
}
#[test]
fn anthropic_stream_parser_rejects_duplicate_active_tool_index() {
let mut parser = AnthropicStreamParser::default();
parser.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"call_1\",\"name\":\"read\"}}\n\n").unwrap();
let error = parser
.push_chunk_outcome("data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"call_2\",\"name\":\"write\"}}\n\n")
.unwrap_err()
.to_string();
assert!(error.contains("duplicate active tool content block index 1"));
}
#[test]
fn anthropic_stream_parser_saturates_usage_total() {
let mut parser = AnthropicStreamParser::default();
let event = json!({"type":"message_start","message":{"usage":{"input_tokens":u64::MAX,"output_tokens":1,"cache_read_input_tokens":1,"cache_creation_input_tokens":1}}});
let events = parser
.push_chunk_outcome(&format!("data: {event}\n\n"))
.unwrap()
.events;
assert!(events.contains(&ProviderEvent::Usage(Usage {
input: u64::MAX,
output: 1,
total: u64::MAX,
cache_read: 1,
cache_write: 1,
..Usage::default()
})));
}
#[test]
fn anthropic_sse_data_strips_only_one_optional_space() {
assert_eq!(
sse_data("data: hello \ndata:\tthere\t").as_deref(),
Some(" hello \n\tthere\t")
);
}
#[test]
fn anthropic_model_catalog_request_uses_models_endpoint_and_headers() {
let provider =
AnthropicProvider::new("claude-test", "secret", CapturingTransport::default());
let http = provider.build_model_catalog_request();
assert_eq!(http.method, "GET");
assert_eq!(http.url, format!("{ANTHROPIC_MODELS_URL}?limit=1000"));
assert_eq!(http.headers["x-api-key"], "secret");
assert_eq!(http.headers["accept"], "application/json");
}
#[test]
fn anthropic_model_catalog_parser_maps_model_info() {
let entries = parse_anthropic_model_catalog_response(r#"{"data":[{"id":"claude-test","display_name":"Claude Test","max_input_tokens":200000,"max_tokens":8192},{"id":""}]}"#).unwrap();
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].id, "anthropic/claude-test");
assert_eq!(entries[0].display_name.as_deref(), Some("Claude Test"));
assert_eq!(entries[0].context_window, Some(200000));
assert_eq!(entries[0].max_output_tokens, Some(8192));
}
#[test]
fn anthropic_stream_parser_rejects_empty_or_missing_tool_metadata() {
let events = [
json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","name":"read"}}),
json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1"}}),
json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"","name":"read"}}),
json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1","name":" "}}),
];
for event in events {
let mut parser = AnthropicStreamParser::default();
let error = parser
.push_chunk_outcome(&format!("data: {event}\n\n"))
.unwrap_err()
.to_string();
assert!(error.contains("missing non-empty"), "{error}");
assert!(parser.pending_tools.is_empty());
}
}
}