use crate::error::AiError;
use crate::event_stream::AssistantMessageEventStreamProducer;
use crate::providers::anthropic::json_parse::parse_streaming_json;
use crate::providers::anthropic::sse::{AnthropicEvent, SseEventStream};
use crate::types::{
AssistantMessage, AssistantMessageEvent, Content, DoneReason, ErrorReason, StopReason,
TextContent, TextContentType, ThinkingContent, ThinkingContentType, ToolCall, ToolCallType,
Usage, UsageCost,
};
use std::sync::Arc;
#[derive(Debug, Clone)]
struct BlockScratch {
anthropic_index: i64,
content_index: usize,
kind: BlockKind,
partial_json: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum BlockKind {
Text,
Thinking { redacted: bool },
ToolCall,
}
pub struct MapperState {
pub output: AssistantMessage,
blocks: Vec<BlockScratch>,
started: bool,
}
impl MapperState {
pub fn new(
api: crate::types::Api,
provider: impl Into<String>,
model: impl Into<String>,
timestamp: i64,
) -> Self {
Self {
output: AssistantMessage::empty(api, provider, model, timestamp),
blocks: Vec::new(),
started: false,
}
}
fn ensure_started(&mut self, prod: &mut AssistantMessageEventStreamProducer) {
if !self.started {
self.started = true;
prod.push(AssistantMessageEvent::Start {
partial: Arc::new(self.output.clone()),
});
}
}
fn find_block_by_anthropic_index(&self, anthropic_index: i64) -> Option<usize> {
self.blocks
.iter()
.position(|b| b.anthropic_index == anthropic_index)
}
pub fn apply(
&mut self,
event: &AnthropicEvent,
prod: &mut AssistantMessageEventStreamProducer,
) -> Result<(), AiError> {
let AnthropicEvent::Message { event_type, payload } = event else {
return Ok(()); };
match event_type.as_str() {
"message_start" => self.apply_message_start(payload, prod),
"content_block_start" => self.apply_content_block_start(payload, prod),
"content_block_delta" => self.apply_content_block_delta(payload, prod)?,
"content_block_stop" => self.apply_content_block_stop(payload, prod),
"message_delta" => self.apply_message_delta(payload, prod),
"message_stop" => { }
_ => {}
}
Ok(())
}
fn apply_message_start(
&mut self,
payload: &serde_json::Value,
prod: &mut AssistantMessageEventStreamProducer,
) {
self.ensure_started(prod);
if let Some(id) = payload
.pointer("/message/id")
.and_then(|v| v.as_str())
.map(|s| s.to_string())
{
self.output.response_id = Some(id);
}
if let Some(usage) = payload.pointer("/message/usage") {
let input = usage
.get("input_tokens")
.and_then(|v| v.as_i64())
.unwrap_or(0);
let output = usage
.get("output_tokens")
.and_then(|v| v.as_i64())
.unwrap_or(0);
let cache_read = usage
.get("cache_read_input_tokens")
.and_then(|v| v.as_i64())
.unwrap_or(0);
let cache_write = usage
.get("cache_creation_input_tokens")
.and_then(|v| v.as_i64())
.unwrap_or(0);
let cache_write_1h = usage
.pointer("/cache_creation/ephemeral_1h_input_tokens")
.and_then(|v| v.as_i64())
.unwrap_or(0);
self.output.usage.input = input;
self.output.usage.output = output;
self.output.usage.cache_read = cache_read;
self.output.usage.cache_write = cache_write;
self.output.usage.cache_write_1h = if cache_write_1h > 0 {
Some(cache_write_1h)
} else {
None
};
self.output.usage.total_tokens =
input + output + cache_read + cache_write;
self.output.usage.cost = UsageCost::default();
}
}
fn apply_content_block_start(
&mut self,
payload: &serde_json::Value,
prod: &mut AssistantMessageEventStreamProducer,
) {
self.ensure_started(prod);
let anthropic_index = payload
.get("index")
.and_then(|v| v.as_i64())
.unwrap_or(0);
let Some(block) = payload.get("content_block") else {
return;
};
let kind = block.get("type").and_then(|v| v.as_str()).unwrap_or("");
match kind {
"text" => {
let text = block.get("text").and_then(|v| v.as_str()).unwrap_or("").to_string();
let content = Content::Text(TextContent {
kind: TextContentType,
text,
text_signature: None,
});
self.output.content.push(content);
let content_index = self.output.content.len() - 1;
self.blocks.push(BlockScratch {
anthropic_index,
content_index,
kind: BlockKind::Text,
partial_json: String::new(),
});
prod.push(AssistantMessageEvent::TextStart {
content_index,
partial: Arc::new(self.output.clone()),
});
}
"thinking" => {
let thinking = block
.get("thinking")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
let signature = block
.get("signature")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
let content = Content::Thinking(ThinkingContent {
kind: ThinkingContentType,
thinking,
thinking_signature: if signature.is_empty() {
None
} else {
Some(signature)
},
redacted: false,
});
self.output.content.push(content);
let content_index = self.output.content.len() - 1;
self.blocks.push(BlockScratch {
anthropic_index,
content_index,
kind: BlockKind::Thinking { redacted: false },
partial_json: String::new(),
});
prod.push(AssistantMessageEvent::ThinkingStart {
content_index,
partial: Arc::new(self.output.clone()),
});
}
"redacted_thinking" => {
let data = block
.get("data")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
let content = Content::Thinking(ThinkingContent {
kind: ThinkingContentType,
thinking: "[Reasoning redacted]".to_string(),
thinking_signature: Some(data),
redacted: true,
});
self.output.content.push(content);
let content_index = self.output.content.len() - 1;
self.blocks.push(BlockScratch {
anthropic_index,
content_index,
kind: BlockKind::Thinking { redacted: true },
partial_json: String::new(),
});
prod.push(AssistantMessageEvent::ThinkingStart {
content_index,
partial: Arc::new(self.output.clone()),
});
}
"tool_use" => {
let id = block
.get("id")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
let name = block
.get("name")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
let args = block.get("input").cloned().unwrap_or(serde_json::Value::Object(
serde_json::Map::new(),
));
let content = Content::tool_call(id, name, args);
self.output.content.push(content);
let content_index = self.output.content.len() - 1;
self.blocks.push(BlockScratch {
anthropic_index,
content_index,
kind: BlockKind::ToolCall,
partial_json: String::new(),
});
prod.push(AssistantMessageEvent::ToolCallStart {
content_index,
partial: Arc::new(self.output.clone()),
});
}
_ => {}
}
}
fn apply_content_block_delta(
&mut self,
payload: &serde_json::Value,
prod: &mut AssistantMessageEventStreamProducer,
) -> Result<(), AiError> {
let anthropic_index = payload
.get("index")
.and_then(|v| v.as_i64())
.unwrap_or(0);
let Some(delta) = payload.get("delta") else {
return Ok(());
};
let delta_type = delta.get("type").and_then(|v| v.as_str()).unwrap_or("");
let Some(scratch_index) = self.find_block_by_anthropic_index(anthropic_index) else {
return Ok(());
};
let scratch = &mut self.blocks[scratch_index];
let content_index = scratch.content_index;
match (scratch.kind, delta_type) {
(BlockKind::Text, "text_delta") => {
let text = delta
.get("text")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
if let Some(Content::Text(slot)) = self.output.content.get_mut(content_index) {
slot.text.push_str(&text);
}
prod.push(AssistantMessageEvent::TextDelta {
content_index,
delta: text,
partial: Arc::new(self.output.clone()),
});
}
(BlockKind::Thinking { redacted: false }, "thinking_delta") => {
let thinking = delta
.get("thinking")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
if let Some(Content::Thinking(slot)) = self.output.content.get_mut(content_index)
{
slot.thinking.push_str(&thinking);
}
prod.push(AssistantMessageEvent::ThinkingDelta {
content_index,
delta: thinking,
partial: Arc::new(self.output.clone()),
});
}
(BlockKind::Thinking { redacted: false }, "signature_delta") => {
let sig = delta
.get("signature")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
if let Some(Content::Thinking(slot)) = self.output.content.get_mut(content_index)
{
let current = slot.thinking_signature.take().unwrap_or_default();
slot.thinking_signature = Some(format!("{current}{sig}"));
}
}
(BlockKind::ToolCall, "input_json_delta") => {
let partial = delta
.get("partial_json")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
scratch.partial_json.push_str(&partial);
if let Some(Content::ToolCall(slot)) = self.output.content.get_mut(content_index)
{
slot.arguments = parse_streaming_json(Some(&scratch.partial_json));
}
prod.push(AssistantMessageEvent::ToolCallDelta {
content_index,
delta: partial,
partial: Arc::new(self.output.clone()),
});
}
_ => {}
}
Ok(())
}
fn apply_content_block_stop(
&mut self,
payload: &serde_json::Value,
prod: &mut AssistantMessageEventStreamProducer,
) {
let anthropic_index = payload
.get("index")
.and_then(|v| v.as_i64())
.unwrap_or(0);
let Some(scratch_index) = self.find_block_by_anthropic_index(anthropic_index) else {
return;
};
let scratch = &mut self.blocks[scratch_index];
let content_index = scratch.content_index;
match scratch.kind {
BlockKind::Text => {
let content = match &self.output.content[content_index] {
Content::Text(t) => t.text.clone(),
_ => String::new(),
};
prod.push(AssistantMessageEvent::TextEnd {
content_index,
content,
partial: Arc::new(self.output.clone()),
});
}
BlockKind::Thinking { redacted } => {
let thinking_text = match &self.output.content[content_index] {
Content::Thinking(t) => t.thinking.clone(),
_ => String::new(),
};
let _ = redacted;
prod.push(AssistantMessageEvent::ThinkingEnd {
content_index,
content: thinking_text,
partial: Arc::new(self.output.clone()),
});
}
BlockKind::ToolCall => {
let final_args =
parse_streaming_json(Some(&scratch.partial_json));
let tool_call = match &self.output.content[content_index] {
Content::ToolCall(tc) => ToolCall {
kind: ToolCallType,
id: tc.id.clone(),
name: tc.name.clone(),
arguments: final_args,
thought_signature: tc.thought_signature.clone(),
namespace: tc.namespace.clone(),
},
_ => return,
};
if let Some(Content::ToolCall(slot)) = self.output.content.get_mut(content_index)
{
slot.arguments = tool_call.arguments.clone();
}
prod.push(AssistantMessageEvent::ToolCallEnd {
content_index,
tool_call,
partial: Arc::new(self.output.clone()),
});
}
}
}
fn apply_message_delta(
&mut self,
payload: &serde_json::Value,
prod: &mut AssistantMessageEventStreamProducer,
) {
let _ = prod; if let Some(stop_reason) = payload
.pointer("/delta/stop_reason")
.and_then(|v| v.as_str())
{
self.output.raw_stop_reason = Some(stop_reason.to_string());
let stop_details = payload.pointer("/delta/stop_details");
match map_stop_reason(stop_reason, stop_details) {
Ok(MappedStop { stop_reason, error_message }) => {
self.output.stop_reason = stop_reason;
if let Some(msg) = error_message {
self.output.error_message = Some(msg);
}
}
Err(e) => {
self.output.stop_reason = StopReason::Error;
self.output.error_message = Some(e.to_string());
}
}
}
if let Some(usage) = payload.get("usage") {
if let Some(v) = usage.get("input_tokens").and_then(|v| v.as_i64()) {
self.output.usage.input = v;
}
if let Some(v) = usage.get("output_tokens").and_then(|v| v.as_i64()) {
self.output.usage.output = v;
}
if let Some(v) = usage.get("cache_read_input_tokens").and_then(|v| v.as_i64()) {
self.output.usage.cache_read = v;
}
if let Some(v) = usage.get("cache_creation_input_tokens").and_then(|v| v.as_i64()) {
self.output.usage.cache_write = v;
}
if let Some(thinking_tokens) = usage
.pointer("/output_tokens_details/thinking_tokens")
.and_then(|v| v.as_i64())
{
self.output.usage.reasoning = Some(thinking_tokens);
}
}
self.output.usage.total_tokens = self.output.usage.input
+ self.output.usage.output
+ self.output.usage.cache_read
+ self.output.usage.cache_write;
}
}
pub struct MappedStop {
pub stop_reason: StopReason,
pub error_message: Option<String>,
}
pub fn map_stop_reason(
reason: &str,
stop_details: Option<&serde_json::Value>,
) -> Result<MappedStop, AiError> {
match reason {
"end_turn" => Ok(MappedStop {
stop_reason: StopReason::Stop,
error_message: None,
}),
"max_tokens" => Ok(MappedStop {
stop_reason: StopReason::Length,
error_message: None,
}),
"tool_use" => Ok(MappedStop {
stop_reason: StopReason::ToolUse,
error_message: None,
}),
"refusal" => {
let explanation = stop_details
.and_then(|d| d.get("explanation"))
.and_then(|v| v.as_str())
.map(|s| s.to_string())
.unwrap_or_else(|| "The model refused to complete the request".to_string());
Ok(MappedStop {
stop_reason: StopReason::Error,
error_message: Some(explanation),
})
}
"pause_turn" | "stop_sequence" => Ok(MappedStop {
stop_reason: StopReason::Stop,
error_message: None,
}),
"sensitive" => Ok(MappedStop {
stop_reason: StopReason::Error,
error_message: Some("Provider stopped with: sensitive".to_string()),
}),
other => Err(AiError::Provider {
code: "unhandled_stop_reason".to_string(),
message: format!("Unhandled stop reason: {other}"),
}),
}
}
pub async fn run_mapper<F>(
stream: &mut SseEventStream,
prod: &mut AssistantMessageEventStreamProducer,
state: &mut MapperState,
cost_fn: F,
) where
F: Fn(&Usage) -> UsageCost,
{
loop {
match stream.next_event().await {
Ok(None) => break,
Ok(Some(sse_frame)) => {
let event = match crate::providers::anthropic::sse::parse_anthropic_event(&sse_frame)
{
Ok(e) => e,
Err(err) => {
emit_terminal_error(prod, state, err.to_string(), false);
return;
}
};
if let Err(err) = state.apply(&event, prod) {
emit_terminal_error(prod, state, err.to_string(), false);
return;
}
}
Err(err) => {
let aborted = matches!(err, AiError::Abort { .. });
emit_terminal_error(prod, state, err.to_string(), aborted);
return;
}
}
}
finalize_mapper(prod, state, cost_fn);
}
pub fn finalize_mapper<F>(
prod: &mut AssistantMessageEventStreamProducer,
state: &mut MapperState,
cost_fn: F,
) where
F: Fn(&Usage) -> UsageCost,
{
if state.output.stop_reason == StopReason::Pending {
emit_terminal_error(
prod,
state,
"Anthropic stream ended without a stop reason".to_string(),
false,
);
return;
}
if matches!(
state.output.stop_reason,
StopReason::Aborted | StopReason::Error
) {
let aborted = matches!(state.output.stop_reason, StopReason::Aborted);
let msg = state
.output
.error_message
.clone()
.unwrap_or_else(|| "An unknown error occurred".to_string());
emit_terminal_error(prod, state, msg, aborted);
return;
}
state.output.usage.cost = cost_fn(&state.output.usage);
let reason = match state.output.stop_reason {
StopReason::Stop => DoneReason::Stop,
StopReason::Length => DoneReason::Length,
StopReason::ToolUse => DoneReason::ToolUse,
StopReason::Deferred => DoneReason::Deferred,
_ => DoneReason::Stop, };
prod.push(AssistantMessageEvent::Done {
reason,
message: state.output.clone(),
});
}
pub fn emit_terminal_error(
prod: &mut AssistantMessageEventStreamProducer,
state: &mut MapperState,
message: String,
aborted: bool,
) {
state.output.stop_reason = if aborted {
StopReason::Aborted
} else {
StopReason::Error
};
state.output.error_message = Some(message);
let reason = if aborted {
ErrorReason::Aborted
} else {
ErrorReason::Error
};
prod.push(AssistantMessageEvent::Error {
reason,
error: state.output.clone(),
});
}
#[cfg(test)]
mod tests {
use super::*;
use crate::event_stream::create_assistant_message_event_stream;
use crate::providers::anthropic::sse::{parse_anthropic_event, ServerSentEvent};
use crate::types::{Api, Content, StopReason};
use serde_json::json;
fn frame(event: &str, data: &str) -> ServerSentEvent {
ServerSentEvent {
event: Some(event.to_string()),
data: data.to_string(),
raw: vec![format!("event: {event}\ndata: {data}")],
}
}
fn j(value: &serde_json::Value) -> String {
serde_json::to_string(value).unwrap()
}
struct Run {
tags: Vec<&'static str>,
result: AssistantMessage,
}
async fn run_fixture(frames: Vec<ServerSentEvent>) -> Run {
let (mut prod, stream) = create_assistant_message_event_stream();
let mut state = MapperState::new(
Api::AnthropicMessages,
"anthropic",
"claude-haiku-4-5",
0,
);
for f in &frames {
let event = parse_anthropic_event(f).expect("event parses");
state.apply(&event, &mut prod).expect("apply");
}
finalize_mapper(&mut prod, &mut state, |_| UsageCost::default());
drop(prod);
let mut stream = stream;
let mut tags = Vec::new();
while let Some(ev) = stream.next().await {
tags.push(ev.type_tag());
}
let result = stream.result().await.expect("terminal result");
Run { tags, result }
}
#[tokio::test]
async fn repairs_malformed_streamed_tool_json() {
let malformed_delta = r#"{"type":"content_block_delta","index":0,"delta":{"type":"input_json_delta","partial_json":"{\"path\":\"A\H\",\"text\":\"col1\tcol2\"}"}}"#;
let frames = vec![
frame(
"message_start",
&j(&json!({
"type": "message_start",
"message": {
"id": "msg_test",
"usage": {
"input_tokens": 12,
"output_tokens": 0,
"cache_read_input_tokens": 0,
"cache_creation_input_tokens": 0,
},
},
})),
),
frame(
"content_block_start",
&j(&json!({
"type": "content_block_start",
"index": 0,
"content_block": { "type": "tool_use", "id": "toolu_test", "name": "edit", "input": {} },
})),
),
frame("content_block_delta", malformed_delta),
frame(
"content_block_stop",
&j(&json!({ "type": "content_block_stop", "index": 0 })),
),
frame(
"message_delta",
&j(&json!({
"type": "message_delta",
"delta": { "stop_reason": "tool_use" },
"usage": {
"input_tokens": 12,
"output_tokens": 5,
"cache_read_input_tokens": 0,
"cache_creation_input_tokens": 0,
},
})),
),
frame("message_stop", &j(&json!({ "type": "message_stop" }))),
];
let run = run_fixture(frames).await;
assert_eq!(run.result.stop_reason, StopReason::ToolUse);
assert!(run.result.error_message.is_none());
let toolcall = run
.result
.content
.iter()
.find_map(|c| match c {
Content::ToolCall(tc) => Some(tc),
_ => None,
})
.expect("a tool call block");
assert_eq!(toolcall.arguments, json!({ "path": "A\\H", "text": "col1\tcol2" }));
assert_eq!(
run.tags,
vec!["start", "toolcall_start", "toolcall_delta", "toolcall_end", "done"]
);
}
#[tokio::test]
async fn preserves_content_from_content_block_start() {
let frames = vec![
frame(
"message_start",
&j(&json!({
"type": "message_start",
"message": {
"id": "msg_initial_content",
"usage": { "input_tokens": 12, "output_tokens": 0, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 },
},
})),
),
frame(
"content_block_start",
&j(&json!({ "type": "content_block_start", "index": 0, "content_block": { "type": "text", "text": "Initial text" } })),
),
frame(
"content_block_delta",
&j(&json!({ "type": "content_block_delta", "index": 0, "delta": { "type": "text_delta", "text": " plus delta" } })),
),
frame("content_block_stop", &j(&json!({ "type": "content_block_stop", "index": 0 }))),
frame(
"content_block_start",
&j(&json!({ "type": "content_block_start", "index": 1, "content_block": { "type": "thinking", "thinking": "Initial thinking", "signature": "initial signature" } })),
),
frame(
"content_block_delta",
&j(&json!({ "type": "content_block_delta", "index": 1, "delta": { "type": "thinking_delta", "thinking": " plus delta" } })),
),
frame(
"content_block_delta",
&j(&json!({ "type": "content_block_delta", "index": 1, "delta": { "type": "signature_delta", "signature": " plus delta" } })),
),
frame("content_block_stop", &j(&json!({ "type": "content_block_stop", "index": 1 }))),
frame(
"message_delta",
&j(&json!({
"type": "message_delta",
"delta": { "stop_reason": "end_turn" },
"usage": { "input_tokens": 12, "output_tokens": 5, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 },
})),
),
frame("message_stop", &j(&json!({ "type": "message_stop" }))),
];
let run = run_fixture(frames).await;
assert_eq!(run.result.stop_reason, StopReason::Stop);
let text = match &run.result.content[0] {
Content::Text(t) => t,
_ => panic!("first block is text"),
};
assert_eq!(text.text, "Initial text plus delta");
let thinking = match &run.result.content[1] {
Content::Thinking(t) => t,
_ => panic!("second block is thinking"),
};
assert_eq!(thinking.thinking, "Initial thinking plus delta");
assert_eq!(
thinking.thinking_signature.as_deref(),
Some("initial signature plus delta")
);
assert!(!thinking.redacted);
}
#[tokio::test]
async fn preserves_refusal_stop_details() {
let explanation = "This request triggered restrictions on violative cyber content and was blocked under Anthropic's Usage Policy.";
let frames = vec![
frame(
"message_start",
&j(&json!({
"type": "message_start",
"message": { "id": "msg_01XFUDYJgAACzvnptvVoYEL", "usage": { "input_tokens": 412, "output_tokens": 0, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 } },
})),
),
frame(
"message_delta",
&j(&json!({
"type": "message_delta",
"delta": { "stop_reason": "refusal", "stop_details": { "type": "refusal", "category": "cyber", "explanation": explanation } },
"usage": { "input_tokens": 412, "output_tokens": 0, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 },
})),
),
frame("message_stop", &j(&json!({ "type": "message_stop" }))),
];
let run = run_fixture(frames).await;
assert_eq!(run.result.stop_reason, StopReason::Error);
assert_eq!(run.result.raw_stop_reason.as_deref(), Some("refusal"));
assert_eq!(run.result.error_message.as_deref(), Some(explanation));
assert_eq!(run.tags.last().copied(), Some("error"));
}
#[tokio::test]
async fn preserves_sensitive_stop_reasons() {
let frames = vec![
frame(
"message_start",
&j(&json!({ "type": "message_start", "message": { "id": "msg_sensitive", "usage": { "input_tokens": 12, "output_tokens": 0, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 } } })),
),
frame(
"message_delta",
&j(&json!({ "type": "message_delta", "delta": { "stop_reason": "sensitive" }, "usage": { "input_tokens": 12, "output_tokens": 0, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 } })),
),
frame("message_stop", &j(&json!({ "type": "message_stop" }))),
];
let run = run_fixture(frames).await;
assert_eq!(run.result.stop_reason, StopReason::Error);
assert_eq!(run.result.raw_stop_reason.as_deref(), Some("sensitive"));
assert_eq!(
run.result.error_message.as_deref(),
Some("Provider stopped with: sensitive")
);
}
#[tokio::test]
async fn message_delta_without_usage_noop() {
let frames = vec![
frame(
"message_start",
&j(&json!({ "type": "message_start", "message": { "id": "msg_test", "usage": { "input_tokens": 12, "output_tokens": 0, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 } } })),
),
frame(
"content_block_start",
&j(&json!({ "type": "content_block_start", "index": 0, "content_block": { "type": "text", "text": "" } })),
),
frame(
"content_block_delta",
&j(&json!({ "type": "content_block_delta", "index": 0, "delta": { "type": "text_delta", "text": "Hello" } })),
),
frame("content_block_stop", &j(&json!({ "type": "content_block_stop", "index": 0 }))),
frame(
"message_delta",
&j(&json!({ "type": "message_delta", "delta": { "stop_reason": "end_turn" } })),
),
frame("message_stop", &j(&json!({ "type": "message_stop" }))),
];
let run = run_fixture(frames).await;
assert_eq!(run.result.stop_reason, StopReason::Stop);
assert!(run.result.error_message.is_none());
let text = match &run.result.content[0] {
Content::Text(t) => t,
_ => panic!("text block"),
};
assert_eq!(text.text, "Hello");
assert_eq!(run.result.usage.input, 12);
assert_eq!(run.result.usage.total_tokens, 12);
}
#[tokio::test]
async fn ignores_unknown_events_after_message_stop() {
let mut frames = vec![
frame(
"message_start",
&j(&json!({ "type": "message_start", "message": { "id": "msg_test", "usage": { "input_tokens": 12, "output_tokens": 0, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 } } })),
),
frame(
"content_block_start",
&j(&json!({ "type": "content_block_start", "index": 0, "content_block": { "type": "text", "text": "" } })),
),
frame(
"content_block_delta",
&j(&json!({ "type": "content_block_delta", "index": 0, "delta": { "type": "text_delta", "text": "Hello" } })),
),
frame("content_block_stop", &j(&json!({ "type": "content_block_stop", "index": 0 }))),
frame(
"message_delta",
&j(&json!({ "type": "message_delta", "delta": { "stop_reason": "end_turn" }, "usage": { "input_tokens": 12, "output_tokens": 5, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 } })),
),
frame("message_stop", &j(&json!({ "type": "message_stop" }))),
];
frames.push(frame("done", "[DONE]"));
frames.push(frame("proxy.stats", "not json"));
let run = run_fixture(frames).await;
assert_eq!(run.result.stop_reason, StopReason::Stop);
assert!(run.result.error_message.is_none());
let text = match &run.result.content[0] {
Content::Text(t) => t,
_ => panic!("text block"),
};
assert_eq!(text.text, "Hello");
}
#[test]
fn map_stop_reason_table() {
let stop = |r: &str| map_stop_reason(r, None).unwrap();
assert_eq!(stop("end_turn").stop_reason, StopReason::Stop);
assert_eq!(stop("max_tokens").stop_reason, StopReason::Length);
assert_eq!(stop("tool_use").stop_reason, StopReason::ToolUse);
assert_eq!(stop("pause_turn").stop_reason, StopReason::Stop);
assert_eq!(stop("stop_sequence").stop_reason, StopReason::Stop);
let details = json!({ "type": "refusal", "explanation": "blocked" });
let refusal = map_stop_reason("refusal", Some(&details)).unwrap();
assert_eq!(refusal.stop_reason, StopReason::Error);
assert_eq!(refusal.error_message.as_deref(), Some("blocked"));
let refusal_default = map_stop_reason("refusal", None).unwrap();
assert_eq!(
refusal_default.error_message.as_deref(),
Some("The model refused to complete the request")
);
let sensitive = map_stop_reason("sensitive", None).unwrap();
assert_eq!(sensitive.stop_reason, StopReason::Error);
assert_eq!(
sensitive.error_message.as_deref(),
Some("Provider stopped with: sensitive")
);
assert!(map_stop_reason("nonsense", None).is_err());
}
#[tokio::test]
async fn message_start_records_response_id_and_usage() {
let frames = vec![
frame(
"message_start",
&j(&json!({
"type": "message_start",
"message": {
"id": "msg_abc",
"usage": {
"input_tokens": 10,
"output_tokens": 2,
"cache_read_input_tokens": 3,
"cache_creation_input_tokens": 4,
"cache_creation": { "ephemeral_1h_input_tokens": 1 },
},
},
})),
),
frame(
"message_delta",
&j(&json!({ "type": "message_delta", "delta": { "stop_reason": "end_turn" } })),
),
frame("message_stop", &j(&json!({ "type": "message_stop" }))),
];
let run = run_fixture(frames).await;
assert_eq!(run.result.response_id.as_deref(), Some("msg_abc"));
assert_eq!(run.result.usage.input, 10);
assert_eq!(run.result.usage.output, 2);
assert_eq!(run.result.usage.cache_read, 3);
assert_eq!(run.result.usage.cache_write, 4);
assert_eq!(run.result.usage.cache_write_1h, Some(1));
assert_eq!(run.result.usage.total_tokens, 10 + 2 + 3 + 4);
}
#[tokio::test]
async fn message_delta_records_reasoning_tokens() {
let frames = vec![
frame(
"message_start",
&j(&json!({ "type": "message_start", "message": { "id": "m", "usage": { "input_tokens": 1, "output_tokens": 0, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 } } })),
),
frame(
"message_delta",
&j(&json!({
"type": "message_delta",
"delta": { "stop_reason": "end_turn" },
"usage": {
"input_tokens": 1,
"output_tokens": 50,
"cache_read_input_tokens": 0,
"cache_creation_input_tokens": 0,
"output_tokens_details": { "thinking_tokens": 30 },
},
})),
),
frame("message_stop", &j(&json!({ "type": "message_stop" }))),
];
let run = run_fixture(frames).await;
assert_eq!(run.result.usage.output, 50);
assert_eq!(run.result.usage.reasoning, Some(30));
}
#[tokio::test]
async fn redacted_thinking_block() {
let frames = vec![
frame(
"message_start",
&j(&json!({ "type": "message_start", "message": { "id": "m", "usage": { "input_tokens": 1, "output_tokens": 0, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 } } })),
),
frame(
"content_block_start",
&j(&json!({ "type": "content_block_start", "index": 0, "content_block": { "type": "redacted_thinking", "data": "opaque-base64" } })),
),
frame("content_block_stop", &j(&json!({ "type": "content_block_stop", "index": 0 }))),
frame(
"message_delta",
&j(&json!({ "type": "message_delta", "delta": { "stop_reason": "end_turn" } })),
),
frame("message_stop", &j(&json!({ "type": "message_stop" }))),
];
let run = run_fixture(frames).await;
let thinking = match &run.result.content[0] {
Content::Thinking(t) => t,
_ => panic!("thinking block"),
};
assert!(thinking.redacted);
assert_eq!(thinking.thinking, "[Reasoning redacted]");
assert_eq!(thinking.thinking_signature.as_deref(), Some("opaque-base64"));
assert!(run.tags.contains(&"thinking_start"));
assert!(run.tags.contains(&"thinking_end"));
}
#[tokio::test]
async fn stream_without_stop_reason_is_error() {
let frames = vec![
frame(
"message_start",
&j(&json!({ "type": "message_start", "message": { "id": "m", "usage": { "input_tokens": 1, "output_tokens": 0, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 } } })),
),
frame("message_stop", &j(&json!({ "type": "message_stop" }))),
];
let run = run_fixture(frames).await;
assert_eq!(run.result.stop_reason, StopReason::Error);
assert_eq!(run.tags.last().copied(), Some("error"));
}
#[tokio::test]
async fn tool_call_partial_json_reparse_each_delta() {
let frames = vec![
frame(
"message_start",
&j(&json!({ "type": "message_start", "message": { "id": "m", "usage": { "input_tokens": 1, "output_tokens": 0, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0 } } })),
),
frame(
"content_block_start",
&j(&json!({ "type": "content_block_start", "index": 0, "content_block": { "type": "tool_use", "id": "t1", "name": "write", "input": {} } })),
),
frame(
"content_block_delta",
&j(&json!({ "type": "content_block_delta", "index": 0, "delta": { "type": "input_json_delta", "partial_json": "{\"path\":\"a\"," } })),
),
frame(
"content_block_delta",
&j(&json!({ "type": "content_block_delta", "index": 0, "delta": { "type": "input_json_delta", "partial_json": "\"text\":\"b\"}" } })),
),
frame("content_block_stop", &j(&json!({ "type": "content_block_stop", "index": 0 }))),
frame(
"message_delta",
&j(&json!({ "type": "message_delta", "delta": { "stop_reason": "tool_use" } })),
),
frame("message_stop", &j(&json!({ "type": "message_stop" }))),
];
let run = run_fixture(frames).await;
let tc = match &run.result.content[0] {
Content::ToolCall(t) => t,
_ => panic!("toolcall"),
};
assert_eq!(tc.arguments, json!({ "path": "a", "text": "b" }));
let deltas: Vec<&'static str> = run
.tags
.iter()
.filter(|&&t| t == "toolcall_delta")
.copied()
.collect();
assert_eq!(deltas.len(), 2);
}
#[tokio::test]
async fn sse_error_event_surfaces_error() {
let (prod, stream) = create_assistant_message_event_stream();
let mut prod = prod;
let mut state = MapperState::new(
Api::AnthropicMessages,
"anthropic",
"claude-haiku-4-5",
0,
);
let err_frame = frame("error", "rate limited");
match parse_anthropic_event(&err_frame) {
Err(e) => emit_terminal_error(&mut prod, &mut state, e.to_string(), false),
Ok(_) => panic!("expected error"),
}
let result = stream.result().await.expect("terminal");
assert_eq!(result.stop_reason, StopReason::Error);
assert!(result.error_message.as_deref().unwrap().contains("rate limited"));
}
}