use std::collections::VecDeque;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use af_agent::{
ChatModel, ContextAuthority, ContextContribution, ContextContributor, ContextRequest, Hook,
HookDecision, Tool, ToolConcurrency, ToolExecutionFuture, ToolMeta, ToolRegistry,
};
use af_agent_runtime::{
recovery_events, AgentRuntime, CancellationToken, CompactionResult, Compactor, EventWriter,
RuntimeError, RuntimeLimits, TokenMeter, TurnRequest,
};
use af_agent_session::{text, DeliveryMode, Event, SessionEvent, SessionProjection};
use af_llm::{
AssistantBlock, ChatMessage, Choice, CompletionRequest, CompletionResponse, FinishReason,
FunctionCall, LlmError, ReasoningEffort, Role, ToolCall, Usage,
};
use async_trait::async_trait;
use serde_json::{json, Value};
struct ScriptedModel(Mutex<VecDeque<Result<CompletionResponse, LlmError>>>);
#[async_trait]
impl ChatModel for ScriptedModel {
async fn complete_streaming(
&self,
_request: &CompletionRequest,
delta_tx: tokio::sync::mpsc::UnboundedSender<(String, bool)>,
) -> Result<CompletionResponse, LlmError> {
let response = self.0.lock().unwrap().pop_front().unwrap()?;
if let Some(content) = response.first_content() {
let _ = delta_tx.send((content.into(), response.first_tool_calls().is_some()));
}
Ok(response)
}
}
#[tokio::test]
async fn structured_assistant_output_is_persisted_and_replayed_without_text_parsing() {
let mut result = response(ChatMessage::assistant("provider fallback")).unwrap();
result.choices[0].output_blocks = vec![
AssistantBlock::Text {
text: "answer".into(),
},
AssistantBlock::Citation {
resource_id: "doc-1".into(),
label: "Architecture".into(),
uri: "docs://doc-1".into(),
excerpt: Some("event sourced".into()),
},
AssistantBlock::Data {
slot: "chart".into(),
value: json!({"points":[1,2]}),
},
];
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([Ok(result)]))));
let history = initial();
let writer = MemoryWriter::new(history.clone());
let outcome = AgentRuntime::new(model, "model", "system")
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(outcome.final_text.as_deref(), Some("answer"));
let content = writer
.events()
.into_iter()
.find_map(|event| match event.event {
Event::AssistantMessage { content, .. } => Some(content),
_ => None,
})
.unwrap();
assert!(
matches!(&content[1], af_agent_session::ContentBlock::Citation { resource_id, uri, .. } if resource_id == "doc-1" && uri == "docs://doc-1")
);
assert!(
matches!(&content[2], af_agent_session::ContentBlock::Data { slot, value } if slot == "chart" && value["points"][1] == 2)
);
}
#[tokio::test]
async fn malformed_structured_assistant_output_fails_before_persistence() {
let mut result = response(ChatMessage::assistant("ignored")).unwrap();
result.choices[0].output_blocks = vec![AssistantBlock::Data {
slot: " ".into(),
value: json!({}),
}];
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([Ok(result)]))));
let history = initial();
let writer = MemoryWriter::new(history.clone());
assert!(matches!(
AgentRuntime::new(model, "model", "system")
.run(request(history), &writer, CancellationToken::default())
.await,
Err(RuntimeError::Model(message)) if message.contains("invalid structured block")
));
assert!(!writer
.events()
.iter()
.any(|event| matches!(event.event, Event::AssistantMessage { .. })));
}
#[tokio::test]
async fn profile_reasoning_and_output_policy_reach_the_single_model_path() {
let model = Arc::new(CapturingModel {
responses: Mutex::new(VecDeque::from([response(ChatMessage::assistant("answer"))])),
requests: Mutex::new(Vec::new()),
});
let history = initial();
let writer = MemoryWriter::new(history.clone());
AgentRuntime::new(model.clone(), "model", "system")
.with_reasoning_effort(ReasoningEffort::High)
.with_output_policy(json!({
"type":"array",
"items":{"type":"object","required":["type"],"properties":{"type":{"const":"text"}}}
}))
.unwrap()
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(
model.requests.lock().unwrap()[0].reasoning_effort,
Some(ReasoningEffort::High)
);
let history = initial();
let writer = MemoryWriter::new(history.clone());
let rejected = AgentRuntime::new(
Arc::new(ScriptedModel(Mutex::new(VecDeque::from([response(
ChatMessage::assistant("forbidden"),
)])))),
"model",
"system",
)
.with_output_policy(json!({"type":"array","maxItems":0}))
.unwrap()
.run(request(history), &writer, CancellationToken::default())
.await;
assert!(
matches!(rejected, Err(RuntimeError::Model(message)) if message.contains("invalid assistant output"))
);
assert!(!writer
.events()
.iter()
.any(|event| matches!(event.event, Event::AssistantMessage { .. })));
}
struct SlowModel;
#[async_trait]
impl ChatModel for SlowModel {
async fn complete_streaming(
&self,
_request: &CompletionRequest,
_delta_tx: tokio::sync::mpsc::UnboundedSender<(String, bool)>,
) -> Result<CompletionResponse, LlmError> {
tokio::time::sleep(Duration::from_secs(5)).await;
response(ChatMessage::assistant("late"))
}
}
struct SteerModel(AtomicUsize);
#[async_trait]
impl ChatModel for SteerModel {
async fn complete_streaming(
&self,
request: &CompletionRequest,
_delta_tx: tokio::sync::mpsc::UnboundedSender<(String, bool)>,
) -> Result<CompletionResponse, LlmError> {
if self.0.fetch_add(1, Ordering::SeqCst) == 0 {
tokio::time::sleep(Duration::from_secs(5)).await;
return response(ChatMessage::assistant("stale"));
}
assert!(request
.messages
.iter()
.any(|message| message.content.as_deref() == Some("new direction")));
response(ChatMessage::assistant("steered answer"))
}
}
struct CapturingModel {
responses: Mutex<VecDeque<Result<CompletionResponse, LlmError>>>,
requests: Mutex<Vec<CompletionRequest>>,
}
#[async_trait]
impl ChatModel for CapturingModel {
async fn complete_streaming(
&self,
request: &CompletionRequest,
delta_tx: tokio::sync::mpsc::UnboundedSender<(String, bool)>,
) -> Result<CompletionResponse, LlmError> {
self.requests.lock().unwrap().push(request.clone());
let response = self.responses.lock().unwrap().pop_front().unwrap()?;
if let Some(content) = response.first_content() {
let _ = delta_tx.send((content.into(), response.first_tool_calls().is_some()));
}
Ok(response)
}
}
struct DynamicContext(AtomicUsize);
#[async_trait]
impl ContextContributor for DynamicContext {
async fn contribute(
&self,
request: &ContextRequest,
) -> Result<Vec<ContextContribution>, String> {
let call = self.0.fetch_add(1, Ordering::SeqCst) + 1;
Ok(vec![ContextContribution {
id: "dynamic".into(),
source: "test-provider".into(),
version: "1".into(),
authority: ContextAuthority::Untrusted,
form: "test".into(),
content: format!("dynamic-{call}:{}", request.query),
}])
}
}
struct BlockingContext(Arc<std::sync::atomic::AtomicBool>);
#[async_trait]
impl ContextContributor for BlockingContext {
async fn contribute(
&self,
request: &ContextRequest,
) -> Result<Vec<ContextContribution>, String> {
while !request.cancellation.is_cancelled() {
tokio::time::sleep(Duration::from_millis(5)).await;
}
self.0.store(true, Ordering::SeqCst);
Ok(Vec::new())
}
}
#[tokio::test]
async fn context_contributors_observe_deadline_cancellation() {
let observed = Arc::new(std::sync::atomic::AtomicBool::new(false));
let history = initial();
let writer = MemoryWriter::new(history.clone());
let error = AgentRuntime::new(
Arc::new(ScriptedModel(Mutex::new(VecDeque::new()))),
"model",
"system",
)
.with_context_contributor(Arc::new(BlockingContext(Arc::clone(&observed))))
.with_limits(RuntimeLimits {
provider_deadline: Duration::from_millis(30),
..RuntimeLimits::default()
})
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap_err();
assert!(matches!(error, RuntimeError::ExtensionDeadline(_)));
assert!(observed.load(Ordering::SeqCst));
}
fn response(message: ChatMessage) -> Result<CompletionResponse, LlmError> {
let finish_reason = if message
.tool_calls
.as_ref()
.is_some_and(|calls| !calls.is_empty())
{
FinishReason::ToolCalls
} else {
FinishReason::Stop
};
response_with_reason(message, finish_reason)
}
fn response_with_reason(
message: ChatMessage,
finish_reason: FinishReason,
) -> Result<CompletionResponse, LlmError> {
Ok(CompletionResponse {
id: "response".into(),
choices: vec![Choice {
index: 0,
message,
finish_reason: Some(finish_reason),
output_blocks: Vec::new(),
}],
usage: Some(Usage {
prompt_tokens: 1,
completion_tokens: 1,
total_tokens: 2,
}),
})
}
#[tokio::test]
async fn provider_attempt_identity_is_durable_and_passed_to_the_provider() {
let model = Arc::new(CapturingModel {
responses: Mutex::new(VecDeque::from([response(ChatMessage::assistant("done"))])),
requests: Mutex::new(Vec::new()),
});
let history = initial();
let writer = MemoryWriter::new(history.clone());
AgentRuntime::new(model.clone(), "model", "system")
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
let provider_attempt_id = model.requests.lock().unwrap()[0]
.provider_attempt_id
.clone()
.unwrap();
assert_eq!(provider_attempt_id, "run:model:1:attempt:1");
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::ModelRequestPrepared { provider_attempt_id: durable, operation_id, reserved_completion_tokens: 4096, .. }
if durable == &provider_attempt_id && operation_id == "model:1:attempt:1"
)));
}
#[tokio::test]
async fn context_is_injected_per_step_with_provenance_without_accumulating() {
let model = Arc::new(CapturingModel {
responses: Mutex::new(VecDeque::from([
response(calls(&[("missing-call", "missing")])),
response(ChatMessage::assistant("done")),
])),
requests: Mutex::new(Vec::new()),
});
let history = initial();
let writer = MemoryWriter::new(history.clone());
let outcome = AgentRuntime::new(model.clone(), "model", "system")
.with_context_contributor(Arc::new(DynamicContext(AtomicUsize::new(0))))
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(outcome.status, "completed");
let requests = model.requests.lock().unwrap();
assert_eq!(requests.len(), 2);
let first = serde_json::to_string(&requests[0]).unwrap();
let second = serde_json::to_string(&requests[1]).unwrap();
assert!(first.contains("dynamic-1:hello"));
assert!(second.contains("dynamic-2:hello"));
assert!(!second.contains("dynamic-1:hello"));
assert!(requests[0]
.messages
.iter()
.any(|message| message.role == Role::User
&& message
.content
.as_deref()
.is_some_and(|text| text.contains("authority=\"untrusted\""))));
let events = writer.events();
for step in [1, 2] {
let context_seq = events
.iter()
.find(|event| matches!(&event.event, Event::ContextInjected { step: actual, source, .. } if *actual == step && source == "test-provider"))
.unwrap()
.seq;
let request_seq = events
.iter()
.find(|event| matches!(&event.event, Event::ModelRequestPrepared { step: actual, .. } if *actual == step))
.unwrap()
.seq;
assert!(context_seq < request_seq);
}
}
#[tokio::test]
async fn finish_reason_is_the_only_tool_dispatch_authority() {
let count = Arc::new(AtomicUsize::new(0));
for (reason, expected_status) in [
(FinishReason::Stop, "completed"),
(FinishReason::Length, "max_steps_reached"),
] {
let message = ChatMessage {
role: Role::Assistant,
content: Some("partial".into()),
tool_calls: calls(&[("must-not-run", "count")]).tool_calls,
tool_call_id: None,
name: None,
};
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response_with_reason(message, reason),
]))));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(CountingTool(Arc::clone(&count))))
.unwrap();
let history = initial();
let writer = MemoryWriter::new(history.clone());
let outcome = AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(outcome.status, expected_status);
assert!(!writer.events().iter().any(|event| matches!(
event.event,
Event::AssistantToolCalls { .. } | Event::ToolCall { .. }
)));
}
assert_eq!(count.load(Ordering::SeqCst), 0);
}
#[tokio::test]
async fn rejected_finish_reasons_do_not_persist_assistant_or_tool_calls() {
for reason in [
FinishReason::ContentFilter,
FinishReason::Unknown("provider_reason".into()),
] {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response_with_reason(calls(&[("must-not-run", "count")]), reason.clone()),
]))));
let history = initial();
let writer = MemoryWriter::new(history.clone());
assert!(matches!(
AgentRuntime::new(model, "model", "system")
.run(request(history), &writer, CancellationToken::default())
.await,
Err(RuntimeError::FinishReason(actual)) if actual == reason
));
assert!(!writer.events().iter().any(|event| matches!(
event.event,
Event::AssistantMessage { .. }
| Event::AssistantToolCalls { .. }
| Event::ToolCall { .. }
)));
}
}
fn calls(calls: &[(&str, &str)]) -> ChatMessage {
ChatMessage {
role: Role::Assistant,
content: None,
tool_calls: Some(
calls
.iter()
.map(|(id, name)| ToolCall {
id: (*id).into(),
kind: "function".into(),
function: FunctionCall {
name: (*name).into(),
arguments: "{}".into(),
},
})
.collect(),
),
tool_call_id: None,
name: None,
}
}
struct DelayTool {
name: &'static str,
delay_ms: u64,
}
struct ExclusiveTool(&'static str);
#[async_trait]
impl Tool for ExclusiveTool {
fn name(&self) -> &str {
self.0
}
fn description(&self) -> &str {
"exclusive test tool"
}
fn parameters(&self) -> Value {
json!({"type":"object"})
}
fn output_schema(&self) -> Value {
json!({"type":"object","required":["tool"],"properties":{"tool":{"type":"string"}}})
}
fn meta(&self) -> ToolMeta {
ToolMeta {
concurrency: ToolConcurrency::Exclusive,
..ToolMeta::default()
}
}
async fn call(&self, _args: Value) -> Result<Value, String> {
Ok(json!({"tool":self.0}))
}
}
struct CooperativeTimeoutTool(Arc<AtomicUsize>);
struct NeverSettlesTool(Arc<AtomicUsize>);
struct GraceSettlementTool {
name: &'static str,
delay_ms: u64,
error: Option<&'static str>,
invalid_output: bool,
started: Arc<AtomicUsize>,
}
#[async_trait]
impl Tool for GraceSettlementTool {
fn name(&self) -> &str {
self.name
}
fn description(&self) -> &str {
"settles during cancellation grace"
}
fn parameters(&self) -> Value {
json!({"type":"object"})
}
fn output_schema(&self) -> Value {
json!({"type":"object","required":["tool"],"properties":{"tool":{"type":"string"}}})
}
fn meta(&self) -> ToolMeta {
ToolMeta {
concurrency: ToolConcurrency::Concurrent,
..ToolMeta::default()
}
}
async fn call(&self, _: Value) -> Result<Value, String> {
Err("context required".into())
}
async fn call_with_context(
&self,
context: &af_agent::ToolExecutionContext,
_: Value,
) -> Result<Value, String> {
self.started.fetch_add(1, Ordering::SeqCst);
while !context.cancellation.is_cancelled() {
tokio::time::sleep(Duration::from_millis(5)).await;
}
tokio::time::sleep(Duration::from_millis(self.delay_ms)).await;
match self.error {
Some(error) => Err(error.into()),
None if self.invalid_output => Ok(json!({})),
None => Ok(json!({"tool":self.name})),
}
}
}
#[async_trait]
impl Tool for NeverSettlesTool {
fn name(&self) -> &str {
"never-settles"
}
fn description(&self) -> &str {
"ignores cancellation"
}
fn parameters(&self) -> Value {
json!({"type":"object"})
}
fn output_schema(&self) -> Value {
json!({"type":"object"})
}
fn meta(&self) -> ToolMeta {
ToolMeta {
timeout_secs: 1,
..ToolMeta::default()
}
}
async fn call(&self, _: Value) -> Result<Value, String> {
self.0.fetch_add(1, Ordering::SeqCst);
std::future::pending().await
}
}
#[async_trait]
impl Tool for CooperativeTimeoutTool {
fn name(&self) -> &str {
"cooperative-timeout"
}
fn description(&self) -> &str {
"settles after cancellation"
}
fn parameters(&self) -> Value {
json!({"type":"object"})
}
fn output_schema(&self) -> Value {
json!({"type":"object","required":["settled"],"properties":{"settled":{"type":"boolean"}}})
}
fn meta(&self) -> ToolMeta {
ToolMeta {
timeout_secs: 1,
..ToolMeta::default()
}
}
async fn call(&self, _: Value) -> Result<Value, String> {
Err("context required".into())
}
async fn call_with_context(
&self,
context: &af_agent::ToolExecutionContext,
_: Value,
) -> Result<Value, String> {
while !context.cancellation.is_cancelled() {
tokio::time::sleep(Duration::from_millis(5)).await;
}
tokio::time::sleep(Duration::from_millis(20)).await;
self.0.fetch_add(1, Ordering::SeqCst);
Ok(json!({"settled":true}))
}
}
#[tokio::test]
async fn tool_timeout_persists_real_outcome_settled_during_grace() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response(calls(&[("call", "cooperative-timeout")])),
response(ChatMessage::assistant("recovered")),
]))));
let settled = Arc::new(AtomicUsize::new(0));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(CooperativeTimeoutTool(settled.clone())))
.unwrap();
let history = initial();
let writer = MemoryWriter::new(history.clone());
AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(settled.load(Ordering::SeqCst), 1);
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::ToolResult { is_error: false, result, .. } if result == &json!({"settled":true})
)));
}
#[tokio::test]
async fn cancellation_grace_persists_ordered_real_results_and_stays_cancelled() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([response(
calls(&[("first-call", "late-ok"), ("second-call", "early-error")]),
)]))));
let started = Arc::new(AtomicUsize::new(0));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(GraceSettlementTool {
name: "late-ok",
delay_ms: 40,
error: None,
invalid_output: false,
started: Arc::clone(&started),
}))
.unwrap()
.register(Arc::new(GraceSettlementTool {
name: "early-error",
delay_ms: 10,
error: Some("settled failure"),
invalid_output: false,
started: Arc::clone(&started),
}))
.unwrap();
let history = initial();
let writer = Arc::new(MemoryWriter::new(history.clone()));
let cancellation = CancellationToken::default();
let task = tokio::spawn({
let writer = Arc::clone(&writer);
let cancellation = cancellation.clone();
async move {
AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.run(request(history), writer.as_ref(), cancellation)
.await
}
});
tokio::time::timeout(Duration::from_secs(1), async {
while started.load(Ordering::SeqCst) != 2 {
tokio::task::yield_now().await;
}
})
.await
.unwrap();
cancellation.cancel();
assert_eq!(task.await.unwrap().unwrap().status, "cancelled");
let events = writer.events();
let results = events
.iter()
.filter_map(|event| match &event.event {
Event::ToolResult {
call_id,
result,
is_error,
..
} => Some((call_id.as_str(), result, *is_error)),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(
results,
vec![
("first-call", &json!({"tool":"late-ok"}), false),
("second-call", &json!({"error":"settled failure"}), true),
]
);
assert!(events.iter().any(|event| matches!(
event.event,
Event::RunFinished {
status: af_agent_session::RunStatus::Cancelled,
..
}
)));
assert!(!events.iter().any(|event| match &event.event {
Event::ToolResult { result, .. } => result.to_string().contains("tool_outcome_unknown"),
_ => false,
}));
}
#[tokio::test]
async fn cancellation_grace_validates_settled_output() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([response(
calls(&[("call", "invalid-after-cancel")]),
)]))));
let started = Arc::new(AtomicUsize::new(0));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(GraceSettlementTool {
name: "invalid-after-cancel",
delay_ms: 10,
error: None,
invalid_output: true,
started: Arc::clone(&started),
}))
.unwrap();
let history = initial();
let writer = Arc::new(MemoryWriter::new(history.clone()));
let cancellation = CancellationToken::default();
let task = tokio::spawn({
let writer = Arc::clone(&writer);
let cancellation = cancellation.clone();
async move {
AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.run(request(history), writer.as_ref(), cancellation)
.await
}
});
tokio::time::timeout(Duration::from_secs(1), async {
while started.load(Ordering::SeqCst) != 1 {
tokio::task::yield_now().await;
}
})
.await
.unwrap();
cancellation.cancel();
assert_eq!(task.await.unwrap().unwrap().status, "cancelled");
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::ToolResult { is_error: true, result, .. }
if result.to_string().contains("invalid tool output")
)));
}
#[tokio::test]
async fn tool_timeout_drops_non_cooperative_execution_after_bounded_grace() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response(calls(&[("call", "never-settles")])),
response(calls(&[("call", "never-settles")])),
]))));
let executions = Arc::new(AtomicUsize::new(0));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(NeverSettlesTool(Arc::clone(&executions))))
.unwrap();
let history = initial();
let writer = MemoryWriter::new(history.clone());
let started = std::time::Instant::now();
let outcome = AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert!(started.elapsed() < Duration::from_secs(2));
assert_eq!(outcome.status, "failed");
assert_eq!(executions.load(Ordering::SeqCst), 1);
assert_eq!(
writer
.events()
.iter()
.filter(|event| matches!(event.event, Event::ModelRequestPrepared { .. }))
.count(),
1
);
assert!(writer.events().iter().any(|event| matches!(&event.event,
Event::ToolResult { is_error: true, result, .. }
if result.to_string().contains("tool_outcome_unknown")
)));
}
#[async_trait]
impl Tool for DelayTool {
fn name(&self) -> &str {
self.name
}
fn description(&self) -> &str {
"test"
}
fn parameters(&self) -> Value {
json!({"type":"object"})
}
fn output_schema(&self) -> Value {
json!({"type":"object","required":["tool"],"properties":{"tool":{"type":"string"},"wrapped":{"type":"boolean"}}})
}
fn meta(&self) -> ToolMeta {
ToolMeta {
concurrency: ToolConcurrency::Concurrent,
..ToolMeta::default()
}
}
async fn call(&self, _args: Value) -> Result<Value, String> {
tokio::time::sleep(Duration::from_millis(self.delay_ms)).await;
Ok(json!({"tool":self.name}))
}
}
struct MemoryWriter(Mutex<Vec<SessionEvent>>);
impl MemoryWriter {
fn new(events: Vec<SessionEvent>) -> Self {
Self(Mutex::new(events))
}
fn events(&self) -> Vec<SessionEvent> {
self.0.lock().unwrap().clone()
}
}
struct TerminalRaceWriter {
inner: MemoryWriter,
injected: std::sync::atomic::AtomicBool,
}
#[async_trait]
impl EventWriter for TerminalRaceWriter {
async fn append(&self, events: Vec<Event>) -> Result<Vec<SessionEvent>, RuntimeError> {
if events
.iter()
.any(|event| matches!(event, Event::RunFinished { .. }))
&& !self.injected.swap(true, Ordering::SeqCst)
{
self.inner
.append(vec![Event::InputQueued {
input_id: "racing-input".into(),
run_id: "run".into(),
mode: DeliveryMode::Inject,
content: text("racing direction"),
explicit_skill: None,
}])
.await?;
}
self.inner.append(events).await
}
async fn load_after(&self, seq: u64) -> Result<Vec<SessionEvent>, RuntimeError> {
self.inner.load_after(seq).await
}
}
struct CompactionRaceWriter {
inner: MemoryWriter,
injected: std::sync::atomic::AtomicBool,
}
struct CrashAfterCompactionStartWriter {
inner: MemoryWriter,
crashed: std::sync::atomic::AtomicBool,
}
#[async_trait]
impl EventWriter for CrashAfterCompactionStartWriter {
async fn append(&self, events: Vec<Event>) -> Result<Vec<SessionEvent>, RuntimeError> {
let starts_compaction = events
.iter()
.any(|event| matches!(event, Event::CompactionStarted { .. }));
let appended = self.inner.append(events).await?;
if starts_compaction && !self.crashed.swap(true, Ordering::SeqCst) {
return Err(RuntimeError::Event(
"simulated crash after compaction start".into(),
));
}
Ok(appended)
}
async fn load_after(&self, seq: u64) -> Result<Vec<SessionEvent>, RuntimeError> {
self.inner.load_after(seq).await
}
}
#[async_trait]
impl EventWriter for CompactionRaceWriter {
async fn append(&self, events: Vec<Event>) -> Result<Vec<SessionEvent>, RuntimeError> {
let starts_compaction = events
.iter()
.any(|event| matches!(event, Event::CompactionStarted { .. }));
let appended = self.inner.append(events).await?;
if starts_compaction && !self.injected.swap(true, Ordering::SeqCst) {
self.inner
.append(vec![Event::InputQueued {
input_id: "during-compaction".into(),
run_id: "run".into(),
mode: DeliveryMode::Inject,
content: text("new direction"),
explicit_skill: None,
}])
.await?;
}
Ok(appended)
}
async fn load_after(&self, seq: u64) -> Result<Vec<SessionEvent>, RuntimeError> {
self.inner.load_after(seq).await
}
}
#[async_trait]
impl EventWriter for MemoryWriter {
async fn append(
&self,
events: Vec<Event>,
) -> Result<Vec<SessionEvent>, af_agent_runtime::RuntimeError> {
let mut history = self.0.lock().unwrap();
let start = history.len() as u64;
let appended = events
.into_iter()
.enumerate()
.map(|(index, event)| SessionEvent {
session_id: "session".into(),
seq: start + index as u64 + 1,
occurred_at: chrono::Utc::now(),
event,
})
.collect::<Vec<_>>();
let mut candidate = history.clone();
candidate.extend(appended.clone());
SessionProjection::replay(&candidate)
.map_err(|error| af_agent_runtime::RuntimeError::Event(error.to_string()))?;
history.extend(appended.clone());
Ok(appended)
}
async fn load_after(
&self,
seq: u64,
) -> Result<Vec<SessionEvent>, af_agent_runtime::RuntimeError> {
Ok(self
.0
.lock()
.unwrap()
.iter()
.filter(|event| event.seq > seq)
.cloned()
.collect())
}
}
fn initial() -> Vec<SessionEvent> {
vec![
SessionEvent {
session_id: "session".into(),
seq: 1,
occurred_at: chrono::Utc::now(),
event: Event::SessionCreated {
profile_revision_id: "profile-r1".into(),
},
},
SessionEvent {
session_id: "session".into(),
seq: 2,
occurred_at: chrono::Utc::now(),
event: Event::InputQueued {
input_id: "input".into(),
run_id: "run".into(),
mode: DeliveryMode::Followup,
content: text("hello"),
explicit_skill: None,
},
},
]
}
fn request(history: Vec<SessionEvent>) -> TurnRequest {
TurnRequest {
context: af_agent::RequestContext {
tenant_id: "tenant".into(),
subject_id: "subject".into(),
roles: Default::default(),
locale: "en".into(),
request_id: "request".into(),
entitlements: Default::default(),
},
session_id: "session".into(),
run_id: "run".into(),
input_id: "input".into(),
content: text("hello"),
history,
}
}
#[tokio::test]
async fn steer_is_bound_to_the_active_run_and_interrupts_the_model_attempt() {
let history = initial();
let writer = Arc::new(MemoryWriter::new(history.clone()));
let runtime = AgentRuntime::new(Arc::new(SteerModel(AtomicUsize::new(0))), "model", "system");
let run_writer = writer.clone();
let task = tokio::spawn(async move {
runtime
.run(
request(history),
run_writer.as_ref(),
CancellationToken::default(),
)
.await
});
for _ in 0..100 {
if writer
.events()
.iter()
.any(|event| matches!(event.event, Event::ModelRequestPrepared { step: 1, .. }))
{
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
writer
.append(vec![Event::InputQueued {
input_id: "steer-input".into(),
run_id: "run".into(),
mode: DeliveryMode::Steer,
content: text("new direction"),
explicit_skill: None,
}])
.await
.unwrap();
let outcome = task.await.unwrap().unwrap();
assert_eq!(outcome.final_text.as_deref(), Some("steered answer"));
let events = writer.events();
assert!(events.iter().any(|event| matches!(
&event.event,
Event::ModelAttemptFailed { step: 1, error, .. } if error == "steered"
)));
assert!(events.iter().any(|event| matches!(
&event.event,
Event::InputClaimed { input_id, run_id } if input_id == "steer-input" && run_id == "run"
)));
}
#[tokio::test]
async fn recovery_claims_durable_steer_queued_before_restart_once() {
let writer = MemoryWriter::new(initial());
writer
.append(vec![
Event::InputClaimed {
input_id: "input".into(),
run_id: "run".into(),
},
Event::RunStarted {
run_id: "run".into(),
input_id: "input".into(),
},
Event::TurnStarted {
run_id: "run".into(),
turn: 1,
},
Event::UserMessage {
run_id: "run".into(),
content: text("hello"),
},
Event::InputQueued {
input_id: "durable-steer".into(),
run_id: "run".into(),
mode: DeliveryMode::Steer,
content: text("new direction"),
explicit_skill: None,
},
])
.await
.unwrap();
let model = Arc::new(CapturingModel {
responses: Mutex::new(VecDeque::from([response(ChatMessage::assistant("done"))])),
requests: Mutex::new(Vec::new()),
});
AgentRuntime::new(model.clone(), "model", "system")
.run(
request(writer.events()),
&writer,
CancellationToken::default(),
)
.await
.unwrap();
let events = writer.events();
assert_eq!(
events
.iter()
.filter(|event| matches!(&event.event, Event::InputClaimed { input_id, .. } if input_id == "durable-steer"))
.count(),
1
);
assert!(model.requests.lock().unwrap()[0]
.messages
.iter()
.any(|message| message.content.as_deref() == Some("new direction")));
}
#[tokio::test]
async fn input_racing_with_terminal_commit_is_claimed_before_the_run_finishes() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response(ChatMessage::assistant("stale answer")),
response(ChatMessage::assistant("answer after injection")),
]))));
let history = initial();
let writer = TerminalRaceWriter {
inner: MemoryWriter::new(history.clone()),
injected: std::sync::atomic::AtomicBool::new(false),
};
let outcome = AgentRuntime::new(model, "model", "system")
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(
outcome.final_text.as_deref(),
Some("answer after injection")
);
let events = writer.inner.events();
assert!(events.iter().any(|event| matches!(
&event.event,
Event::InputClaimed { input_id, run_id } if input_id == "racing-input" && run_id == "run"
)));
let projection = SessionProjection::replay(&events).unwrap();
assert!(projection.queued_inputs.is_empty());
assert_eq!(
projection.run_status.get("run").map(String::as_str),
Some("completed")
);
}
struct ZeroMeter;
impl TokenMeter for ZeroMeter {
fn count(&self, _: &str, _: &[ChatMessage]) -> u64 {
0
}
}
#[tokio::test]
async fn missing_provider_usage_is_estimated_and_enforces_the_run_budget() {
let mut completion = response(ChatMessage::assistant("usage missing")).unwrap();
completion.usage = None;
let history = initial();
let writer = MemoryWriter::new(history.clone());
let outcome = AgentRuntime::new(
Arc::new(ScriptedModel(Mutex::new(VecDeque::from([Ok(completion)])))),
"model",
"system",
)
.with_token_meter(Arc::new(ZeroMeter))
.with_limits(RuntimeLimits {
max_tokens: 0,
..RuntimeLimits::default()
})
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(outcome.status, "max_steps_reached");
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::UsageRecorded { operation_id, completion_tokens: 1, .. }
if operation_id == "model:1:attempt:1"
)));
}
#[tokio::test]
async fn runtime_guards_close_runs_without_starting_unsafe_work() {
let unused_model = || Arc::new(ScriptedModel(Mutex::new(VecDeque::new())));
let history = initial();
let writer = MemoryWriter::new(history.clone());
let mut empty = request(history);
empty.content.clear();
assert!(matches!(
AgentRuntime::new(unused_model(), "model", "system")
.run(empty, &writer, CancellationToken::default())
.await,
Err(RuntimeError::InvalidInput(_))
));
let mut history = initial();
if let Event::InputQueued { run_id, .. } = &mut history[1].event {
*run_id = "other-run".into();
}
let writer = MemoryWriter::new(history.clone());
writer
.append(vec![
Event::InputClaimed {
input_id: "input".into(),
run_id: "other-run".into(),
},
Event::RunStarted {
run_id: "other-run".into(),
input_id: "input".into(),
},
])
.await
.unwrap();
assert!(matches!(
AgentRuntime::new(unused_model(), "model", "system")
.run(
request(writer.events()),
&writer,
CancellationToken::default()
)
.await,
Err(RuntimeError::SessionBusy)
));
let history = initial();
let writer = MemoryWriter::new(history.clone());
let cancelled = CancellationToken::default();
cancelled.cancel();
let outcome = AgentRuntime::new(unused_model(), "model", "system")
.run(request(history), &writer, cancelled)
.await
.unwrap();
assert_eq!(outcome.status, "cancelled");
let history = initial();
let writer = MemoryWriter::new(history.clone());
let outcome = AgentRuntime::new(
Arc::new(ScriptedModel(Mutex::new(VecDeque::from([response(
calls(&[("call", "missing")]),
)])))),
"model",
"system",
)
.with_limits(RuntimeLimits {
max_tool_calls: 0,
..RuntimeLimits::default()
})
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(outcome.status, "max_steps_reached");
let history = initial();
let writer = MemoryWriter::new(history.clone());
let outcome = AgentRuntime::new(
Arc::new(ScriptedModel(Mutex::new(VecDeque::from([response(
ChatMessage::assistant("too expensive"),
)])))),
"model",
"system",
)
.with_token_meter(Arc::new(ZeroMeter))
.with_limits(RuntimeLimits {
max_tokens: 1,
..RuntimeLimits::default()
})
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(outcome.status, "max_steps_reached");
assert_eq!(
RuntimeError::InvalidInput(String::new()).code(),
"invalid_input"
);
assert_eq!(RuntimeError::SessionBusy.code(), "session_busy");
assert_eq!(
RuntimeError::CompactionConflict.code(),
"compaction_conflict"
);
assert_eq!(RuntimeError::Cancelled.code(), "cancelled");
assert_eq!(
RuntimeError::ToolDenied("tool".into(), "reason".into()).code(),
"tool_denied"
);
assert_eq!(RuntimeError::EmptyModelResponse.code(), "agent_run_failed");
}
#[tokio::test]
async fn cancellation_interrupts_an_inflight_model_attempt() {
let history = initial();
let writer = MemoryWriter::new(history.clone());
let cancellation = CancellationToken::default();
let trigger = cancellation.clone();
tokio::spawn(async move {
tokio::time::sleep(Duration::from_millis(5)).await;
trigger.cancel();
});
let outcome = AgentRuntime::new(Arc::new(SlowModel), "model", "system")
.run(request(history), &writer, cancellation)
.await
.unwrap();
assert_eq!(outcome.status, "cancelled");
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::ModelAttemptFailed {
error,
retryable: false,
..
} if error == "cancelled"
)));
}
#[tokio::test]
async fn concurrent_tools_commit_results_in_model_order_and_replay() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response(calls(&[("slow-call", "slow"), ("fast-call", "fast")])),
response(ChatMessage::assistant("done")),
]))));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(DelayTool {
name: "slow",
delay_ms: 20,
}))
.unwrap()
.register(Arc::new(DelayTool {
name: "fast",
delay_ms: 1,
}))
.unwrap();
let history = initial();
let writer = MemoryWriter::new(history.clone());
let outcome = AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(outcome.status, "completed");
let events = writer.events();
let result_ids = events
.iter()
.filter_map(|event| match &event.event {
Event::ToolResult { call_id, .. } => Some(call_id.as_str()),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(result_ids, ["slow-call", "fast-call"]);
assert!(events
.iter()
.any(|event| matches!(event.event, Event::ModelRequestPrepared { .. })));
assert!(events
.iter()
.any(|event| matches!(event.event, Event::AssistantDelta { .. })));
assert_eq!(
events
.iter()
.filter(|event| matches!(
event.event,
Event::ToolAuthorization {
status: af_agent_session::ToolAuthorizationStatus::Allowed,
..
}
))
.count(),
2
);
assert!(SessionProjection::replay(&writer.events())
.unwrap()
.active_run_id
.is_none());
}
#[tokio::test]
async fn exclusive_tool_is_marked_started_only_when_its_batch_dispatches() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response(calls(&[("first-call", "first"), ("second-call", "second")])),
response(ChatMessage::assistant("done")),
]))));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(ExclusiveTool("first")))
.unwrap()
.register(Arc::new(ExclusiveTool("second")))
.unwrap();
let history = initial();
let writer = MemoryWriter::new(history.clone());
AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
let events = writer.events();
let first_result = events
.iter()
.position(|event| matches!(&event.event, Event::ToolResult { call_id, .. } if call_id == "first-call"))
.unwrap();
let second_started = events
.iter()
.position(|event| matches!(&event.event, Event::ToolExecutionStarted { call_id, .. } if call_id == "second-call"))
.unwrap();
assert!(first_result < second_started);
}
struct WaitHook;
#[async_trait]
impl Hook for WaitHook {
async fn before_tool(
&self,
_context: &af_agent::ToolExecutionContext,
_tool: &str,
_arguments: &Value,
) -> HookDecision {
HookDecision::WaitForInput {
kind: "action".into(),
payload: json!({"reason":"confirm"}),
}
}
}
struct BlockingBeforeHook(Arc<std::sync::atomic::AtomicBool>);
#[async_trait]
impl Hook for BlockingBeforeHook {
async fn before_tool(
&self,
context: &af_agent::ToolExecutionContext,
_: &str,
_: &Value,
) -> HookDecision {
while !context.cancellation.is_cancelled() {
tokio::time::sleep(Duration::from_millis(5)).await;
}
self.0.store(true, Ordering::SeqCst);
HookDecision::Continue
}
}
struct BlockingAfterHook(Arc<std::sync::atomic::AtomicBool>);
struct FailingAfterHook;
#[async_trait]
impl Hook for FailingAfterHook {
async fn after_tool(
&self,
_: &af_agent::ToolExecutionContext,
_: &str,
_: &Value,
_: &Value,
) -> Result<(), af_agent::PluginError> {
Err(af_agent::PluginError::Mount("observer failed".into()))
}
}
#[async_trait]
impl Hook for BlockingAfterHook {
async fn after_tool(
&self,
context: &af_agent::ToolExecutionContext,
_: &str,
_: &Value,
_: &Value,
) -> Result<(), af_agent::PluginError> {
while !context.cancellation.is_cancelled() {
tokio::time::sleep(Duration::from_millis(5)).await;
}
self.0.store(true, Ordering::SeqCst);
Ok(())
}
}
#[tokio::test]
async fn tool_hooks_share_run_cancellation_and_finish_before_results_commit() {
for (before, after) in [(true, false), (false, true)] {
let observed = Arc::new(std::sync::atomic::AtomicBool::new(false));
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([response(
calls(&[("call", "fast")]),
)]))));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(DelayTool {
name: "fast",
delay_ms: 0,
}))
.unwrap();
let hooks: Vec<Arc<dyn Hook>> = if before {
vec![Arc::new(BlockingBeforeHook(Arc::clone(&observed)))]
} else {
vec![Arc::new(BlockingAfterHook(Arc::clone(&observed)))]
};
let history = initial();
let writer = Arc::new(MemoryWriter::new(history.clone()));
let cancellation = CancellationToken::default();
let task = tokio::spawn({
let writer = Arc::clone(&writer);
let cancellation = cancellation.clone();
async move {
AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.with_hooks(hooks)
.run(request(history), writer.as_ref(), cancellation)
.await
}
});
tokio::time::sleep(Duration::from_millis(30)).await;
cancellation.cancel();
let outcome = task.await.unwrap().unwrap();
assert_eq!(outcome.status, "cancelled");
assert!(observed.load(Ordering::SeqCst));
if after {
let result_seq = writer
.events()
.iter()
.find(|event| matches!(event.event, Event::ToolResult { .. }))
.map(|event| event.seq)
.unwrap();
assert!(observed.load(Ordering::SeqCst));
assert!(result_seq > 0);
}
}
}
#[tokio::test]
async fn after_hook_failure_cannot_replace_a_committed_tool_result() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response(calls(&[("call", "fast")])),
response(ChatMessage::assistant("done")),
]))));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(DelayTool {
name: "fast",
delay_ms: 0,
}))
.unwrap();
let history = initial();
let writer = MemoryWriter::new(history.clone());
AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.with_hooks(vec![Arc::new(FailingAfterHook)])
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
let events = writer.events();
let result = events
.iter()
.find(|event| matches!(event.event, Event::ToolResult { .. }))
.unwrap();
assert!(matches!(
result.event,
Event::ToolResult {
is_error: false,
..
}
));
let observer = events
.iter()
.find(|event| matches!(&event.event, Event::Extension { event_type, .. } if event_type == "after_tool_failed"))
.unwrap();
assert!(observer.seq > result.seq);
}
struct WaitForSlow;
#[async_trait]
impl Hook for WaitForSlow {
async fn before_tool(
&self,
_context: &af_agent::ToolExecutionContext,
tool: &str,
_arguments: &Value,
) -> HookDecision {
if tool == "slow" {
HookDecision::WaitForInput {
kind: "action".into(),
payload: json!({"reason":"confirm"}),
}
} else {
HookDecision::Continue
}
}
}
#[tokio::test]
async fn approval_wait_closes_every_non_waiting_sibling_call() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response(calls(&[("fast-call", "fast"), ("slow-call", "slow")])),
response(ChatMessage::assistant("approved")),
]))));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(DelayTool {
name: "fast",
delay_ms: 0,
}))
.unwrap()
.register(Arc::new(DelayTool {
name: "slow",
delay_ms: 0,
}))
.unwrap();
let writer = MemoryWriter::new(initial());
let runtime = AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.with_hooks(vec![Arc::new(WaitForSlow)]);
runtime
.run(
request(writer.events()),
&writer,
CancellationToken::default(),
)
.await
.unwrap();
let waiting = SessionProjection::replay(&writer.events()).unwrap();
assert_eq!(
waiting.open_tool_calls.keys().cloned().collect::<Vec<_>>(),
["slow-call"]
);
let interaction_id = waiting.waiting_interaction_id.unwrap();
writer
.append(vec![
Event::InteractionResolved {
run_id: "run".into(),
interaction_id: interaction_id.clone(),
resolution: af_agent_session::InteractionResolution::Confirmed,
payload: Value::Null,
},
Event::RunResumed {
run_id: "run".into(),
interaction_id,
},
])
.await
.unwrap();
let outcome = runtime
.run(
request(writer.events()),
&writer,
CancellationToken::default(),
)
.await
.unwrap();
assert_eq!(outcome.final_text.as_deref(), Some("approved"));
assert!(SessionProjection::replay(&writer.events())
.unwrap()
.open_tool_calls
.is_empty());
}
struct AroundHook;
#[async_trait]
impl Hook for AroundHook {
fn around_tool<'a>(
&'a self,
_context: &'a af_agent::ToolExecutionContext,
_tool: &'a str,
_arguments: &'a Value,
next: ToolExecutionFuture<'a>,
) -> ToolExecutionFuture<'a> {
Box::pin(async move {
let mut result = next.await?;
result["wrapped"] = json!(true);
Ok(result)
})
}
}
#[tokio::test]
async fn around_hook_wraps_the_actual_tool_execution() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response(calls(&[("call", "fast")])),
response(ChatMessage::assistant("done")),
]))));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(DelayTool {
name: "fast",
delay_ms: 1,
}))
.unwrap();
let history = initial();
let writer = MemoryWriter::new(history.clone());
AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.with_hooks(vec![Arc::new(AroundHook)])
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::ToolResult { result, .. } if result["wrapped"] == json!(true)
)));
}
struct InvalidAroundHook;
#[async_trait]
impl Hook for InvalidAroundHook {
fn around_tool<'a>(
&'a self,
_context: &'a af_agent::ToolExecutionContext,
_tool: &'a str,
_arguments: &'a Value,
next: ToolExecutionFuture<'a>,
) -> ToolExecutionFuture<'a> {
Box::pin(async move {
let mut result = next.await?;
result["tool"] = json!(7);
Ok(result)
})
}
}
#[tokio::test]
async fn final_hook_output_must_match_the_tool_contract() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response(calls(&[("call", "fast")])),
response(ChatMessage::assistant("recovered")),
]))));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(DelayTool {
name: "fast",
delay_ms: 1,
}))
.unwrap();
let history = initial();
let writer = MemoryWriter::new(history.clone());
AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.with_hooks(vec![Arc::new(InvalidAroundHook)])
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::ToolResult { is_error: true, result, .. } if result.to_string().contains("invalid tool output")
)));
}
struct QuestionHook;
#[async_trait]
impl Hook for QuestionHook {
async fn before_tool(
&self,
_context: &af_agent::ToolExecutionContext,
_tool: &str,
_arguments: &Value,
) -> HookDecision {
HookDecision::WaitForInput {
kind: "question".into(),
payload: json!({"question":"Which value?"}),
}
}
}
struct CountingTool(Arc<AtomicUsize>);
#[async_trait]
impl Tool for CountingTool {
fn name(&self) -> &str {
"count"
}
fn description(&self) -> &str {
"must not run before replanning"
}
fn parameters(&self) -> Value {
json!({"type":"object"})
}
fn output_schema(&self) -> Value {
json!({"type":"object","required":["ok"],"additionalProperties":false,"properties":{"ok":{"type":"boolean"}}})
}
async fn call(&self, _args: Value) -> Result<Value, String> {
self.0.fetch_add(1, Ordering::SeqCst);
Ok(json!({"ok":true}))
}
}
struct CostTool;
#[async_trait]
impl Tool for CostTool {
fn name(&self) -> &str {
"cost"
}
fn description(&self) -> &str {
"cost test"
}
fn parameters(&self) -> Value {
json!({"type":"object"})
}
fn output_schema(&self) -> Value {
json!({"type":"object","required":["ok"],"properties":{"ok":{"type":"boolean"}}})
}
fn meta(&self) -> ToolMeta {
ToolMeta {
cost_units: 7,
..ToolMeta::default()
}
}
async fn call(&self, _args: Value) -> Result<Value, String> {
Ok(json!({"ok":true}))
}
}
#[tokio::test]
async fn terminal_tool_cost_is_recorded_once() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response(calls(&[("cost-call", "cost")])),
response(ChatMessage::assistant("done")),
]))));
let mut tools = ToolRegistry::new();
tools.register(Arc::new(CostTool)).unwrap();
let history = initial();
let writer = MemoryWriter::new(history.clone());
AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
let events = writer.events();
assert_eq!(
events
.iter()
.filter(|event| matches!(
&event.event,
Event::UsageRecorded { operation_id, cost_units: 7, .. }
if operation_id == "tool:cost-call"
))
.count(),
1
);
assert_eq!(
SessionProjection::replay(&events)
.unwrap()
.billable_units_for("run"),
11
);
}
#[tokio::test]
async fn answered_question_is_model_visible_and_does_not_execute_stale_arguments() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response(calls(&[("question-call", "count")])),
response(ChatMessage::assistant("replanned")),
]))));
let count = Arc::new(AtomicUsize::new(0));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(CountingTool(Arc::clone(&count))))
.unwrap();
let history = initial();
let writer = MemoryWriter::new(history.clone());
let runtime = AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.with_hooks(vec![Arc::new(QuestionHook)]);
runtime
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
let interaction_id = SessionProjection::replay(&writer.events())
.unwrap()
.waiting_interaction_id
.unwrap();
writer
.append(vec![
Event::InteractionResolved {
run_id: "run".into(),
interaction_id: interaction_id.clone(),
resolution: af_agent_session::InteractionResolution::Answered,
payload: json!({"answer":"blue"}),
},
Event::RunResumed {
run_id: "run".into(),
interaction_id,
},
])
.await
.unwrap();
let outcome = runtime
.run(
request(writer.events()),
&writer,
CancellationToken::default(),
)
.await
.unwrap();
assert_eq!(outcome.final_text.as_deref(), Some("replanned"));
assert_eq!(count.load(Ordering::SeqCst), 0);
assert!(writer.events().iter().any(|event| match &event.event {
Event::ModelRequestPrepared { request, .. } =>
request
.to_string()
.contains("User answered the pending question")
&& request.to_string().contains("blue"),
_ => false,
}));
}
#[tokio::test]
async fn retry_wait_and_crash_recovery_are_event_sourced() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
Err(LlmError::Api {
status: 503,
body: "retry".into(),
}),
response(calls(&[("call", "slow")])),
response(ChatMessage::assistant("approved")),
]))));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(DelayTool {
name: "slow",
delay_ms: 0,
}))
.unwrap();
let history = initial();
let writer = MemoryWriter::new(history.clone());
let runtime = AgentRuntime::new(model, "model", "system")
.with_tools(tools)
.with_hooks(vec![Arc::new(WaitHook)])
.with_limits(RuntimeLimits {
provider_attempts: 2,
..RuntimeLimits::default()
});
let outcome = runtime
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(outcome.status, "waiting_for_input");
let waiting = SessionProjection::replay(&writer.events()).unwrap();
assert!(waiting.usage_for("run").0 > 1);
assert!(writer
.events()
.iter()
.any(|event| matches!(event.event, Event::RetryScheduled { .. })));
let recovery_writer = MemoryWriter::new(writer.events());
recovery_writer
.append(recovery_events(&waiting))
.await
.unwrap();
assert!(SessionProjection::replay(&recovery_writer.events())
.unwrap()
.active_run_id
.is_none());
let interaction_id = waiting.waiting_interaction_id.clone().unwrap();
writer
.append(vec![
Event::InteractionResolved {
run_id: "run".into(),
interaction_id: interaction_id.clone(),
resolution: af_agent_session::InteractionResolution::Confirmed,
payload: Value::Null,
},
Event::RunResumed {
run_id: "run".into(),
interaction_id,
},
])
.await
.unwrap();
let resumed = runtime
.run(
request(writer.events()),
&writer,
CancellationToken::default(),
)
.await
.unwrap();
assert_eq!(resumed.final_text.as_deref(), Some("approved"));
assert!(SessionProjection::replay(&writer.events())
.unwrap()
.active_run_id
.is_none());
}
#[tokio::test]
async fn confirmation_is_bound_to_the_exact_tool_call_event() {
let count = Arc::new(AtomicUsize::new(0));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(CountingTool(Arc::clone(&count))))
.unwrap();
let writer = MemoryWriter::new(initial());
writer
.append(vec![
Event::InputClaimed {
input_id: "input".into(),
run_id: "run".into(),
},
Event::RunStarted {
run_id: "run".into(),
input_id: "input".into(),
},
Event::TurnStarted {
run_id: "run".into(),
turn: 1,
},
Event::UserMessage {
run_id: "run".into(),
content: text("hello"),
},
Event::StepStarted {
run_id: "run".into(),
step: 1,
},
Event::ToolCall {
run_id: "run".into(),
step: 1,
call_id: "new-call".into(),
tool: "count".into(),
arguments: json!({}),
},
Event::ToolAuthorization {
run_id: "run".into(),
step: 1,
call_id: "new-call".into(),
status: af_agent_session::ToolAuthorizationStatus::Waiting,
reason: None,
},
Event::InteractionRequested {
run_id: "run".into(),
interaction_id: "old-confirmation".into(),
kind: af_agent_session::InteractionKind::Action,
payload: json!({"call_id":"old-call","source_event_seq":1}),
},
Event::RunWaiting {
run_id: "run".into(),
interaction_id: "old-confirmation".into(),
},
Event::InteractionResolved {
run_id: "run".into(),
interaction_id: "old-confirmation".into(),
resolution: af_agent_session::InteractionResolution::Confirmed,
payload: Value::Null,
},
Event::RunResumed {
run_id: "run".into(),
interaction_id: "old-confirmation".into(),
},
])
.await
.unwrap();
let outcome = AgentRuntime::new(
Arc::new(ScriptedModel(Mutex::new(VecDeque::new()))),
"model",
"system",
)
.with_tools(tools)
.run(
request(writer.events()),
&writer,
CancellationToken::default(),
)
.await
.unwrap();
assert_eq!(outcome.status, "failed");
assert_eq!(count.load(Ordering::SeqCst), 0);
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::RunFinished { error_code: Some(code), .. } if code == "tool_outcome_unknown"
)));
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::ToolResult { result, .. }
if result.to_string().contains("Verify external state")
)));
}
#[tokio::test]
async fn crash_after_tool_execution_start_fails_without_repeating_the_effect() {
let count = Arc::new(AtomicUsize::new(0));
let mut tools = ToolRegistry::new();
tools
.register(Arc::new(CountingTool(Arc::clone(&count))))
.unwrap();
let writer = MemoryWriter::new(initial());
writer
.append(vec![
Event::InputClaimed {
input_id: "input".into(),
run_id: "run".into(),
},
Event::RunStarted {
run_id: "run".into(),
input_id: "input".into(),
},
Event::TurnStarted {
run_id: "run".into(),
turn: 1,
},
Event::UserMessage {
run_id: "run".into(),
content: text("hello"),
},
Event::StepStarted {
run_id: "run".into(),
step: 1,
},
Event::ToolCall {
run_id: "run".into(),
step: 1,
call_id: "call".into(),
tool: "count".into(),
arguments: json!({}),
},
Event::ToolAuthorization {
run_id: "run".into(),
step: 1,
call_id: "call".into(),
status: af_agent_session::ToolAuthorizationStatus::Allowed,
reason: None,
},
Event::ToolExecutionStarted {
run_id: "run".into(),
step: 1,
call_id: "call".into(),
},
])
.await
.unwrap();
let outcome = AgentRuntime::new(
Arc::new(ScriptedModel(Mutex::new(VecDeque::new()))),
"model",
"system",
)
.with_tools(tools)
.run(
request(writer.events()),
&writer,
CancellationToken::default(),
)
.await
.unwrap();
assert_eq!(outcome.status, "failed");
assert_eq!(count.load(Ordering::SeqCst), 0);
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::RunFinished { error_code: Some(code), .. } if code == "tool_outcome_unknown"
)));
}
#[tokio::test]
async fn unavailable_tool_on_resume_becomes_a_durable_error_result() {
let writer = MemoryWriter::new(initial());
writer
.append(vec![
Event::InputClaimed {
input_id: "input".into(),
run_id: "run".into(),
},
Event::RunStarted {
run_id: "run".into(),
input_id: "input".into(),
},
Event::TurnStarted {
run_id: "run".into(),
turn: 1,
},
Event::UserMessage {
run_id: "run".into(),
content: text("hello"),
},
Event::StepStarted {
run_id: "run".into(),
step: 1,
},
Event::ToolCall {
run_id: "run".into(),
step: 1,
call_id: "call".into(),
tool: "removed-tool".into(),
arguments: json!({}),
},
Event::ToolAuthorization {
run_id: "run".into(),
step: 1,
call_id: "call".into(),
status: af_agent_session::ToolAuthorizationStatus::Allowed,
reason: None,
},
])
.await
.unwrap();
let outcome = AgentRuntime::new(
Arc::new(ScriptedModel(Mutex::new(VecDeque::from([response(
ChatMessage::assistant("recovered"),
)])))),
"model",
"system",
)
.run(
request(writer.events()),
&writer,
CancellationToken::default(),
)
.await
.unwrap();
assert_eq!(outcome.final_text.as_deref(), Some("recovered"));
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::ToolResult { call_id, is_error: true, result, .. }
if call_id == "call" && result.to_string().contains("unavailable")
)));
}
struct Pressure;
impl TokenMeter for Pressure {
fn count(&self, _: &str, _: &[ChatMessage]) -> u64 {
10
}
}
struct PressureUntilSummary;
impl TokenMeter for PressureUntilSummary {
fn count(&self, _: &str, messages: &[ChatMessage]) -> u64 {
if messages.iter().any(|message| {
message
.content
.as_deref()
.is_some_and(|content| content.starts_with("Conversation summary:"))
}) {
1
} else {
10
}
}
}
#[tokio::test]
async fn default_model_compactor_is_durable_and_uses_an_idempotent_attempt() {
let model = Arc::new(CapturingModel {
responses: Mutex::new(VecDeque::from([
response(ChatMessage::assistant("Objective: finish the task")),
response(ChatMessage::assistant("done")),
])),
requests: Mutex::new(Vec::new()),
});
let history = initial();
let writer = MemoryWriter::new(history.clone());
let outcome = AgentRuntime::new(model.clone(), "model", "system")
.with_token_meter(Arc::new(PressureUntilSummary))
.with_limits(RuntimeLimits {
max_tokens: 5,
..RuntimeLimits::default()
})
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(outcome.final_text.as_deref(), Some("done"));
let requests = model.requests.lock().unwrap();
assert_eq!(requests.len(), 2);
assert_eq!(
requests[0].provider_attempt_id.as_deref(),
Some("run:compaction:1:attempt:1")
);
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::ModelRequestPrepared { provider_attempt_id, operation_id, .. }
if provider_attempt_id == "run:compaction:1:attempt:1"
&& operation_id == "compaction:1:attempt:1"
)));
assert!(writer.events().iter().any(|event| matches!(
&event.event,
Event::SummaryReplaced { summary, .. } if summary.contains("Objective")
)));
}
#[tokio::test]
async fn compaction_crash_takeover_closes_run_and_keeps_reserved_usage_billable() {
let model = Arc::new(CapturingModel {
responses: Mutex::new(VecDeque::new()),
requests: Mutex::new(Vec::new()),
});
let history = initial();
let writer = CrashAfterCompactionStartWriter {
inner: MemoryWriter::new(history.clone()),
crashed: std::sync::atomic::AtomicBool::new(false),
};
let runtime = AgentRuntime::new(model.clone(), "model", "system")
.with_token_meter(Arc::new(Pressure))
.with_limits(RuntimeLimits {
max_tokens: 1,
..RuntimeLimits::default()
});
assert!(matches!(
runtime
.run(request(history), &writer, CancellationToken::default())
.await,
Err(RuntimeError::Event(error)) if error.contains("simulated crash")
));
assert!(model.requests.lock().unwrap().is_empty());
let crashed = writer.inner.events();
let (reserved_prompt_tokens, reserved_completion_tokens) = crashed
.iter()
.find_map(|event| match &event.event {
Event::ModelRequestPrepared {
provider_attempt_id,
operation_id,
reserved_prompt_tokens,
reserved_completion_tokens,
..
} if provider_attempt_id == "run:compaction:1:attempt:1"
&& operation_id == "compaction:1:attempt:1" =>
{
Some((*reserved_prompt_tokens, *reserved_completion_tokens))
}
_ => None,
})
.unwrap();
assert!(reserved_prompt_tokens > 0);
assert_eq!(reserved_completion_tokens, 2_048);
assert_eq!(
SessionProjection::replay(&crashed)
.unwrap()
.billable_units_for("run"),
reserved_prompt_tokens + reserved_completion_tokens
);
let outcome = runtime
.run(
request(writer.inner.events()),
&writer,
CancellationToken::default(),
)
.await
.unwrap();
assert_eq!(outcome.status, "failed");
assert_eq!(outcome.prompt_tokens, reserved_prompt_tokens);
assert_eq!(outcome.completion_tokens, reserved_completion_tokens);
assert!(model.requests.lock().unwrap().is_empty());
let events = writer.inner.events();
let projection = SessionProjection::replay(&events).unwrap();
assert!(projection.active_run_id.is_none());
assert!(projection.open_compaction.is_none());
assert!(projection.open_steps.is_empty());
assert!(projection.open_turn.is_none());
assert_eq!(
projection.billable_units_for("run"),
reserved_prompt_tokens + reserved_completion_tokens
);
assert_eq!(
projection.usage_for("run"),
(reserved_prompt_tokens, reserved_completion_tokens)
);
assert_eq!(
events
.iter()
.filter(|event| matches!(
&event.event,
Event::UsageRecorded { operation_id, .. }
if operation_id == "compaction:1:attempt:1"
))
.count(),
1
);
assert!(events.iter().any(|event| matches!(
&event.event,
Event::RunFinished { error_code: Some(error), .. } if error == "worker_restarted"
)));
}
struct Summary;
#[async_trait]
impl Compactor for Summary {
async fn summarize(
&self,
_: &str,
_: &[ChatMessage],
_: &str,
_: CancellationToken,
_: std::time::Instant,
) -> Result<CompactionResult, RuntimeError> {
Ok(CompactionResult {
summary: "stable summary".into(),
prompt_tokens: 4,
completion_tokens: 2,
})
}
}
struct FailingSummary;
#[async_trait]
impl Compactor for FailingSummary {
async fn summarize(
&self,
_: &str,
_: &[ChatMessage],
_: &str,
_: CancellationToken,
_: std::time::Instant,
) -> Result<CompactionResult, RuntimeError> {
Err(RuntimeError::Model("compactor unavailable".into()))
}
}
struct BlockingSummary(Arc<std::sync::atomic::AtomicBool>);
#[async_trait]
impl Compactor for BlockingSummary {
async fn summarize(
&self,
_: &str,
_: &[ChatMessage],
_: &str,
cancellation: CancellationToken,
_: std::time::Instant,
) -> Result<CompactionResult, RuntimeError> {
while !cancellation.is_cancelled() {
tokio::time::sleep(Duration::from_millis(5)).await;
}
self.0.store(true, Ordering::SeqCst);
Err(RuntimeError::Cancelled)
}
}
#[tokio::test]
async fn compactor_observes_run_cancellation() {
let observed = Arc::new(std::sync::atomic::AtomicBool::new(false));
let history = initial();
let writer = Arc::new(MemoryWriter::new(history.clone()));
let cancellation = CancellationToken::default();
let task = tokio::spawn({
let writer = Arc::clone(&writer);
let cancellation = cancellation.clone();
let observed = Arc::clone(&observed);
async move {
AgentRuntime::new(
Arc::new(ScriptedModel(Mutex::new(VecDeque::new()))),
"model",
"system",
)
.with_token_meter(Arc::new(Pressure))
.with_compactor(Arc::new(BlockingSummary(observed)))
.with_limits(RuntimeLimits {
max_tokens: 1,
..RuntimeLimits::default()
})
.run(request(history), writer.as_ref(), cancellation)
.await
}
});
tokio::time::sleep(Duration::from_millis(30)).await;
cancellation.cancel();
assert_eq!(task.await.unwrap().unwrap().status, "cancelled");
assert!(observed.load(Ordering::SeqCst));
}
#[tokio::test]
async fn failed_compaction_closes_its_event_pair() {
let history = initial();
let writer = MemoryWriter::new(history.clone());
let error = AgentRuntime::new(
Arc::new(ScriptedModel(Mutex::new(VecDeque::new()))),
"model",
"system",
)
.with_token_meter(Arc::new(Pressure))
.with_compactor(Arc::new(FailingSummary))
.with_limits(RuntimeLimits {
max_tokens: 1,
..RuntimeLimits::default()
})
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap_err();
assert!(matches!(error, RuntimeError::Model(_)));
let events = writer.events();
assert_eq!(
events
.iter()
.filter(|event| matches!(event.event, Event::CompactionStarted { .. }))
.count(),
1
);
assert!(events.iter().any(|event| matches!(
&event.event,
Event::CompactionFinished { status, error: Some(error), .. }
if status == "failed" && error.contains("compactor unavailable")
)));
assert!(SessionProjection::replay(&events)
.unwrap()
.open_compaction
.is_none());
}
#[tokio::test]
async fn input_during_compaction_retries_at_the_next_boundary_without_sticking() {
let history = initial();
let writer = CompactionRaceWriter {
inner: MemoryWriter::new(history.clone()),
injected: std::sync::atomic::AtomicBool::new(false),
};
let outcome = AgentRuntime::new(
Arc::new(ScriptedModel(Mutex::new(VecDeque::from([response(
ChatMessage::assistant("done"),
)])))),
"model",
"system",
)
.with_token_meter(Arc::new(Pressure))
.with_compactor(Arc::new(Summary))
.with_limits(RuntimeLimits {
max_tokens: 1,
..RuntimeLimits::default()
})
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(outcome.status, "max_steps_reached");
let events = writer.inner.events();
assert!(events.iter().any(|event| matches!(
&event.event,
Event::CompactionFinished { status, error: Some(error), .. }
if status == "failed" && error == "surface_changed"
)));
let projection = SessionProjection::replay(&events).unwrap();
assert!(projection.active_run_id.is_none());
assert!(projection.queued_inputs.is_empty());
}
#[tokio::test]
async fn cancellation_repairs_open_tool_step_turn_and_compaction() {
let writer = MemoryWriter::new(initial());
writer
.append(vec![
Event::InputClaimed {
input_id: "input".into(),
run_id: "run".into(),
},
Event::RunStarted {
run_id: "run".into(),
input_id: "input".into(),
},
Event::TurnStarted {
run_id: "run".into(),
turn: 1,
},
Event::StepStarted {
run_id: "run".into(),
step: 1,
},
Event::ToolCall {
run_id: "run".into(),
step: 1,
call_id: "call".into(),
tool: "slow".into(),
arguments: json!({}),
},
Event::CompactionStarted {
run_id: "run".into(),
compaction_id: "compaction".into(),
source_through_seq: 7,
},
])
.await
.unwrap();
let projection = SessionProjection::replay(&writer.events()).unwrap();
let terminal = af_agent_runtime::cancel_events(&projection);
assert!(terminal.iter().any(|event| matches!(
event,
Event::CompactionFinished { status, .. } if status == "failed"
)));
writer.append(terminal).await.unwrap();
let projection = SessionProjection::replay(&writer.events()).unwrap();
assert!(projection.active_run_id.is_none());
assert!(projection.open_tool_calls.is_empty());
assert!(projection.open_steps.is_empty());
assert!(projection.open_compaction.is_none());
}
#[tokio::test]
async fn token_pressure_compacts_with_real_event_sequence() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([response(
ChatMessage::assistant("done"),
)]))));
let history = initial();
let writer = MemoryWriter::new(history.clone());
let outcome = AgentRuntime::new(model, "model", "system")
.with_token_meter(Arc::new(Pressure))
.with_compactor(Arc::new(Summary))
.with_limits(RuntimeLimits {
max_tokens: 1,
..RuntimeLimits::default()
})
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(outcome.status, "max_steps_reached");
let events = writer.events();
let started_seq = events
.iter()
.find(|event| matches!(event.event, Event::CompactionStarted { .. }))
.unwrap()
.seq;
let replacement = events
.iter()
.find_map(|event| match &event.event {
Event::SummaryReplaced { through_seq, .. } => Some(*through_seq),
_ => None,
})
.unwrap();
assert_eq!(replacement + 1, started_seq);
assert!(
SessionProjection::replay(&events)
.unwrap()
.usage_for("run")
.0
>= 4
);
}
#[tokio::test]
async fn non_retryable_provider_error_is_not_retried() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
Err(LlmError::Api {
status: 400,
body: "bad request".into(),
}),
response(ChatMessage::assistant("must not run")),
]))));
let history = initial();
let writer = MemoryWriter::new(history.clone());
let error = AgentRuntime::new(model, "model", "system")
.with_limits(RuntimeLimits {
provider_attempts: 2,
..RuntimeLimits::default()
})
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap_err();
assert!(matches!(error, RuntimeError::Model(_)));
assert!(!writer
.events()
.iter()
.any(|event| matches!(event.event, Event::RetryScheduled { .. })));
assert!(
SessionProjection::replay(&writer.events())
.unwrap()
.usage_for("run")
.0
> 0
);
}
#[tokio::test]
async fn unknown_tool_is_a_model_visible_result_instead_of_a_failed_run() {
let model = Arc::new(ScriptedModel(Mutex::new(VecDeque::from([
response(calls(&[("unknown-call", "missing.tool")])),
response(ChatMessage::assistant("recovered")),
]))));
let history = initial();
let writer = MemoryWriter::new(history.clone());
let outcome = AgentRuntime::new(model, "model", "system")
.run(request(history), &writer, CancellationToken::default())
.await
.unwrap();
assert_eq!(outcome.final_text.as_deref(), Some("recovered"));
assert!(writer.events().iter().any(|event| matches!(
event.event,
Event::ToolResult {
ref call_id,
is_error: true,
..
} if call_id == "unknown-call"
)));
}