use async_trait::async_trait;
use serde::Deserialize;
use serde_json::{Value, json};
use crate::core::Secret;
#[cfg(test)]
use super::ModelCall;
use super::wire::{RESPOND_TOOL, classify_status, classify_transport, structured};
use super::{
Completion, ModelError, ModelId, ModelProvider, Request, SchemaMode, Usage, anthropic_stream,
sse,
};
pub struct Anthropic {
http: reqwest::Client,
key: Secret,
base: String,
version: String,
default_schema_mode: SchemaMode,
schema_modes: std::collections::BTreeMap<String, SchemaMode>,
stream: bool,
egress: Option<crate::core::Egress>,
timeout: std::time::Duration,
}
impl std::fmt::Debug for Anthropic {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Anthropic")
.field("base", &self.base)
.field("version", &self.version)
.field("key", &"<redacted>")
.finish_non_exhaustive()
}
}
impl Anthropic {
pub const DEFAULT_TIMEOUT: std::time::Duration = std::time::Duration::from_mins(5);
pub const VERSION: &'static str = "2023-06-01";
pub fn new(key: impl Into<String>) -> Result<Self, ModelError> {
let http = reqwest::Client::builder()
.build()
.map_err(|e| ModelError::Unreachable {
model: ModelId::new("anthropic", "*"),
detail: format!("could not build an HTTP client: {e}"),
})?;
Ok(Self {
http,
key: Secret::new(key),
base: "https://api.anthropic.com".to_owned(),
version: Self::VERSION.to_owned(),
default_schema_mode: SchemaMode::Native,
schema_modes: std::collections::BTreeMap::new(),
stream: true,
egress: None,
timeout: Self::DEFAULT_TIMEOUT,
})
}
#[must_use]
pub const fn timeout(mut self, timeout: std::time::Duration) -> Self {
self.timeout = timeout;
self
}
#[must_use]
pub fn base(mut self, base: impl Into<String>) -> Self {
self.base = base.into();
self
}
#[must_use]
pub fn structured_via(mut self, mode: SchemaMode) -> Self {
self.default_schema_mode = mode;
self
}
#[must_use]
pub fn structured_via_for(mut self, model: impl Into<String>, mode: SchemaMode) -> Self {
self.schema_modes.insert(model.into(), mode);
self
}
#[must_use]
pub fn egress(mut self, egress: crate::core::Egress) -> Self {
self.egress = Some(egress);
self
}
#[must_use]
pub const fn buffered(mut self) -> Self {
self.stream = false;
self
}
fn check_egress(&self, model: &ModelId) -> Result<(), ModelError> {
let Some(egress) = &self.egress else {
return Ok(());
};
let host = reqwest::Url::parse(&self.base)
.ok()
.and_then(|u| u.host_str().map(ToOwned::to_owned));
egress
.permits(host.as_deref())
.map_err(|e| ModelError::Refused {
model: model.clone(),
detail: e.to_string(),
})
}
fn mode_for(&self, model: &ModelId) -> SchemaMode {
self.schema_modes
.get(&model.model)
.copied()
.unwrap_or(self.default_schema_mode)
}
}
#[allow(clippy::struct_field_names)]
#[derive(Debug, Deserialize)]
struct ApiUsage {
#[serde(default)]
input_tokens: u64,
#[serde(default)]
output_tokens: u64,
#[serde(default)]
cache_creation_input_tokens: u64,
#[serde(default)]
cache_read_input_tokens: u64,
}
#[derive(Debug, Deserialize)]
struct ApiResponse {
#[serde(default)]
content: Vec<Value>,
#[serde(default)]
usage: Option<ApiUsage>,
#[serde(default)]
stop_reason: Option<String>,
}
impl ApiResponse {
fn usage(&self) -> Usage {
let u = self.usage.as_ref();
let write = u.map_or(0, |u| u.cache_creation_input_tokens);
let read = u.map_or(0, |u| u.cache_read_input_tokens);
Usage {
input_tokens: u.map_or(0, |u| u.input_tokens) + write + read,
output_tokens: u.map_or(0, |u| u.output_tokens),
cache_write_tokens: write,
cache_read_tokens: read,
minor_units: 0,
}
}
fn text(&self) -> String {
self.content
.iter()
.filter(|b| b.get("type").and_then(Value::as_str) == Some("text"))
.filter_map(|b| b.get("text").and_then(Value::as_str))
.collect::<Vec<_>>()
.join("")
}
fn tool_calls(&self) -> Result<Vec<super::ToolCall>, String> {
self.content
.iter()
.filter(|block| block.get("type").and_then(Value::as_str) == Some("tool_use"))
.filter(|block| block.get("name").and_then(Value::as_str) != Some(RESPOND_TOOL))
.map(|block| {
let id = block
.get("id")
.and_then(Value::as_str)
.filter(|value| !value.is_empty())
.ok_or_else(|| {
"Anthropic returned a tool_use block without an id".to_owned()
})?;
let name = block
.get("name")
.and_then(Value::as_str)
.filter(|value| !value.is_empty())
.ok_or_else(|| {
"Anthropic returned a tool_use block without a name".to_owned()
})?;
let arguments = block
.get("input")
.cloned()
.ok_or_else(|| format!("Anthropic tool_use '{id}' has no input arguments"))?;
Ok(super::ToolCall {
id: id.to_owned(),
name: name.to_owned(),
arguments,
})
})
.collect()
}
fn forced_tool_input(&self) -> Option<&Value> {
self.content
.iter()
.find(|b| {
b.get("type").and_then(Value::as_str) == Some("tool_use")
&& b.get("name").and_then(Value::as_str) == Some(RESPOND_TOOL)
})
.and_then(|b| b.get("input"))
}
}
fn continue_with(
messages: Value,
exchanges: &[super::ToolExchange],
continuation: Option<&super::ProviderContinuation>,
) -> Value {
if exchanges.is_empty() {
return messages;
}
let mut out = match messages {
Value::Array(v) => v,
other => vec![json!({ "role": "user", "content": other })],
};
if let Some(state) = continuation.and_then(|state| state.state.as_array()) {
out.extend(state.iter().cloned());
} else {
out.push(json!({
"role": "assistant",
"content": exchanges
.iter()
.map(|e| json!({
"type": "tool_use",
"id": e.call.id,
"name": e.call.name,
"input": e.call.arguments,
}))
.collect::<Vec<_>>(),
}));
}
out.push(tool_results(exchanges));
Value::Array(out)
}
fn tool_results(exchanges: &[super::ToolExchange]) -> Value {
json!({
"role": "user",
"content": exchanges
.iter()
.map(|exchange| json!({
"type": "tool_result",
"tool_use_id": exchange.call.id,
"content": match &exchange.output {
Value::String(value) => value.clone(),
other => other.to_string(),
},
"is_error": exchange.failed,
}))
.collect::<Vec<_>>(),
})
}
fn accumulate_continuation(
completion: &mut Completion,
prior: Option<&super::ProviderContinuation>,
exchanges: &[super::ToolExchange],
) {
let Some(current) = completion.continuation.as_mut() else {
return;
};
let mut transcript = prior
.and_then(|value| value.state.as_array())
.cloned()
.unwrap_or_default();
if !exchanges.is_empty() {
transcript.push(tool_results(exchanges));
}
if let Some(messages) = current.state.as_array() {
transcript.extend(messages.iter().cloned());
}
current.state = Value::Array(transcript);
}
fn messages(prompt: &Value) -> Value {
match prompt {
Value::String(s) => json!([{ "role": "user", "content": s }]),
Value::Array(_) => prompt.clone(),
other => other.get("messages").cloned().unwrap_or_else(|| {
let mut rest = other.clone();
if let Some(map) = rest.as_object_mut() {
map.remove("system");
}
json!([{ "role": "user", "content": rest.to_string() }])
}),
}
}
fn system(prompt: &Value) -> Option<Value> {
prompt.get("system").cloned().filter(|s| !s.is_null())
}
struct Assembled {
text: String,
forced: Option<Value>,
tool_calls: Vec<super::ToolCall>,
usage: Usage,
stop_reason: Option<String>,
continuation: Value,
}
fn interpret(
model: &ModelId,
schema: Option<&Value>,
emulating: bool,
assembled: Assembled,
) -> Result<Completion, ModelError> {
let Assembled {
text,
forced,
tool_calls,
usage,
stop_reason,
continuation,
} = assembled;
if stop_reason.as_deref() == Some("refusal") {
return Err(ModelError::Unusable {
model: model.clone(),
usage,
detail: "the model declined to answer".to_owned(),
});
}
let (text, structured_value) = if emulating {
let Some(value) = forced else {
return Err(ModelError::Unusable {
model: model.clone(),
usage,
detail: "a tool call was forced and no usable arguments came back — \
the model did not honour `tool_choice`, or its streamed \
fragments did not reassemble into JSON"
.to_owned(),
});
};
if let Some(schema) = schema {
super::validate_schema(schema, &value).map_err(|detail| ModelError::Unusable {
model: model.clone(),
usage,
detail,
})?;
}
(value.to_string(), Some(value))
} else {
if text.is_empty() {
return Err(ModelError::Unusable {
model: model.clone(),
usage,
detail: "the answer carried no text content".to_owned(),
});
}
let parsed_schema = structured(schema, &text, model, usage)?;
(text, parsed_schema)
};
let continuation = (!tool_calls.is_empty()).then(|| {
super::ProviderContinuation::new(
"anthropic",
json!([{ "role": "assistant", "content": continuation }]),
)
});
Ok(Completion {
structured: structured_value,
tool_calls,
text,
usage,
truncated: stop_reason.as_deref() == Some("max_tokens"),
stop_reason,
continuation,
})
}
impl Anthropic {
#[cfg(test)]
fn body(
&self,
model: &ModelId,
prompt: &Value,
schema: Option<&Value>,
tools: &[super::ToolDeclaration],
exchanges: &[super::ToolExchange],
) -> Result<Value, ModelError> {
self.body_with_max(
model,
prompt,
ModelCall::DEFAULT_MAX_OUTPUT_TOKENS,
None,
schema,
tools,
exchanges,
None,
)
}
#[allow(clippy::too_many_arguments)]
fn body_with_max(
&self,
model: &ModelId,
prompt: &Value,
max_output_tokens: u32,
reasoning_effort: Option<super::ReasoningEffort>,
schema: Option<&Value>,
tools: &[super::ToolDeclaration],
exchanges: &[super::ToolExchange],
continuation: Option<&super::ProviderContinuation>,
) -> Result<Value, ModelError> {
if let Some(state) = continuation
&& (state.provider != "anthropic" || !state.state.is_array())
{
return Err(ModelError::Refused {
model: model.clone(),
detail: "the continuation was not an Anthropic assistant-content array".to_owned(),
});
}
if reasoning_effort.is_some() && !exchanges.is_empty() && continuation.is_none() {
return Err(ModelError::Refused {
model: model.clone(),
detail: "reasoning-enabled tool continuation requires the signed thinking \
and assistant blocks from the prior Anthropic response"
.to_owned(),
});
}
let mut body = json!({
"model": model.model,
"max_tokens": max_output_tokens,
"messages": continue_with(messages(prompt), exchanges, continuation),
});
if let Some(effort) = reasoning_effort {
if matches!(
effort,
super::ReasoningEffort::None
| super::ReasoningEffort::Minimal
| super::ReasoningEffort::XHigh
) {
return Err(ModelError::Refused {
model: model.clone(),
detail: format!(
"Anthropic adaptive thinking does not support reasoning effort '{}'",
effort.as_str()
),
});
}
body["thinking"] = json!({ "type": "adaptive" });
body["output_config"]["effort"] = json!(effort.as_str());
}
if let Some(system) = system(prompt) {
body["system"] = system;
}
if !tools.is_empty() {
body["tools"] = Value::Array(
tools
.iter()
.map(|t| {
json!({
"name": t.name,
"description": t.description,
"input_schema": t.parameters,
"strict": true,
})
})
.collect(),
);
}
if let Some(schema) = schema {
match self.mode_for(model) {
SchemaMode::Native => {
body["output_config"]["format"] =
json!({ "type": "json_schema", "schema": schema });
}
SchemaMode::ForcedTool => {
if !tools.is_empty() {
return Err(ModelError::Refused {
model: model.clone(),
detail: format!(
"model '{}' has no native structured output here, so a declared \
response schema is obtained by forcing a synthetic tool — which \
cannot be combined with the {} tool(s) this request declares. \
Use a model with native structured output, or drop the schema \
and validate the answer yourself",
model.model,
tools.len()
),
});
}
body["tools"] = json!([{
"name": RESPOND_TOOL,
"description": "Return the answer in the required shape.",
"input_schema": schema,
}]);
body["tool_choice"] = json!({ "type": "tool", "name": RESPOND_TOOL });
}
}
}
if self.stream {
body["stream"] = json!(true);
}
Ok(body)
}
async fn read_buffered(
&self,
response: reqwest::Response,
model: &ModelId,
schema: Option<&Value>,
) -> Result<Completion, ModelError> {
let parsed: ApiResponse = response.json().await.map_err(|e| ModelError::Unusable {
model: model.clone(),
usage: Usage::default(),
detail: format!("the response body did not parse: {e}"),
})?;
let emulating = schema.is_some() && self.mode_for(model) == SchemaMode::ForcedTool;
let usage = parsed.usage();
let tool_calls = parsed.tool_calls().map_err(|detail| ModelError::Unusable {
model: model.clone(),
usage,
detail,
})?;
interpret(
model,
schema,
emulating,
Assembled {
text: parsed.text(),
forced: parsed.forced_tool_input().cloned(),
tool_calls,
usage,
stop_reason: parsed.stop_reason.clone(),
continuation: Value::Array(parsed.content.clone()),
},
)
}
async fn read_streamed(
&self,
response: reqwest::Response,
model: &ModelId,
schema: Option<&Value>,
observer: Option<(&dyn super::ModelStreamObserver, &crate::core::Label)>,
) -> Result<Completion, ModelError> {
use futures_util::StreamExt;
let mut decoder = sse::Decoder::new();
let mut acc = anthropic_stream::Accumulator::new();
let mut body = response.bytes_stream();
while let Some(chunk) = body.next().await {
let chunk = match chunk {
Ok(chunk) => chunk,
Err(e) => return Err(severed(model, &acc, &e.to_string())),
};
let events = decoder
.push(&chunk)
.map_err(|error| ModelError::Unaccounted {
model: model.clone(),
detail: error.to_string(),
})?;
for event in events {
if event.name == "content_block_delta"
&& let Ok(value) = serde_json::from_str::<Value>(&event.data)
&& value
.get("delta")
.and_then(|delta| delta.get("type"))
.and_then(Value::as_str)
== Some("text_delta")
&& let Some(delta) = value
.get("delta")
.and_then(|delta| delta.get("text"))
.and_then(Value::as_str)
&& let Some((observer, label)) = observer
{
observer.event(crate::core::Tainted::with_label(
super::ModelStreamEvent::TextDelta(delta.to_owned()),
label.clone(),
));
}
acc.event(&event.name, &event.data);
}
if let Some(err) = acc.error() {
return Err(stream_error(model, &acc, err));
}
}
if !acc.complete() {
return Err(severed(
model,
&acc,
"the stream ended before `message_stop`",
));
}
let emulating = schema.is_some() && self.mode_for(model) == SchemaMode::ForcedTool;
let usage = acc.billed();
let tool_calls = acc.tool_calls().map_err(|detail| ModelError::Unusable {
model: model.clone(),
usage,
detail,
})?;
let completion = interpret(
model,
schema,
emulating,
Assembled {
text: acc.text().to_owned(),
forced: acc.forced_tool_input(),
tool_calls,
usage,
stop_reason: acc.stop_reason().map(ToOwned::to_owned),
continuation: acc.continuation_content(),
},
)?;
if let Some((observer, label)) = observer {
observer.event(crate::core::Tainted::with_label(
super::ModelStreamEvent::Usage(completion.usage),
label.clone(),
));
}
Ok(completion)
}
}
fn severed(model: &ModelId, acc: &anthropic_stream::Accumulator, detail: &str) -> ModelError {
if acc.started() {
return ModelError::Interrupted {
model: model.clone(),
usage: acc.billed(),
detail: detail.to_owned(),
};
}
ModelError::Unavailable {
model: model.clone(),
detail: format!("the stream ended before it began: {detail}"),
}
}
fn stream_error(
model: &ModelId,
acc: &anthropic_stream::Accumulator,
err: &anthropic_stream::StreamError,
) -> ModelError {
let detail = format!("{}: {}", err.kind, err.message);
if acc.started() {
return ModelError::Interrupted {
model: model.clone(),
usage: acc.billed(),
detail,
};
}
match err.kind.as_str() {
"overloaded_error" | "rate_limit_error" => ModelError::RateLimited {
model: model.clone(),
detail,
},
"invalid_request_error"
| "authentication_error"
| "permission_error"
| "not_found_error" => ModelError::Refused {
model: model.clone(),
detail,
},
_ => ModelError::Unavailable {
model: model.clone(),
detail,
},
}
}
#[async_trait]
impl ModelProvider for Anthropic {
fn request_profile(&self, model: &ModelId) -> Value {
let schema_mode = match self.mode_for(model) {
SchemaMode::Native => "native",
SchemaMode::ForcedTool => "forced-tool",
};
json!({
"driver": "anthropic-messages/v1",
"base": self.base,
"api_version": self.version,
"schema_mode": schema_mode,
"stream": self.stream,
"timeout_ms": self.timeout.as_millis(),
})
}
async fn complete(&self, request: Request<'_>) -> Result<Completion, ModelError> {
let Request {
model,
prompt,
max_output_tokens,
reasoning_effort,
schema,
tools,
exchanges,
continuation,
stream,
} = request;
super::refuse_provider_side_media(prompt, model)?;
self.check_egress(model)?;
let response = self
.http
.post(format!("{}/v1/messages", self.base))
.timeout(self.timeout)
.header("x-api-key", self.key.expose())
.header("anthropic-version", &self.version)
.json(&self.body_with_max(
model,
prompt,
max_output_tokens,
reasoning_effort,
schema,
tools,
exchanges,
continuation,
)?)
.send()
.await
.map_err(|e| classify_transport(model, &e))?;
let status = response.status();
if !status.is_success() {
let detail = response.text().await.unwrap_or_default();
return Err(classify_status(model, status.as_u16(), &detail));
}
let mut completion = if self.stream {
self.read_streamed(response, model, schema, stream).await?
} else {
self.read_buffered(response, model, schema).await?
};
if !self.stream
&& let Some((observer, label)) = stream
{
observer.event(crate::core::Tainted::with_label(
super::ModelStreamEvent::Usage(completion.usage),
label.clone(),
));
}
accumulate_continuation(&mut completion, continuation, exchanges);
Ok(completion)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn driver() -> Anthropic {
Anthropic::new("test-key").expect("build the driver")
}
#[test]
fn the_forced_tool_is_found_by_name_not_by_position() {
let parsed: ApiResponse = serde_json::from_value(json!({
"content": [
{ "type": "tool_use", "id": "c1", "name": "refund",
"input": { "amount": 999 } },
{ "type": "tool_use", "id": "c2", "name": RESPOND_TOOL,
"input": { "verdict": "ship" } },
],
}))
.expect("parse");
assert_eq!(
parsed.forced_tool_input(),
Some(&json!({ "verdict": "ship" })),
"the structured answer was taken from the caller's tool call"
);
let calls = parsed.tool_calls().expect("well-formed tool calls");
assert_eq!(calls.len(), 1, "the forced tool is not a caller tool call");
assert_eq!(calls[0].name, "refund");
assert_eq!(calls[0].id, "c1");
}
#[test]
fn a_malformed_buffered_tool_call_is_loud() {
let parsed: ApiResponse = serde_json::from_value(json!({
"content": [
{ "type": "tool_use", "name": "refund", "input": { "amount": 999 } }
],
"usage": { "input_tokens": 10, "output_tokens": 4 }
}))
.expect("response envelope");
assert!(
parsed.tool_calls().is_err(),
"a missing call id was silently turned into no tool calls"
);
}
#[test]
fn a_system_instruction_rides_beside_the_messages() {
let body = driver().body(
&ModelId::new("anthropic", "claude-x"),
&json!({ "system": "answer only in French", "messages": [{"role": "user", "content": "hi"}] }),
None,
&[],
&[],
)
.expect("a body without tools");
assert_eq!(
body["system"], "answer only in French",
"the system instruction must be a top-level parameter: {body}"
);
assert_eq!(
body["messages"],
json!([{ "role": "user", "content": "hi" }]),
"it must not also be pushed into the conversation: {body}"
);
}
#[test]
fn a_system_instruction_is_not_shown_as_the_question() {
let body = driver()
.body(
&ModelId::new("anthropic", "claude-x"),
&json!({ "system": "be terse", "ticket": "printer on fire" }),
None,
&[],
&[],
)
.expect("a body without tools");
let asked = body["messages"][0]["content"].as_str().unwrap_or_default();
assert!(
!asked.contains("be terse"),
"the instruction leaked into the question: {asked}"
);
assert!(
asked.contains("printer on fire"),
"the actual content went missing: {asked}"
);
}
#[test]
fn a_multimodal_message_is_passed_through_verbatim() {
let parts = json!([{
"role": "user",
"content": [
{ "type": "text", "text": "what is in this image?" },
{ "type": "image", "source": {
"type": "base64", "media_type": "image/png", "data": "iVBORw0KGgo=" } }
]
}]);
let body = driver()
.body(
&ModelId::new("anthropic", "claude-x"),
&json!({ "messages": parts }),
None,
&[],
&[],
)
.expect("a body without tools");
assert_eq!(
body["messages"], parts,
"content blocks must survive untouched: {body}"
);
}
#[test]
fn a_prompt_without_a_system_sends_no_system() {
let body = driver()
.body(
&ModelId::new("anthropic", "claude-x"),
&json!("hi"),
None,
&[],
&[],
)
.expect("a body without tools");
assert!(
body.get("system").is_none(),
"an unset instruction must not become an empty one: {body}"
);
}
}
#[cfg(test)]
mod tool_tests {
use super::*;
use crate::model::ToolDeclaration;
fn driver() -> Anthropic {
Anthropic::new("test-key").expect("build the driver")
}
fn decl() -> ToolDeclaration {
ToolDeclaration::new(
"ledger.read",
"Read a ledger entry.",
json!({ "type": "object", "properties": { "id": { "type": "string" } } }),
)
}
#[test]
fn a_declared_tool_is_rendered_in_anthropics_shape() {
let body = driver()
.body(
&ModelId::new("anthropic", "claude-x"),
&json!({ "messages": [] }),
None,
&[decl()],
&[],
)
.expect("a body with tools");
let tool = &body["tools"][0];
assert_eq!(tool["name"], "ledger.read");
assert_eq!(
tool["strict"], true,
"strict tool use must enforce the declared input schema during generation: {body}"
);
assert_eq!(
tool["input_schema"]["type"], "object",
"Anthropic names the argument schema `input_schema`; `parameters` is OpenAI's spelling and this request would be rejected: {body}"
);
assert!(
tool.get("function").is_none(),
"the OpenAI `function` wrapper must not appear: {body}"
);
}
#[test]
fn a_forced_schema_and_declared_tools_are_refused_together() {
let model = ModelId::new("anthropic", "claude-x");
let forced = driver().structured_via(SchemaMode::ForcedTool);
let schema = json!({ "type": "object" });
let out = forced.body(
&model,
&json!({ "messages": [] }),
Some(&schema),
&[decl()],
&[],
);
match out {
Err(ModelError::Refused { detail, .. }) => assert!(
detail.contains("structured output") && detail.contains("tool"),
"the refusal must say which two things collided: {detail}"
),
Err(e) => panic!("wrong refusal: {e}"),
Ok(body) => panic!(
"a schema and declared tools were combined, so one silently \
replaced the other: {body}"
),
}
assert!(
forced
.body(&model, &json!({ "messages": [] }), Some(&schema), &[], &[])
.is_ok(),
"a schema alone must still work"
);
assert!(
forced
.body(&model, &json!({ "messages": [] }), None, &[decl()], &[])
.is_ok(),
"tools alone must still work"
);
let native = driver()
.body(
&model,
&json!({ "messages": [] }),
Some(&schema),
&[decl()],
&[],
)
.expect("native structured output coexists with tools");
assert_eq!(native["tools"][0]["name"], "ledger.read");
assert!(native["output_config"].is_object());
}
}
#[cfg(test)]
mod continuation_tests {
use super::*;
use crate::model::{
ProviderContinuation, ReasoningEffort, ToolCall as ModelToolCall, ToolExchange,
};
fn driver() -> Anthropic {
Anthropic::new("test-key").expect("build the driver")
}
fn exchange(failed: bool) -> ToolExchange {
let call = ModelToolCall {
id: "toolu_01".to_owned(),
name: "ledger.read".to_owned(),
arguments: json!({ "id": "AC-1" }),
};
if failed {
ToolExchange::failed(call, "the ledger was unreachable")
} else {
ToolExchange::ok(call, json!({ "balance": 42 }))
}
}
#[test]
fn a_continuation_echoes_the_call_beside_its_result() {
let body = driver()
.body(
&ModelId::new("anthropic", "claude-x"),
&json!({ "messages": [{"role": "user", "content": "balance?"}] }),
None,
&[],
&[exchange(false)],
)
.expect("a continuation body");
let msgs = body["messages"].as_array().expect("messages");
assert_eq!(msgs.len(), 3, "question, the call, the result: {body}");
let used = &msgs[1];
assert_eq!(used["role"], "assistant");
assert_eq!(used["content"][0]["type"], "tool_use");
assert_eq!(
used["content"][0]["input"],
json!({ "id": "AC-1" }),
"Anthropic takes arguments as an object, not a JSON string: {body}"
);
let result = &msgs[2];
assert_eq!(result["role"], "user");
assert_eq!(result["content"][0]["type"], "tool_result");
assert_eq!(
result["content"][0]["tool_use_id"], used["content"][0]["id"],
"the result must carry the id of the call it answers, or the API \
rejects it: {body}"
);
}
#[test]
fn signed_thinking_round_trips_unchanged_before_tool_results() {
let content = json!([
{"type": "thinking", "thinking": "private reasoning", "signature": "signed-value"},
{"type": "text", "text": "checking"},
{"type": "tool_use", "id": "toolu_01", "name": "ledger.read", "input": {"id": "AC-1"}}
]);
let state = ProviderContinuation::new(
"anthropic",
json!([{ "role": "assistant", "content": content.clone() }]),
);
let body = driver()
.body_with_max(
&ModelId::new("anthropic", "claude-x"),
&json!({"messages": [{"role": "user", "content": "balance?"}]}),
4096,
Some(ReasoningEffort::High),
None,
&[],
&[exchange(false)],
Some(&state),
)
.expect("lossless thinking continuation");
assert_eq!(body["messages"][1]["content"], content);
assert_eq!(body["messages"][2]["content"][0]["type"], "tool_result");
}
#[test]
fn continuation_accumulates_every_prior_message() {
let prior = ProviderContinuation::new(
"anthropic",
json!([{ "role": "assistant", "content": [{"type": "thinking", "signature": "first"}] }]),
);
let mut completion = Completion {
text: String::new(),
tool_calls: vec![ModelToolCall {
id: "toolu_02".to_owned(),
name: "ledger.read".to_owned(),
arguments: json!({}),
}],
usage: Usage::default(),
stop_reason: Some("tool_use".to_owned()),
truncated: false,
structured: None,
continuation: Some(ProviderContinuation::new(
"anthropic",
json!([{ "role": "assistant", "content": [{"type": "tool_use", "id": "toolu_02"}] }]),
)),
};
accumulate_continuation(&mut completion, Some(&prior), &[exchange(false)]);
let state = completion.continuation.unwrap().state;
assert_eq!(state[0]["content"][0]["signature"], "first");
assert_eq!(state[1]["role"], "user");
assert_eq!(state[2]["content"][0]["id"], "toolu_02");
}
#[test]
fn a_failed_tool_is_marked_is_error() {
let body = driver()
.body(
&ModelId::new("anthropic", "claude-x"),
&json!({ "messages": [] }),
None,
&[],
&[exchange(true)],
)
.expect("body");
let msgs = body["messages"].as_array().expect("messages");
let (used, result) = (&msgs[msgs.len() - 2], &msgs[msgs.len() - 1]);
assert_eq!(
used["content"][0]["type"], "tool_use",
"the call is echoed even when it failed: {body}"
);
assert_eq!(
result["content"][0]["is_error"], true,
"a failure rendered as an ordinary result teaches the model the \
operation succeeded and returned something strange: {body}"
);
}
#[test]
fn a_first_turn_is_untouched() {
let prompt = json!({ "messages": [{"role": "user", "content": "hi"}] });
let plain = driver()
.body(
&ModelId::new("anthropic", "claude-x"),
&prompt,
None,
&[],
&[],
)
.expect("body");
assert_eq!(plain["messages"].as_array().expect("messages").len(), 1);
}
}