use super::*;
use crate::providers::error::provider_stream_trace_from_error;
fn push(parser: &mut AnthropicStreamParser, event: Value) -> StreamParseOutcome {
parser
.push_chunk_outcome(&format!("data: {event}\n\n"))
.unwrap()
}
#[test]
fn fragmented_frames_emit_text_usage_and_terminal_done_in_order() {
let mut parser = AnthropicStreamParser::default();
let chunks = [
"data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":2}}}\n",
"\ndata: {\"type\":\"content_block_delta\",\"delta\":{\"type\":\"text_delta\",\"text\":\"hi\"}}\n\nda",
"ta: {\"type\":\"message_delta\",\"usage\":{\"output_tokens\":3}}\n\ndata: {\"type\":\"message_stop\"}\n\n",
];
let mut events = Vec::new();
for chunk in chunks {
events.extend(parser.push_chunk_outcome(chunk).unwrap().events);
}
assert_eq!(
events,
vec![
ProviderEvent::UsageObserved(crate::providers::stream::UsageObservation {
usage: Usage {
input: 2,
total: 2,
..Usage::default()
},
presence: crate::providers::stream::UsagePresence {
input: true,
..Default::default()
},
}),
ProviderEvent::TextDelta("hi".to_string()),
ProviderEvent::UsageObserved(crate::providers::stream::UsageObservation {
usage: Usage {
input: 2,
output: 3,
total: 5,
..Usage::default()
},
presence: crate::providers::stream::UsagePresence {
output: true,
..Default::default()
},
}),
ProviderEvent::Done,
]
);
assert!(parser.finish().unwrap().is_empty());
}
#[test]
fn assembles_tool_arguments_and_marks_unsafe_progress() {
let mut parser = AnthropicStreamParser::default();
let start = push(
&mut parser,
json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1","name":"read"}}),
);
let first = push(
&mut parser,
json!({"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":"{\"path\":"}}),
);
let second = push(
&mut parser,
json!({"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":"\"a.txt\"}"}}),
);
let stop = push(&mut parser, json!({"type":"content_block_stop","index":1}));
assert!(start.semantic_progress && start.unsafe_recovery_progress);
assert!(first.semantic_progress && first.unsafe_recovery_progress);
assert!(second.semantic_progress && second.unsafe_recovery_progress);
assert_eq!(
stop.events,
vec![
ProviderEvent::ResponseItem(
json!({"type":"function_call","call_id":"call_1","name":"read","arguments":"{\"path\":\"a.txt\"}","status":"completed"})
),
ProviderEvent::ToolCall(ToolCall {
id: "call_1".into(),
name: "read".into(),
arguments: json!({"path":"a.txt"}),
}),
]
);
}
#[test]
fn emits_thinking_and_redacted_thinking_blocks() {
let mut parser = AnthropicStreamParser::default();
assert!(
push(
&mut parser,
json!({"type":"content_block_start","index":0,"content_block":{"type":"thinking"}}),
)
.semantic_progress
);
push(
&mut parser,
json!({"type":"content_block_delta","index":0,"delta":{"type":"thinking_delta","thinking":"private"}}),
);
push(
&mut parser,
json!({"type":"content_block_delta","index":0,"delta":{"type":"signature_delta","signature":"sig_1"}}),
);
assert_eq!(
push(&mut parser, json!({"type":"content_block_stop","index":0}),).events,
vec![ProviderEvent::ResponseItem(
json!({"type":"thinking","thinking":"private","signature":"sig_1"})
)]
);
push(
&mut parser,
json!({"type":"content_block_start","index":2,"content_block":{"type":"redacted_thinking","data":"encrypted_1"}}),
);
assert_eq!(
push(&mut parser, json!({"type":"content_block_stop","index":2}),).events,
vec![ProviderEvent::ResponseItem(
json!({"type":"redacted_thinking","data":"encrypted_1"})
)]
);
}
#[test]
fn ping_has_no_semantic_progress_but_text_does() {
let mut parser = AnthropicStreamParser::default();
let ping = push(&mut parser, json!({"type":"ping"}));
assert!(ping.events.is_empty());
assert!(!ping.semantic_progress);
let text = push(
&mut parser,
json!({"type":"content_block_delta","delta":{"type":"text_delta","text":"hi"}}),
);
assert!(text.semantic_progress);
assert!(!text.unsafe_recovery_progress);
}
#[test]
fn usage_total_saturates() {
let mut parser = AnthropicStreamParser::default();
let events = push(
&mut parser,
json!({"type":"message_start","message":{"usage":{"input_tokens":u64::MAX,"output_tokens":1,"cache_read_input_tokens":1,"cache_creation_input_tokens":1}}}),
)
.events;
assert!(matches!(
events.as_slice(),
[ProviderEvent::UsageObserved(
crate::providers::stream::UsageObservation {
usage: Usage {
input: u64::MAX,
output: 1,
total: u64::MAX,
cache_read: 1,
cache_write: 1,
..
},
..
}
)]
));
}
#[test]
fn completion_and_incomplete_eof_are_validated() {
let mut complete = AnthropicStreamParser::default();
assert_eq!(
push(&mut complete, json!({"type":"message_stop"})).events,
vec![ProviderEvent::Done]
);
assert!(complete.finish().unwrap().is_empty());
let mut incomplete = AnthropicStreamParser::default();
push(
&mut incomplete,
json!({"type":"content_block_delta","delta":{"type":"text_delta","text":"hi"}}),
);
assert!(
incomplete
.finish()
.unwrap_err()
.to_string()
.contains("missing provider stream completion")
);
}
#[test]
fn bounded_diagnostics_redact_and_escape_sensitive_fragments() {
let secret = "sk-ant-api03-secret";
let long = format!("{}tail-marker", "x".repeat(240));
let mut malformed = AnthropicStreamParser::default();
let error = malformed
.push_chunk_outcome(&format!(
"data: {{\"type\":\"message_start\",\"token\":\"{secret}\",\"text\":\"line\\n{long}\"\n\n"
))
.unwrap_err()
.to_string();
assert!(
error.contains("malformed provider SSE data JSON"),
"{error}"
);
assert!(!error.contains(secret), "{error}");
assert!(error.contains("<redacted>"), "{error}");
assert!(error.contains("\\n"), "{error}");
assert!(error.contains("..."), "{error}");
assert!(!error.contains("tail-marker"), "{error}");
let mut incomplete = AnthropicStreamParser::default();
incomplete
.push_chunk_outcome(&format!("data: token={secret}\n{long}"))
.unwrap();
let error = incomplete.finish().unwrap_err().to_string();
assert!(!error.contains(secret), "{error}");
assert!(error.contains("<redacted>"), "{error}");
assert!(error.contains("\\n"), "{error}");
assert!(error.contains("..."), "{error}");
assert!(!error.contains("tail-marker"), "{error}");
}
#[test]
fn provider_error_and_bad_tool_arguments_use_bounded_redacted_snippets() {
let secret = "sk-ant-api03-secret";
let long = format!("{}tail-marker", "x".repeat(240));
let mut parser = AnthropicStreamParser::default();
let error = push_error(
&mut parser,
json!({"type":"error","error":{"message":format!("{secret}\n{long}")}}),
);
assert!(error.contains("anthropic provider stream error"), "{error}");
assert!(!error.contains(secret), "{error}");
assert!(error.contains("<redacted>"), "{error}");
assert!(error.contains("\\n"), "{error}");
assert!(error.contains("..."), "{error}");
assert!(!error.contains("tail-marker"), "{error}");
let mut parser = AnthropicStreamParser::default();
push(
&mut parser,
json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1","name":"read"}}),
);
push(
&mut parser,
json!({"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":format!("{{\"token\":\"{secret}\",\"text\":\"{long}")}}),
);
let error = push_error(&mut parser, json!({"type":"content_block_stop","index":1}));
assert!(error.contains("malformed non-empty provider tool call arguments"));
assert!(!error.contains(secret), "{error}");
assert!(error.contains("<redacted>"), "{error}");
assert!(error.contains("..."), "{error}");
assert!(!error.contains("tail-marker"), "{error}");
}
fn push_error(parser: &mut AnthropicStreamParser, event: Value) -> String {
parser
.push_chunk_outcome(&format!("data: {event}\n\n"))
.unwrap_err()
.to_string()
}
#[test]
fn rejects_oversized_event_and_tool_arguments() {
let mut event_parser = AnthropicStreamParser::default();
let error = event_parser
.push_chunk_outcome(&"x".repeat(MAX_SSE_EVENT_BUFFER_BYTES + 1))
.unwrap_err()
.to_string();
assert!(error.contains("maximum buffered size"), "{error}");
let mut tool_parser = AnthropicStreamParser::default();
push(
&mut tool_parser,
json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1","name":"read"}}),
);
let error = push_error(
&mut tool_parser,
json!({"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":"x".repeat(MAX_TOOL_ARGUMENT_BYTES + 1)}}),
);
assert!(error.contains("maximum size"), "{error}");
}
#[test]
fn validates_tool_indexes_metadata_and_active_state() {
let invalid = [
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 invalid {
let mut parser = AnthropicStreamParser::default();
assert!(push_error(&mut parser, event).contains("missing non-empty"));
assert!(parser.pending_tools.is_empty());
}
let mut missing_index = AnthropicStreamParser::default();
assert!(
push_error(
&mut missing_index,
json!({"type":"content_block_start","content_block":{"type":"tool_use","id":"call_1","name":"read"}}),
)
.contains("missing required index")
);
let mut duplicate = AnthropicStreamParser::default();
push(
&mut duplicate,
json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_1","name":"read"}}),
);
assert!(
push_error(
&mut duplicate,
json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"call_2","name":"write"}}),
)
.contains("duplicate active tool content block index 1")
);
}
#[test]
fn pending_tool_error_has_bounded_hashed_trace_without_raw_arguments() {
let mut parser = AnthropicStreamParser::default();
push(
&mut parser,
json!({"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"toolu_trace_full_id","name":"read"}}),
);
push(
&mut parser,
json!({"type":"content_block_delta","index":1,"delta":{"type":"input_json_delta","partial_json":"{\"path\":\"/tmp"}}),
);
push(
&mut parser,
json!({"type":"message_delta","delta":{"stop_reason":"tool_use"},"usage":{"output_tokens":7}}),
);
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"));
let pending = trace.pending_tools.first().unwrap();
assert_eq!(pending.argument_bytes, 13);
assert_eq!(
pending.argument_sha256,
"ac2a055274ea75011d1fc6958da49b826e7a101c25258bb537ba6c36fed6a56c"
);
let delta = trace
.recent_events
.iter()
.find(|event| event.delta_type.as_deref() == Some("input_json_delta"))
.unwrap();
assert_eq!(delta.partial_json_bytes, Some(13));
assert_eq!(
delta.partial_json_sha256.as_deref(),
Some("ac2a055274ea75011d1fc6958da49b826e7a101c25258bb537ba6c36fed6a56c")
);
let serialized = serde_json::to_string(&trace).unwrap();
assert!(!serialized.contains("/tmp"));
assert!(!serialized.contains(r#"{\"path\":\"/tmp"#));
}
#[test]
fn stream_trace_bounds_events_and_pending_tools() {
let mut parser = AnthropicStreamParser::default();
for index in 0..20 {
push(
&mut parser,
json!({"type":"content_block_start","index":index,"content_block":{"type":"tool_use","id":format!("toolu_{index}"),"name":"read"}}),
);
push(
&mut parser,
json!({"type":"content_block_delta","index":index,"delta":{"type":"input_json_delta","partial_json":format!("raw-/tmp-marker-{index}")}}),
);
}
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 sse_data_strips_only_one_optional_space() {
assert_eq!(
sse_data("data: hello \ndata:\tthere\t").as_deref(),
Some(" hello \n\tthere\t")
);
}