use super::super::completion::{CLAUDE_OPUS_4_8, CLAUDE_OPUS_5_5, CLAUDE_SONNET_4_6};
use super::*;
use crate::completion::CompletionRequest;
use crate::completion::Message as RigMessage;
use crate::completion::request::Document as RigDocument;
use crate::driver::{Decoded, decode_events};
use crate::message::{
AssistantContent, DocumentRange, Opaque, Reasoning, Source, SourceLocation, StopReason,
};
use crate::providers::anthropic::wire::AnthropicConfig;
use crate::wire::Mode;
use serde_json::json;
fn adapter() -> MessagesDecoder {
MessagesDecoder::new(false)
}
fn decode(events: impl IntoIterator<Item = MessagesEvent>) -> Decoded<Completion> {
decode_events!(MessagesDecoder::new(false), "anthropic", events)
}
fn classified(frame: &str) -> MessagesEvent {
let crate::wire::WireEvent::Known(event) = adapter().classify(WireFrame::Text(frame.into()))
else {
panic!("{frame} must classify Known");
};
event
}
fn streamed(frames: &[Value]) -> Result<crate::completion::CompletionResponse, ProviderError> {
let wire = AnthropicConfig::new("test-key").completion(CLAUDE_SONNET_4_6);
crate::test_utils::decode_reply(
&wire,
&CompletionRequest::new("hello"),
Mode::Streaming,
frames
.iter()
.map(|frame| WireFrame::Text(frame.to_string())),
Value::Null,
)
}
fn block(index: usize, start: Value, deltas: &[Value]) -> Vec<Value> {
std::iter::once(json!({"type": "content_block_start", "index": index, "content_block": start}))
.chain(
deltas.iter().map(
|delta| json!({"type": "content_block_delta", "index": index, "delta": delta}),
),
)
.chain([json!({"type": "content_block_stop", "index": index})])
.collect()
}
fn reply(blocks: Vec<Vec<Value>>, stop_reason: &str) -> Vec<Value> {
std::iter::once(json!({"type": "message_start", "message": {
"id": "msg_1", "type": "message", "role": "assistant", "model": CLAUDE_SONNET_4_6,
"content": [], "stop_reason": null, "stop_sequence": null,
"usage": {"input_tokens": 3, "output_tokens": 1}
}}))
.chain(blocks.into_iter().flatten())
.chain([json!({"type": "message_delta",
"delta": {"stop_reason": stop_reason, "stop_sequence": null},
"usage": {"output_tokens": 5}})])
.collect()
}
fn item(content: &AssistantContent) -> Option<&Value> {
match content {
AssistantContent::Opaque(Opaque { item, .. }) => Some(item),
content => content.native_item(),
}
}
fn built_streaming_body(
model: &str,
request: CompletionRequest,
strict_tools: bool,
) -> Result<Value, ProviderError> {
use crate::wire::{Body, Mode, Wire};
let wire =
crate::providers::anthropic::wire::AnthropicConfig::new("test-key").completion(model);
let wire = if strict_tools {
wire.with_strict_tools()
} else {
wire
};
let request = <Completion as crate::wire::Operation>::prepare(request, &wire.describe())?;
let encoded = wire.encode(request, Mode::Streaming)?;
match encoded.request.body() {
Body::Bytes(bytes) => Ok(serde_json::from_slice(bytes)?),
Body::Multipart(_) => Err(ProviderError::request("the Messages endpoint takes JSON")),
}
}
#[test]
fn test_streaming_tool_build_marks_final_combined_tool() {
let request = CompletionRequest::new(RigMessage::user("Hi"))
.max_tokens(64)
.tools(vec![crate::completion::ToolDefinition {
name: crate::message::ToolName::new("rig_tool").expect("tool name"),
description: "Rig tool".to_string(),
parameters: json!({"type": "object", "properties": {}}),
}])
.additional_params(json!({
"tools": [{
"name": "provider_tool",
"description": "Provider tool",
"input_schema": {"type": "object"}
}]
}));
let wire = AnthropicConfig::new("test-key")
.completion(CLAUDE_SONNET_4_6)
.with_prompt_caching();
let body = streamed_body(&wire, request);
let tools = body["tools"].as_array().expect("tools");
assert_eq!(tools.len(), 2);
assert!(tools[0].get("cache_control").is_none());
assert_eq!(tools[1]["name"], "provider_tool");
assert_eq!(tools[1]["cache_control"]["type"], "ephemeral");
}
fn streamed_body(
wire: &crate::providers::anthropic::Messages,
request: CompletionRequest,
) -> Value {
use crate::wire::Wire;
crate::test_utils::json_body(
&wire
.encode(request, Mode::Streaming)
.expect("the request encodes")
.request,
)
}
#[test]
fn streaming_request_keeps_documents_after_leading_system_messages() {
let request = CompletionRequest::from(vec![
RigMessage::system("System prompt"),
RigMessage::assistant("Earlier assistant turn"),
RigMessage::system("Mid-conversation instruction"),
RigMessage::user("Prompt"),
])
.max_tokens(64)
.documents(vec![RigDocument {
id: "doc1".to_string(),
text: "Document text.".to_string(),
additional_props: Default::default(),
}]);
let body = built_streaming_body(CLAUDE_OPUS_4_8, request, false)
.expect("streaming request body should build");
assert_eq!(body["system"][0]["text"], "System prompt");
assert_eq!(body["system"].as_array().map(Vec::len), Some(1));
let messages = body["messages"]
.as_array()
.expect("messages should be array");
assert_eq!(messages.len(), 4);
assert_eq!(messages[3]["role"], "system");
assert_eq!(messages[0]["role"], "user");
assert!(
messages[0].to_string().contains("<file id: doc1>"),
"document message should follow top-level system: {messages:?}"
);
assert_eq!(messages[1]["role"], "assistant");
assert_eq!(messages[2]["role"], "user");
assert_eq!(
messages
.iter()
.filter(|message| message.to_string().contains("<file id: doc1>"))
.count(),
1,
"document message should appear exactly once: {messages:?}"
);
}
#[test]
fn streaming_body_is_blocking_body_plus_stream_flag_and_carries_output_schema() {
let schema: schemars::Schema = serde_json::from_value(json!({
"title": "WeatherResponse",
"type": "object",
"properties": { "city": { "type": "string" } }
}))
.expect("schema should deserialize");
let request = CompletionRequest::from(vec![
RigMessage::system("You are helpful"),
RigMessage::user("What's the weather?"),
])
.temperature(0.5)
.max_tokens(64)
.output_schema(schema);
let streaming_body = built_streaming_body(CLAUDE_OPUS_4_8, request.clone(), false)
.expect("streaming request body should build");
assert_eq!(streaming_body["stream"], serde_json::Value::Bool(true));
assert_eq!(
streaming_body["output_config"]["format"]["type"],
"json_schema"
);
assert!(
streaming_body["output_config"]["format"]["schema"].is_object(),
"streaming body must carry the structured-output schema: {streaming_body}"
);
let wire = crate::providers::anthropic::wire::AnthropicConfig::new("test-key")
.completion(CLAUDE_OPUS_4_8);
let mut expected = crate::test_utils::json_body(
&crate::wire::Wire::encode(&wire, request, Mode::Unary)
.expect("blocking request body should build")
.request,
);
expected
.as_object_mut()
.expect("body is an object")
.insert("stream".to_string(), serde_json::Value::Bool(true));
assert_eq!(streaming_body, expected);
}
#[test]
fn streaming_body_drops_tool_choice_when_no_tools_are_advertised() {
let request = CompletionRequest::new(RigMessage::user("Hi"))
.max_tokens(64)
.tool_choice(crate::message::ToolChoice::Auto);
let body = built_streaming_body(CLAUDE_OPUS_4_8, request, false)
.expect("streaming request body should build");
assert!(
body.get("tool_choice").is_none(),
"tool_choice must be omitted when no tools are advertised: {body}"
);
assert!(body.get("tools").is_none());
}
#[test]
fn thinking_assembles_its_text_and_signature_into_its_item() {
let frames = reply(
vec![block(
0,
json!({"type": "thinking", "thinking": "", "signature": "open-"}),
&[
json!({"type": "thinking_delta", "thinking": "Let me "}),
json!({"type": "thinking_delta", "thinking": "think."}),
json!({"type": "signature_delta", "signature": "sig_"}),
json!({"type": "signature_delta", "signature": "end"}),
],
)],
"end_turn",
);
let response = streamed(&frames).expect("the reply folds");
let [AssistantContent::Reasoning(reasoning)] = response.choice.as_slice() else {
panic!("one reasoning block: {:?}", response.choice);
};
assert_eq!(reasoning.text, "Let me think.");
assert!(!reasoning.redacted);
assert_eq!(
item(&response.choice[0]),
Some(
&json!({"type": "thinking", "thinking": "Let me think.", "signature": "open-sig_end"})
)
);
}
#[test]
fn redacted_thinking_is_a_redacted_block_holding_its_payload() {
let frames = reply(
vec![block(
0,
json!({"type": "redacted_thinking", "data": "redacted_blob"}),
&[],
)],
"end_turn",
);
let response = streamed(&frames).expect("the reply folds");
let [AssistantContent::Reasoning(Reasoning { text, redacted, .. })] =
response.choice.as_slice()
else {
panic!("one reasoning block: {:?}", response.choice);
};
assert!(text.is_empty() && *redacted);
assert_eq!(
item(&response.choice[0]),
Some(&json!({"type": "redacted_thinking", "data": "redacted_blob"}))
);
}
#[test]
fn a_call_assembles_its_input_into_its_item() {
let frames = reply(
vec![block(
0,
json!({"type": "tool_use", "id": "toolu_1", "name": "lookup", "input": {},
"caller": {"type": "direct"}}),
&[
json!({"type": "input_json_delta", "partial_json": "{\"location\":"}),
json!({"type": "input_json_delta", "partial_json": "\"Paris\"}"}),
],
)],
"tool_use",
);
let response = streamed(&frames).expect("the reply folds");
let [AssistantContent::ToolCall(call)] = response.choice.as_slice() else {
panic!("one call: {:?}", response.choice);
};
assert_eq!(call.id.to_string(), "toolu_1");
assert_eq!(
call.function.arguments_value(),
json!({"location": "Paris"})
);
assert_eq!(
item(&response.choice[0]),
Some(
&json!({"type": "tool_use", "id": "toolu_1", "name": "lookup",
"input": {"location": "Paris"}, "caller": {"type": "direct"}})
)
);
assert_eq!(response.stop(), StopReason::ToolUse);
}
#[test]
fn a_stream_without_its_stop_reason_never_ends_successfully() {
let mut frames = reply(
vec![block(0, json!({"type": "text", "text": "partial"}), &[])],
"end_turn",
);
if let Some(delta) = frames.last_mut() {
delta["delta"]["stop_reason"] = Value::Null;
}
frames.push(json!({"type": "message_stop"}));
assert!(matches!(streamed(&frames), Err(ProviderError::Truncated)));
}
#[test]
fn input_to_a_text_block_fails_the_reply() {
let frames = reply(
vec![block(
0,
json!({"type": "text", "text": ""}),
&[json!({"type": "input_json_delta", "partial_json": "{}"})],
)],
"end_turn",
);
assert!(streamed(&frames).is_err());
}
#[test]
fn citation_deltas_land_on_the_text_item() {
let citation = json!({"type": "char_location", "cited_text": "The grass is green.",
"document_index": 0, "document_title": "Example", "start_char_index": 0,
"end_char_index": 20});
for opening in [json!(null), json!([])] {
let frames = reply(
vec![block(
0,
json!({"type": "text", "text": "", "citations": opening}),
&[
json!({"type": "citations_delta", "citation": citation}),
json!({"type": "text_delta", "text": "the grass is green"}),
],
)],
"end_turn",
);
let response = streamed(&frames).expect("the reply folds");
let [AssistantContent::Text(text)] = response.choice.as_slice() else {
panic!("one text block: {:?}", response.choice);
};
assert_eq!(text.text, "the grass is green");
assert_eq!(
response.choice[0]
.native_item()
.and_then(|item| item.get("citations")),
Some(&json!([citation]))
);
}
}
#[test]
fn a_fallback_block_is_kept_first_and_an_error_after_output() {
let fallback = json!({"type": "fallback", "model": "claude-opus-4-8"});
let text = json!({"type": "text", "text": ""});
let leading = streamed(&reply(
vec![
block(0, fallback.clone(), &[]),
block(
1,
text.clone(),
&[json!({"type": "text_delta", "text": "hi"})],
),
],
"end_turn",
))
.expect("a leading fallback is legal");
assert!(matches!(
leading.choice.first(),
Some(AssistantContent::Opaque(Opaque { replay: false, .. }))
));
let late = streamed(&reply(
vec![
block(0, text, &[json!({"type": "text_delta", "text": "hi"})]),
block(1, fallback, &[]),
],
"end_turn",
));
assert!(late.is_err(), "{late:?}");
}
#[test]
fn classify_dispatches_on_the_known_event_list() {
let adapter = adapter();
let frame = WireFrame::Text(r#"{"type":"something_new_from_anthropic","field":"x"}"#.into());
assert!(matches!(
adapter.classify(frame),
crate::wire::WireEvent::Unknown { event_type, .. }
if event_type == "something_new_from_anthropic"
));
let frame = WireFrame::Text(r#"{"type":"ping"}"#.into());
assert!(matches!(
adapter.classify(frame),
crate::wire::WireEvent::Known(event) if event.kind() == "ping"
));
let frame = WireFrame::Text("{not json".into());
assert!(matches!(
adapter.classify(frame),
crate::wire::WireEvent::Corrupt(_)
));
}
#[test]
fn cache_usage_from_message_start_survives_output_only_terminal_delta() {
let frames = [
r#"{"type":"message_start","message":{"id":"msg_1","role":"assistant","content":[],"model":"claude-sonnet-4-6","stop_reason":null,"stop_sequence":null,"usage":{"input_tokens":10,"output_tokens":0,"cache_creation_input_tokens":4,"cache_read_input_tokens":6,"cache_creation":{"ephemeral_1h_input_tokens":3,"ephemeral_5m_input_tokens":1}}}}"#,
r#"{"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":3}}"#,
];
let decoded = crate::driver::feed_frames!(
MessagesDecoder::new(false),
document::Message::default(),
"anthropic",
frames.map(|frame| WireFrame::Text(frame.to_owned()))
);
let response = decoded.outcome.expect("the message_delta ends the reply");
assert_eq!(response.usage.input_tokens, Some(10 + 6 + 4));
assert_eq!(response.usage.cache_creation_input_tokens, Some(4));
assert_eq!(response.usage.cached_input_tokens, Some(6));
assert_eq!(response.usage.output_tokens, Some(3));
assert_eq!(response.usage.total_tokens, Some(23));
let usage = &response.raw["usage"];
assert_eq!(usage["cache_creation_input_tokens"], 4);
assert_eq!(usage["cache_read_input_tokens"], 6);
assert_eq!(usage["cache_creation"]["ephemeral_1h_input_tokens"], 3);
}
#[test]
fn a_delta_without_its_type_or_text_fails_the_reply() {
for delta in [
json!({"text": "hello"}),
json!({"type": "text_delta", "text": 42}),
] {
let frames = reply(
vec![block(
0,
json!({"type": "text", "text": ""}),
std::slice::from_ref(&delta),
)],
"end_turn",
);
assert!(streamed(&frames).is_err(), "{delta}");
}
}
#[test]
fn message_start_with_null_message_is_a_known_noop() {
let decoded = decode([classified(r#"{"type":"message_start","message":null}"#)]);
assert!(
decoded.events().is_empty(),
"a message-less message_start is a no-op"
);
}
#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
mod terminal_emission {
use super::super::super::completion::CLAUDE_SONNET_4_6;
use crate::providers::anthropic::wire::AnthropicConfig;
use crate::streaming::{Item, StreamEvent};
use crate::test_utils::MockStreamingClient;
use futures::StreamExt;
const MESSAGE_START: &str = r#"{"type":"message_start","message":{"id":"msg_1","role":"assistant","content":[],"model":"claude-sonnet-4-6","stop_reason":null,"stop_sequence":null,"usage":{"input_tokens":5,"output_tokens":0}}}"#;
const TEXT_START: &str =
r#"{"type":"content_block_start","index":0,"content_block":{"type":"text","text":""}}"#;
const TEXT_DELTA: &str =
r#"{"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"hi"}}"#;
const MESSAGE_DELTA: &str = r#"{"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":3}}"#;
fn sse(frames: &[&str]) -> bytes::Bytes {
bytes::Bytes::from(
frames
.iter()
.map(|frame| format!("data: {frame}\n\n"))
.collect::<String>(),
)
}
async fn collect(
sse_bytes: bytes::Bytes,
) -> (
Vec<String>,
bool,
Result<crate::completion::CompletionResponse, crate::error::ProviderError>,
) {
let bound = crate::driver::Model::new(
AnthropicConfig::new("test-key").completion(CLAUDE_SONNET_4_6),
MockStreamingClient { sse_bytes },
);
let request = crate::completion::CompletionRequest::new("hello");
let stream = bound.stream(request).expect("stream should open");
drain(stream).await
}
async fn drain(
mut stream: crate::streaming::CompletionStream,
) -> (
Vec<String>,
bool,
Result<crate::completion::CompletionResponse, crate::error::ProviderError>,
) {
let mut texts = Vec::new();
let mut saw_error = false;
while let Some(item) = stream.next().await {
match item {
Ok(Item::Event(StreamEvent::Text { text, .. })) => texts.push(text),
Ok(_) => {}
Err(_) => saw_error = true,
}
}
(texts, saw_error, stream.finish().await)
}
#[tokio::test]
async fn errored_stream_forwards_the_error_and_no_end() {
use crate::test_utils::SequencedStreamingHttpClient;
let bound = crate::driver::Model::new(
AnthropicConfig::new("test-key").completion(CLAUDE_SONNET_4_6),
SequencedStreamingHttpClient::new(vec![
Ok(sse(&[MESSAGE_START, TEXT_START, TEXT_DELTA])),
Err(crate::http_client::Error::non_success_with_details(
http::StatusCode::BAD_GATEWAY,
http::HeaderMap::new(),
"connection reset".to_string(),
)),
]),
);
let request = crate::completion::CompletionRequest::new("hello");
let stream = bound.stream(request).expect("stream should open");
let (texts, saw_error, finished) = drain(stream).await;
assert_eq!(texts, ["hi"]);
assert!(saw_error, "the transport failure must reach the consumer");
assert!(finished.is_err(), "a failed stream has no response");
}
#[tokio::test]
async fn input_tokens_prefer_the_terminal_delta_and_fall_back_to_message_start() {
fn message_start(input_tokens: usize) -> String {
format!(
r#"{{"type":"message_start","message":{{"id":"msg_1","role":"assistant","content":[],"model":"claude-sonnet-4-6","stop_reason":null,"stop_sequence":null,"usage":{{"input_tokens":{input_tokens},"output_tokens":0}}}}}}"#
)
}
fn message_delta(input_tokens: usize) -> String {
format!(
r#"{{"type":"message_delta","delta":{{"stop_reason":"end_turn","stop_sequence":null}},"usage":{{"input_tokens":{input_tokens},"output_tokens":3}}}}"#
)
}
for (start, delta, expected, case) in [
(
message_start(0),
message_delta(9),
9,
"a gateway reporting the prompt size on message_delta must reach the consumer",
),
(
message_start(5),
MESSAGE_DELTA.to_owned(),
5,
"a delta without input_tokens falls back to message_start",
),
(
message_start(5),
message_delta(5),
5,
"agreeing frames report that count",
),
(
message_start(5),
message_delta(0),
5,
"a zero on the delta must not erase the message_start count",
),
] {
let (_texts, _saw_error, finished) =
collect(sse(&[&start, TEXT_START, TEXT_DELTA, &delta])).await;
let response = finished.expect("the turn must complete");
assert_eq!(response.usage.input_tokens, Some(expected), "{case}");
}
}
}
mod projection {
use std::sync::Arc;
use crate::completion::CompletionRequest;
use crate::observe::{
Action, AdapterContext, AdapterEnding, AdapterErrorBoundary, AdapterErrorEnvelope,
AdapterEvent, ObservationLog, Subject,
};
use crate::providers::anthropic::wire::{AnthropicConfig, Messages};
use crate::test_utils::RecordingHttpClient;
fn adapter_events(log: &ObservationLog) -> Vec<AdapterEvent> {
log.trace()
.observations
.iter()
.filter_map(|o| match &o.action {
Action::Adapter { observation } => Some(observation.event.clone()),
_ => None,
})
.collect()
}
fn context(log: &Arc<ObservationLog>) -> AdapterContext {
AdapterContext::new(log.clone(), Subject::default(), "call")
}
fn wire() -> Messages {
AnthropicConfig::new("test-key").completion("claude-test")
}
fn request() -> CompletionRequest {
CompletionRequest::new("hello").max_tokens(64)
}
#[tokio::test]
async fn messages_rejection_projects_the_envelope() {
let http = RecordingHttpClient::with_error(
http::StatusCode::SERVICE_UNAVAILABLE,
r#"{"type":"error","error":{"type":"overloaded_error","message":"Overloaded"}}"#,
);
let log = Arc::new(ObservationLog::default());
let error = crate::driver::Model::new(wire(), http.clone())
.call_observed(request(), context(&log))
.await
.expect_err("the transport rejects the call");
assert!(error.is_retryable());
let events = adapter_events(&log);
assert!(
matches!(&events[0], AdapterEvent::Started { method, route } if method == "POST" && route == "/v1/messages")
);
assert!(events.contains(&AdapterEvent::Response { status: 503 }));
assert!(events.contains(&AdapterEvent::ErrorEnvelope {
error: AdapterErrorEnvelope {
code: None,
status: Some("overloaded_error".into()),
message: Some("Overloaded".into()),
}
}));
assert_eq!(
events.last(),
Some(&AdapterEvent::Finished {
ending: AdapterEnding::Error {
boundary: AdapterErrorBoundary::ProviderResponse,
kind: "provider_response".into(),
status: Some(503),
retryable: true,
}
})
);
}
}
#[test]
fn every_messages_event_has_a_sample() {
let samples: Vec<MessagesEvent> = [
r#"{"type":"message_start","message":null}"#,
r#"{"type":"message","id":"msg_1","role":"assistant","model":"m","content":[],"stop_reason":"end_turn","stop_sequence":null,"usage":{"input_tokens":1,"output_tokens":1}}"#,
r#"{"type":"content_block_start","index":0,"content_block":{"type":"text","text":""}}"#,
r#"{"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"x"}}"#,
r#"{"type":"content_block_stop","index":0}"#,
r#"{"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":1}}"#,
r#"{"type":"message_stop"}"#,
r#"{"type":"ping"}"#,
r#"{"type":"error","error":{"type":"overloaded_error","message":"Overloaded"}}"#,
]
.into_iter()
.map(classified)
.collect();
let index = |event: &MessagesEvent| {
KNOWN_EVENT_TYPES
.iter()
.position(|kind| *kind == event.kind())
.unwrap_or(usize::MAX)
};
crate::test_utils::history::assert_every_variant(&samples, index, KNOWN_EVENT_TYPES.len());
}
#[test]
fn an_invented_item_and_field_survive_decode_and_same_model_replay() {
use crate::wire::{Operation, Wire};
let invented = json!({"type": "frobnicate", "payload": {"x": 1}});
let text = json!({"type": "text", "text": "hi", "sparkle": true});
let wire = AnthropicConfig::new("test-key").completion(CLAUDE_SONNET_4_6);
let whole = json!({
"type": "message", "id": "msg_1", "role": "assistant", "model": CLAUDE_SONNET_4_6,
"content": [invented, text], "stop_reason": "end_turn", "stop_sequence": null,
"usage": {"input_tokens": 3, "output_tokens": 1}
});
let unary = crate::test_utils::decode_reply(
&wire,
&CompletionRequest::new("hello"),
Mode::Unary,
[WireFrame::Text(whole.to_string())],
whole.clone(),
)
.expect("the whole reply folds");
let stream = streamed(&reply(
vec![
block(0, invented.clone(), &[]),
block(
1,
json!({"type": "text", "text": "", "sparkle": true}),
&[json!({"type": "text_delta", "text": "hi"})],
),
],
"end_turn",
))
.expect("the stream folds");
assert_eq!(unary.message(), stream.message());
for response in [unary, stream] {
let turn = response.message().expect("an assistant turn");
let request = CompletionRequest::from(vec![
RigMessage::user("hello"),
turn,
RigMessage::user("again"),
]);
let request = Completion::prepare(request, &wire.describe()).expect("the history adapts");
let encoded = wire
.encode(request, Mode::Unary)
.expect("the request encodes");
let body = crate::test_utils::json_body(&encoded.request);
assert_eq!(body["messages"][1]["content"], json!([invented, text]));
}
}
fn whole(message: Value) -> Result<crate::completion::CompletionResponse, ProviderError> {
let wire = AnthropicConfig::new("test-key").completion(CLAUDE_SONNET_4_6);
crate::test_utils::decode_reply(
&wire,
&CompletionRequest::new("hello"),
Mode::Unary,
[WireFrame::Text(message.to_string())],
Value::Null,
)
}
#[test]
fn an_empty_reply_folds_the_same_whole_or_streamed() {
for reason in [
"end_turn",
"stop_sequence",
"max_tokens",
"model_context_window_exceeded",
"tool_use",
"pause_turn",
"refusal",
"sensitive",
] {
let streamed = streamed(&reply(vec![], reason)).map(|response| response.stop());
let whole = whole(json!({
"type": "message", "id": "msg_1", "role": "assistant", "model": CLAUDE_SONNET_4_6,
"content": [], "stop_reason": reason, "stop_sequence": null,
"usage": {"input_tokens": 3, "output_tokens": 5}
}))
.map(|response| response.stop());
assert!(streamed.is_ok(), "{reason}: {streamed:?}");
assert_eq!(
whole.as_ref().ok(),
streamed.as_ref().ok(),
"{reason}: {whole:?}"
);
}
}
#[test]
fn a_stop_reason_states_every_open_block_complete() {
let without_stops = |reason: &str, input: &str| {
let mut frames = reply(Vec::new(), reason);
let end = frames.pop().expect("the message_delta");
frames.extend([
json!({"type": "content_block_start", "index": 0,
"content_block": {"type": "thinking", "thinking": "", "signature": ""}}),
json!({"type": "content_block_delta", "index": 0,
"delta": {"type": "thinking_delta", "thinking": "plan"}}),
json!({"type": "content_block_delta", "index": 0,
"delta": {"type": "signature_delta", "signature": "sig"}}),
json!({"type": "content_block_start", "index": 1,
"content_block": {"type": "tool_use", "id": "toolu_a", "name": "f", "input": {}}}),
json!({"type": "content_block_delta", "index": 1,
"delta": {"type": "input_json_delta", "partial_json": input}}),
end,
]);
frames
};
let response = streamed(&without_stops("tool_use", "{\"x\":1}")).expect("the reply folds");
assert_eq!(response.stop(), StopReason::ToolUse);
assert_eq!(
item(&response.choice[0]),
Some(&json!({"type": "thinking", "thinking": "plan", "signature": "sig"}))
);
assert_eq!(
item(&response.choice[1]),
Some(&json!({"type": "tool_use", "id": "toolu_a", "name": "f", "input": {"x": 1}}))
);
let cut = streamed(&without_stops("max_tokens", "{\"x\":")).expect("the reply folds");
assert_eq!(cut.stop(), StopReason::Length);
assert!(item(&cut.choice[0]).is_some(), "{:?}", cut.choice);
assert!(item(&cut.choice[1]).is_none(), "{:?}", cut.choice);
let unfinished = streamed(&without_stops("tool_use", "{\"x\":")).expect("the reply folds");
assert!(unfinished.stop().is_failure(), "{:?}", unfinished.stop());
}
#[test]
fn a_delta_for_a_block_that_never_started_opens_it() {
let frames = reply(
vec![vec![
json!({"type": "content_block_delta", "index": 0,
"delta": {"type": "text_delta", "text": "hi"}}),
json!({"type": "content_block_stop", "index": 0}),
]],
"end_turn",
);
let response = streamed(&frames).expect("the reply decodes");
assert!(
matches!(response.choice.as_slice(), [AssistantContent::Text(text)] if text.text == "hi"),
"{:?}",
response.choice
);
}
const PROGRESS_UPDATE_TEXT: &str = "Confirmed the retry path never refreshes the expired token. Editing auth.py to add the refresh call.";
fn progress_update_content(update_text: &str) -> Value {
json!([
{"type": "thinking", "thinking": "", "signature": "EqMBCkYICxIM-reasoning"},
{"type": "thinking", "thinking": update_text, "signature": "Es8CCkYICxIM-update"},
{
"type": "tool_use",
"id": "toolu_01D7FLrfh4GYq7yT1ULFeyMV",
"name": "edit_file",
"input": {"path": "auth.py", "content": "..."}
}
])
}
fn progress_update_reply(update_text: &str, mode: Mode) -> crate::completion::CompletionResponse {
let wire = AnthropicConfig::new("test-key").completion(CLAUDE_OPUS_5_5);
let frames: Vec<Value> = match mode {
Mode::Unary => vec![json!({
"type": "message", "id": "msg_progress", "model": CLAUDE_OPUS_5_5,
"role": "assistant", "content": progress_update_content(update_text),
"stop_reason": "tool_use", "stop_sequence": null,
"usage": {"input_tokens": 120, "output_tokens": 80}
})],
_ => {
let thinking = |index, text: &str, signature: &str| {
block(
index,
json!({"type": "thinking", "thinking": "", "signature": ""}),
&[
json!({"type": "thinking_delta", "thinking": text}),
json!({"type": "signature_delta", "signature": signature}),
],
)
};
std::iter::once(json!({"type": "message_start", "message": {
"type": "message", "id": "msg_progress", "model": CLAUDE_OPUS_5_5,
"role": "assistant", "content": [], "stop_reason": null, "stop_sequence": null,
"usage": {"input_tokens": 120, "output_tokens": 1}
}}))
.chain(thinking(0, "", "EqMBCkYICxIM-reasoning"))
.chain(thinking(1, update_text, "Es8CCkYICxIM-update"))
.chain(block(
2,
json!({"type": "tool_use", "id": "toolu_01D7FLrfh4GYq7yT1ULFeyMV",
"name": "edit_file", "input": {}}),
&[json!({"type": "input_json_delta",
"partial_json": "{\"path\": \"auth.py\", \"content\": \"...\"}"})],
))
.chain([
json!({"type": "message_delta",
"delta": {"stop_reason": "tool_use", "stop_sequence": null},
"usage": {"output_tokens": 80}}),
json!({"type": "message_stop"}),
])
.collect()
}
};
crate::test_utils::decode_reply(
&wire,
&CompletionRequest::new("hello"),
mode,
frames
.iter()
.map(|frame| WireFrame::Text(frame.to_string())),
Value::Null,
)
.expect("the reply folds")
}
#[test]
fn progress_update_blocks_survive_decoding_and_replay_under_every_display() {
use crate::wire::{Operation, Wire};
let wire = AnthropicConfig::new("test-key").completion(CLAUDE_OPUS_5_5);
for update_text in ["", PROGRESS_UPDATE_TEXT] {
for mode in [Mode::Unary, Mode::Streaming] {
let response = progress_update_reply(update_text, mode);
let items: Vec<Value> = response.choice.iter().filter_map(item).cloned().collect();
assert_eq!(
Value::Array(items),
progress_update_content(update_text),
"{mode:?}: each block decodes as received"
);
let turn = response.message().expect("an assistant turn");
let mut request = CompletionRequest::from(vec![
RigMessage::user("The login test fails after an hour of uptime."),
turn,
]);
request.chat_history.push(RigMessage::tool_result(
crate::message::CallId::from_wire("toolu_01D7FLrfh4GYq7yT1ULFeyMV"),
crate::message::ToolName::new("edit_file").expect("a tool name"),
"saved",
));
request.tools = vec![crate::completion::ToolDefinition {
name: crate::message::ToolName::new("edit_file").expect("a tool name"),
description: "Edit a file.".to_owned(),
parameters: json!({"type": "object", "properties": {}}),
}];
let request =
Completion::prepare(request, &wire.describe()).expect("the history adapts");
let encoded = wire
.encode(request, Mode::Unary)
.expect("the request encodes");
let body = crate::test_utils::json_body(&encoded.request);
assert_eq!(
body["messages"][1]["content"],
progress_update_content(update_text),
"the {mode:?} turn replays exactly as received"
);
}
}
}
fn cited_blocks() -> Vec<(&'static str, Value)> {
vec![
(
"The grass is green.",
json!({"type": "char_location", "cited_text": "The grass is green.",
"document_index": 0, "document_title": "Lawn", "start_char_index": 0,
"end_char_index": 20, "file_id": "file_1"}),
),
(
"Pages two and three.",
json!({"type": "page_location", "cited_text": "Two. Three.",
"document_index": 1, "document_title": null, "start_page_number": 2,
"end_page_number": 4}),
),
(
"Block one.",
json!({"type": "content_block_location", "cited_text": "One.",
"document_index": 2, "document_title": "Blocks", "start_block_index": 1,
"end_block_index": 2}),
),
(
"Café opens at nine.",
json!({"type": "search_result_location", "cited_text": "Opens 9am.",
"source": "https://example.com/cafe", "title": "Café", "search_result_index": 3,
"start_block_index": 0, "end_block_index": 1}),
),
(
"Rust 2.0 shipped.",
json!({"type": "web_search_result_location", "cited_text": "Rust 2.0 is out.",
"url": "https://example.com/rust", "title": "Rust", "encrypted_index": "Eo8B"}),
),
(
"Unknown.",
json!({"type": "frobnicate_location", "cited_text": "?"}),
),
]
}
fn expected_sources() -> Vec<Vec<Source>> {
let document = |index, within| SourceLocation::Document {
index: Some(index),
id: None,
within: Some(within),
};
vec![
vec![
Source::new(SourceLocation::Document {
index: Some(0),
id: Some("file_1".to_owned()),
within: Some(DocumentRange::Chars(0..20)),
})
.title("Lawn")
.cited_text("The grass is green."),
],
vec![Source::new(document(1, DocumentRange::Pages(2..4))).cited_text("Two. Three.")],
vec![
Source::new(document(2, DocumentRange::Blocks(1..2)))
.title("Blocks")
.cited_text("One."),
],
vec![
Source::new(SourceLocation::SearchResult {
index: 3,
source: "https://example.com/cafe".to_owned(),
blocks: Some(0..1),
})
.title("Café")
.cited_text("Opens 9am."),
],
vec![
Source::new(SourceLocation::Url {
url: "https://example.com/rust".to_owned(),
})
.title("Rust")
.cited_text("Rust 2.0 is out."),
],
vec![],
]
}
fn citations_of(
response: &crate::completion::CompletionResponse,
) -> Vec<(String, Vec<(Option<std::ops::Range<usize>>, Vec<Source>)>)> {
response
.choice
.iter()
.filter_map(|content| match content {
AssistantContent::Text(text) => Some((
text.text.clone(),
text.citations()
.iter()
.map(|citation| {
(
citation.span.map(|span| span.range()),
citation.sources.clone(),
)
})
.collect(),
)),
_ => None,
})
.collect()
}
#[test]
fn every_citation_kind_resolves_the_same_unary_and_streamed() {
use crate::providers::anthropic::wire::{ANTHROPIC, MINIMAX, MOONSHOT, XIAOMIMIMO, ZAI};
let blocks = cited_blocks();
let expected: Vec<_> = blocks
.iter()
.zip(expected_sources())
.map(|((text, _), sources)| {
let cited = (!sources.is_empty()).then_some((None, sources));
((*text).to_owned(), cited.into_iter().collect::<Vec<_>>())
})
.collect();
let frames = reply(
blocks
.iter()
.enumerate()
.map(|(index, (text, citation))| {
block(
index,
json!({"type": "text", "text": "", "citations": null}),
&[
json!({"type": "citations_delta", "citation": citation}),
json!({"type": "text_delta", "text": text}),
],
)
})
.collect(),
"end_turn",
);
for (dialect, model) in [
(&ANTHROPIC, CLAUDE_SONNET_4_6),
(&ZAI, "glm-5"),
(&MINIMAX, "MiniMax-M2.7"),
(&MOONSHOT, "kimi-k3"),
(&XIAOMIMIMO, "mimo-v2.5"),
] {
let whole = json!({
"type": "message", "id": "msg_1", "role": "assistant", "model": model,
"content": blocks.iter().map(|(text, citation)| json!({
"type": "text", "text": text, "citations": [citation]
})).collect::<Vec<_>>(),
"stop_reason": "end_turn", "stop_sequence": null,
"usage": {"input_tokens": 3, "output_tokens": 1}
});
let wire = AnthropicConfig::with_key(dialect, "test-key").completion(model);
let decoded = |mode, frames: Vec<WireFrame>| {
crate::test_utils::decode_reply(
&wire,
&CompletionRequest::new("hello"),
mode,
frames,
whole.clone(),
)
.expect("the reply folds")
};
let unary = decoded(Mode::Unary, vec![WireFrame::Text(whole.to_string())]);
let stream = decoded(
Mode::Streaming,
frames
.iter()
.map(|frame| WireFrame::Text(frame.to_string()))
.collect(),
);
assert_eq!(citations_of(&unary), expected, "{}", dialect.name);
assert_eq!(citations_of(&stream), expected, "{}", dialect.name);
for response in [&unary, &stream] {
for (content, (_, citation)) in response.choice.iter().zip(&blocks) {
assert_eq!(
content.native_item().and_then(|item| item.get("citations")),
Some(&json!([citation])),
"{}",
dialect.name
);
}
}
}
}
#[test]
fn a_citation_without_its_source_stays_native() {
for citation in [
json!({"type": "web_search_result_location", "cited_text": "x"}),
json!({"type": "search_result_location", "cited_text": "x", "search_result_index": 0}),
json!("not an object"),
] {
let response = streamed(&reply(
vec![block(
0,
json!({"type": "text", "text": "Claim."}),
&[json!({"type": "citations_delta", "citation": citation})],
)],
"end_turn",
))
.expect("the stream folds");
let [AssistantContent::Text(text)] = response.choice.as_slice() else {
panic!("one text block: {:?}", response.choice);
};
assert!(text.citations().is_empty(), "{citation}");
assert_eq!(
response.choice[0]
.native_item()
.and_then(|item| item.get("citations")),
Some(&json!([citation]))
);
}
}
#[test]
fn a_citation_without_its_passage_or_title_still_cites_its_source() {
let citation = json!({"type": "web_search_result_location",
"url": "https://example.com/rust"});
let expected = vec![(
"Rust 2.0 shipped.".to_owned(),
vec![(
None,
vec![Source::new(SourceLocation::Url {
url: "https://example.com/rust".to_owned(),
})],
)],
)];
let opened = reply(
vec![block(
0,
json!({"type": "text", "text": "", "citations": [citation]}),
&[json!({"type": "text_delta", "text": "Rust 2.0 shipped."})],
)],
"end_turn",
);
let delta = reply(
vec![block(
0,
json!({"type": "text", "text": ""}),
&[
json!({"type": "citations_delta", "citation": citation}),
json!({"type": "text_delta", "text": "Rust 2.0 shipped."}),
],
)],
"end_turn",
);
for frames in [opened, delta] {
let response = streamed(&frames).expect("the stream folds");
assert_eq!(citations_of(&response), expected);
}
}