use super::*;
#[test]
fn session_persistence_warning_is_emitted_once_per_run() {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let provider = SessionBreakingProvider {
session_path: session.path().to_path_buf(),
};
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"prompt",
None,
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(output.text, "still runs");
let warnings = sink
.outputs
.iter()
.filter(|event| {
matches!(
event,
OutputEvent::Diagnostic { level, message }
if level == "warning"
&& message.contains("session persistence failed")
&& message.contains("future resume may be incomplete")
)
})
.count();
assert_eq!(warnings, 1);
}
#[test]
fn session_replay_read_diagnostics_emit_local_warning_without_provider_leak() {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.open("safe").unwrap();
let previous_user = SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"previous user text"}),
);
let previous_assistant = SessionEvent::new(
"assistant_output",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"previous assistant text"}),
);
std::fs::create_dir_all(session.path().parent().unwrap()).unwrap();
std::fs::write(
session.path(),
format!(
"{}\n{{\"access_token\":\"MALFORMED_LEAK_MARKER\",\n{}\n{{\"event_type\":\"truncated\"\n",
serde_json::to_string(&previous_user).unwrap(),
serde_json::to_string(&previous_assistant).unwrap()
),
)
.unwrap();
crate::sessions::secure_test_session_root(session.path().parent().unwrap());
let provider = ScriptedProvider::new(vec![text_done("done")]);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
agent
.run_print_with_tools_streaming_output(
&provider,
"current prompt",
None,
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
let requests = provider.requests();
assert_eq!(requests.len(), 1);
let request_debug = format!("{:?}", requests[0]);
assert!(request_debug.contains("previous user text"));
assert!(request_debug.contains("previous assistant text"));
assert!(request_debug.contains("current prompt"));
let diagnostic_messages = sink
.outputs
.iter()
.filter_map(|event| match event {
OutputEvent::Diagnostic { level, message } if level == "warning" => {
Some(message.as_str())
}
_ => None,
})
.collect::<Vec<_>>();
assert!(
diagnostic_messages
.iter()
.any(|message| message.contains("session replay") && message.contains("2"))
);
for forbidden in [
"MALFORMED_LEAK_MARKER",
"failed to parse",
"line 2",
"line 4",
"session replay",
] {
assert!(
!request_debug.contains(forbidden),
"request leaked {forbidden}"
);
}
}
#[test]
fn cancellation_after_text_delta_flushes_pending_assistant_chunk_to_session() {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let cancel = Arc::new(AtomicBool::new(false));
let provider = ActiveCancelingProvider::new(Arc::clone(&cancel), 100);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let error = agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
initial_instructions: &[],
prompt: "stream forever",
prompt_origin: crate::output::UserPromptOrigin::User,
effective_prompt: None,
tools: None,
hooks: None,
session: Some(&session),
cwd: temp.path(),
output_sink: None,
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));
let events = session.read_events().unwrap();
assert_eq!(assistant_chunk_text(&events), "chunk-0");
assert!(
events
.iter()
.all(|event| event.event_type != "assistant_output")
);
let turn_statuses = events
.iter()
.filter(|event| event.event_type == "turn_status")
.collect::<Vec<_>>();
assert_eq!(turn_statuses.len(), 1);
let turn_status = turn_statuses[0];
assert_eq!(turn_status.payload["status"], "cancelled");
assert_eq!(turn_status.payload["assistant_text"], "chunk-0");
}
#[test]
fn provider_stream_auto_continue_recovers_once_without_failed_status() {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let provider = AutoContinueProvider::new(vec![
AutoContinueStep::TextThenEligibleTimeout("partial "),
AutoContinueStep::TextThenDone("done"),
]);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"will recover",
None,
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(output.text, "partial done");
assert!(output.recovered_incomplete_stream);
let requests = provider.requests();
assert_eq!(requests.len(), 2);
let continuation_items = requests[1].conversation_items();
assert!(continuation_items.windows(2).any(|pair| matches!(
(&pair[0], &pair[1]),
(
ProviderConversationItem::Message(assistant),
ProviderConversationItem::Message(user),
) if assistant.role == crate::providers::MessageRole::Assistant
&& assistant.content == "partial "
&& user.role == crate::providers::MessageRole::User
&& user.content == "Continue"
)));
let events = session.read_events().unwrap();
assert_eq!(assistant_chunk_text(&events), "partial done");
assert!(!events.iter().any(|event| event.event_type == "turn_status"));
assert!(
!events.iter().any(|event| {
event.event_type == "user_input" && event.payload["text"] == "Continue"
})
);
assert!(sink.outputs.iter().any(|event| matches!(
event,
OutputEvent::AutomaticUserPrompt { text } if text == "Continue"
)));
assert!(sink.outputs.iter().any(|event| matches!(
event,
OutputEvent::Diagnostic { level, message }
if level == "info" && message.contains("ended before completion")
)));
assert!(sink.outputs.iter().any(|event| matches!(
event,
OutputEvent::Diagnostic { level, message }
if level == "info" && message.contains("auto-continuing once")
)));
}
#[test]
fn incomplete_stream_cancellation_wins_before_ephemeral_continue() {
let temp = tempfile::TempDir::new().unwrap();
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let cancel = Arc::new(AtomicBool::new(false));
let provider = CancelingIncompleteProvider {
cancel: Arc::clone(&cancel),
requests: RecordedRequests::default(),
};
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let error = agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
session: Some(&session),
output_sink: Some(&mut sink),
cancellation: AgentCancellation::new(cancel),
..run_request("cancel recovery", temp.path())
},
)
.unwrap_err();
assert!(is_run_canceled(&error));
assert_eq!(provider.requests.snapshot().len(), 1);
assert!(
!sink
.outputs
.iter()
.any(|event| matches!(event, OutputEvent::AutomaticUserPrompt { .. }))
);
let events = session.read_events().unwrap();
assert!(events.iter().any(|event| {
event.event_type == "turn_status" && event.payload["status"] == "cancelled"
}));
assert!(
!events.iter().any(|event| {
event.event_type == "user_input" && event.payload["text"] == "Continue"
})
);
}
#[test]
fn provider_stream_auto_continue_completed_tool_progress_blocks_recovery() {
for step in [
AutoContinueStep::TextThenToolThenEligibleTimeout,
AutoContinueStep::TextThenFunctionItemThenEligibleTimeout,
] {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let provider = AutoContinueProvider::new(vec![step]);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let error = agent
.run_print_with_tools_streaming_output(
&provider,
"will not recover",
None,
Some(&session),
temp.path(),
None,
)
.unwrap_err();
assert!(
error
.to_string()
.contains("provider stream no semantic progress before timeout")
);
assert_eq!(provider.requests().len(), 1);
let events = session.read_events().unwrap();
let failed = events
.iter()
.filter(|event| {
event.event_type == "turn_status" && event.payload["status"] == "failed"
})
.count();
assert_eq!(failed, 1);
}
}
#[test]
fn provider_stream_auto_continue_hidden_tool_progress_recovers_without_execution() {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let provider = AutoContinueProvider::new(vec![
AutoContinueStep::TextThenHiddenToolProgressTimeout,
AutoContinueStep::TextThenDone("done"),
]);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"will recover hidden progress",
None,
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(output.text, "done");
let requests = provider.requests();
assert_eq!(requests.len(), 2);
assert!(matches!(
requests[1].conversation_items().last(),
Some(ProviderConversationItem::Message(user))
if user.role == crate::providers::MessageRole::User
&& user.content == "Continue"
));
assert!(
!output
.tool_results
.iter()
.any(|result| result.tool_name == "read")
);
assert!(
!session.read_events().unwrap().iter().any(|event| {
event.event_type == "user_input" && event.payload["text"] == "Continue"
})
);
assert!(sink.outputs.iter().any(|event| matches!(
event,
OutputEvent::AutomaticUserPrompt { text } if text == "Continue"
)));
}
#[test]
fn openai_body_error_after_partial_tool_arguments_recovers_ephemerally() {
let temp = tempfile::TempDir::new().unwrap();
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let transport = OpenAiPartialToolBodyErrorTransport::new();
let requests = transport.requests_handle();
let provider = OpenAiCompatibleProvider::custom(
"custom",
"model",
None,
"https://example.invalid/v1",
false,
transport,
);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"recover body error",
None,
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(output.text, "done");
assert!(output.recovered_incomplete_stream);
assert!(output.tool_results.is_empty());
let requests = requests.lock().unwrap();
assert_eq!(requests.len(), 2);
let continuation = requests[1].body.to_string();
assert!(continuation.contains("Continue"));
assert!(!continuation.contains("UNSAFE_PARTIAL_MARKER"));
drop(requests);
let events = session.read_events().unwrap();
assert!(!events.iter().any(|event| event.event_type == "tool_call"));
assert!(
!events.iter().any(|event| {
event.event_type == "user_input" && event.payload["text"] == "Continue"
})
);
assert!(
!serde_json::to_string(&events)
.unwrap()
.contains("UNSAFE_PARTIAL_MARKER")
);
assert!(events.iter().any(|event| {
event.event_type == "diagnostic"
&& event.payload["message"]
.as_str()
.is_some_and(|message| message.contains("partial tool-call progress observed"))
}));
let replay = crate::context::build_conversation_replay(Some(&session)).unwrap();
let replay_debug = format!("{:?}", replay.items);
assert!(replay_debug.contains("recover body error"));
assert!(replay_debug.contains("done"));
assert!(!replay_debug.contains("Continue"));
assert!(!replay_debug.contains("UNSAFE_PARTIAL_MARKER"));
assert!(sink.outputs.iter().any(|event| matches!(
event,
OutputEvent::AutomaticUserPrompt { text } if text == "Continue"
)));
}
#[test]
fn anthropic_pending_tool_message_stop_retries_once_then_records_provider_stream_trace() {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let transport = AnthropicPendingToolMessageStopTransport::new();
let requests = transport.requests_handle();
let provider = AnthropicProvider::new("claude-test", "secret", transport);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let error = agent
.run_print_with_tools_streaming_output(
&provider,
"trace pending tool",
None,
Some(&session),
temp.path(),
None,
)
.unwrap_err();
let error = error.to_string();
assert!(error.contains("incomplete tool_use"), "{error}");
assert!(error.contains("unsafe tool-call progress"), "{error}");
assert_eq!(requests.lock().unwrap().len(), 2);
let events = session.read_events().unwrap();
let event_types = events
.iter()
.map(|event| event.event_type.as_str())
.collect::<Vec<_>>();
assert!(event_types.contains(&"user_input"));
assert!(event_types.contains(&"turn_status"));
assert!(event_types.contains(&"provider_stream_trace"));
assert!(event_types.contains(&"abort_recovery"));
assert!(!event_types.contains(&"tool_call"));
assert!(!event_types.contains(&"provider_response_item"));
assert_eq!(
events
.iter()
.filter(|event| event.event_type == "user_input")
.count(),
1
);
let turn_status_index = event_types
.iter()
.position(|kind| *kind == "turn_status")
.unwrap();
let trace_index = event_types
.iter()
.position(|kind| *kind == "provider_stream_trace")
.unwrap();
let recovery_index = event_types
.iter()
.position(|kind| *kind == "abort_recovery")
.unwrap();
assert!(turn_status_index < trace_index && trace_index < recovery_index);
let trace = &events[trace_index].payload;
assert_eq!(trace["provider"], "anthropic");
assert_eq!(trace["failure_context"], "message_stop");
assert_eq!(trace["message_delta_stop_reason"], "tool_use");
assert_eq!(
trace["recent_events"]
.as_array()
.unwrap()
.iter()
.map(|event| event["event_type"].as_str().unwrap())
.collect::<Vec<_>>(),
vec![
"content_block_start",
"content_block_delta",
"message_delta",
"message_stop"
]
);
let pending = &trace["pending_tools"][0];
assert_eq!(pending["index"], 1);
assert_eq!(pending["id"], "toolu_trace_full_id");
assert_eq!(pending["name"], "read");
assert_eq!(pending["argument_bytes"], 13);
let hash = pending["argument_sha256"].as_str().unwrap();
assert_eq!(hash.len(), 64);
assert!(
hash.chars()
.all(|ch| ch.is_ascii_hexdigit() && !ch.is_ascii_uppercase())
);
let serialized = serde_json::to_string(trace).unwrap();
assert!(!serialized.contains("/tmp"));
assert!(!serialized.contains(r#"{\"path\":\"/tmp"#));
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 == "trace pending tool"
)));
let replay_debug = format!("{replay:?}");
assert!(!replay_debug.contains("Session recovery summary for failed historical turn"));
assert!(!replay_debug.contains("Abort recovery facts"));
assert!(replay.replay_diagnostics.iter().any(|diagnostic| {
diagnostic.contains("without_recovery_summary unsafe_tool_call_progress=true")
}));
}
#[test]
fn unsafe_tool_call_replay_preserves_completed_pre_failure_session_replay_items() {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let append = |event_type: &str, payload: serde_json::Value| {
session
.append(&SessionEvent::new(
event_type,
session.id().to_string(),
temp.path().to_path_buf(),
payload,
))
.unwrap();
};
append("user_input", json!({"text":"inspect file before bad tool"}));
append("assistant_chunk", json!({"text":"I will inspect first."}));
append(
"tool_call",
json!({"id":"call_1","name":"read","arguments":{"path":"file.txt"}}),
);
append(
"tool_result",
json!({"call_id":"call_1","result":{"tool_name":"read","success":true,"content":"file body"}}),
);
append(
"turn_status",
json!({"status":"failed","assistant_text":""}),
);
append(
"abort_recovery",
json!({
"observed_reasoning_delta": true,
"unsafe_tool_call_progress": true,
"reasoning_preview": "private repair preview",
"partial_tool_call_summary": "partial invalid tool call omitted"
}),
);
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 == "inspect file before bad tool"
)));
assert!(replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::Assistant
&& message.content == "I will inspect first."
)));
assert!(replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::ResponseItem(value)
if value.get("type").and_then(serde_json::Value::as_str) == Some("function_call")
&& value.get("call_id").and_then(serde_json::Value::as_str) == Some("call_1")
&& value.get("name").and_then(serde_json::Value::as_str) == Some("read")
)));
assert!(replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::ToolResult(ProviderToolResult {
call_id,
tool_name: _,
success: _,
output,
}) if call_id == "call_1" && output == "file body"
)));
let replay_debug = format!("{replay:?}");
assert!(!replay_debug.contains("Session recovery summary for failed historical turn"));
assert!(!replay_debug.contains("Abort recovery facts"));
assert!(!replay_debug.contains("partial invalid tool call omitted"));
assert!(replay.replay_diagnostics.iter().any(|diagnostic| {
diagnostic.contains("without_recovery_summary unsafe_tool_call_progress=true")
}));
}
#[test]
fn provider_stream_auto_continue_retry_cap_records_one_final_failure() {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let provider = AutoContinueProvider::new(vec![
AutoContinueStep::TextThenEligibleTimeout("partial one "),
AutoContinueStep::TextThenEligibleTimeout("partial two"),
]);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let error = agent
.run_print_with_tools_streaming_output(
&provider,
"will fail once after recovery",
None,
Some(&session),
temp.path(),
None,
)
.unwrap_err();
assert!(
error
.to_string()
.contains("provider stream no semantic progress before timeout")
);
assert_eq!(provider.requests().len(), 2);
let events = session.read_events().unwrap();
let failed = events
.iter()
.filter(|event| event.event_type == "turn_status" && event.payload["status"] == "failed")
.collect::<Vec<_>>();
assert_eq!(failed.len(), 1);
assert_eq!(
failed[0].payload["assistant_text"],
"partial one partial two"
);
}
#[test]
fn provider_error_without_text_records_failed_terminal_status_for_replay() {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let provider = FailingProvider { events: vec![] };
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let error = agent
.run_print_with_tools_streaming_output(
&provider,
"reasoning-only failure objective",
None,
Some(&session),
temp.path(),
None,
)
.unwrap_err();
assert_eq!(error.to_string(), "provider broke");
let events = session.read_events().unwrap();
let turn_status = events
.iter()
.find(|event| event.event_type == "turn_status")
.unwrap();
assert_eq!(turn_status.payload["status"], "failed");
assert_eq!(turn_status.payload["assistant_text"], "");
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 == "reasoning-only failure objective"
)));
}
#[test]
fn provider_error_after_reasoning_delta_records_abort_recovery_without_raw_tool_args() {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let provider = FailingProvider {
events: vec![ProviderEvent::ReasoningSummaryDelta(
"consider sk-1234567890abcdef".to_string(),
)],
};
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let error = agent
.run_print_with_tools_streaming_output(
&provider,
"recover reasoning",
None,
Some(&session),
temp.path(),
None,
)
.unwrap_err();
assert_eq!(error.to_string(), "provider broke");
let events = session.read_events().unwrap();
let recovery = events
.iter()
.find(|event| event.event_type == "abort_recovery")
.unwrap();
assert_eq!(recovery.payload["observed_reasoning_delta"], true);
assert_eq!(recovery.payload["unsafe_tool_call_progress"], false);
assert!(!recovery.payload.to_string().contains("sk-1234567890abcdef"));
let replay = crate::context::build_conversation_replay(Some(&session)).unwrap();
let replay_debug = format!("{replay:?}");
assert!(replay_debug.contains("Abort recovery facts"));
assert!(!replay_debug.contains("sk-1234567890abcdef"));
}
#[test]
fn provider_error_after_text_records_failed_terminal_status_for_replay() {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let provider = FailingProvider {
events: vec![text("partial before error")],
};
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let error = agent
.run_print_with_tools_streaming_output(
&provider,
"will fail",
None,
Some(&session),
temp.path(),
None,
)
.unwrap_err();
assert_eq!(error.to_string(), "provider broke");
let events = session.read_events().unwrap();
assert_eq!(assistant_chunk_text(&events), "partial before error");
let turn_status = events
.iter()
.find(|event| event.event_type == "turn_status")
.unwrap();
assert_eq!(turn_status.payload["status"], "failed");
assert_eq!(
turn_status.payload["assistant_text"],
"partial before error"
);
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::Assistant
&& message.content == "partial before error"
)));
}
#[test]
fn visible_successful_hooks_do_not_persist_hook_diagnostics() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().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,
show_in_tui: true,
before_tool: vec![crate::config::HookDefinition {
label: Some("successful".into()),
command: "true".into(),
..crate::config::HookDefinition::default()
}],
..crate::config::HookSettings::default()
},
true,
)
.unwrap();
let mut sink = ProbeActivitySink::new(|_| {});
let agent = AgentSession::new("model", &[], &SkillDiscovery::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: 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!(sink.sender_events.lock().unwrap().iter().any(|event| {
matches!(
event,
ActivityEvent::Finished {
status: ActivityStatus::Success,
..
}
)
}));
assert!(
!sink
.outputs
.iter()
.any(|event| matches!(event, OutputEvent::HookDiagnostic { .. }))
);
let event_types = session
.read_events()
.unwrap()
.iter()
.map(|event| event.event_type.clone())
.collect::<Vec<_>>();
assert_eq!(
event_types,
vec![
"user_input",
"tool_call",
"hook_lifecycle",
"hook_lifecycle",
"tool_result",
"assistant_output",
]
);
assert!(
!event_types
.iter()
.any(|event_type| event_type == "hook_diagnostic")
);
assert_eq!(provider.requests()[1].tool_results()[0].output, "hello");
}
#[test]
fn after_hook_provider_context_reaches_next_request_and_session_replay() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
std::fs::write(
temp.path().join("hook.out"),
r#"{"context_items":[{"role":"user","content":"hook memory note"}]}"#,
)
.unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let provider = ScriptedProvider::new(vec![read_done("call_1"), text_done("done")]);
let 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("memory".into()),
command: "cat hook.out".into(),
..crate::config::HookDefinition::default()
}],
..crate::config::HookSettings::default()
},
true,
)
.unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::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: 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();
let requests = provider.requests();
assert_eq!(requests.len(), 2);
let continuation = requests[1].conversation_items();
assert!(continuation.windows(2).any(|pair| matches!(
(&pair[0], &pair[1]),
(
ProviderConversationItem::ToolResult(result),
ProviderConversationItem::Message(message),
) if result.call_id == "call_1" && message.content == "hook memory note"
)));
let events = session.read_events().unwrap();
assert!(
events
.iter()
.any(|event| event.event_type == "hook_context_injection")
);
assert!(
events
.iter()
.any(|event| event.event_type == "provider_context_item"
&& event.payload["content"] == "hook memory note")
);
let replay = crate::context::build_conversation_replay(Some(&session)).unwrap();
assert!(replay.items.windows(2).any(|pair| matches!(
(&pair[0], &pair[1]),
(
ProviderConversationItem::ToolResult(result),
ProviderConversationItem::Message(message),
) if result.call_id == "call_1" && message.content == "hook memory note"
)));
}
#[test]
fn before_hook_fail_emits_local_diagnostic_and_skips_tool_dispatch() {
let temp = tempfile::TempDir::new().unwrap();
let target = temp.path().join("blocked.txt");
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
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(),
crate::config::HookSettings {
enabled: true,
before_tool: vec![crate::config::HookDefinition {
label: Some("strict".into()),
command: "printf before-secret; exit 9".into(),
failure_policy: Some(crate::config::HookFailurePolicy::Fail),
..crate::config::HookDefinition::default()
}],
..crate::config::HookSettings::default()
},
true,
)
.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: "write file",
prompt_origin: crate::output::UserPromptOrigin::User,
effective_prompt: None,
tools: Some(&tools),
hooks: Some(&hooks),
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_err()
.to_string();
assert!(error.contains("before_tool"));
assert!(!error.contains("before-secret"));
assert!(!target.exists());
assert_eq!(provider.requests().len(), 1);
assert_eq!(
sink.outputs
.iter()
.filter(|event| matches!(event, OutputEvent::HookDiagnostic { .. }))
.count(),
1
);
assert!(sink.outputs.iter().any(|event| matches!(
event,
OutputEvent::ToolResult { call, result, .. }
if call.id == "call_write"
&& !result.success
&& result.metadata["hook_failed"] == true
&& !result.content.contains("before-secret")
)));
assert!(sink.activities.iter().any(|event| matches!(
event,
ActivityEvent::Finished {
status: ActivityStatus::Failed,
..
}
)));
let events = session.read_events().unwrap();
assert_eq!(
events
.iter()
.map(|event| event.event_type.as_str())
.collect::<Vec<_>>(),
vec![
"user_input",
"tool_call",
"hook_lifecycle",
"hook_lifecycle",
"hook_diagnostic",
"tool_display_result",
"turn_status",
]
);
let display_event = events
.iter()
.find(|event| event.event_type == "tool_display_result")
.unwrap();
assert_eq!(display_event.payload["call_id"], "call_write");
assert_eq!(display_event.payload["result"]["success"], false);
assert_eq!(
display_event.payload["result"]["metadata"]["hook_failed"],
true
);
let hook_event = events
.iter()
.find(|event| event.event_type == "hook_diagnostic")
.unwrap();
assert_eq!(hook_event.payload["phase"], "before_tool");
assert_eq!(hook_event.payload["policy"], "fail");
assert_eq!(hook_event.payload["target_ran"], false);
assert!(
!hook_event.payload["message"]
.as_str()
.unwrap()
.contains("before-secret")
);
let request_debug = format!("{:?}", provider.requests());
let cache_material = conversation_cache_material(
"provider",
"model",
agent.system_prompt(),
provider.requests()[0].conversation_items().as_ref(),
);
for forbidden in [
"before-secret",
"hook_diagnostic",
"before_tool",
"hook exited with status",
] {
assert!(
!request_debug.contains(forbidden),
"request leaked {forbidden}"
);
assert!(
!cache_material.contains(forbidden),
"material leaked {forbidden}"
);
}
}
#[test]
fn reasoning_summary_event_reaches_sink_and_session_without_answer_pollution() {
let temp = tempfile::TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
std::fs::write(temp.path().join("file.txt"), "file body").unwrap();
let session = manager.create().unwrap();
let provider = ScriptedProvider::new(vec![
vec![
ProviderEvent::ReasoningSummaryDelta("provider ".to_string()),
ProviderEvent::ReasoningSummaryCompleteIdentified(crate::providers::ReasoningSummary {
text: "provider summary".to_string(),
item_id: Some("rs_1".to_string()),
turn_id: Some("turn_1".to_string()),
}),
ProviderEvent::ResponseItem(json!({
"id":"rs_1",
"type":"reasoning",
"summary":[{"type":"summary_text","text":"provider summary"}],
"encrypted_content":{"nested":[{"encrypted_content":"opaque"}]}
})),
text("answer"),
read_call("call_reasoning"),
done(),
],
vec![done()],
]);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut sink = CapturingOutputSink::default();
let output = agent
.run_print_with_tools_streaming_output(
&provider,
"say hi",
Some(&ToolRuntime::new(temp.path()).unwrap()),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(output.text, "answer");
assert!(sink.outputs.iter().any(|event| matches!(
event,
OutputEvent::ThinkingSummaryCompleteIdentified { text, .. } if text == "provider summary"
)));
assert!(sink.outputs.iter().any(|event| matches!(
event,
OutputEvent::ThinkingSummaryCompleteIdentified {
item_id: Some(item_id),
turn_id: Some(turn_id),
..
} if item_id == "rs_1" && turn_id == "turn_1"
)));
let events = session.read_events().unwrap();
let reasoning_events = events
.iter()
.filter(|event| event.event_type == "reasoning_summary")
.collect::<Vec<_>>();
assert_eq!(reasoning_events.len(), 1);
assert_eq!(reasoning_events[0].payload["text"], "provider summary");
assert_eq!(reasoning_events[0].payload["item_id"], "rs_1");
assert_eq!(reasoning_events[0].payload["turn_id"], "turn_1");
assert!(events.iter().any(|event| {
event.event_type == "provider_response_item"
&& event.payload["item"].get("encrypted_content").is_none()
}));
assert!(
!std::fs::read_to_string(session.path())
.unwrap()
.contains("opaque")
);
assert!(format!("{:?}", provider.requests()[1]).contains("opaque"));
}
#[test]
fn assistant_delta_reaches_sink_before_best_effort_session_chunk_recording() {
let temp = tempfile::TempDir::new().unwrap();
let provider = ScriptedProvider::new(vec![text_done("hello")]);
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let mut event_types_seen_by_sink = Vec::new();
let mut sink = |delta: &str| {
assert_eq!(delta, "hello");
event_types_seen_by_sink = session
.read_events()
.unwrap()
.into_iter()
.map(|event| event.event_type)
.collect::<Vec<_>>();
Ok(())
};
let output = agent
.run_print_with_tools_streaming(
&provider,
"say hi",
None,
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(output.text, "hello");
assert_eq!(event_types_seen_by_sink, vec!["user_input"]);
let final_event_types = session
.read_events()
.unwrap()
.into_iter()
.map(|event| event.event_type)
.collect::<Vec<_>>();
assert_eq!(
final_event_types,
vec!["user_input", "assistant_chunk", "assistant_output"]
);
}
#[test]
fn assistant_delta_sink_failure_does_not_advance_output_or_session_state() {
let temp = tempfile::TempDir::new().unwrap();
let provider = ScriptedProvider::new(vec![text_done("not accepted")]);
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 = |_delta: &str| anyhow::bail!("sink rejected delta");
let error = agent
.run_print_with_tools_streaming(
&provider,
"say hi",
None,
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap_err();
assert_eq!(error.to_string(), "sink rejected delta");
let events = session.read_events().unwrap();
let event_types = events
.iter()
.map(|event| event.event_type.as_str())
.collect::<Vec<_>>();
assert_eq!(event_types, vec!["user_input", "turn_status"]);
assert_eq!(events[1].payload["status"], "failed");
assert_eq!(events[1].payload["assistant_text"], "");
}
#[test]
fn agent_streaming_batches_assistant_chunk_persistence() {
let temp = tempfile::TempDir::new().unwrap();
let deltas = (0..100)
.map(|index| text(format!("{index},")))
.chain(std::iter::once(done()))
.collect::<Vec<_>>();
let provider = ScriptedProvider::new(vec![deltas]);
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 output = agent
.run_print_with_tools_streaming_output(
&provider,
"count",
None,
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
let expected = (0..100)
.map(|index| format!("{index},"))
.collect::<String>();
assert_eq!(output.text, expected);
let streamed = sink
.outputs
.iter()
.filter_map(|event| match event {
OutputEvent::AssistantDelta { text } => Some(text.as_str()),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(streamed.len(), 100);
assert_eq!(streamed.concat(), expected);
let events = session.read_events().unwrap();
let assistant_chunks = events
.iter()
.filter(|event| event.event_type == "assistant_chunk")
.count();
assert!(assistant_chunks < 100);
assert_eq!(assistant_chunk_text(&events), expected);
assert_eq!(
events
.iter()
.filter(|event| event.event_type == "assistant_output")
.count(),
1
);
}
#[test]
fn agent_skips_whitespace_only_assistant_persistence() {
let temp = tempfile::TempDir::new().unwrap();
let provider = ScriptedProvider::new(vec![text_done("\n\n")]);
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 output = agent
.run_print_with_tools_streaming_output(
&provider,
"blank",
None,
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(output.text, "");
assert!(
sink.outputs
.iter()
.any(|event| matches!(event, OutputEvent::AssistantDelta { text } if text == "\n\n"))
);
assert!(
sink.outputs.iter().any(
|event| matches!(event, OutputEvent::AssistantComplete { text } if text == "\n\n")
)
);
let events = session.read_events().unwrap();
assert_eq!(
events
.iter()
.map(|event| event.event_type.as_str())
.collect::<Vec<_>>(),
vec!["user_input"]
);
}
#[test]
fn streaming_tool_turn_preserves_session_event_order() {
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("read "), text("complete"), 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 streamed = String::new();
let mut sink = |delta: &str| {
streamed.push_str(delta);
Ok(())
};
let output = agent
.run_print_with_tools_streaming(
&provider,
"read file",
Some(&tools),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
assert_eq!(streamed, "read complete");
assert_eq!(output.text, "read complete");
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",
"assistant_chunk",
"assistant_output"
]
);
}
#[test]
fn write_start_activity_happens_before_file_mutation_and_session_order_is_unchanged() {
enum CapturedEvent {
Output(OutputEvent),
Activity(ActivityEvent),
}
struct FileAbsentOnStartSink {
target: std::path::PathBuf,
events: Vec<CapturedEvent>,
}
impl AgentOutputSink for FileAbsentOnStartSink {
fn assistant_delta(&mut self, _text: &str) -> anyhow::Result<()> {
Ok(())
}
fn output_event(&mut self, event: OutputEvent) -> anyhow::Result<()> {
self.events.push(CapturedEvent::Output(event));
Ok(())
}
fn activity_event(&mut self, event: ActivityEvent) -> anyhow::Result<()> {
if matches!(event, ActivityEvent::Started { .. }) {
assert!(
!self.target.exists(),
"write executed before started activity event"
);
}
self.events.push(CapturedEvent::Activity(event));
Ok(())
}
fn tool_block(&mut self, _block: &str) -> anyhow::Result<()> {
Ok(())
}
}
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", "hello"), 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());
let mut sink = FileAbsentOnStartSink {
target,
events: Vec::new(),
};
agent
.run_print_with_tools_streaming_output(
&provider,
"write file",
Some(&tools),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
let tool_event_positions = sink
.events
.iter()
.enumerate()
.filter_map(|(index, event)| match event {
CapturedEvent::Activity(ActivityEvent::Started { id, .. })
if id.as_str() == "call_write" =>
{
Some(("activity_started", index))
}
CapturedEvent::Output(OutputEvent::ToolStarted { call, .. })
if call.id == "call_write" =>
{
Some(("tool_started", index))
}
CapturedEvent::Output(OutputEvent::ToolResult { call, .. })
if call.id == "call_write" =>
{
Some(("tool_result", index))
}
CapturedEvent::Activity(ActivityEvent::Finished { id, .. })
if id.as_str() == "call_write" =>
{
Some(("activity_finished", index))
}
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(
tool_event_positions
.iter()
.map(|(name, _)| *name)
.collect::<Vec<_>>(),
vec![
"activity_started",
"tool_started",
"tool_result",
"activity_finished"
]
);
assert!(tool_event_positions[0].1 < tool_event_positions[1].1);
assert!(sink.events.iter().any(|event| matches!(
event,
CapturedEvent::Output(OutputEvent::ToolStarted { call, .. }) if call.id == "call_write"
)));
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", "assistant_output"]
);
}
#[test]
fn lifecycle_events_wrap_tool_result_without_persistence_or_provider_changes() {
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 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();
agent
.run_print_with_tools_streaming_output(
&provider,
"read file",
Some(&tools),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
assert!(matches!(sink.outputs[0], OutputEvent::UserPrompt { .. }));
assert!(
matches!(sink.activities[0], ActivityEvent::Started { ref id, .. } if id.as_str() == "call_1")
);
assert!(
matches!(sink.activities[1], ActivityEvent::Finished { ref id, status: ActivityStatus::Success, .. } if id.as_str() == "call_1")
);
assert_eq!(sink.activities.len(), 2);
assert!(matches!(
sink.outputs
.iter()
.find(|event| matches!(event, OutputEvent::ToolResult { .. })),
Some(OutputEvent::ToolResult { .. })
));
let requests = provider.requests();
assert_eq!(requests[1].tool_results()[0].output, "hello");
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",
"assistant_chunk",
"assistant_output"
]
);
}
#[test]
fn agent_continues_after_one_tool_result_and_records_session_events() {
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("read complete")]);
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,
"read file",
Some(&tools),
Some(&session),
temp.path(),
)
.unwrap();
assert_eq!(output.text, "read complete");
assert_eq!(output.tool_results[0].content, "hello");
let requests = provider.requests();
assert_eq!(requests.len(), 2);
assert!(requests[0].tool_results().is_empty());
assert_eq!(requests[1].response_items()[0]["type"], "function_call");
assert_eq!(requests[1].response_items()[0]["call_id"], "call_1");
assert_eq!(
requests[1].response_items()[1]["type"],
"function_call_output"
);
assert_eq!(requests[1].response_items()[1]["output"], "hello");
assert_eq!(requests[1].tool_results()[0].call_id, "call_1");
assert_eq!(requests[1].tool_results()[0].output, "hello");
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",
"assistant_chunk",
"assistant_output"
]
);
}
#[test]
fn visible_tool_blocks_preserve_session_tool_result_content() {
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 = CapturingOutputSink::default();
agent
.run_print_with_tools_streaming_output(
&provider,
"read file",
Some(&tools),
Some(&session),
temp.path(),
Some(&mut sink),
)
.unwrap();
let events = session.read_events().unwrap();
let tool_result = events
.iter()
.find(|event| event.event_type == "tool_result")
.unwrap();
assert_eq!(tool_result.payload["result"]["content"], "hello");
assert!(!tool_result.payload.to_string().contains("TOOL START"));
}
#[test]
fn provider_response_items_are_persisted_for_replay() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "hello").unwrap();
let response_item = json!({
"type": "function_call",
"call_id": "call_1",
"name": "read",
"arguments": "{\"path\":\"file.txt:raw\"}",
"status": "completed"
});
let provider = ScriptedProvider::new(vec![
vec![
ProviderEvent::ResponseItem(response_item.clone()),
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 events = session.read_events().unwrap();
assert!(events.iter().any(|event| {
event.event_type == "provider_response_item" && event.payload["item"] == response_item
}));
}
#[test]
fn current_turn_replay_preserves_text_response_item_tool_result_order() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "file body").unwrap();
let response_item = json!({
"type": "function_call",
"call_id": "call_1",
"name": "read",
"arguments": "{\"path\":\"file.txt:raw\"}",
"status": "completed"
});
let provider = ScriptedProvider::new(vec![
vec![
text("I'll read."),
ProviderEvent::ResponseItem(response_item.clone()),
read_call("call_1"),
done(),
],
vec![done()],
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
agent
.run_print_with_tools_streaming_output(
&provider,
"read file",
Some(&tools),
None,
temp.path(),
None,
)
.unwrap();
let requests = provider.requests();
assert_eq!(requests.len(), 2);
let continuation = requests[1].conversation_items();
let assistant_index = continuation
.iter()
.position(|item| {
matches!(
item,
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::Assistant
&& message.content == "I'll read."
)
})
.unwrap();
let response_index = continuation
.iter()
.position(|item| item == &ProviderConversationItem::ResponseItem(response_item.clone()))
.unwrap();
let tool_result_index = continuation
.iter()
.position(|item| {
matches!(
item,
ProviderConversationItem::ToolResult(result)
if result.call_id == "call_1" && result.output == "file body"
)
})
.unwrap();
assert!(assistant_index < response_index);
assert!(response_index < tool_result_index);
}
#[test]
fn agent_session_normalizes_empty_provider_and_model() {
let agent = AgentSession::new(" ", &[], &SkillDiscovery::default()).with_provider_id(" ");
assert_eq!(agent.model, "model");
assert_eq!(agent.provider_id, "provider");
}
#[test]
fn bash_mode_runner_emits_transient_tool_events_without_auth_or_session_persistence() {
let temp = tempfile::TempDir::new().unwrap();
let paths = crate::config::McPaths::from_root(temp.path().join("mc"));
let config = crate::config::EffectiveConfig {
provider: Some(crate::providers::OPENAI_CODEX_PROVIDER.to_string()),
model: Some("test-model".to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: std::collections::BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: None,
auth: None,
paths,
};
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let skills = SkillDiscovery::default();
let mut sink = CapturingOutputSink::default();
let output = runner::run_bash_mode_once(
&config,
&skills,
runner::ProviderRunOptions {
settings: Some(crate::config::Settings::default()),
prompt: "!printf hi",
session: Some(&session),
cwd: temp.path(),
output_sink: Some(&mut sink),
selected_primary_agent: None,
cancellation: None,
session_title_notifier: None,
invocation_mode: crate::output::InvocationMode::Shell,
disabled_tools: None,
disabled_subagent_profiles: None,
subagent_profile_discovery: None,
mcp: None,
},
)
.unwrap();
assert_eq!(output.tool_results.len(), 1);
assert!(
output.tool_results[0].success,
"{}",
output.tool_results[0].content
);
assert_eq!(output.tool_results[0].metadata["stdout"], "hi");
assert!(matches!(sink.outputs[0], OutputEvent::SessionHeader { .. }));
assert!(
matches!(sink.outputs[1], OutputEvent::BashCommand { ref command } if command == "printf hi")
);
assert!(
!sink
.outputs
.iter()
.any(|event| matches!(event, OutputEvent::UserPrompt { .. }))
);
assert!(sink.outputs.iter().any(
|event| matches!(event, OutputEvent::ToolStarted { call, .. } if call.name == "bash")
));
assert!(sink.outputs.iter().any(|event| matches!(event, OutputEvent::ToolResult { call, result, .. } if call.name == "bash" && result.metadata["stdout"] == "hi")));
let events = session.read_events().unwrap();
assert!(events.is_empty(), "events: {events:?}");
}
#[test]
fn bash_mode_runner_failure_stays_transient_to_session() {
let temp = tempfile::TempDir::new().unwrap();
let paths = crate::config::McPaths::from_root(temp.path().join("mc"));
crate::config::write_settings(&paths, &crate::config::Settings::default()).unwrap();
let config = crate::config::EffectiveConfig {
provider: Some(crate::providers::OPENAI_CODEX_PROVIDER.to_string()),
model: Some("test-model".to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: std::collections::BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: None,
auth: None,
paths,
};
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let skills = SkillDiscovery::default();
let mut sink = CapturingOutputSink::default();
let output = runner::run_bash_mode_once(
&config,
&skills,
runner::ProviderRunOptions {
settings: Some(crate::config::Settings::default()),
prompt: "!printf err >&2; exit 7",
session: Some(&session),
cwd: temp.path(),
output_sink: Some(&mut sink),
selected_primary_agent: None,
cancellation: None,
session_title_notifier: None,
invocation_mode: crate::output::InvocationMode::Shell,
disabled_tools: None,
disabled_subagent_profiles: None,
subagent_profile_discovery: None,
mcp: None,
},
)
.unwrap();
assert_eq!(output.tool_results.len(), 1);
assert!(!output.tool_results[0].success);
assert_eq!(output.tool_results[0].metadata["exit_code"], 7);
assert!(sink.outputs.iter().any(|event| matches!(event, OutputEvent::ToolResult { call, result, .. } if call.name == "bash" && result.metadata["exit_code"] == 7)));
let events = session.read_events().unwrap();
assert!(events.is_empty(), "events: {events:?}");
}
#[test]
fn bash_mode_hook_diagnostic_stays_transient_to_session() {
let temp = tempfile::TempDir::new().unwrap();
let paths = crate::config::McPaths::from_root(temp.path().join("mc"));
let settings = crate::config::Settings {
hooks: crate::config::HookSettings {
enabled: true,
before_tool: vec![crate::config::HookDefinition {
label: Some("warn-hook".to_string()),
command: "printf hookfail >&2; exit 2".to_string(),
..Default::default()
}],
..Default::default()
},
..Default::default()
};
let config = crate::config::EffectiveConfig {
provider: Some(crate::providers::OPENAI_CODEX_PROVIDER.to_string()),
model: Some("test-model".to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: std::collections::BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
api_key: None,
auth: None,
paths,
};
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let skills = SkillDiscovery::default();
let mut sink = CapturingOutputSink::default();
runner::run_bash_mode_once(
&config,
&skills,
runner::ProviderRunOptions {
settings: Some(settings),
prompt: "!printf hi",
session: Some(&session),
cwd: temp.path(),
output_sink: Some(&mut sink),
selected_primary_agent: None,
cancellation: None,
session_title_notifier: None,
invocation_mode: crate::output::InvocationMode::Shell,
disabled_tools: None,
disabled_subagent_profiles: None,
subagent_profile_discovery: None,
mcp: None,
},
)
.unwrap();
assert!(sink.outputs.iter().any(
|event| matches!(event, OutputEvent::HookDiagnostic { diagnostic } if diagnostic.label == "warn-hook")
));
let events = session.read_events().unwrap();
assert!(events.is_empty(), "events: {events:?}");
}
#[test]
fn oversized_failed_turn_replays_completed_tool_pairs_and_excludes_unpaired_calls() {
let temp = tempfile::TempDir::new().unwrap();
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let append = |event_type: &str, payload: serde_json::Value| {
session
.append(&SessionEvent::new(
event_type,
session.id().to_string(),
temp.path().to_path_buf(),
payload,
))
.unwrap();
};
let output = "x".repeat(70_000);
append("user_input", json!({"text":"inspect large files"}));
for (call_id, tool_name) in [("call_1", "read"), ("call_2", "grep")] {
append(
"tool_call",
json!({"id":call_id,"name":tool_name,"arguments":{"path":"file.txt"}}),
);
append(
"tool_result",
json!({"call_id":call_id,"result":{"tool_name":tool_name,"success":true,"content":output}}),
);
}
append(
"tool_call",
json!({"id":"call_unpaired","name":"write","arguments":{"path":"danger.txt"}}),
);
append(
"turn_status",
json!({"status":"failed","assistant_text":""}),
);
let replay = crate::context::build_conversation_replay(Some(&session)).unwrap();
let call_ids = replay
.items
.iter()
.filter_map(|item| match item {
ProviderConversationItem::ResponseItem(value) if value["type"] == "function_call" => {
value["call_id"].as_str()
}
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(call_ids, vec!["call_1", "call_2"]);
assert_eq!(
replay
.items
.iter()
.filter(|item| matches!(item, ProviderConversationItem::ToolResult(_)))
.count(),
2
);
assert!(
!replay
.replay_diagnostics
.iter()
.any(|diagnostic| diagnostic.contains("skipped_partial_historical_turn"))
);
}