use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use wiremock::matchers::{header, method, path};
use wiremock::{Mock, MockServer, Request, ResponseTemplate};
use yoagent::provider::{
AnthropicCompat, AnthropicProvider, ModelConfig, StreamConfig, StreamProvider,
};
use yoagent::types::*;
struct HeaderAbsent(&'static str);
impl wiremock::Match for HeaderAbsent {
fn matches(&self, request: &Request) -> bool {
!request.headers.contains_key(self.0)
}
}
fn sse_empty_with_stop(stop_reason: &str) -> String {
format!(
"event: message_start\n\
data: {{\"type\":\"message_start\",\"message\":{{\"usage\":{{\"input_tokens\":10,\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}}}}\n\n\
event: message_delta\n\
data: {{\"type\":\"message_delta\",\"delta\":{{\"stop_reason\":\"{stop_reason}\"}},\"usage\":{{\"output_tokens\":0}}}}\n\n\
event: message_stop\n\
data: {{\"type\":\"message_stop\"}}\n\n"
)
}
fn stream_config(base_url: &str, anthropic: Option<AnthropicCompat>) -> StreamConfig {
let mut mc = ModelConfig::anthropic("claude-sonnet-5", "Claude Sonnet 5");
mc.base_url = base_url.to_string();
mc.anthropic = anthropic;
let mut config = StreamConfig::new("claude-sonnet-5", "test-key");
config.system_prompt = "test".into();
config.messages = vec![Message::user("hi")];
config.max_tokens = Some(256);
config.model_config = Some(mc);
config
}
async fn run_stream(config: StreamConfig) -> Result<Message, yoagent::provider::ProviderError> {
let (tx, _rx) = mpsc::unbounded_channel();
AnthropicProvider
.stream(config, tx, CancellationToken::new())
.await
}
#[tokio::test]
async fn refusal_stop_reason_maps_to_refusal_with_error_message() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(
ResponseTemplate::new(200)
.set_body_raw(sse_empty_with_stop("refusal"), "text/event-stream"),
)
.mount(&server)
.await;
let message = run_stream(stream_config(&server.uri(), None))
.await
.expect("stream should succeed");
let Message::Assistant {
stop_reason,
error_message,
..
} = &message
else {
panic!("expected assistant message");
};
assert_eq!(*stop_reason, StopReason::Refusal);
assert!(
error_message.as_deref().unwrap_or("").contains("refusal"),
"error_message should explain the refusal, got {error_message:?}"
);
}
#[tokio::test]
async fn context_window_exceeded_maps_to_overflow_error() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(ResponseTemplate::new(200).set_body_raw(
sse_empty_with_stop("model_context_window_exceeded"),
"text/event-stream",
))
.mount(&server)
.await;
let message = run_stream(stream_config(&server.uri(), None))
.await
.expect("stream should succeed");
assert!(
message.is_context_overflow(),
"in-stream overflow must trigger the documented recovery hook"
);
}
#[tokio::test]
async fn bearer_auth_sends_authorization_and_no_x_api_key() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/messages"))
.and(header("authorization", "Bearer test-key"))
.and(HeaderAbsent("x-api-key"))
.respond_with(
ResponseTemplate::new(200)
.set_body_raw(sse_empty_with_stop("end_turn"), "text/event-stream"),
)
.expect(1)
.mount(&server)
.await;
let config = stream_config(
&server.uri(),
Some(AnthropicCompat {
adaptive_thinking: true,
bearer_auth: true,
}),
);
run_stream(config).await.expect("stream should succeed");
}
#[tokio::test]
async fn default_auth_sends_x_api_key() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/messages"))
.and(header("x-api-key", "test-key"))
.respond_with(
ResponseTemplate::new(200)
.set_body_raw(sse_empty_with_stop("end_turn"), "text/event-stream"),
)
.expect(1)
.mount(&server)
.await;
run_stream(stream_config(&server.uri(), None))
.await
.expect("stream should succeed");
}
#[tokio::test]
async fn user_authorization_header_suppresses_x_api_key() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/messages"))
.and(header("authorization", "Bearer custom-token"))
.and(HeaderAbsent("x-api-key"))
.respond_with(
ResponseTemplate::new(200)
.set_body_raw(sse_empty_with_stop("end_turn"), "text/event-stream"),
)
.expect(1)
.mount(&server)
.await;
let mut config = stream_config(&server.uri(), None);
if let Some(mc) = &mut config.model_config {
mc.headers
.insert("Authorization".into(), "Bearer custom-token".into());
}
run_stream(config).await.expect("stream should succeed");
}
#[tokio::test]
async fn rate_limit_carries_retry_after_from_header() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(
ResponseTemplate::new(429)
.insert_header("retry-after", "7")
.set_body_string(
r#"{"type":"error","error":{"type":"rate_limit_error","message":"rate limited"}}"#,
),
)
.mount(&server)
.await;
let err = run_stream(stream_config(&server.uri(), None))
.await
.expect_err("429 must surface as an error");
match err {
yoagent::provider::ProviderError::RateLimited { retry_after_ms } => {
assert_eq!(retry_after_ms, Some(7000));
}
other => panic!("expected RateLimited, got: {:?}", other),
}
}
#[tokio::test]
async fn provider_comes_from_model_config_not_hardcoded() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(
ResponseTemplate::new(200)
.set_body_raw(sse_empty_with_stop("end_turn"), "text/event-stream"),
)
.mount(&server)
.await;
let mut mc = ModelConfig::opencode_zen("claude-sonnet-5");
mc.base_url = server.uri();
assert_eq!(mc.provider, "opencode-zen", "preset sets the provider name");
let mut config = stream_config(&server.uri(), None);
config.model_config = Some(mc);
let message = run_stream(config).await.expect("stream should succeed");
let Message::Assistant { provider, .. } = &message else {
panic!("expected assistant message");
};
assert_eq!(
provider, "opencode-zen",
"provider must be propagated from ModelConfig, not hardcoded"
);
}
#[tokio::test]
async fn stream_ended_without_stop_reason_is_retryable_network_error() {
let server = MockServer::start().await;
let truncated = "event: message_start\n\
data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":10,\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}\n\n";
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(ResponseTemplate::new(200).set_body_raw(truncated, "text/event-stream"))
.mount(&server)
.await;
let err = run_stream(stream_config(&server.uri(), None))
.await
.expect_err("truncation before stop_reason must be an error");
assert!(
matches!(err, yoagent::provider::ProviderError::Network(_)),
"expected retryable Network, got: {err:?}"
);
assert!(err.is_retryable(), "truncation must be retryable");
}
#[tokio::test]
async fn stream_ended_after_stop_reason_is_clean_eof() {
let server = MockServer::start().await;
let no_terminator = "event: message_start\n\
data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":10,\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}\n\n\
event: content_block_start\n\
data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n\n\
event: content_block_delta\n\
data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"hello\"}}\n\n\
event: message_delta\n\
data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"},\"usage\":{\"output_tokens\":5}}\n\n";
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(ResponseTemplate::new(200).set_body_raw(no_terminator, "text/event-stream"))
.mount(&server)
.await;
let message = run_stream(stream_config(&server.uri(), None))
.await
.expect("close after message_delta must not be an error");
let Message::Assistant {
stop_reason,
content,
usage,
..
} = &message
else {
panic!("expected assistant message");
};
assert_eq!(*stop_reason, StopReason::Stop);
assert_eq!(usage.output, 5, "usage from message_delta must survive");
let text: String = content
.iter()
.filter_map(|c| match c {
Content::Text { text } => Some(text.clone()),
_ => None,
})
.collect();
assert_eq!(
text, "hello",
"content must survive the terminator-less close"
);
}
#[tokio::test]
async fn intermediate_message_delta_does_not_arm_the_clean_eof_guard() {
let server = MockServer::start().await;
let body = "event: message_start\n\
data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":10,\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}\n\n\
event: content_block_start\n\
data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n\n\
event: content_block_delta\n\
data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"partial\"}}\n\n\
event: message_delta\n\
data: {\"type\":\"message_delta\",\"delta\":{},\"usage\":{\"output_tokens\":3}}\n\n";
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(ResponseTemplate::new(200).set_body_raw(body, "text/event-stream"))
.mount(&server)
.await;
let err = run_stream(stream_config(&server.uri(), None))
.await
.expect_err("a delta without stop_reason must not mark the response complete");
assert!(
err.is_retryable(),
"expected retryable truncation, got: {err:?}"
);
}
#[tokio::test]
async fn clean_eof_guard_rejects_unterminated_tool_call() {
let server = MockServer::start().await;
let body = "event: message_start\n\
data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":10,\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}\n\n\
event: content_block_start\n\
data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tu_1\",\"name\":\"bash\",\"input\":{}}}\n\n\
event: content_block_delta\n\
data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"command\\\": \\\"rm -rf /tm\"}}\n\n\
event: message_delta\n\
data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"tool_use\"},\"usage\":{\"output_tokens\":9}}\n\n";
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(ResponseTemplate::new(200).set_body_raw(body, "text/event-stream"))
.mount(&server)
.await;
let err = run_stream(stream_config(&server.uri(), None))
.await
.expect_err("an unterminated tool_use block must not be returned as success");
assert!(
err.is_retryable(),
"expected retryable truncation, got: {err:?}"
);
}
#[tokio::test]
async fn message_delta_without_usage_still_yields_its_stop_reason() {
let server = MockServer::start().await;
let body = "event: message_start\n\
data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":10,\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}\n\n\
event: message_delta\n\
data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"max_tokens\"}}\n\n\
event: message_stop\n\
data: {\"type\":\"message_stop\"}\n\n";
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(ResponseTemplate::new(200).set_body_raw(body, "text/event-stream"))
.mount(&server)
.await;
let message = run_stream(stream_config(&server.uri(), None))
.await
.expect("a usage-less message_delta must still parse");
let Message::Assistant { stop_reason, .. } = &message else {
panic!("expected assistant message");
};
assert_eq!(
*stop_reason,
StopReason::Length,
"stop_reason must survive a message_delta with no usage field"
);
}
#[tokio::test]
async fn malformed_tool_arguments_fail_the_turn_instead_of_defaulting() {
let server = MockServer::start().await;
let body = "event: message_start\n\
data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":10,\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}\n\n\
event: content_block_start\n\
data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tu_1\",\"name\":\"bash\",\"input\":{}}}\n\n\
event: content_block_delta\n\
data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{not valid json\"}}\n\n\
event: content_block_stop\n\
data: {\"type\":\"content_block_stop\",\"index\":0}\n\n\
event: message_delta\n\
data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"tool_use\"},\"usage\":{\"output_tokens\":9}}\n\n\
event: message_stop\n\
data: {\"type\":\"message_stop\"}\n\n";
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(ResponseTemplate::new(200).set_body_raw(body, "text/event-stream"))
.mount(&server)
.await;
let message = run_stream(stream_config(&server.uri(), None))
.await
.expect("the stream itself is well-formed");
let Message::Assistant {
stop_reason,
error_message,
content,
..
} = &message
else {
panic!("expected assistant message");
};
assert_eq!(
*stop_reason,
StopReason::Error,
"a tool call we cannot parse must not be presented as a usable turn"
);
assert!(
error_message.as_deref().unwrap_or("").contains("bash"),
"the error must name the tool, got: {error_message:?}"
);
assert!(
!content
.iter()
.any(|c| matches!(c, Content::ToolCall { .. })),
"the unusable tool call must not remain in the message"
);
assert!(
!format!("{content:?}").contains("__partial_json"),
"the accumulator sentinel must not escape: {content:?}"
);
let Some(Content::Text { text }) = content.first() else {
panic!("the dropped tool call must leave a text block behind: {content:?}");
};
assert!(
text.contains("bash"),
"the replacement must record which tool was dropped, got: {text}"
);
assert_eq!(content.len(), 1, "no other blocks expected: {content:?}");
assert!(
error_message
.as_deref()
.unwrap_or("")
.contains("{not valid json"),
"the error must quote the unparseable input, got: {error_message:?}"
);
}
#[tokio::test]
async fn content_block_stop_without_an_index_does_not_close_block_zero() {
let server = MockServer::start().await;
let body = "event: message_start\n\
data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":10,\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}\n\n\
event: content_block_start\n\
data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tu_1\",\"name\":\"bash\",\"input\":{}}}\n\n\
event: content_block_delta\n\
data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"cmd\\\":\"}}\n\n\
event: content_block_stop\n\
data: {\"type\":\"content_block_stop\"}\n\n\
event: content_block_delta\n\
data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"\\\"ls\\\"}\"}}\n\n\
event: content_block_stop\n\
data: {\"type\":\"content_block_stop\",\"index\":0}\n\n\
event: message_delta\n\
data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"tool_use\"},\"usage\":{\"output_tokens\":9}}\n\n\
event: message_stop\n\
data: {\"type\":\"message_stop\"}\n\n";
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(ResponseTemplate::new(200).set_body_raw(body, "text/event-stream"))
.mount(&server)
.await;
let message = run_stream(stream_config(&server.uri(), None))
.await
.expect("stream should succeed");
let Message::Assistant {
content,
stop_reason,
..
} = &message
else {
panic!("expected assistant message");
};
assert_eq!(
*stop_reason,
StopReason::ToolUse,
"an index-less stop must not close block 0 mid-accumulation and fail the turn"
);
let Some(Content::ToolCall { arguments, .. }) = content
.iter()
.find(|c| matches!(c, Content::ToolCall { .. }))
else {
panic!("the tool call must survive: {content:?}");
};
assert_eq!(
arguments["cmd"], "ls",
"arguments must assemble fully: {arguments:?}"
);
}
#[derive(Clone, Default)]
struct CapturedLogs(std::sync::Arc<std::sync::Mutex<Vec<String>>>);
impl CapturedLogs {
fn contains(&self, needle: &str) -> bool {
self.0.lock().unwrap().iter().any(|l| l.contains(needle))
}
}
impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for CapturedLogs {
fn on_event(
&self,
event: &tracing::Event<'_>,
_ctx: tracing_subscriber::layer::Context<'_, S>,
) {
struct Msg(String);
impl tracing::field::Visit for Msg {
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
if field.name() == "message" {
self.0 = format!("{value:?}");
}
}
}
let mut msg = Msg(String::new());
event.record(&mut msg);
self.0
.lock()
.unwrap()
.push(format!("{}: {}", event.metadata().level(), msg.0));
}
}
async fn run_stream_capturing_logs(stop_reason: &str) -> CapturedLogs {
use tracing_subscriber::layer::SubscriberExt;
let logs = CapturedLogs::default();
let _guard =
tracing::subscriber::set_default(tracing_subscriber::registry().with(logs.clone()));
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(
ResponseTemplate::new(200)
.set_body_raw(sse_empty_with_stop(stop_reason), "text/event-stream"),
)
.mount(&server)
.await;
run_stream(stream_config(&server.uri(), None))
.await
.unwrap_or_else(|e| panic!("[{stop_reason}] stream should succeed: {e}"));
logs
}
#[tokio::test]
async fn a_healthy_turn_logs_no_stop_reason_warning() {
for reason in ["end_turn", "stop_sequence"] {
let logs = run_stream_capturing_logs(reason).await;
assert!(
!logs.contains("unrecognized Anthropic stop_reason"),
"[{reason}] a recognized stop reason must not warn; captured: {:?}",
logs.0.lock().unwrap()
);
}
}
#[tokio::test]
async fn an_unrecognized_stop_reason_is_logged() {
let logs = run_stream_capturing_logs("reason_that_does_not_exist").await;
assert!(
logs.contains("unrecognized Anthropic stop_reason"),
"an unrecognized stop reason must be logged; captured: {:?}",
logs.0.lock().unwrap()
);
}
#[tokio::test]
async fn end_turn_and_stop_sequence_are_recognized_stop_reasons() {
for reason in ["end_turn", "stop_sequence"] {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(
ResponseTemplate::new(200)
.set_body_raw(sse_empty_with_stop(reason), "text/event-stream"),
)
.mount(&server)
.await;
let message = run_stream(stream_config(&server.uri(), None))
.await
.unwrap_or_else(|e| panic!("[{reason}] stream should succeed: {e}"));
let Message::Assistant { stop_reason, .. } = &message else {
panic!("[{reason}] expected assistant message");
};
assert_eq!(*stop_reason, StopReason::Stop, "[{reason}]");
}
}
#[tokio::test]
async fn pause_turn_is_reported_as_incomplete_rather_than_finished() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(
ResponseTemplate::new(200)
.set_body_raw(sse_empty_with_stop("pause_turn"), "text/event-stream"),
)
.mount(&server)
.await;
let message = run_stream(stream_config(&server.uri(), None))
.await
.expect("stream should succeed");
let Message::Assistant {
stop_reason,
error_message,
..
} = &message
else {
panic!("expected assistant message");
};
assert_eq!(
*stop_reason,
StopReason::Error,
"a paused turn is not a finished one"
);
assert!(
error_message
.as_deref()
.unwrap_or("")
.contains("pause_turn"),
"the error must name the cause, got: {error_message:?}"
);
}
#[tokio::test]
async fn unrecognized_stop_reason_falls_back_to_a_normal_stop() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(ResponseTemplate::new(200).set_body_raw(
sse_empty_with_stop("reason_that_does_not_exist"),
"text/event-stream",
))
.mount(&server)
.await;
let message = run_stream(stream_config(&server.uri(), None))
.await
.expect("an unrecognized stop reason must not fail the stream");
let Message::Assistant { stop_reason, .. } = &message else {
panic!("expected assistant message");
};
assert_eq!(*stop_reason, StopReason::Stop);
}
#[tokio::test]
async fn an_unfinalized_tool_call_never_reaches_the_caller() {
for (label, stop_line) in [
("index-less", r#"data: {"type":"content_block_stop"}"#),
("non-JSON body", "data: not-json-at-all"),
] {
let server = MockServer::start().await;
let body = format!(
"event: message_start\n\
data: {{\"type\":\"message_start\",\"message\":{{\"usage\":{{\"input_tokens\":10,\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}}}}\n\n\
event: content_block_start\n\
data: {{\"type\":\"content_block_start\",\"index\":0,\"content_block\":{{\"type\":\"tool_use\",\"id\":\"tu_1\",\"name\":\"bash\",\"input\":{{}}}}}}\n\n\
event: content_block_delta\n\
data: {{\"type\":\"content_block_delta\",\"index\":0,\"delta\":{{\"type\":\"input_json_delta\",\"partial_json\":\"{{\\\"cmd\\\":\\\"ls\\\"}}\"}}}}\n\n\
event: content_block_stop\n\
{stop_line}\n\n\
event: message_delta\n\
data: {{\"type\":\"message_delta\",\"delta\":{{\"stop_reason\":\"tool_use\"}},\"usage\":{{\"output_tokens\":9}}}}\n\n\
event: message_stop\n\
data: {{\"type\":\"message_stop\"}}\n\n"
);
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(ResponseTemplate::new(200).set_body_raw(body, "text/event-stream"))
.mount(&server)
.await;
let message = run_stream(stream_config(&server.uri(), None))
.await
.unwrap_or_else(|e| panic!("[{label}] stream should not error: {e}"));
let Message::Assistant {
content,
stop_reason,
..
} = &message
else {
panic!("[{label}] expected assistant message");
};
assert!(
!format!("{content:?}").contains("__partial_json"),
"[{label}] the accumulator escaped as tool arguments: {content:?}"
);
assert_ne!(
*stop_reason,
StopReason::ToolUse,
"[{label}] an unfinalized tool call must not be presented as runnable"
);
}
}
#[tokio::test]
async fn a_healthy_sibling_tool_call_does_not_dangle() {
let server = MockServer::start().await;
let body = "event: message_start\n\
data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":10,\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}\n\n\
event: content_block_start\n\
data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tu_0\",\"name\":\"read\",\"input\":{}}}\n\n\
event: content_block_delta\n\
data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{\\\"p\\\":\\\"a\\\"}\"}}\n\n\
event: content_block_stop\n\
data: {\"type\":\"content_block_stop\",\"index\":0}\n\n\
event: content_block_start\n\
data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"tool_use\",\"id\":\"tu_1\",\"name\":\"bash\",\"input\":{}}}\n\n\
event: content_block_delta\n\
data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"input_json_delta\",\"partial_json\":\"{not valid\"}}\n\n\
event: content_block_stop\n\
data: {\"type\":\"content_block_stop\",\"index\":1}\n\n\
event: message_delta\n\
data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"tool_use\"},\"usage\":{\"output_tokens\":9}}\n\n\
event: message_stop\n\
data: {\"type\":\"message_stop\"}\n\n";
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(ResponseTemplate::new(200).set_body_raw(body, "text/event-stream"))
.mount(&server)
.await;
let message = run_stream(stream_config(&server.uri(), None))
.await
.expect("stream should succeed");
let Message::Assistant {
content,
stop_reason,
error_message,
..
} = &message
else {
panic!("expected assistant message");
};
assert_eq!(*stop_reason, StopReason::Error);
assert!(error_message.as_deref().unwrap_or("").contains("bash"));
assert!(
!content
.iter()
.any(|c| matches!(c, Content::ToolCall { .. })),
"no tool_use may survive an errored turn unanswered: {content:?}"
);
}
#[tokio::test]
async fn a_trailing_empty_delta_preserves_stop_reason_and_usage() {
let server = MockServer::start().await;
let body = "event: message_start\n\
data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":10,\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}\n\n\
event: message_delta\n\
data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"refusal\"},\"usage\":{\"output_tokens\":42}}\n\n\
event: message_delta\n\
data: {\"type\":\"message_delta\",\"delta\":{}}\n\n\
event: message_stop\n\
data: {\"type\":\"message_stop\"}\n\n";
Mock::given(method("POST"))
.and(path("/messages"))
.respond_with(ResponseTemplate::new(200).set_body_raw(body, "text/event-stream"))
.mount(&server)
.await;
let message = run_stream(stream_config(&server.uri(), None))
.await
.expect("stream should succeed");
let Message::Assistant {
stop_reason, usage, ..
} = &message
else {
panic!("expected assistant message");
};
assert_eq!(
*stop_reason,
StopReason::Refusal,
"a trailing delta must not downgrade a terminal stop reason"
);
assert_eq!(
usage.output, 42,
"a usage-less trailing delta must not zero the count"
);
}