use crate::providers::{
ProviderEvent, ToolCall, Usage,
error::{
ProviderError, ProviderStreamTrace, ProviderStreamTraceEvent,
ProviderStreamTracePendingTool, ProviderStreamTraceUsage,
},
sse::{diagnostic_snippet, next_sse_event_boundary, sse_data},
stream::{StreamParseOutcome, is_semantic_progress_event},
};
use serde_json::{Value, json};
use sha2::{Digest, Sha256};
use std::collections::{BTreeMap, VecDeque};
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(Default)]
pub(super) 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 {
fn push_chunk_outcome(&mut self, chunk: &str) -> anyhow::Result<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.iter().any(is_semantic_progress_event);
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(StreamParseOutcome {
events,
semantic_progress,
unsafe_recovery_progress,
})
}
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(ProviderError::stream_terminal(
"missing provider stream completion before EOF",
)
.into());
}
Ok(Vec::new())
}
fn parse_data_event(&mut self, data: &str) -> anyhow::Result<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(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(value) = input {
self.usage.input = value;
}
if let Some(value) = output {
self.usage.output = value;
}
if let Some(value) = cache_read {
self.usage.cache_read = value;
}
if let Some(value) = cache_write {
self.usage.cache_write = value;
}
self.usage.total = self
.usage
.input
.saturating_add(self.usage.output)
.saturating_add(self.usage.cache_read)
.saturating_add(self.usage.cache_write);
Some(ProviderEvent::UsageObserved(
crate::providers::stream::UsageObservation {
usage: self.usage.clone(),
presence: crate::providers::stream::UsagePresence {
input: input.is_some(),
output: output.is_some(),
cache_read: cache_read.is_some(),
cache_write: cache_write.is_some(),
total: false,
reasoning: false,
},
},
))
}
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,
}
}
}
impl crate::providers::openai_stream::ProviderStreamParser for AnthropicStreamParser {
fn push_chunk_outcome(&mut self, chunk: &str) -> anyhow::Result<StreamParseOutcome> {
self.push_chunk_outcome(chunk)
}
fn finish(&mut self) -> anyhow::Result<Vec<ProviderEvent>> {
self.finish()
}
}
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")
})
}
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;