use super::*;
#[test]
fn reasoning_only_turn_auto_continues_once() {
let temp = tempfile::TempDir::new().unwrap();
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let provider = ScriptedProvider::new(vec![
vec![
ProviderEvent::ReasoningSummaryComplete("thinking only".to_string()),
done(),
],
text_done("done"),
]);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"answer directly",
None,
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(output.text, "done");
assert_eq!(provider.requests().len(), 2);
assert!(provider.requests()[1].messages().iter().any(|message| {
message.role == crate::providers::MessageRole::User && message.content == "Continue"
}));
assert!(sink.outputs.iter().any(|event| matches!(
event,
OutputEvent::Diagnostic { message, .. }
if message == "Provider returned reasoning with no text or tool calls; auto-continuing once."
)));
let auto_continue = session
.read_events()
.unwrap()
.into_iter()
.find(|event| event.payload["reason"] == "reasoning_only_no_output")
.expect("reasoning-only auto-continue event");
assert_eq!(auto_continue.event_type, "user_input");
assert_eq!(auto_continue.payload["text"], "Continue");
assert_eq!(auto_continue.payload["auto_recovery"], true);
}
#[test]
fn reasoning_only_recovery_blocks_soft_compaction_after_tool_growth() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "large enough").unwrap();
let provider = ScriptedProvider::new(vec![
vec![
ProviderEvent::ReasoningSummaryComplete("thinking only".to_string()),
done(),
],
read_done("recovery_read"),
text_done("done"),
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let output = agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
tools: Some(&tools),
continuation_auto_compaction_policy: Some((1, "1 token".to_string())),
..run_request("read after reasoning", temp.path())
},
)
.unwrap();
assert_eq!(provider.requests().len(), 3);
assert_eq!(output.text, "done");
assert!(output.auto_compaction_blocked_by_recovery);
}
#[test]
fn incomplete_stream_recovery_blocks_soft_compaction_after_tool_growth() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "large enough").unwrap();
let provider = AutoContinueProvider::new(vec![
AutoContinueStep::TextThenEligibleTimeout("partial "),
AutoContinueStep::ReadThenDone,
AutoContinueStep::TextThenDone("done"),
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let output = agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
tools: Some(&tools),
continuation_auto_compaction_policy: Some((1, "1 token".to_string())),
..run_request("read after incomplete stream", temp.path())
},
)
.unwrap();
assert_eq!(provider.requests().len(), 3);
assert!(output.text.contains("partial"));
assert!(output.text.contains("done"));
assert!(output.recovered_incomplete_stream);
assert!(output.auto_compaction_blocked_by_recovery);
}
#[test]
fn optimization_harness_executes_large_synthetic_tool_set_before_continuation() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "file text").unwrap();
let tool_call_count = 128;
let first_turn = (0..tool_call_count)
.map(|index| read_call(&format!("call_{index}")))
.chain(std::iter::once(done()))
.collect::<Vec<_>>();
let provider = ScriptedProvider::new(vec![first_turn, text_done("done")]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"read many files",
Some(&tools),
None,
temp.path(),
None,
)
.unwrap();
let requests = provider.requests();
assert_eq!(output.tool_results.len(), tool_call_count);
assert!(output.tool_results.iter().all(|result| result.success));
assert_eq!(requests.len(), 2);
assert_eq!(requests[1].tool_results().len(), tool_call_count);
assert!(
requests[1]
.tool_results()
.iter()
.any(|result| result.call_id == format!("call_{}", tool_call_count - 1))
);
}
#[test]
fn steering_update_injects_after_tool_result_and_persists_as_user_input() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "file text").unwrap();
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let provider = ScriptedProvider::new(vec![read_done("call_steer"), text_done("done")]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let steering = AgentSteering::new();
steering
.try_enqueue("prefer concise answer".to_string())
.unwrap();
let mut sink = CapturingOutputSink::default();
agent
.run_print_with_tools_streaming_output_cancellable_with_steering(
&provider,
AgentRunRequest {
tools: Some(&tools),
session: Some(&session),
output_sink: Some(&mut sink),
..run_request("read file", temp.path())
},
steering,
)
.unwrap();
let requests = provider.requests();
assert_eq!(requests.len(), 2);
let continuation = requests[1].conversation_items();
let tool_index = continuation
.iter()
.position(|item| {
matches!(
item,
ProviderConversationItem::ToolResult(result)
if result.call_id == "call_steer" && result.output == "file text"
)
})
.expect("tool result in continuation");
let steering_index = continuation
.iter()
.position(|item| matches!(
item,
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User
&& message.content == "Steering update from user while current run was active:\n\nprefer concise answer"
))
.expect("steering user message in continuation");
assert!(tool_index < steering_index);
let user_prompt_events = sink
.outputs
.iter()
.filter_map(|event| match event {
OutputEvent::UserPrompt { text } => Some(text.as_str()),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(
user_prompt_events,
vec![
"read file",
"Steering update from user while current run was active:\n\nprefer concise answer",
]
);
let tool_result_event = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::ToolResult { call, .. } if call.id == "call_steer"))
.expect("tool result event");
let steering_prompt_event = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::UserPrompt { text } if text.contains("prefer concise answer")))
.expect("steering prompt event");
assert!(tool_result_event < steering_prompt_event);
let user_inputs = session
.read_events()
.unwrap()
.into_iter()
.filter(|event| event.event_type == "user_input")
.map(|event| event.payload["text"].as_str().unwrap().to_string())
.collect::<Vec<_>>();
assert_eq!(
user_inputs,
vec![
"read file".to_string(),
"Steering update from user while current run was active:\n\nprefer concise answer"
.to_string(),
]
);
}
#[cfg(unix)]
#[test]
fn steering_persistence_failure_retains_queue_and_blocks_continuation() {
use std::os::unix::fs::PermissionsExt;
let temp = tempfile::TempDir::new().unwrap();
let sessions_root = temp.path().join("sessions");
let session = crate::sessions::SessionManager::new(sessions_root.clone())
.create()
.unwrap();
let steering = AgentSteering::new();
steering.try_enqueue("must survive".to_string()).unwrap();
std::fs::set_permissions(&sessions_root, std::fs::Permissions::from_mode(0o500)).unwrap();
let persistence = SessionPersistence::new(Some(&session), temp.path());
let mut turn_state = AgentTurnState::default();
let mut output_sink = None;
let error = super::super::inject_steering_update_at_continuation_boundary(
Some(&steering),
steering.observe_collapsed(),
&mut turn_state,
&persistence,
&mut output_sink,
)
.unwrap_err()
.to_string();
std::fs::set_permissions(&sessions_root, std::fs::Permissions::from_mode(0o700)).unwrap();
assert!(error.contains("failed to persist steering input before provider continuation"));
assert_eq!(steering.pending_count(), 1);
assert_eq!(
steering.observe_collapsed().unwrap().text,
"Steering update from user while current run was active:\n\nmust survive"
);
assert!(
!session
.read_events()
.unwrap()
.iter()
.any(|event| event.event_type == "user_input")
);
}
#[test]
fn steering_boundary_injects_only_observed_batch_and_leaves_later_input_queued() {
let temp = tempfile::TempDir::new().unwrap();
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let steering = AgentSteering::new();
steering.try_enqueue("projected input".to_string()).unwrap();
let projected_batch = steering.observe_collapsed();
steering.try_enqueue("later input".to_string()).unwrap();
let persistence = SessionPersistence::new(Some(&session), temp.path());
let mut turn_state = AgentTurnState::default();
let mut output_sink = None;
assert!(
super::super::inject_steering_update_at_continuation_boundary(
Some(&steering),
projected_batch,
&mut turn_state,
&persistence,
&mut output_sink,
)
.unwrap()
);
assert_eq!(steering.pending_count(), 1);
assert_eq!(
steering.observe_collapsed().unwrap().text,
"Steering update from user while current run was active:\n\nlater input"
);
let persisted = session.read_events().unwrap();
assert_eq!(persisted.len(), 1);
assert_eq!(
persisted[0].payload["text"],
"Steering update from user while current run was active:\n\nprojected input"
);
}
#[test]
fn empty_steering_queue_does_not_inject_or_persist_mid_turn_input() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "file text").unwrap();
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let provider = ScriptedProvider::new(vec![read_done("call_empty_steer"), text_done("done")]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
agent
.run_print_with_tools_streaming_output_cancellable_with_steering(
&provider,
AgentRunRequest {
tools: Some(&tools),
session: Some(&session),
..run_request("read file", temp.path())
},
AgentSteering::new(),
)
.unwrap();
let requests = provider.requests();
assert_eq!(requests.len(), 2);
assert!(!requests[1].messages().iter().any(|message| {
message
.content
.contains("Steering update from user while current run was active")
}));
let user_input_count = session
.read_events()
.unwrap()
.into_iter()
.filter(|event| event.event_type == "user_input")
.count();
assert_eq!(user_input_count, 1);
}
#[test]
fn absent_steering_handle_does_not_inject_and_continuation_still_runs() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "file text").unwrap();
let provider = ScriptedProvider::new(vec![read_done("call_no_steer"), text_done("done")]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"read file",
Some(&tools),
None,
temp.path(),
None,
)
.unwrap();
let requests = provider.requests();
assert_eq!(output.text, "done");
assert_eq!(requests.len(), 2);
assert_eq!(requests[1].tool_results().len(), 1);
assert!(!requests[1].messages().iter().any(|message| {
message
.content
.contains("Steering update from user while current run was active")
}));
}
#[test]
fn text_only_steering_continues_after_assistant_action() {
let temp = tempfile::TempDir::new().unwrap();
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let provider = ScriptedProvider::new(vec![text_done("first"), text_done("second")]);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let steering = AgentSteering::new();
steering
.try_enqueue("prefer concise answer".to_string())
.unwrap();
let steering_probe = steering.clone();
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output_cancellable_with_steering(
&provider,
AgentRunRequest {
session: Some(&session),
output_sink: Some(&mut sink),
invocation_mode: InvocationMode::MissionControl,
..run_request("answer directly", temp.path())
},
steering,
)
.unwrap();
assert_eq!(output.text, "first\n\nsecond");
assert_eq!(steering_probe.pending_count(), 0);
let requests = provider.requests();
assert_eq!(requests.len(), 2);
let second_messages = requests[1].messages();
assert_eq!(second_messages.len(), 4);
assert_eq!(second_messages[1].content, "answer directly");
assert_eq!(second_messages[2].content, "first");
assert_eq!(
second_messages[3].content,
"Steering update from user while current run was active:\n\nprefer concise answer"
);
let first_complete = sink
.outputs
.iter()
.position(
|event| matches!(event, OutputEvent::AssistantComplete { text } if text == "first"),
)
.expect("first assistant complete");
let steering_prompt = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::UserPrompt { text } if text.starts_with("Steering update from user")))
.expect("steering prompt");
assert!(first_complete < steering_prompt);
let user_inputs = session
.read_events()
.unwrap()
.into_iter()
.filter(|event| event.event_type == "user_input")
.map(|event| event.payload["text"].as_str().unwrap().to_string())
.collect::<Vec<_>>();
assert_eq!(
user_inputs,
vec![
"answer directly".to_string(),
"Steering update from user while current run was active:\n\nprefer concise answer"
.to_string()
]
);
let assistant_outputs = session
.read_events()
.unwrap()
.into_iter()
.filter(|event| event.event_type == "assistant_output")
.map(|event| event.payload["text"].as_str().unwrap().to_string())
.collect::<Vec<_>>();
assert_eq!(assistant_outputs, vec!["first", "second"]);
}
#[test]
fn after_assistant_hook_context_follows_text_before_steering_continuation() {
let temp = tempfile::TempDir::new().unwrap();
let provider = ScriptedProvider::new(vec![text_done("first"), text_done("second")]);
let mut hook_settings = after_assistant_hook(
"memory",
"printf '%s' '{\"context_items\":[{\"role\":\"user\",\"content\":\"HOOK_CONTEXT\"}]}'",
None,
);
hook_settings.after_assistant[0].provider_context_injection = Some(true);
let hooks = HookRuntime::new(temp.path(), hook_settings, true).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let steering = AgentSteering::new();
steering
.try_enqueue("prefer concise answer".to_string())
.unwrap();
agent
.run_print_with_tools_streaming_output_cancellable_with_steering(
&provider,
AgentRunRequest {
initial_instructions: &[],
prompt: "answer directly",
prompt_origin: crate::output::UserPromptOrigin::User,
effective_prompt: None,
tools: None,
hooks: Some(&hooks),
session: None,
cwd: temp.path(),
output_sink: None,
cancellation: AgentCancellation::default(),
session_title_job: None,
semantic_progress_timeout: None,
invocation_mode: InvocationMode::MissionControl,
agent_id: None,
herdr_reporter: None,
ttsr: crate::config::TtsrSettings::default(),
continuation_auto_compaction_policy: None,
},
steering,
)
.unwrap();
let requests = provider.requests();
assert_eq!(requests.len(), 2);
let second_messages = requests[1].messages();
assert_eq!(second_messages.len(), 5);
assert_eq!(second_messages[1].content, "answer directly");
assert_eq!(second_messages[2].content, "first");
assert_eq!(second_messages[3].content, "HOOK_CONTEXT");
assert_eq!(
second_messages[4].content,
"Steering update from user while current run was active:\n\nprefer concise answer"
);
}
#[test]
fn final_assistant_boundary_without_pending_steering_emits_no_diagnostic() {
let temp = tempfile::TempDir::new().unwrap();
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let provider = ScriptedProvider::new(vec![text_done("done")]);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output_cancellable_with_steering(
&provider,
AgentRunRequest {
initial_instructions: &[],
prompt: "answer directly",
prompt_origin: crate::output::UserPromptOrigin::User,
effective_prompt: None,
tools: None,
hooks: None,
session: Some(&session),
cwd: temp.path(),
output_sink: Some(&mut sink),
cancellation: AgentCancellation::default(),
session_title_job: None,
semantic_progress_timeout: None,
invocation_mode: InvocationMode::Print,
agent_id: None,
herdr_reporter: None,
ttsr: crate::config::TtsrSettings::default(),
continuation_auto_compaction_policy: None,
},
AgentSteering::new(),
)
.unwrap();
assert_eq!(output.text, "done");
assert_eq!(provider.requests().len(), 1);
let user_input_count = session
.read_events()
.unwrap()
.into_iter()
.filter(|event| event.event_type == "user_input")
.count();
assert_eq!(user_input_count, 1);
}
#[test]
fn tool_turn_assistant_complete_precedes_chunk_flush_and_tool_lifecycle() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "tool text").unwrap();
let provider = ScriptedProvider::new(vec![
vec![text("I will inspect."), read_call("call_order"), done()],
vec![done()],
]);
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = SessionOrderProbeSink::new(&session);
agent
.run_print_with_tools_streaming_output(
&provider,
"read file",
Some(&tools),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(sink.events_at_assistant_complete, vec!["user_input"]);
let assistant_complete = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::AssistantComplete { text } if text == "I will inspect."))
.expect("assistant complete");
let tool_started = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::ToolStarted { call, .. } if call.id == "call_order"))
.expect("tool started");
assert!(assistant_complete < tool_started);
let event_types = session
.read_events()
.unwrap()
.into_iter()
.map(|event| event.event_type)
.collect::<Vec<_>>();
assert_eq!(event_types[1], "assistant_chunk");
assert_eq!(event_types[2], "tool_call");
}
#[test]
fn mixed_text_tool_turn_emits_assistant_complete_before_tools() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "tool text").unwrap();
let provider = ScriptedProvider::new(vec![
vec![text("I will inspect."), read_call("call_mixed"), done()],
text_done(" done"),
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"read file",
Some(&tools),
None,
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(output.text, "I will inspect.\n\n done");
let assistant_delta = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::AssistantDelta { text } if text == "I will inspect."))
.expect("live assistant delta");
let tool_started = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::ToolStarted { call, .. } if call.id == "call_mixed"))
.expect("tool started");
let tool_result = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::ToolResult { call, .. } if call.id == "call_mixed"))
.expect("tool result");
let first_complete = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::AssistantComplete { text } if text == "I will inspect."))
.expect("mixed-turn assistant complete");
assert!(assistant_delta < first_complete);
assert!(first_complete < tool_started);
assert!(tool_started < tool_result);
}
#[test]
fn post_tool_continuation_no_progress_error_preserves_tool_result_output() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "final tool text").unwrap();
let provider = PostToolContinuationFailingProvider {
requests: RecordedRequests::default(),
};
let tools = ToolRuntime::new(temp.path()).unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let error = agent
.run_print_with_tools_streaming_output(
&provider,
"read final file",
Some(&tools),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap_err()
.to_string();
assert!(error.contains("no semantic progress"), "{error}");
assert!(sink.outputs.iter().any(|event| matches!(
event,
OutputEvent::ToolResult { result, .. }
if result.tool_name == "read" && result.content == "final tool text"
)));
let requests = provider.requests();
assert_eq!(requests.len(), 2);
assert_eq!(requests[1].tool_results().len(), 1);
assert_eq!(requests[1].tool_results()[0].call_id, "call_final");
assert_eq!(requests[1].tool_results()[0].output, "final tool text");
assert!(requests[1].conversation_items().iter().any(|item| matches!(
item,
ProviderConversationItem::ToolResult(result)
if result.call_id == "call_final" && result.output == "final tool text"
)));
let events = session.read_events().unwrap();
let statuses = events
.iter()
.filter(|event| event.event_type == "turn_status")
.collect::<Vec<_>>();
assert_eq!(statuses.len(), 1);
assert_eq!(statuses[0].payload["status"], "failed");
let replay = crate::context::build_conversation_replay(Some(&session)).unwrap();
assert!(replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User
&& message.content == "read final file"
)));
assert!(replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::ToolResult(result) if result.call_id == "call_final"
)));
}
#[test]
fn subagents_final_result_is_sent_to_parent_continuation() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(ParentSubagentProvider {
requests: RecordedRequests::default(),
});
let parent_agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let base_tools = ToolRuntime::new(temp.path()).unwrap();
let subagent_config = crate::subagents::SubagentRunConfig {
parent_agent: parent_agent.clone(),
provider: provider.clone(),
provider_override: None,
parent_tools: base_tools.clone(),
parent_cwd: temp.path().to_path_buf(),
cancellation: AgentCancellation::default(),
profiles: Default::default(),
subagent_profiles_prompt: None,
sessions_root: None,
depth: 0,
parent_activity_id: None,
activity_sender: None,
inherited_hooks: None,
semantic_progress_timeout: std::time::Duration::from_millis(40),
schema_validation_max_retries: 2,
};
let tools = base_tools.with_subagents(move |arguments, context| {
let mut config = subagent_config.clone();
config.parent_activity_id = context.parent_activity_id;
config.activity_sender = context.activity_sender;
crate::subagents::dispatch_subagents(arguments, config)
});
parent_agent
.run_print_with_tools_streaming_output(
provider.as_ref(),
"run subagent",
Some(&tools),
None,
temp.path(),
None,
)
.unwrap();
let requests = provider.requests();
let parent_continuation = requests
.iter()
.find(|request| {
request
.tool_results()
.iter()
.any(|result| result.call_id == "parent_subagents")
})
.expect("parent continuation with subagents result");
let parent_tool_results = parent_continuation.tool_results();
let result = parent_tool_results
.iter()
.find(|result| result.call_id == "parent_subagents")
.unwrap();
assert_eq!(result.tool_name, "subagents");
assert!(
result.output.contains("child raw output"),
"{}",
result.output
);
for display_only in ["--- TOOL START", "STATUS:", "TASK_COUNT:"] {
assert!(
!result.output.contains(display_only),
"{display_only} leaked"
);
}
}
#[test]
fn subagent_activity_finishes_tools_before_assistant_finished() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "child file text").unwrap();
let provider = Arc::new(ParentSubagentMixedToolProvider {
requests: RecordedRequests::default(),
});
let parent_agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let base_tools = ToolRuntime::new(temp.path()).unwrap();
let subagent_config = crate::subagents::SubagentRunConfig {
parent_agent: parent_agent.clone(),
provider: provider.clone(),
provider_override: None,
parent_tools: base_tools.clone(),
parent_cwd: temp.path().to_path_buf(),
cancellation: AgentCancellation::default(),
profiles: Default::default(),
subagent_profiles_prompt: None,
sessions_root: None,
depth: 0,
parent_activity_id: None,
activity_sender: None,
inherited_hooks: None,
semantic_progress_timeout: std::time::Duration::from_millis(40),
schema_validation_max_retries: 2,
};
let tools = base_tools.with_subagents(move |arguments, context| {
let mut config = subagent_config.clone();
config.parent_activity_id = context.parent_activity_id;
config.activity_sender = context.activity_sender;
crate::subagents::dispatch_subagents(arguments, config)
});
let mut sink = ProbeActivitySink::new(|_| {});
parent_agent
.run_print_with_tools_streaming_output(
provider.as_ref(),
"run subagent",
Some(&tools),
None,
temp.path(),
Some(&mut sink),
)
.unwrap();
let events = sink.sender_events.lock().unwrap();
let assistant_started = events
.iter()
.position(|event| {
matches!(
event,
ActivityEvent::Started {
kind: ActivityKind::Assistant,
..
}
)
})
.expect("assistant started");
let assistant_delta = events
.iter()
.position(|event| matches!(event, ActivityEvent::Delta { preview, .. } if preview == "child will read"))
.expect("assistant delta");
let tool_preview = events
.iter()
.position(|event| matches!(event, ActivityEvent::ToolResultDetail { id, .. } if id.as_str().contains("child_read")))
.expect("child tool final detail");
let tool_finished = events
.iter()
.position(|event| matches!(event, ActivityEvent::Finished { id, .. } if id.as_str().contains("child_read")))
.expect("child tool finished");
let assistant_finished = events
.iter()
.position(|event| matches!(event, ActivityEvent::Finished { id, .. } if id.as_str().ends_with("/assistant")))
.expect("assistant finished");
assert!(assistant_started < assistant_delta);
assert!(assistant_delta < tool_preview);
assert!(tool_finished < assistant_finished);
assert!(tool_preview < assistant_finished);
}
#[test]
fn inherited_child_hook_context_stays_child_local() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "child file text").unwrap();
std::fs::write(
temp.path().join("hook.out"),
r#"{"context_items":[{"role":"user","content":"CHILD_HOOK_CONTEXT_SENTINEL"}]}"#,
)
.unwrap();
let provider = Arc::new(ParentSubagentMixedToolProvider {
requests: RecordedRequests::default(),
});
let parent_agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let session_root = temp.path().join("sessions");
let parent_session = crate::sessions::SessionManager::new(session_root.clone())
.create()
.unwrap();
let base_tools = ToolRuntime::new(temp.path()).unwrap();
let hooks = HookRuntime::new(
temp.path(),
crate::config::HookSettings {
enabled: true,
provider_context_injection: true,
after_tool: vec![crate::config::HookDefinition {
label: Some("child-memory".into()),
command: "cat hook.out".into(),
include_tools: vec!["read".into()],
..crate::config::HookDefinition::default()
}],
..crate::config::HookSettings::default()
},
true,
)
.unwrap();
let subagent_config = crate::subagents::SubagentRunConfig {
parent_agent: parent_agent.clone(),
provider: provider.clone(),
provider_override: None,
parent_tools: base_tools.clone(),
parent_cwd: temp.path().to_path_buf(),
cancellation: AgentCancellation::default(),
profiles: Default::default(),
subagent_profiles_prompt: None,
sessions_root: Some(session_root.clone()),
depth: 0,
parent_activity_id: None,
activity_sender: None,
inherited_hooks: Some(hooks),
semantic_progress_timeout: std::time::Duration::from_millis(40),
schema_validation_max_retries: 2,
};
let tools = base_tools.with_subagents(move |arguments, context| {
let mut config = subagent_config.clone();
config.parent_activity_id = context.parent_activity_id;
config.activity_sender = context.activity_sender;
crate::subagents::dispatch_subagents(arguments, config)
});
parent_agent
.run_print_with_tools_streaming_output(
provider.as_ref(),
"run subagent",
Some(&tools),
Some(&parent_session),
temp.path(),
None,
)
.unwrap();
let requests = provider.requests();
let child_continuation = requests
.iter()
.find(|request| {
request
.tool_results()
.iter()
.any(|result| result.call_id == "child_read")
})
.expect("child continuation with read result");
assert!(child_continuation.conversation_items().windows(2).any(|pair| matches!(
(&pair[0], &pair[1]),
(
ProviderConversationItem::ToolResult(result),
ProviderConversationItem::Message(message),
) if result.call_id == "child_read" && message.content == "CHILD_HOOK_CONTEXT_SENTINEL"
)));
let parent_continuation = requests
.iter()
.find(|request| {
request
.tool_results()
.iter()
.any(|result| result.call_id == "parent_subagents")
})
.expect("parent continuation with subagents result");
assert!(!format!("{parent_continuation:?}").contains("CHILD_HOOK_CONTEXT_SENTINEL"));
let parent_result = parent_continuation
.tool_results()
.into_iter()
.find(|result| result.call_id == "parent_subagents")
.unwrap();
assert!(!parent_result.output.contains("CHILD_HOOK_CONTEXT_SENTINEL"));
let parent_session_text = std::fs::read_to_string(parent_session.path()).unwrap();
assert!(!parent_session_text.contains("CHILD_HOOK_CONTEXT_SENTINEL"));
assert!(!parent_session_text.contains("hook_context_injection"));
let mut session_files = Vec::new();
let mut pending = vec![session_root.clone()];
while let Some(dir) = pending.pop() {
for entry in std::fs::read_dir(dir).unwrap() {
let path = entry.unwrap().path();
if path.is_dir() {
pending.push(path);
} else if path != parent_session.path() {
session_files.push(path);
}
}
}
let child_session_texts = session_files
.into_iter()
.map(|path| std::fs::read_to_string(path).unwrap())
.collect::<Vec<_>>();
assert!(child_session_texts.iter().any(|text| {
text.contains("CHILD_HOOK_CONTEXT_SENTINEL")
&& text.contains("hook_context_injection")
&& text.contains("provider_context_item")
}));
}
#[test]
fn inherited_child_hook_failure_is_sanitized_in_parent_provider_continuation() {
let temp = tempfile::TempDir::new().unwrap();
let provider = Arc::new(ParentSubagentHookFailureProvider {
requests: RecordedRequests::default(),
});
let parent_agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let base_tools = ToolRuntime::new(temp.path()).unwrap();
let hooks = HookRuntime::new(
temp.path(),
crate::config::HookSettings {
enabled: true,
before_tool: vec![crate::config::HookDefinition {
label: Some("PARENT_PROVIDER_LEAK_MARKER".into()),
command: "printf CHILD_STDOUT_PROVIDER_POISON; printf CHILD_STDERR_PROVIDER_POISON >&2; exit 7".into(),
failure_policy: Some(crate::config::HookFailurePolicy::Fail),
include_tools: vec!["write".into()],
..crate::config::HookDefinition::default()
}],
..crate::config::HookSettings::default()
},
true,
)
.unwrap();
let subagent_config = crate::subagents::SubagentRunConfig {
parent_agent: parent_agent.clone(),
provider: provider.clone(),
provider_override: None,
parent_tools: base_tools.clone(),
parent_cwd: temp.path().to_path_buf(),
cancellation: AgentCancellation::default(),
profiles: Default::default(),
subagent_profiles_prompt: None,
sessions_root: Some(temp.path().join("sessions")),
depth: 0,
parent_activity_id: None,
activity_sender: None,
inherited_hooks: Some(hooks),
semantic_progress_timeout: std::time::Duration::from_millis(40),
schema_validation_max_retries: 2,
};
let tools = base_tools.with_subagents(move |arguments, context| {
let mut config = subagent_config.clone();
config.parent_activity_id = context.parent_activity_id;
config.activity_sender = context.activity_sender;
crate::subagents::dispatch_subagents(arguments, config)
});
parent_agent
.run_print_with_tools_streaming_output(
provider.as_ref(),
"run subagent",
Some(&tools),
None,
temp.path(),
None,
)
.unwrap();
let requests = provider.requests();
let parent_continuation = requests
.iter()
.find(|request| {
request
.tool_results()
.iter()
.any(|result| result.call_id == "parent_subagents")
})
.expect("parent continuation with subagents result");
let result = parent_continuation
.tool_results()
.into_iter()
.find(|result| result.call_id == "parent_subagents")
.unwrap();
assert!(
result
.output
.contains("child subagent stopped by local hook policy"),
"{}",
result.output
);
for marker in [
"PARENT_PROVIDER_LEAK_MARKER",
"CHILD_STDOUT_PROVIDER_POISON",
"CHILD_STDERR_PROVIDER_POISON",
"hook_lifecycle",
"hook_diagnostic",
] {
assert!(
!result.output.contains(marker),
"{marker} leaked: {}",
result.output
);
assert!(
!format!("{parent_continuation:?}").contains(marker),
"{marker} leaked in request"
);
}
}
#[test]
fn tool_output_compression_changes_provider_output_only() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("before.txt"), "before\n").unwrap();
std::fs::write(temp.path().join("after.txt"), "RAW_PROVIDER_POISON\n").unwrap();
let provider = ScriptedProvider::new(vec![
vec![bash_call(
"call_diff",
"git --no-pager diff --no-index before.txt after.txt",
)],
vec![done()],
]);
let tools = ToolRuntime::new_with_settings(
temp.path(),
crate::tools::ToolSettings {
output_compression: crate::tools::ToolOutputCompressionSettings { enabled: true },
..crate::tools::ToolSettings::default()
},
)
.unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
initial_instructions: &[],
prompt: "diff files",
prompt_origin: crate::output::UserPromptOrigin::User,
effective_prompt: None,
tools: Some(&tools),
hooks: None,
session: Some(&session),
cwd: temp.path(),
output_sink: Some(&mut sink),
cancellation: AgentCancellation::default(),
session_title_job: None,
semantic_progress_timeout: None,
invocation_mode: crate::output::InvocationMode::Print,
agent_id: None,
herdr_reporter: None,
ttsr: crate::config::TtsrSettings::default(),
continuation_auto_compaction_policy: None,
},
)
.unwrap();
assert_eq!(output.tool_results.len(), 1);
assert!(
output.tool_results[0]
.content
.contains("RAW_PROVIDER_POISON")
);
assert!(
sink.blocks
.iter()
.any(|block| block.contains("RAW_PROVIDER_POISON"))
);
let session_tool_result = session
.read_events()
.unwrap()
.into_iter()
.find(|event| event.event_type == "tool_result")
.expect("session tool_result event");
assert!(
session_tool_result.payload["result"]["content"]
.as_str()
.unwrap()
.contains("RAW_PROVIDER_POISON")
);
let requests = provider.requests();
assert_eq!(requests[1].tool_results().len(), 1);
let provider_output = &requests[1].tool_results()[0].output;
assert!(provider_output.contains("[tool_output_compression]"));
assert!(provider_output.contains("rule: bash.git_diff"));
assert!(!provider_output.contains("RAW_PROVIDER_POISON"));
}
#[test]
fn cancellation_during_provider_callback_stops_before_tool_dispatch() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "secret").unwrap();
let cancel = Arc::new(AtomicBool::new(false));
let provider = CancelingProvider::new(Arc::clone(&cancel));
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let error = agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
initial_instructions: &[],
prompt: "read file",
prompt_origin: crate::output::UserPromptOrigin::User,
effective_prompt: None,
tools: Some(&tools),
hooks: None,
session: None,
cwd: temp.path(),
output_sink: Some(&mut sink),
cancellation: AgentCancellation::new(cancel),
session_title_job: None,
semantic_progress_timeout: None,
invocation_mode: crate::output::InvocationMode::Print,
agent_id: None,
herdr_reporter: None,
ttsr: crate::config::TtsrSettings::default(),
continuation_auto_compaction_policy: None,
},
)
.unwrap_err();
assert!(is_run_canceled(&error));
assert!(sink.blocks.is_empty());
assert!(sink.outputs.iter().all(|event| !matches!(
event,
OutputEvent::ToolStarted { .. } | OutputEvent::ToolResult { .. }
)));
}
#[test]
fn hooks_warn_continue_and_do_not_change_provider_tool_output() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![read_done("call_1"), vec![done()]]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let hooks = HookRuntime::new(
temp.path(),
crate::config::HookSettings {
enabled: true,
before_tool: vec![crate::config::HookDefinition {
label: Some("noisy".into()),
command: "printf hook-secret; exit 7".into(),
failure_policy: Some(crate::config::HookFailurePolicy::Warn),
..crate::config::HookDefinition::default()
}],
..crate::config::HookSettings::default()
},
true,
)
.unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
initial_instructions: &[],
prompt: "read file",
prompt_origin: crate::output::UserPromptOrigin::User,
effective_prompt: None,
tools: Some(&tools),
hooks: Some(&hooks),
session: None,
cwd: temp.path(),
output_sink: Some(&mut sink),
cancellation: AgentCancellation::default(),
session_title_job: None,
semantic_progress_timeout: None,
invocation_mode: crate::output::InvocationMode::Print,
agent_id: None,
herdr_reporter: None,
ttsr: crate::config::TtsrSettings::default(),
continuation_auto_compaction_policy: None,
},
)
.unwrap();
let requests = provider.requests();
assert_eq!(requests[1].tool_results()[0].output, "hello");
assert!(!format!("{:?}", requests[1]).contains("hook-secret"));
assert!(
sink.outputs
.iter()
.any(|event| matches!(event, OutputEvent::HookDiagnostic { .. }))
);
}
#[test]
fn before_hook_block_prevents_target_tool_dispatch() {
let temp = tempfile::TempDir::new().unwrap();
let target = temp.path().join("blocked.txt");
let provider = ScriptedProvider::new(vec![
vec![
write_call("call_write", "blocked.txt", "should-not-write"),
done(),
],
vec![done()],
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let hooks = HookRuntime::new(
temp.path(),
before_tool_hook(
"gate",
"exit 9",
Some(crate::config::HookFailurePolicy::Block),
),
true,
)
.unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
initial_instructions: &[],
prompt: "write file",
prompt_origin: crate::output::UserPromptOrigin::User,
effective_prompt: None,
tools: Some(&tools),
hooks: Some(&hooks),
session: None,
cwd: temp.path(),
output_sink: None,
cancellation: AgentCancellation::default(),
session_title_job: None,
semantic_progress_timeout: None,
invocation_mode: crate::output::InvocationMode::Print,
agent_id: None,
herdr_reporter: None,
ttsr: crate::config::TtsrSettings::default(),
continuation_auto_compaction_policy: None,
},
)
.unwrap();
assert!(!target.exists());
let requests = provider.requests();
assert!(
requests[1].tool_results()[0]
.output
.contains("blocked by local before_tool hook policy")
);
assert!(!requests[1].tool_results()[0].output.contains("exit 9"));
}
#[test]
fn after_hook_fail_preserves_tool_result_and_stops_before_provider_continuation() {
let temp = tempfile::TempDir::new().unwrap();
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![
read_done("call_1"),
vec![text("should-not-continue"), done()],
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let hooks = HookRuntime::new(
temp.path(),
after_tool_hook(
"strict",
"printf after-secret; exit 4",
Some(crate::config::HookFailurePolicy::Fail),
),
true,
)
.unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let error = agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
initial_instructions: &[],
prompt: "read file",
prompt_origin: crate::output::UserPromptOrigin::User,
effective_prompt: None,
tools: Some(&tools),
hooks: Some(&hooks),
session: Some(&session),
cwd: temp.path(),
output_sink: None,
cancellation: AgentCancellation::default(),
session_title_job: None,
semantic_progress_timeout: None,
invocation_mode: crate::output::InvocationMode::Print,
agent_id: None,
herdr_reporter: None,
ttsr: crate::config::TtsrSettings::default(),
continuation_auto_compaction_policy: None,
},
)
.unwrap_err()
.to_string();
assert!(error.contains("after_tool"));
assert!(!error.contains("after-secret"));
assert_eq!(provider.requests().len(), 1);
let events = session.read_events().unwrap();
assert_eq!(
events
.iter()
.filter(|event| event.event_type == "turn_status" && event.payload["status"] == "failed")
.count(),
1
);
assert!(events.iter().any(|event| {
event.event_type == "tool_result" && event.payload["result"]["content"] == "hello"
}));
}
#[test]
fn after_hook_block_action_stops_before_provider_continuation() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![
read_done("call_1"),
vec![text("should-not-continue"), done()],
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let hooks = HookRuntime::new(
temp.path(),
crate::config::HookSettings {
enabled: true,
after_tool: vec![crate::config::HookDefinition {
label: Some("invalid-after-block".into()),
command: "exit 4".into(),
failure_policy: Some(crate::config::HookFailurePolicy::Block),
..crate::config::HookDefinition::default()
}],
..crate::config::HookSettings::default()
},
true,
)
.unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let error = agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
initial_instructions: &[],
prompt: "read file",
prompt_origin: crate::output::UserPromptOrigin::User,
effective_prompt: None,
tools: Some(&tools),
hooks: Some(&hooks),
session: None,
cwd: temp.path(),
output_sink: None,
cancellation: AgentCancellation::default(),
session_title_job: None,
semantic_progress_timeout: None,
invocation_mode: crate::output::InvocationMode::Print,
agent_id: None,
herdr_reporter: None,
ttsr: crate::config::TtsrSettings::default(),
continuation_auto_compaction_policy: None,
},
)
.unwrap_err()
.to_string();
assert!(error.contains("after_tool"));
assert_eq!(provider.requests().len(), 1);
}
#[test]
fn thinking_level_is_attached_to_initial_and_continuation_requests() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![read_done("call_1"), vec![done()]]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("gpt-5", &[], &SkillDiscovery::default())
.with_thinking_level(crate::thinking::ThinkingLevel::High);
agent
.run_print_with_tools(&provider, "read file", Some(&tools), None, temp.path())
.unwrap();
let requests = provider.requests();
assert_eq!(requests.len(), 2);
assert!(
requests
.iter()
.all(|request| request.thinking_level == crate::thinking::ThinkingLevel::High)
);
}
#[test]
fn text_then_tool_flushes_assistant_chunk_before_tool_call() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![
vec![text("before "), text("tool"), read_call("call_1"), done()],
vec![done()],
]);
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
agent
.run_print_with_tools(&provider, "read", Some(&tools), Some(&session), temp.path())
.unwrap();
let event_types = session
.read_events()
.unwrap()
.into_iter()
.map(|event| event.event_type)
.collect::<Vec<_>>();
assert_eq!(event_types[1], "assistant_chunk");
assert_eq!(event_types[2], "tool_call");
}
#[test]
fn sink_error_during_activity_finished_stops_without_continuation() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![read_done("call_1"), vec![done()]]);
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = FailingOutputSink::new(SinkFailurePoint::ActivityFinished);
let error = agent
.run_print_with_tools_streaming_output(
&provider,
"read file",
Some(&tools),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap_err()
.to_string();
assert_eq!(error, "sink activity broke");
assert_eq!(provider.requests().len(), 1);
let event_types = session
.read_events()
.unwrap()
.into_iter()
.map(|event| event.event_type)
.collect::<Vec<_>>();
assert_eq!(
event_types,
vec!["user_input", "tool_call", "tool_result", "turn_status"]
);
}
#[test]
fn sink_error_during_tool_result_stops_without_continuation() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![read_done("call_1"), vec![done()]]);
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = FailingOutputSink::new(SinkFailurePoint::ToolResult);
let error = agent
.run_print_with_tools_streaming_output(
&provider,
"read file",
Some(&tools),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap_err()
.to_string();
assert_eq!(error, "sink tool result broke");
assert_eq!(provider.requests().len(), 1);
let event_types = session
.read_events()
.unwrap()
.into_iter()
.map(|event| event.event_type)
.collect::<Vec<_>>();
assert_eq!(
event_types,
vec!["user_input", "tool_call", "tool_result", "turn_status"]
);
let events = session.read_events().unwrap();
let statuses = events
.iter()
.filter(|event| event.event_type == "turn_status")
.collect::<Vec<_>>();
assert_eq!(statuses.len(), 1);
assert_eq!(statuses[0].payload["status"], "failed");
let replay = crate::context::build_conversation_replay(Some(&session)).unwrap();
assert!(replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User
&& message.content == "read file"
)));
assert!(replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::ToolResult(result) if result.call_id == "call_1"
)));
}
#[test]
fn sink_error_during_tool_turn_assistant_complete_flushes_assistant_chunk() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider =
ScriptedProvider::new(vec![vec![text("before tool"), read_call("call_1"), done()]]);
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = FailingOutputSink::new(SinkFailurePoint::AssistantComplete);
let error = agent
.run_print_with_tools_streaming_output(
&provider,
"read file",
Some(&tools),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap_err()
.to_string();
assert_eq!(error, "sink assistant complete broke");
assert_eq!(provider.requests().len(), 1);
let events = session.read_events().unwrap();
assert_eq!(assistant_chunk_text(&events), "before tool");
assert_eq!(
events
.iter()
.map(|event| event.event_type.as_str())
.collect::<Vec<_>>(),
vec!["user_input", "assistant_chunk", "turn_status"]
);
}
#[test]
fn assistant_complete_after_tool_interleave_contains_current_segment_only() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "file contents").unwrap();
let provider = ScriptedProvider::new(vec![
vec![text("I'll read."), read_call("call_1"), done()],
text_done("Done."),
]);
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"inspect file",
Some(&tools),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(output.text, "I'll read.\n\nDone.");
assert_eq!(
assistant_chunk_text(&session.read_events().unwrap()),
output.text
);
assert_eq!(sink.text, "I'll read.\n\nDone.");
let assistant_deltas = sink
.outputs
.iter()
.filter_map(|event| match event {
OutputEvent::AssistantDelta { text } => Some(text.as_str()),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(assistant_deltas, vec!["I'll read.", "\n\n", "Done."]);
let assistant_completes = sink
.outputs
.iter()
.filter_map(|event| match event {
OutputEvent::AssistantComplete { text } => Some(text.as_str()),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(assistant_completes, vec!["I'll read.", "Done."]);
assert!(matches!(
provider.requests()[1].conversation_items()[2],
ProviderConversationItem::Message(ref message)
if message.role == crate::providers::MessageRole::Assistant
&& message.content == "I'll read."
));
assert!(matches!(
provider.requests()[1].conversation_items()[3],
ProviderConversationItem::ResponseItem(ref item)
if item["type"] == "function_call" && item["call_id"] == "call_1"
));
assert!(matches!(
provider.requests()[1].conversation_items()[4],
ProviderConversationItem::ToolResult(ref result)
if result.call_id == "call_1" && result.output == "file contents"
));
}
#[test]
fn assistant_complete_emitted_after_tool_result_for_pre_tool_segment_recovery() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "file contents").unwrap();
let provider = ScriptedProvider::new(vec![
vec![text("I'll read."), read_call("call_1"), done()],
text_done("Done."),
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"inspect file",
Some(&tools),
None,
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(output.text, "I'll read.\n\nDone.");
let assistant_deltas = sink
.outputs
.iter()
.filter_map(|event| match event {
OutputEvent::AssistantDelta { text } => Some(text.as_str()),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(assistant_deltas, vec!["I'll read.", "\n\n", "Done."]);
let assistant_completes = sink
.outputs
.iter()
.filter_map(|event| match event {
OutputEvent::AssistantComplete { text } => Some(text.as_str()),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(assistant_completes, vec!["I'll read.", "Done."]);
}
#[test]
fn write_tool_start_event_precedes_result_and_provider_gets_raw_output() {
let temp = tempfile::TempDir::new().unwrap();
let target = temp.path().join("new.txt");
let provider = ScriptedProvider::new(vec![
vec![
write_call("call_write", "new.txt", "secret payload"),
done(),
],
vec![done()],
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
agent
.run_print_with_tools_streaming_output(
&provider,
"write file",
Some(&tools),
None,
temp.path(),
Some(&mut sink),
)
.unwrap();
let start_position = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::ToolStarted { call, .. } if call.id == "call_write"))
.unwrap();
let result_position = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::ToolResult { call, .. } if call.id == "call_write"))
.unwrap();
assert!(start_position < result_position);
assert_eq!(std::fs::read_to_string(target).unwrap(), "secret payload");
let requests = provider.requests();
assert!(requests[1].tool_results()[0].output.contains("wrote"));
assert!(!requests[1].tool_results()[0].output.contains("RUNNING"));
assert!(!requests[1].tool_results()[0].output.contains("TOOL START"));
}
#[test]
fn edit_tool_start_event_precedes_failure_result_and_provider_gets_raw_error() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![
vec![
edit_call("call_edit", "file.txt", "missing", "secret replacement"),
done(),
],
vec![done()],
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
agent
.run_print_with_tools_streaming_output(
&provider,
"edit file",
Some(&tools),
None,
temp.path(),
Some(&mut sink),
)
.unwrap();
let start_position = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::ToolStarted { call, .. } if call.id == "call_edit"))
.unwrap();
let result_position = sink
.outputs
.iter()
.position(|event| matches!(event, OutputEvent::ToolResult { call, result, .. } if call.id == "call_edit" && !result.success))
.unwrap();
assert!(start_position < result_position);
let requests = provider.requests();
assert_eq!(
requests[1].tool_results()[0].output,
sink.outputs
.iter()
.find_map(|event| match event {
OutputEvent::ToolResult { result, .. } => Some(result.content.clone()),
_ => None,
})
.unwrap()
);
assert!(!requests[1].tool_results()[0].output.contains("RUNNING"));
assert!(!requests[1].tool_results()[0].output.contains("TOOL START"));
}
#[test]
fn tool_result_reconciliation_does_not_emit_redundant_final_activity_delta() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![read_done("call_1"), vec![done()]]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
agent
.run_print_with_tools_streaming_output(
&provider,
"read file",
Some(&tools),
None,
temp.path(),
Some(&mut sink),
)
.unwrap();
assert!(
sink.activities
.iter()
.all(|event| !matches!(event, ActivityEvent::Delta { .. }))
);
assert!(
sink.outputs
.iter()
.any(|event| matches!(event, OutputEvent::ToolResult { .. }))
);
}
#[test]
fn fallback_tool_result_uses_display_activity_id_without_provider_continuation_change() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![vec![read_call(""), done()], vec![done()]]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
agent
.run_print_with_tools_streaming_output(
&provider,
"read file",
Some(&tools),
None,
temp.path(),
Some(&mut sink),
)
.unwrap();
assert!(matches!(
sink.outputs.iter().find(|event| matches!(event, OutputEvent::ToolResult { .. })),
Some(OutputEvent::ToolResult { call, .. }) if call.id == "turn-0/tool-0-read"
));
let requests = provider.requests();
assert_eq!(requests[1].tool_results()[0].call_id, "");
}
#[test]
fn visible_tool_blocks_do_not_change_provider_continuation() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![read_done("call_1"), text_done("done")]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"read file",
Some(&tools),
None,
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(sink.blocks.len(), 1);
assert!(!sink.blocks[0].contains("TOOL START"));
assert!(sink.blocks[0].contains("OUTPUT:\nhello"));
assert_eq!(output.tool_results[0].content, "hello");
assert!(!output.tool_results[0].content.contains("TOOL START"));
let requests = provider.requests();
assert_eq!(requests[1].tool_results()[0].output, "hello");
assert_eq!(requests[1].response_items()[1]["output"], "hello");
assert!(!requests[1].tool_results()[0].output.contains("TOOL START"));
}
#[test]
fn visible_tool_blocks_capture_two_iterations_in_order() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
std::fs::write(temp.path().join("other.txt"), "world").unwrap();
let provider = ScriptedProvider::new(vec![
read_done("call_1"),
vec![
ProviderEvent::ToolCall(ToolCall {
id: "call_2".to_string(),
name: "read".to_string(),
arguments: json!({"path":"other.txt:raw"}),
}),
done(),
],
text_done("done"),
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
agent
.run_print_with_tools_streaming_output(
&provider,
"read twice",
Some(&tools),
None,
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(sink.blocks.len(), 2);
assert!(sink.blocks[0].contains("PATH: file.txt"));
assert!(sink.blocks[1].contains("PATH: other.txt"));
assert!(sink.blocks[0].contains("OUTPUT:\nhello"));
assert!(sink.blocks[1].contains("OUTPUT:\nworld"));
}
#[test]
fn visible_tool_blocks_display_failures_and_preserve_raw_failure_content() {
let temp = tempfile::TempDir::new().unwrap();
let provider = ScriptedProvider::new(vec![
vec![
ProviderEvent::ToolCall(ToolCall {
id: "call_bad".to_string(),
name: "read".to_string(),
arguments: json!({"path":"missing.txt"}),
}),
done(),
],
vec![text("handled"), done()],
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"read missing",
Some(&tools),
None,
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(sink.blocks.len(), 1);
assert!(sink.blocks[0].contains("STATUS: failure"));
assert!(sink.blocks[0].contains("OUTPUT:"));
assert_eq!(output.tool_results.len(), 1);
let requests = provider.requests();
assert_eq!(
requests[1].tool_results()[0].output,
output.tool_results[0].content
);
assert!(!requests[1].tool_results()[0].output.contains("TOOL START"));
}
#[test]
fn agent_executes_subagents_tool_call_and_continues() {
let temp = tempfile::TempDir::new().unwrap();
let provider = ScriptedProvider::new(vec![
vec![
ProviderEvent::ToolCall(ToolCall {
id: "call_subagents".to_string(),
name: "subagents".to_string(),
arguments: json!({"tasks":[{"intent":"inspect"}],"concurrency":1}),
}),
done(),
],
vec![text("aggregated"), done()],
]);
let captured_args = Arc::new(Mutex::new(Vec::new()));
let captured_args_for_runner = Arc::clone(&captured_args);
let tools =
ToolRuntime::new(temp.path())
.unwrap()
.with_subagents(move |arguments, _context| {
captured_args_for_runner.lock().unwrap().push(arguments);
ToolResult {
tool_name: "subagents".to_string(),
success: true,
content: "{\"summary\":{\"total\":1,\"completed\":1}}".to_string(),
metadata: json!({}),
display: crate::tools::ToolResultDisplay::default(),
}
});
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let output = agent
.run_print_with_tools(&provider, "delegate", Some(&tools), None, temp.path())
.unwrap();
assert_eq!(output.text, "aggregated");
assert_eq!(
captured_args.lock().unwrap()[0]["tasks"][0]["intent"],
"inspect"
);
let requests = provider.requests();
assert_eq!(requests.len(), 2);
assert_eq!(requests[1].tool_results()[0].call_id, "call_subagents");
assert_eq!(requests[1].tool_results()[0].tool_name, "subagents");
}
#[test]
fn duplicate_tool_call_is_blocked_before_second_dispatch() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![
read_done("call_1"),
read_done("call_1"),
vec![text("unreachable"), done()],
]);
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let error = agent
.run_print_with_tools_streaming_output(
&provider,
"read twice",
Some(&tools),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap_err()
.to_string();
assert_eq!(error, "duplicate tool call suppressed");
assert_eq!(provider.requests().len(), 2);
assert_eq!(sink.activities.len(), 2);
assert_eq!(sink.blocks.len(), 1);
let events = session.read_events().unwrap();
assert_eq!(
events
.iter()
.filter(|event| event.event_type == "tool_call")
.count(),
1
);
assert_eq!(
events
.iter()
.filter(|event| event.event_type == "tool_result")
.count(),
1
);
assert!(!error.contains("file.txt"));
assert!(!error.contains("hello"));
}
#[test]
fn same_turn_duplicate_tool_call_is_blocked() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider =
ScriptedProvider::new(vec![vec![read_call("call_1"), read_call("call_1"), done()]]);
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let error = agent
.run_print_with_tools_streaming_output(
&provider,
"read twice",
Some(&tools),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap_err()
.to_string();
assert_eq!(error, "duplicate tool call suppressed");
assert_eq!(provider.requests().len(), 1);
assert_eq!(sink.activities.len(), 2);
let events = session.read_events().unwrap();
assert_eq!(
events
.iter()
.filter(|event| event.event_type == "tool_call")
.count(),
1
);
assert_eq!(
events
.iter()
.filter(|event| event.event_type == "tool_result")
.count(),
1
);
}
#[test]
fn tool_calls_with_different_ids_are_not_duplicates() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![
vec![
ProviderEvent::ToolCall(ToolCall {
id: "call_1".to_string(),
name: "read".to_string(),
arguments: json!({"path":"file.txt","offset":0}),
}),
done(),
],
vec![
ProviderEvent::ToolCall(ToolCall {
id: "call_2".to_string(),
name: "read".to_string(),
arguments: json!({"path":"file.txt","offset":1}),
}),
done(),
],
vec![done()],
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let output = agent
.run_print_with_tools(&provider, "read offsets", Some(&tools), None, temp.path())
.unwrap();
assert_eq!(output.tool_results.len(), 2);
assert_eq!(provider.requests().len(), 3);
}
#[test]
fn conflicting_duplicate_tool_call_id_stops_before_second_dispatch() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![
vec![
ProviderEvent::ToolCall(ToolCall {
id: "same_id".to_string(),
name: "read".to_string(),
arguments: json!({"path":"file.txt","offset":0}),
}),
done(),
],
vec![
ProviderEvent::ToolCall(ToolCall {
id: "same_id".to_string(),
name: "read".to_string(),
arguments: json!({"path":"file.txt","offset":1}),
}),
done(),
],
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let error = agent
.run_print_with_tools(&provider, "read offsets", Some(&tools), None, temp.path())
.unwrap_err()
.to_string();
assert!(error.contains("conflicting duplicate tool call id"));
assert_eq!(provider.requests().len(), 2);
}
#[test]
fn agent_supports_multiple_tool_iterations() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let provider = ScriptedProvider::new(vec![
read_done("call_1"),
read_done("call_2"),
text_done("done"),
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let output = agent
.run_print_with_tools(&provider, "read twice", Some(&tools), None, temp.path())
.unwrap();
assert_eq!(output.text, "done");
assert_eq!(output.tool_results.len(), 2);
assert_eq!(provider.requests().len(), 3);
}
#[test]
fn agent_continues_past_eight_tool_iterations() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let mut responses = (0..9)
.map(|index| vec![read_call(&format!("call_{index}")), done()])
.collect::<Vec<_>>();
responses.push(text_done("done"));
let provider = ScriptedProvider::new(responses);
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let output = agent
.run_print_with_tools(&provider, "loop", Some(&tools), Some(&session), temp.path())
.unwrap();
assert_eq!(output.text, "done");
assert_eq!(output.tool_results.len(), 9);
assert_eq!(provider.requests().len(), 10);
let events = session.read_events().unwrap();
assert!(!events.iter().any(|event| event.event_type == "diagnostic"));
}
#[test]
fn continuation_over_window_fails_before_provider_dispatch() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "x".repeat(60_000)).unwrap();
let provider = ScriptedProvider::new(vec![read_done("call_1"), vec![done()]]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default()).with_context_budget(
crate::context::ContextBudget {
max_tokens: 20_000,
reserve_tokens: 10_000,
..crate::context::ContextBudget::default()
},
);
let error = agent
.run_print_with_tools(&provider, "read", Some(&tools), None, temp.path())
.unwrap_err()
.to_string();
assert!(error.contains("no history was omitted"));
assert_eq!(provider.requests().len(), 1);
}
#[test]
fn before_hook_fail_records_one_failed_terminal_status() {
let temp = tempfile::TempDir::new().unwrap();
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let provider = ScriptedProvider::new(vec![vec![read_call("call_1"), done()], vec![done()]]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let hooks = HookRuntime::new(
temp.path(),
before_tool_hook(
"strict",
"exit 4",
Some(crate::config::HookFailurePolicy::Fail),
),
true,
)
.unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
prompt: "read file",
tools: Some(&tools),
hooks: Some(&hooks),
session: Some(&session),
cwd: temp.path(),
..run_request("read file", temp.path())
},
)
.unwrap_err();
let events = session.read_events().unwrap();
assert_eq!(
events
.iter()
.filter(|event| event.event_type == "turn_status" && event.payload["status"] == "failed")
.count(),
1
);
assert!(!events.iter().any(|event| event.event_type == "tool_result"));
}