use super::*;
#[test]
fn conversation_replay_drops_incomplete_marked_user_only_turn() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for (kind, payload) in [
(
SessionEventKind::UserInput,
json!({"text":"unreliable prompt"}),
),
(SessionEventKind::TurnStatus, json!({"status":"incomplete"})),
] {
session
.append(&SessionEvent::new_kind(
kind,
session.id().to_string(),
temp.path().to_path_buf(),
payload,
))
.unwrap();
}
let replay = build_conversation_replay(Some(&session)).unwrap();
assert!(replay.items.is_empty());
assert!(
replay
.replay_diagnostics
.iter()
.any(|message| { message.contains("dropped_incomplete_historical_user_only_turn") })
);
}
#[test]
fn full_conversation_replay_reads_all_provider_visible_events_in_order() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let call_item = json!({
"type": "function_call",
"call_id": "call_1",
"name": "read",
"arguments": "{\"path\":\"a.txt\"}",
"status": "completed"
});
for (event_type, payload) in [
("user_input", json!({"text":"old request"})),
("assistant_chunk", json!({"text":"I'll read."})),
("assistant_output", json!({"text":"I'll read."})),
("provider_response_item", json!({"item": call_item.clone()})),
(
"tool_call",
json!({"id":"call_1","name":"read","arguments":{"path":"a.txt"}}),
),
(
"tool_result",
json!({"call_id":"call_1","result":{"tool_name":"read","success":true,"content":"file contents"}}),
),
("assistant_chunk", json!({"text":"Found it."})),
("assistant_output", json!({"text":"I'll read.Found it."})),
("user_input", json!({"text":"next request"})),
] {
session
.append(&SessionEvent::new(
event_type,
session.id().to_string(),
temp.path().to_path_buf(),
payload,
))
.unwrap();
}
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(replay.legacy_lossy_events, Vec::<String>::new());
assert!(matches!(
&replay.items[0],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User && message.content == "old request"
));
assert!(matches!(
&replay.items[1],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::Assistant && message.content == "I'll read."
));
assert_eq!(
replay.items[2],
ProviderConversationItem::ResponseItem(call_item)
);
assert!(matches!(
&replay.items[3],
ProviderConversationItem::ToolResult(result)
if result.call_id == "call_1" && result.output == "file contents"
));
assert!(matches!(
&replay.items[4],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::Assistant && message.content == "Found it."
));
assert!(matches!(
&replay.items[5],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User && message.content == "next request"
));
}
#[test]
fn auto_recovery_continue_replays_inside_active_turn() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for (event_type, payload) in [
("user_input", json!({"text":"initial request"})),
("assistant_chunk", json!({"text":"partial answer"})),
(
"user_input",
json!({"text":"Continue", "auto_recovery": true, "reason":"incomplete_semantic_progress_timeout"}),
),
("assistant_chunk", json!({"text":" final answer"})),
(
"assistant_output",
json!({"text":"partial answer final answer"}),
),
] {
append_test_event(&session, temp.path(), event_type, payload);
}
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(
replay.items,
vec![
user_message("initial request"),
assistant_message("partial answer"),
user_message("Continue"),
assistant_message(" final answer"),
]
);
}
#[test]
fn normal_continue_does_not_replay_inside_active_turn() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for (event_type, payload) in [
("user_input", json!({"text":"initial request"})),
("assistant_chunk", json!({"text":"partial answer"})),
("user_input", json!({"text":"Continue"})),
("assistant_output", json!({"text":"later answer"})),
] {
append_test_event(&session, temp.path(), event_type, payload);
}
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(
replay.items,
vec![user_message("Continue"), assistant_message("later answer")]
);
}
#[test]
fn mid_run_steering_user_input_preserves_preceding_tool_lifecycle_on_replay() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let call_item = function_call_item("call_1", "read");
for (event_type, payload) in [
("user_input", json!({"text":"initial request"})),
("assistant_chunk", json!({"text":"I'll read."})),
("provider_response_item", json!({"item": call_item.clone()})),
(
"tool_result",
json!({"call_id":"call_1","result":{"tool_name":"read","success":true,"content":"file contents"}}),
),
(
"user_input",
json!({"text":"Steering update from user while current run was active:\n\nstay focused"}),
),
("assistant_chunk", json!({"text":"Done."})),
("assistant_output", json!({"text":"I'll read.Done."})),
] {
append_test_event(&session, temp.path(), event_type, payload);
}
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(replay.legacy_lossy_events, Vec::<String>::new());
assert_eq!(
replay.items,
vec![
user_message("initial request"),
assistant_message("I'll read."),
ProviderConversationItem::ResponseItem(call_item),
tool_result("call_1", "read", "file contents"),
user_message("Steering update from user while current run was active:\n\nstay focused"),
assistant_message("Done."),
]
);
}
#[test]
fn text_only_steering_after_completed_assistant_continues_active_turn() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for (event_type, payload) in [
("user_input", json!({"text":"initial request"})),
("assistant_chunk", json!({"text":"first text-only answer"})),
("assistant_output", json!({"text":"first text-only answer"})),
(
"user_input",
json!({"text":"Steering update from user while current run was active:\n\nrevise with detail"}),
),
("assistant_chunk", json!({"text":"continuation answer"})),
("assistant_output", json!({"text":"continuation answer"})),
] {
append_test_event(&session, temp.path(), event_type, payload);
}
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(
replay.items,
vec![
user_message("initial request"),
assistant_message("first text-only answer"),
user_message(
"Steering update from user while current run was active:\n\nrevise with detail"
),
assistant_message("continuation answer"),
]
);
}
#[test]
fn normal_user_input_after_completed_tool_turn_still_starts_new_replay_turn() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let call_item = function_call_item("call_1", "read");
for (event_type, payload) in [
("user_input", json!({"text":"initial request"})),
("assistant_chunk", json!({"text":"I'll read."})),
("provider_response_item", json!({"item": call_item.clone()})),
(
"tool_result",
json!({"call_id":"call_1","result":{"tool_name":"read","success":true,"content":"file contents"}}),
),
("assistant_output", json!({"text":"I'll read. Done."})),
("user_input", json!({"text":"next normal request"})),
] {
append_test_event(&session, temp.path(), event_type, payload);
}
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(
replay.items,
vec![
user_message("initial request"),
assistant_message("I'll read."),
ProviderConversationItem::ResponseItem(call_item),
tool_result("call_1", "read", "file contents"),
user_message("next normal request"),
]
);
}
#[test]
fn multiple_mid_run_steering_inputs_across_tool_boundaries_replay_in_provider_order() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let first_call = function_call_item("call_1", "read");
let second_call = function_call_item("call_2", "bash");
for (event_type, payload) in [
("user_input", json!({"text":"initial request"})),
("assistant_chunk", json!({"text":"Read first."})),
(
"provider_response_item",
json!({"item": first_call.clone()}),
),
(
"tool_result",
json!({"call_id":"call_1","result":{"tool_name":"read","success":true,"content":"first output"}}),
),
(
"user_input",
json!({"text":"Steering update from user while current run was active:\n\nsteering one"}),
),
("assistant_chunk", json!({"text":"Run command."})),
(
"provider_response_item",
json!({"item": second_call.clone()}),
),
(
"tool_result",
json!({"call_id":"call_2","result":{"tool_name":"bash","success":true,"content":"second output"}}),
),
(
"user_input",
json!({"text":"Steering update from user while current run was active:\n\nsteering two"}),
),
("assistant_chunk", json!({"text":"Finished."})),
(
"assistant_output",
json!({"text":"Read first.Run command.Finished."}),
),
] {
append_test_event(&session, temp.path(), event_type, payload);
}
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(
replay.items,
vec![
user_message("initial request"),
assistant_message("Read first."),
ProviderConversationItem::ResponseItem(first_call),
tool_result("call_1", "read", "first output"),
user_message("Steering update from user while current run was active:\n\nsteering one"),
assistant_message("Run command."),
ProviderConversationItem::ResponseItem(second_call),
tool_result("call_2", "bash", "second output"),
user_message("Steering update from user while current run was active:\n\nsteering two"),
assistant_message("Finished."),
]
);
}
#[test]
fn steering_input_followed_immediately_by_another_tool_call_replays_as_turn_continuation() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let first_call = function_call_item("call_1", "read");
let second_call = function_call_item("call_2", "bash");
for (event_type, payload) in [
("user_input", json!({"text":"initial request"})),
(
"provider_response_item",
json!({"item": first_call.clone()}),
),
(
"tool_result",
json!({"call_id":"call_1","result":{"tool_name":"read","success":true,"content":"first output"}}),
),
(
"user_input",
json!({"text":"Steering update from user while current run was active:\n\nuse bash next"}),
),
(
"provider_response_item",
json!({"item": second_call.clone()}),
),
(
"tool_result",
json!({"call_id":"call_2","result":{"tool_name":"bash","success":true,"content":"second output"}}),
),
("assistant_output", json!({"text":"done"})),
] {
append_test_event(&session, temp.path(), event_type, payload);
}
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(
replay.items,
vec![
user_message("initial request"),
ProviderConversationItem::ResponseItem(first_call),
tool_result("call_1", "read", "first output"),
user_message(
"Steering update from user while current run was active:\n\nuse bash next"
),
ProviderConversationItem::ResponseItem(second_call),
tool_result("call_2", "bash", "second output"),
assistant_message("done"),
]
);
}
fn user_message(text: &str) -> ProviderConversationItem {
ProviderConversationItem::Message(ChatMessage::user(text))
}
fn assistant_message(text: &str) -> ProviderConversationItem {
ProviderConversationItem::Message(ChatMessage::assistant(text))
}
fn function_call_item(call_id: &str, name: &str) -> Value {
json!({
"type": "function_call",
"call_id": call_id,
"name": name,
"arguments": "{}",
"status": "completed"
})
}
fn tool_result(call_id: &str, tool_name: &str, output: &str) -> ProviderConversationItem {
ProviderConversationItem::ToolResult(crate::providers::ProviderToolResult {
call_id: call_id.to_string(),
tool_name: tool_name.to_string(),
success: true,
output: output.to_string(),
})
}
#[test]
fn conversation_replay_downgrades_orphan_tool_results_without_structured_replay() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for (event_type, payload) in [
("user_input", json!({"text":"old request"})),
(
"tool_result",
json!({"call_id":"orphan_call","result":{"tool_name":"read","success":true,"content":"orphan output"}}),
),
(
"provider_response_item",
json!({"item":{"type":"function_call_output","call_id":"orphan_provider","output":"provider orphan"}}),
),
("assistant_output", json!({"text":"done"})),
] {
session
.append(&SessionEvent::new(
event_type,
session.id().to_string(),
temp.path().to_path_buf(),
payload,
))
.unwrap();
}
let replay = build_conversation_replay(Some(&session)).unwrap();
assert!(replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::LegacyReplayNote { event_type, content }
if event_type == "tool_result" && content.contains("orphan output")
)));
assert!(replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::LegacyReplayNote { event_type, content }
if event_type == "provider_response_item" && content.contains("provider orphan")
)));
assert!(!replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::ToolResult(result) if result.call_id == "orphan_call"
)));
assert!(!replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::ResponseItem(value)
if value.get("type").and_then(Value::as_str) == Some("function_call_output")
)));
assert!(
replay
.replay_diagnostics
.iter()
.any(|message| message.contains("downgraded_orphan_session_tool_result"))
);
assert!(
replay
.replay_diagnostics
.iter()
.any(|message| message.contains("downgraded_orphan_provider_tool_result"))
);
}
#[test]
fn conversation_replay_downgrades_orphan_tool_result_with_call_id() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for (event_type, payload) in [
("user_input", json!({"text":"old request"})),
(
"tool_call",
json!({"id":"call_1","name":"read","arguments":{"path":"a.txt"}}),
),
(
"tool_result",
json!({"call_id":"call_1","result":{"tool_name":"read","success":true,"content":"matched output"}}),
),
(
"tool_result",
json!({"call_id":"call_1","result":{"tool_name":"read","success":true,"content":"orphan duplicate output"}}),
),
("assistant_output", json!({"text":"done"})),
] {
session
.append(&SessionEvent::new(
event_type,
session.id().to_string(),
temp.path().to_path_buf(),
payload,
))
.unwrap();
}
let replay = build_conversation_replay(Some(&session)).unwrap();
let structured_results = replay
.items
.iter()
.filter(|item| {
matches!(
item,
ProviderConversationItem::ToolResult(result) if result.call_id == "call_1"
)
})
.count();
assert_eq!(structured_results, 1);
assert!(replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::ToolResult(result)
if result.output == "matched output"
)));
assert!(replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::LegacyReplayNote { event_type, content }
if event_type == "tool_result" && content.contains("orphan duplicate output")
)));
assert!(replay.replay_diagnostics.iter().any(|message| {
message.contains("downgraded_orphan_session_tool_result") && message.contains("call_1")
}));
}
#[test]
fn full_conversation_replay_excludes_local_only_session_events() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for event_type in [
"hook_diagnostic",
"hook_lifecycle",
"hook_context_injection",
"context_cache",
"session_title",
"diagnostic",
] {
session
.append(&SessionEvent::new(
event_type,
session.id().to_string(),
temp.path().to_path_buf(),
json!({"message":"LOCAL_ONLY_MARKER"}),
))
.unwrap();
}
session
.append(&SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"visible"}),
))
.unwrap();
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = conversation_cache_material("p", "m", "sys", &replay.items);
assert!(material.contains("visible"));
assert!(!material.contains("LOCAL_ONLY_MARKER"));
}
#[test]
fn provider_stream_trace_is_local_only_and_not_replayed() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"visible"}),
);
append_test_event(
&session,
temp.path(),
"provider_stream_trace",
json!({"pending_tools":[{"argument_sha256":"trace_hash_marker_123"}]}),
);
append_test_event(
&session,
temp.path(),
"assistant_output",
json!({"text":"done"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = conversation_cache_material("p", "m", "sys", &replay.items);
let recent =
recent_events_context(&session, ContextBudget::default().keep_recent_tokens).unwrap();
let debug = format!("{replay:?}");
assert!(material.contains("visible"));
assert!(material.contains("done"));
for provider_context in [material, recent, debug] {
assert!(!provider_context.contains("provider_stream_trace"));
assert!(!provider_context.contains("trace_hash_marker_123"));
}
assert!(
!replay
.legacy_lossy_events
.contains(&"provider_stream_trace".to_string())
);
}
#[test]
fn hook_context_injection_audit_is_local_only_but_provider_context_replays() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for (event_type, payload) in [
("user_input", json!({"text":"start"})),
(
"tool_call",
json!({"id":"call_1","name":"read","arguments":{"path":"a.txt"}}),
),
(
"tool_result",
json!({"call_id":"call_1","result":{"tool_name":"read","success":true,"content":"file text"}}),
),
(
"hook_context_injection",
json!({"status":"success","label":"audit","item_count":1,"byte_count":21}),
),
(
"provider_context_item",
json!({"role":"user","content":"injected memory note"}),
),
("assistant_output", json!({"text":"done"})),
] {
append_test_event(&session, temp.path(), event_type, payload);
}
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
let recent =
recent_events_context(&session, ContextBudget::default().keep_recent_tokens).unwrap();
assert!(replay.items.windows(2).any(|pair| matches!(
(&pair[0], &pair[1]),
(
ProviderConversationItem::ToolResult(result),
ProviderConversationItem::Message(message),
) if result.output == "file text" && message.content == "injected memory note"
)));
assert!(material.contains("injected memory note"));
assert!(recent.contains("injected memory note"));
assert!(!material.contains("hook_context_injection"));
assert!(!recent.contains("hook_context_injection"));
}
#[test]
fn recent_context_exposes_only_compaction_summary_not_checkpoint_aggregate() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"compaction",
json!({
"schema_version": 1,
"summary": "intended summary boundary",
"provider": "provider-marker",
"model": "model-marker",
"aggregate": {
"first_user_input_text": "first-user-marker",
"session_path": "session-path-marker",
"latest_title": "title-marker",
"aggregate_marker": "aggregate-marker",
"secret": "sk-compaction-secret-marker"
},
"cutoff_event_count": 0
}),
);
let recent =
recent_events_context(&session, ContextBudget::default().keep_recent_tokens).unwrap();
let replay = build_conversation_replay(Some(&session)).unwrap();
let replay_payload = serde_json::to_string(&replay.items).unwrap();
for value in [recent, replay_payload] {
assert!(value.contains("intended summary boundary"));
for marker in [
"first-user-marker",
"session-path-marker",
"title-marker",
"aggregate-marker",
"sk-compaction-secret-marker",
"provider-marker",
"model-marker",
] {
assert!(!value.contains(marker), "leaked marker: {marker}");
}
}
}
#[test]
fn conversation_replay_ignores_reasoning_summary_but_keeps_reasoning_item() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let reasoning_item = json!({
"id":"rs_1",
"type":"reasoning",
"summary":[{"type":"summary_text","text":"provider summary"}],
"encrypted_content":{"nested":[{"encrypted_content":"CACHE_OPAQUE"}]}
});
for (event_type, payload) in [
("user_input", json!({"text":"hello"})),
("reasoning_summary", json!({"text":"provider summary"})),
(
"provider_response_item",
json!({"item": reasoning_item.clone()}),
),
("assistant_output", json!({"text":"answer"})),
] {
session
.append(&SessionEvent::new(
event_type,
session.id().to_string(),
temp.path().to_path_buf(),
payload,
))
.unwrap();
}
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = conversation_cache_material("p", "m", "sys", &replay.items);
assert!(!material.contains("reasoning_summary"));
let replay_serialized = serde_json::to_string(&replay.items).unwrap();
assert!(replay_serialized.contains("\"type\":\"reasoning"));
assert!(!replay_serialized.contains("CACHE_OPAQUE"));
}
#[test]
fn legacy_text_only_session_replays_best_effort() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for (event_type, payload) in [
("user_input", json!({"text":"hello"})),
("assistant_output", json!({"text":"hi there"})),
] {
session
.append(&SessionEvent::new(
event_type,
session.id().to_string(),
temp.path().to_path_buf(),
payload,
))
.unwrap();
}
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(replay.legacy_lossy_events, Vec::<String>::new());
assert!(matches!(
&replay.items[0],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User && message.content == "hello"
));
assert!(matches!(
&replay.items[1],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::Assistant && message.content == "hi there"
));
}
#[test]
fn conversation_replay_skips_whitespace_only_assistant_messages() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for (event_type, payload) in [
("user_input", json!({"text":"hello"})),
("assistant_chunk", json!({"text":"\n\n"})),
("assistant_output", json!({"text":"\n\n"})),
("user_input", json!({"text":"next"})),
("assistant_output", json!({"text":"\nvisible\n"})),
] {
session
.append(&SessionEvent::new(
event_type,
session.id().to_string(),
temp.path().to_path_buf(),
payload,
))
.unwrap();
}
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(replay.legacy_lossy_events, Vec::<String>::new());
assert_eq!(replay.items.len(), 3);
assert!(matches!(
&replay.items[0],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User && message.content == "hello"
));
assert!(matches!(
&replay.items[1],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User && message.content == "next"
));
assert!(matches!(
&replay.items[2],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::Assistant && message.content == "\nvisible\n"
));
}
#[test]
fn legacy_lossy_tool_event_is_marked_not_exact() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
session
.append(&SessionEvent::new(
"tool_result",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"result":{"tool_name":"read","success":true,"content":"orphan"}}),
))
.unwrap();
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(replay.legacy_lossy_events, vec!["tool_result"]);
assert!(matches!(
&replay.items[0],
ProviderConversationItem::LegacyReplayNote { event_type, content }
if event_type == "tool_result" && content.contains("orphan")
));
}
#[test]
fn unknown_provider_visible_legacy_event_becomes_replay_note() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
session
.append(&SessionEvent::new(
"context_compaction",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"summary":"old compacted facts"}),
))
.unwrap();
session
.append(&SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"new request"}),
))
.unwrap();
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(replay.legacy_lossy_events, vec!["context_compaction"]);
assert!(matches!(
&replay.items[0],
ProviderConversationItem::LegacyReplayNote { event_type, content }
if event_type == "context_compaction" && content.contains("old compacted facts")
));
assert!(matches!(
&replay.items[1],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User && message.content == "new request"
));
}
#[test]
fn conversation_replay_tolerates_truncated_session_tail() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.open("safe").unwrap();
let first = SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"survives"}),
);
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::write(
session.path(),
format!(
"{}\n{{\"event_type\":\"assistant_output\",\"payload\":",
serde_json::to_string(&first).unwrap()
),
)
.unwrap();
crate::sessions::secure_test_session_root(session.path().parent().unwrap());
let replay = build_conversation_replay(Some(&session)).unwrap();
assert_eq!(replay.session_read_diagnostics.len(), 1);
assert_eq!(replay.session_read_diagnostics[0].line, 2);
assert!(matches!(
&replay.items[0],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User && message.content == "survives"
));
}
#[test]
fn recent_events_context_tolerates_malformed_legacy_line() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.open("safe").unwrap();
let first = SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"before"}),
);
let second = SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"after"}),
);
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::write(
session.path(),
format!(
"{}\nlegacy damaged line\n{}\n",
serde_json::to_string(&first).unwrap(),
serde_json::to_string(&second).unwrap()
),
)
.unwrap();
crate::sessions::secure_test_session_root(session.path().parent().unwrap());
let report =
recent_events_context_report(&session, ContextBudget::default().keep_recent_tokens)
.unwrap();
assert!(report.text.contains("before"));
assert!(report.text.contains("after"));
assert_eq!(report.diagnostics[0].line, 2);
assert!(
!report.diagnostics[0]
.message
.contains("legacy damaged line")
);
}
#[test]
fn bounded_context_reads_keep_recent_events() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for index in 0..50 {
session
.append(&SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
serde_json::json!({"text": format!("event-{index}")}),
))
.unwrap();
}
let recent = recent_events_context(&session, 8).unwrap();
assert!(recent.contains("event-49"));
assert!(!recent.contains("event-0"));
}
#[test]
fn recent_events_context_tail_omits_older_large_session_and_reports_diagnostic() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.open("safe").unwrap();
let early = SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"EARLY_SHOULD_BE_OMITTED"}),
);
let recent = SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"TAIL_SHOULD_SURVIVE"}),
);
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::write(
session.path(),
format!(
"{}\n{}\n{}\n",
serde_json::to_string(&early).unwrap(),
"x".repeat(2 * 1024 * 1024),
serde_json::to_string(&recent).unwrap()
),
)
.unwrap();
crate::sessions::secure_test_session_root(session.path().parent().unwrap());
let report = recent_events_context_report(&session, usize::MAX).unwrap();
assert!(report.text.contains("TAIL_SHOULD_SURVIVE"));
assert!(!report.text.contains("EARLY_SHOULD_BE_OMITTED"));
assert!(
report
.diagnostics
.iter()
.any(|diagnostic| diagnostic.message.contains("bounded tail window"))
);
}
#[test]
fn recent_events_context_report_surfaces_malformed_tail_diagnostics_without_text_leak() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.open("safe").unwrap();
let visible = SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"visible after malformed tail"}),
);
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::write(
session.path(),
format!(
"{}\nmalformed tail secret-token-should-not-leak\n{}\n",
"x".repeat(2 * 1024 * 1024),
serde_json::to_string(&visible).unwrap()
),
)
.unwrap();
crate::sessions::secure_test_session_root(session.path().parent().unwrap());
let report = recent_events_context_report(&session, usize::MAX).unwrap();
assert!(report.text.contains("visible after malformed tail"));
assert!(!report.text.contains("secret-token-should-not-leak"));
assert!(report.diagnostics.iter().any(|diagnostic| {
diagnostic.message.contains("failed to parse session JSONL")
&& !diagnostic.message.contains("secret-token-should-not-leak")
}));
}
#[test]
fn recent_events_context_report_matches_text_wrapper_for_small_valid_session() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"small one"}),
);
append_test_event(
&session,
temp.path(),
"assistant_output",
json!({"text":"small two"}),
);
let text = recent_events_context(&session, usize::MAX).unwrap();
let report = recent_events_context_report(&session, usize::MAX).unwrap();
assert_eq!(report.text, text);
assert!(report.text.contains("small one"));
assert!(report.text.contains("small two"));
assert!(report.diagnostics.is_empty());
}
#[test]
fn provider_session_context_excludes_hook_diagnostics() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let marker = "HOOK_DIAGNOSTIC_PROVIDER_LEAK_MARKER";
session
.append(&SessionEvent::new(
"hook_diagnostic",
session.id().to_string(),
temp.path().to_path_buf(),
serde_json::json!({"message": marker, "tool": "bash"}),
))
.unwrap();
session
.append(&SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
serde_json::json!({"text":"provider-visible continuity"}),
))
.unwrap();
let recent =
recent_events_context(&session, ContextBudget::default().keep_recent_tokens).unwrap();
let replay = build_conversation_replay(Some(&session)).unwrap();
let replay_material = conversation_cache_material("p", "m", "sys", &replay.items);
let persisted = session.read_events().unwrap();
assert!(
persisted
.iter()
.any(|event| event.event_type == "hook_diagnostic")
);
for provider_context in [recent, replay_material] {
assert!(provider_context.contains("provider-visible continuity"));
assert!(!provider_context.contains(marker));
assert!(!provider_context.contains("hook_diagnostic"));
}
}
#[test]
fn provider_session_context_excludes_hook_lifecycle() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let marker = "HOOK_LIFECYCLE_PROVIDER_LEAK_MARKER";
session
.append(&SessionEvent::new(
"hook_lifecycle",
session.id().to_string(),
temp.path().to_path_buf(),
serde_json::json!({
"phase":"before_tool",
"label": marker,
"target_tool":"bash",
"status":"success",
"policy":"warn",
"target_ran": false
}),
))
.unwrap();
session
.append(&SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
serde_json::json!({"text":"provider-visible continuity"}),
))
.unwrap();
let recent =
recent_events_context(&session, ContextBudget::default().keep_recent_tokens).unwrap();
let replay = build_conversation_replay(Some(&session)).unwrap();
let replay_material = conversation_cache_material("p", "m", "sys", &replay.items);
for provider_context in [recent, replay_material] {
assert!(provider_context.contains("provider-visible continuity"));
assert!(!provider_context.contains(marker));
assert!(!provider_context.contains("hook_lifecycle"));
}
}
#[test]
fn historical_summary_events_do_not_change_recent_context() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let old_type = format!("{}_{}", "context", "compaction");
let old_heading = ["Latest", " compacted", " session", " summary"].concat();
session
.append(&SessionEvent::new(
old_type.clone(),
session.id().to_string(),
temp.path().to_path_buf(),
serde_json::json!({"summary":"old summary line"}),
))
.unwrap();
session
.append(&SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
serde_json::json!({"text":"newer line"}),
))
.unwrap();
let recent =
recent_events_context(&session, ContextBudget::default().keep_recent_tokens).unwrap();
assert!(!recent.contains(&old_heading));
assert!(recent.contains(&old_type));
assert!(recent.contains("newer line"));
}
#[test]
fn optimization_harness_replays_large_synthetic_session_without_dropping_tool_pairs() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let tool_pair_count = 256;
session
.append(&SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"start large replay"}),
))
.unwrap();
for index in 0..tool_pair_count {
session
.append(&SessionEvent::new(
"tool_call",
session.id().to_string(),
temp.path().to_path_buf(),
json!({
"id": format!("call_{index}"),
"name": "read",
"arguments": {"path": format!("fixtures/{index}.txt")}
}),
))
.unwrap();
session
.append(&SessionEvent::new(
"tool_result",
session.id().to_string(),
temp.path().to_path_buf(),
json!({
"call_id": format!("call_{index}"),
"result": {
"tool_name": "read",
"success": true,
"content": format!("file contents {index}")
}
}),
))
.unwrap();
}
session
.append(&SessionEvent::new(
"assistant_output",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"done"}),
))
.unwrap();
let replay = build_conversation_replay(Some(&session)).unwrap();
let function_calls = replay
.items
.iter()
.filter(|item| {
matches!(
item,
ProviderConversationItem::ResponseItem(value)
if value.get("type").and_then(Value::as_str) == Some("function_call")
)
})
.count();
let tool_results = replay
.items
.iter()
.filter(|item| matches!(item, ProviderConversationItem::ToolResult(_)))
.count();
let material = conversation_cache_material("p", "m", "sys", &replay.items);
assert_eq!(replay.legacy_lossy_events, Vec::<String>::new());
assert_eq!(function_calls, tool_pair_count);
assert_eq!(tool_results, tool_pair_count);
assert!(material.contains(&format!("call_{}", tool_pair_count - 1)));
assert!(material.contains(&format!("file contents {}", tool_pair_count - 1)));
}
#[test]
fn conversation_replay_recovers_failed_provider_body_turn_as_summary() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"latest objective: fix resume recovery after provider body error"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"call_subagents","name":"subagents","arguments":{"tasks":[{"intent":"CODE_WRITING_EXECUTION"}]}}),
);
append_test_event(
&session,
temp.path(),
"tool_result",
json!({"call_id":"call_subagents","result":{"tool_name":"subagents","success":true,"content":"g1 PASS\nVerification: cargo test context::tests::conversation_replay passed\nNext action: run cargo clippy -- -D warnings"}}),
);
append_test_event(
&session,
temp.path(),
"diagnostic",
json!({"message":"provider request failed: request or response body error while streaming response body"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"continue"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
let summary = match &replay.items[0] {
ProviderConversationItem::Message(message) => &message.content,
other => panic!("expected recovery summary message, got {other:?}"),
};
assert_eq!(replay.items.len(), 2);
assert!(summary.contains("Session recovery summary for failed historical turn"));
assert!(summary.contains("not exact structured replay"));
assert!(summary.contains("latest objective: fix resume recovery"));
assert!(summary.contains("CWD:"));
assert!(summary.contains(&temp.path().display().to_string()));
assert!(summary.contains("tool=subagents success=true"));
assert!(summary.contains("g1 PASS"));
assert!(
summary.contains("Verification: cargo test context::tests::conversation_replay passed")
);
assert!(summary.contains("Next action: run cargo clippy -- -D warnings"));
assert!(summary.contains("Last diagnostic/error category: provider_body_error"));
assert!(matches!(
&replay.items[1],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User && message.content == "continue"
));
assert!(!replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::ResponseItem(_) | ProviderConversationItem::ToolResult(_)
)));
assert!(
material.find("Session recovery summary").unwrap() < material.find("continue").unwrap()
);
assert!(replay.replay_diagnostics.iter().any(|diagnostic| {
diagnostic.contains("recovered_failed_historical_turn")
&& diagnostic.contains("summary_chars=")
}));
}
#[test]
fn conversation_replay_recovers_semantic_progress_timeout_turn_as_summary() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"original objective: fix timeout replay recovery"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"timeout_call","name":"subagents","arguments":{"tasks":[{"intent":"IMPLEMENT"}]}}),
);
append_test_event(
&session,
temp.path(),
"tool_result",
json!({"call_id":"timeout_call","result":{"tool_name":"subagents","success":true,"content":"subagent finished narrow fix\nverification ready"}}),
);
append_test_event(
&session,
temp.path(),
"diagnostic",
json!({"level":"error","message":"provider stream no semantic progress before timeout"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"continue"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
let summary = match &replay.items[0] {
ProviderConversationItem::Message(message) => &message.content,
other => panic!("expected recovery summary message, got {other:?}"),
};
assert_eq!(replay.items.len(), 2);
assert!(summary.contains("Session recovery summary for failed historical turn"));
assert!(summary.contains("original objective: fix timeout replay recovery"));
assert!(summary.contains("tool=subagents success=true"));
assert!(summary.contains("subagent finished narrow fix"));
assert!(summary.contains("Last diagnostic/error category: provider_stream_timeout"));
assert!(summary.contains("Diagnostic categories: provider_stream_timeout"));
assert!(matches!(
&replay.items[1],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User && message.content == "continue"
));
assert!(!replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::ResponseItem(_) | ProviderConversationItem::ToolResult(_)
)));
assert!(
material.find("Session recovery summary").unwrap() < material.find("continue").unwrap()
);
}
#[test]
fn conversation_replay_recovers_stream_idle_timeout_turn_as_summary() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"original objective: recover stream idle timeout turn"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"idle_timeout_call","name":"subagents","arguments":{"tasks":[{"intent":"IMPLEMENT"}]}}),
);
append_test_event(
&session,
temp.path(),
"tool_result",
json!({"call_id":"idle_timeout_call","result":{"tool_name":"subagents","success":true,"content":"subagent completed before idle timeout\nverification pending"}}),
);
append_test_event(
&session,
temp.path(),
"diagnostic",
json!({"level":"error","message":"provider stream idle timeout after 60s without bytes"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"continue"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
let summary = match &replay.items[0] {
ProviderConversationItem::Message(message) => &message.content,
other => panic!("expected recovery summary message, got {other:?}"),
};
assert_eq!(replay.items.len(), 2);
assert!(summary.contains("Session recovery summary for failed historical turn"));
assert!(summary.contains("original objective: recover stream idle timeout turn"));
assert!(summary.contains("tool=subagents success=true"));
assert!(summary.contains("subagent completed before idle timeout"));
assert!(summary.contains("Last diagnostic/error category: provider_stream_timeout"));
assert!(summary.contains("Diagnostic categories: provider_stream_timeout"));
assert!(matches!(
&replay.items[1],
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User && message.content == "continue"
));
assert!(!replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::ResponseItem(_) | ProviderConversationItem::ToolResult(_)
)));
assert!(
material.find("Session recovery summary").unwrap() < material.find("continue").unwrap()
);
}
#[test]
fn conversation_replay_recovery_summary_redacts_secrets() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.open("safe").unwrap();
let user_api_key = format!("sk-{}", "x".repeat(24));
let tool_bearer = format!("Bearer {}", "x".repeat(24));
let tool_api_key = format!("sk-{}", "x".repeat(24));
let diagnostic_bearer = format!("Bearer {}", "x".repeat(24));
let raw_events = [
SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":format!("recover without leaking api_key={user_api_key}")}),
),
SessionEvent::new(
"tool_call",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"id":"secret_call","name":"subagents","arguments":{}}),
),
SessionEvent::new(
"tool_result",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"call_id":"secret_call","result":{"tool_name":"subagents","success":true,"content":format!("PASS {tool_bearer} api_key={tool_api_key} code=SECRETCODE")}}),
),
SessionEvent::new(
"diagnostic",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"message":format!("request body failed Authorization: {diagnostic_bearer}")}),
),
SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"continue safely"}),
),
];
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::write(
session.path(),
raw_events
.iter()
.map(|event| serde_json::to_string(event).unwrap())
.collect::<Vec<_>>()
.join("\n"),
)
.unwrap();
crate::sessions::secure_test_session_root(session.path().parent().unwrap());
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("<redacted>"));
for secret in [
user_api_key.as_str(),
tool_bearer.as_str(),
tool_api_key.as_str(),
diagnostic_bearer.as_str(),
] {
assert!(!material.contains(secret), "leaked generated fixture");
}
assert!(material.contains("code=SECRETCODE"));
}
#[test]
fn conversation_replay_recovery_summary_tolerates_malformed_jsonl() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.open("safe").unwrap();
let events = [
SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"objective survives malformed line"}),
),
SessionEvent::new(
"tool_result",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"call_id":"malformed_safe","result":{"tool_name":"subagents","success":true,"content":"PASS after malformed line"}}),
),
SessionEvent::new(
"diagnostic",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"message":"failed to read provider error body"}),
),
SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"continue after malformed"}),
),
];
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::write(
session.path(),
format!(
"{}\nmalformed sk-malformedSecret999 line\n{}\n{}\n{}\n",
serde_json::to_string(&events[0]).unwrap(),
serde_json::to_string(&events[1]).unwrap(),
serde_json::to_string(&events[2]).unwrap(),
serde_json::to_string(&events[3]).unwrap()
),
)
.unwrap();
crate::sessions::secure_test_session_root(session.path().parent().unwrap());
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert_eq!(replay.session_read_diagnostics.len(), 1);
assert_eq!(replay.session_read_diagnostics[0].line, 2);
assert!(
!replay.session_read_diagnostics[0]
.message
.contains("sk-malformedSecret999")
);
assert!(material.contains("objective survives malformed line"));
assert!(material.contains("PASS after malformed line"));
assert!(material.contains("continue after malformed"));
}
#[test]
fn conversation_replay_detects_provider_body_diagnostics() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
for diagnostic in [
"request body failed before provider response",
"response body stream ended early",
"failed to read provider error body",
] {
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text": diagnostic}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id": format!("call_{diagnostic}"), "name":"read", "arguments":{}}),
);
append_test_event(
&session,
temp.path(),
"diagnostic",
json!({"message": diagnostic}),
);
}
append_test_event(&session, temp.path(), "user_input", json!({"text":"after"}));
let replay = build_conversation_replay(Some(&session)).unwrap();
let summary_count = replay
.items
.iter()
.filter(|item| {
matches!(
item,
ProviderConversationItem::Message(message)
if message.content.contains("Session recovery summary")
)
})
.count();
assert_eq!(summary_count, 3);
assert!(
replay
.items
.iter()
.filter_map(|item| match item {
ProviderConversationItem::Message(message)
if message.content.contains("Session recovery summary") =>
{
Some(&message.content)
}
_ => None,
})
.all(|summary| summary.contains("provider_body_error"))
);
}
#[test]
fn conversation_replay_does_not_recover_unrelated_local_diagnostic() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"local shell failure should not become context"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"local_diag_call","name":"bash","arguments":{}}),
);
append_test_event(
&session,
temp.path(),
"diagnostic",
json!({"message":"shell command exited with status 1"}),
);
append_test_event(&session, temp.path(), "user_input", json!({"text":"after"}));
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("after"));
assert!(!material.contains("Session recovery summary"));
assert!(!material.contains("local shell failure should not become context"));
assert!(!material.contains("local_diag_call"));
}
#[test]
fn conversation_replay_uses_latest_valid_compaction_as_hard_boundary() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"RAW_BEFORE_BOUNDARY"}),
);
crate::sessions::record_session_compaction(
&session,
temp.path(),
"SUMMARY_ONCE",
"provider-a",
"model-a",
1,
)
.unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"AFTER_BOUNDARY"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("SUMMARY_ONCE"));
assert_eq!(material.matches("SUMMARY_ONCE").count(), 1);
assert!(material.contains("AFTER_BOUNDARY"));
assert!(!material.contains("RAW_BEFORE_BOUNDARY"));
}
#[test]
fn conversation_replay_uses_compaction_cutoff_not_checkpoint_index() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"RAW_BEFORE_CUTOFF"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"CONCURRENT_AFTER_CUTOFF"}),
);
crate::sessions::record_session_compaction(
&session,
temp.path(),
"SUMMARY_CUTOFF",
"provider-a",
"model-a",
1,
)
.unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"AFTER_CHECKPOINT"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("SUMMARY_CUTOFF"));
assert!(material.contains("CONCURRENT_AFTER_CUTOFF"));
assert!(material.contains("AFTER_CHECKPOINT"));
assert!(!material.contains("RAW_BEFORE_CUTOFF"));
}
#[test]
fn conversation_replay_malformed_latest_compaction_falls_back_to_previous_valid() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"RAW_BEFORE_FIRST"}),
);
crate::sessions::record_session_compaction(
&session,
temp.path(),
"VALID_SUMMARY",
"provider-a",
"model-a",
1,
)
.unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"AFTER_VALID"}),
);
let malformed = crate::sessions::SessionEvent::new(
"compaction",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"schema_version": 1, "summary":" "}),
);
session.append(&malformed).unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"AFTER_MALFORMED"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("VALID_SUMMARY"));
assert_eq!(material.matches("VALID_SUMMARY").count(), 1);
assert!(material.contains("AFTER_VALID"));
assert!(material.contains("AFTER_MALFORMED"));
assert!(!material.contains("RAW_BEFORE_FIRST"));
assert!(
replay
.replay_diagnostics
.iter()
.any(|message| message.contains("empty_summary"))
);
}
#[test]
fn conversation_replay_ignores_compaction_with_cutoff_after_checkpoint_index() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"RAW_BEFORE_FIRST"}),
);
crate::sessions::record_session_compaction(
&session,
temp.path(),
"VALID_SUMMARY",
"provider-a",
"model-a",
1,
)
.unwrap();
let malformed = crate::sessions::SessionEvent::new(
"compaction",
session.id().to_string(),
temp.path().to_path_buf(),
json!({
"schema_version": 1,
"summary":"BAD_SUMMARY",
"provider":"provider-a",
"model":"model-a",
"cutoff_event_count": 99,
}),
);
session.append(&malformed).unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"AFTER_BAD"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("VALID_SUMMARY"));
assert!(material.contains("AFTER_BAD"));
assert!(!material.contains("BAD_SUMMARY"));
assert!(!material.contains("RAW_BEFORE_FIRST"));
assert!(
replay
.replay_diagnostics
.iter()
.any(|message| message.contains("cutoff_event_count_out_of_bounds"))
);
}
#[test]
fn conversation_replay_repeated_compaction_uses_only_latest_summary_and_suffix() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"ERA_ONE_RAW"}),
);
crate::sessions::record_session_compaction(
&session,
temp.path(),
"FIRST_SUMMARY",
"provider-a",
"model-a",
1,
)
.unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"ERA_TWO_RAW"}),
);
crate::sessions::record_session_compaction(
&session,
temp.path(),
"SECOND_SUMMARY",
"provider-a",
"model-a",
3,
)
.unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"ERA_THREE_RAW"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("SECOND_SUMMARY"));
assert_eq!(material.matches("SECOND_SUMMARY").count(), 1);
assert!(material.contains("ERA_THREE_RAW"));
assert!(!material.contains("FIRST_SUMMARY"));
assert!(!material.contains("ERA_ONE_RAW"));
assert!(!material.contains("ERA_TWO_RAW"));
}
#[test]
fn conversation_replay_preserves_post_compaction_tool_pairs() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"before"}),
);
crate::sessions::record_session_compaction(
&session,
temp.path(),
"SUMMARY",
"provider-a",
"model-a",
1,
)
.unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"use tool"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"post_call","name":"read","arguments":{"path":"a.txt"}}),
);
append_test_event(
&session,
temp.path(),
"tool_result",
json!({"call_id":"post_call","result":{"tool_name":"read","success":true,"content":"POST_TOOL_OUTPUT"}}),
);
append_test_event(
&session,
temp.path(),
"assistant_output",
json!({"text":"done"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("SUMMARY"));
assert!(material.contains("post_call"));
assert!(material.contains("POST_TOOL_OUTPUT"));
assert!(replay.items.iter().any(|item| matches!(
item,
ProviderConversationItem::ToolResult(result) if result.call_id == "post_call"
)));
}
#[test]
fn conversation_replay_preserves_consecutive_historical_user_only_turns() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"first historical user-only turn"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"second historical user-only turn"}),
);
append_test_event(
&session,
temp.path(),
"assistant_output",
json!({"text":"answer to second turn"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("first historical user-only turn"));
assert!(material.contains("second historical user-only turn"));
assert!(material.contains("answer to second turn"));
}
#[test]
fn conversation_replay_skips_balanced_tool_turn_with_provider_failure_and_no_assistant_answer() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"safe prior"}),
);
append_test_event(
&session,
temp.path(),
"assistant_output",
json!({"text":"done"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"failed turn"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"call_failed","name":"read","arguments":{"path":"x"}}),
);
append_test_event(
&session,
temp.path(),
"tool_result",
json!({"call_id":"call_failed","result":{"tool_name":"read","success":true,"content":"failed output"}}),
);
append_test_event(
&session,
temp.path(),
"diagnostic",
json!({"message":"provider stream ended with failed or incomplete response: response.failed"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"latest pending"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("safe prior"));
assert!(material.contains("done"));
assert!(material.contains("latest pending"));
assert!(material.contains("failed turn"));
assert!(!material.contains("call_failed"));
assert!(material.contains("failed output"));
assert!(
replay
.replay_diagnostics
.iter()
.any(|d| d.contains("recovered_failed_historical_turn"))
);
}
#[test]
fn conversation_replay_uses_local_diagnostic_as_skip_evidence_without_replaying_it() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let diagnostic = "response.failed LOCAL_DIAGNOSTIC_MARKER";
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"do tool"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"call_diag","name":"read","arguments":{}}),
);
append_test_event(
&session,
temp.path(),
"diagnostic",
json!({"message": diagnostic}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"pending"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("pending"));
assert!(material.contains("do tool"));
assert!(material.contains("tool=read"));
assert!(!material.contains("call_diag"));
assert!(!material.contains("LOCAL_DIAGNOSTIC_MARKER"));
assert!(!material.contains("response.failed"));
}
#[test]
fn conversation_replay_keeps_only_latest_pending_user_input_after_failed_retries() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"old ok"}),
);
append_test_event(
&session,
temp.path(),
"assistant_output",
json!({"text":"old answer"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"retry one"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"retry_one","name":"read","arguments":{}}),
);
append_test_event(
&session,
temp.path(),
"diagnostic",
json!({"message":"failed or incomplete response"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"retry two"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"retry_two","name":"read","arguments":{}}),
);
append_test_event(
&session,
temp.path(),
"diagnostic",
json!({"message":"response.incomplete"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"final pending"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("old ok"));
assert!(material.contains("old answer"));
assert!(material.contains("final pending"));
assert!(material.contains("retry one"));
assert!(material.contains("retry two"));
let summary_count = replay
.items
.iter()
.filter(|item| {
matches!(
item,
ProviderConversationItem::Message(message)
if message.content.contains("Session recovery summary")
)
})
.count();
assert_eq!(summary_count, 2);
}
#[test]
fn conversation_replay_skips_failed_oversized_turn_and_keeps_later_user_input() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let poison = "POISON".repeat(100_000);
append_test_event(&session, temp.path(), "user_input", json!({"text":"safe"}));
append_test_event(
&session,
temp.path(),
"assistant_output",
json!({"text":"answer"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"oversized failed"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"huge_call","name":"read","arguments":{}}),
);
append_test_event(
&session,
temp.path(),
"tool_result",
json!({"call_id":"huge_call","result":{"tool_name":"read","success":true,"content":poison}}),
);
append_test_event(
&session,
temp.path(),
"diagnostic",
json!({"message":"response.failed"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"Just respond HI"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("safe"));
assert!(material.contains("answer"));
assert!(material.contains("Just respond HI"));
assert!(material.contains("oversized failed"));
assert!(!material.contains("huge_call"));
assert!(!material.contains("POISON"));
}
#[test]
fn conversation_replay_drops_eof_uncommitted_tool_call_but_keeps_user() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"pending with partial tool"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"eof_call","name":"read","arguments":{"path":"x"}}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("pending with partial tool"));
assert!(!material.contains("eof_call"));
assert!(!material.contains("function_call"));
}
#[test]
fn conversation_replay_drops_eof_uncommitted_tool_pair_but_keeps_user() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"pending with partial tool pair"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"eof_pair","name":"read","arguments":{}}),
);
append_test_event(
&session,
temp.path(),
"tool_result",
json!({"call_id":"eof_pair","result":{"tool_name":"read","success":true,"content":"partial result"}}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("pending with partial tool pair"));
assert!(!material.contains("eof_pair"));
assert!(!material.contains("partial result"));
}
#[test]
fn conversation_replay_replays_failed_partial_assistant_chunk() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"chunk failed"}),
);
append_test_event(
&session,
temp.path(),
"assistant_chunk",
json!({"text":"partial answer"}),
);
append_test_event(
&session,
temp.path(),
"turn_status",
json!({"status":"failed", "assistant_text":"partial answer"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"new pending"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("new pending"));
assert!(material.contains("chunk failed"));
assert!(material.contains("partial answer"));
assert!(replay.replay_diagnostics.iter().any(|diagnostic| {
diagnostic.contains("replayed_partial_historical_turn")
&& diagnostic.contains("failed=true")
}));
}
#[test]
fn conversation_replay_preserves_compaction_required_tool_progress() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"run tool before compaction"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"compact_boundary_call","name":"read","arguments":{}}),
);
append_test_event(
&session,
temp.path(),
"tool_result",
json!({"call_id":"compact_boundary_call","result":{"tool_name":"read","success":true,"content":"COMPACTION_BOUNDARY_RESULT"}}),
);
append_test_event(
&session,
temp.path(),
"turn_status",
json!({"status":"compaction_required", "assistant_text":"tool phase"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("run tool before compaction"));
assert!(material.contains("compact_boundary_call"));
assert!(material.contains("COMPACTION_BOUNDARY_RESULT"));
assert!(material.contains("tool phase"));
}
#[test]
fn conversation_replay_replays_cancelled_partial_assistant_chunk_at_eof() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"cancelled prompt"}),
);
append_test_event(
&session,
temp.path(),
"assistant_chunk",
json!({"text":"partial cancelled answer"}),
);
append_test_event(
&session,
temp.path(),
"turn_status",
json!({"status":"cancelled", "assistant_text":"partial cancelled answer"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("cancelled prompt"));
assert!(material.contains("partial cancelled answer"));
assert!(replay.replay_diagnostics.iter().any(|diagnostic| {
diagnostic.contains("replayed_partial_historical_turn")
&& diagnostic.contains("cancelled=true")
}));
}
#[test]
fn conversation_replay_oversized_detection_counts_tool_result_original_once() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let large_but_below_uncommitted_limit = "B".repeat(120_000);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"large failed"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"large_once","name":"read","arguments":{}}),
);
append_test_event(
&session,
temp.path(),
"tool_result",
json!({"call_id":"large_once","result":{"tool_name":"read","success":true,"content":large_but_below_uncommitted_limit}}),
);
append_test_event(
&session,
temp.path(),
"diagnostic",
json!({"message":"response.failed"}),
);
append_test_event(&session, temp.path(), "user_input", json!({"text":"after"}));
let replay = build_conversation_replay(Some(&session)).unwrap();
assert!(replay.replay_diagnostics.iter().any(|diagnostic| {
diagnostic.contains("recovered_failed_historical_turn")
&& diagnostic.contains("oversized=false")
}));
}
#[test]
fn conversation_replay_preserves_successful_large_completed_tool_turn() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"large success"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"large_call","name":"read","arguments":{}}),
);
append_test_event(
&session,
temp.path(),
"tool_result",
json!({"call_id":"large_call","result":{"tool_name":"read","success":true,"content":"A".repeat(REPLAY_TOOL_RESULT_OUTPUT_CHAR_LIMIT + 100)}}),
);
append_test_event(
&session,
temp.path(),
"assistant_output",
json!({"text":"completed"}),
);
append_test_event(&session, temp.path(), "user_input", json!({"text":"next"}));
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("large success"));
assert!(material.contains("large_call"));
assert!(material.contains("historical replay compacted tool result"));
assert!(material.contains("completed"));
assert!(material.contains("next"));
}
#[test]
fn conversation_replay_compacts_oversized_historical_tool_result_output() {
let temp = TempDir::new().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"compact"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"compact_call","name":"read","arguments":{}}),
);
append_test_event(
&session,
temp.path(),
"tool_result",
json!({"call_id":"compact_call","result":{"tool_name":"read","success":true,"content":"0123456789".repeat(3_000)}}),
);
append_test_event(
&session,
temp.path(),
"assistant_output",
json!({"text":"ok"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let result = replay
.items
.iter()
.find_map(|item| match item {
ProviderConversationItem::ToolResult(result) => Some(result),
_ => None,
})
.unwrap();
assert!(
result
.output
.contains("historical replay compacted tool result")
);
assert!(result.output.chars().count() < 30_000);
assert!(!result.output.contains(&"0123456789".repeat(3_000)));
assert!(
replay
.replay_diagnostics
.iter()
.any(|d| d.contains("compacted_historical_tool_result_output"))
);
}
#[test]
fn historical_tool_output_compaction_preserves_multibyte_suffix_boundary() {
let edge_chars = REPLAY_TOOL_RESULT_OUTPUT_CHAR_LIMIT / 4;
let suffix = format!("{}TAIL", "🦀".repeat(edge_chars - 4));
let output = format!(
"{}{}",
"A".repeat(REPLAY_TOOL_RESULT_OUTPUT_CHAR_LIMIT),
suffix
);
let compacted = compact_historical_tool_output(&output, "read", true);
assert!(compacted.is_char_boundary(compacted.len()));
assert!(compacted.ends_with(&suffix));
assert!(compacted.contains("historical replay compacted tool result"));
}
#[test]
fn failed_child_recovery_prefers_authoritative_assistant_output() {
let temp = TempDir::new().unwrap();
let sessions_root = temp.path().join("sessions");
let manager = crate::sessions::SessionManager::new(sessions_root.clone());
let session = manager.create().unwrap();
let child_id = "child-replay";
let child_dir = sessions_root.join("subagents");
fs::create_dir_all(&child_dir).unwrap();
let child_path = child_dir.join(format!("{child_id}.jsonl"));
let child_events = [
SessionEvent::new(
"assistant_chunk",
child_id.to_string(),
temp.path().to_path_buf(),
json!({"text":"first "}),
),
SessionEvent::new(
"assistant_chunk",
child_id.to_string(),
temp.path().to_path_buf(),
json!({"text":"second"}),
),
SessionEvent::new(
"assistant_output",
child_id.to_string(),
temp.path().to_path_buf(),
json!({"text":"first second"}),
),
];
fs::write(
&child_path,
child_events
.iter()
.map(|event| serde_json::to_string(event).unwrap())
.collect::<Vec<_>>()
.join("\n"),
)
.unwrap();
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"parent"}),
);
append_test_event(
&session,
temp.path(),
"tool_call",
json!({"id":"subagents","name":"subagents","arguments":{}}),
);
append_test_event(
&session,
temp.path(),
"tool_result",
json!({
"call_id":"subagents",
"result": {
"tool_name":"subagents",
"success":true,
"content": serde_json::to_string(&json!({
"results":[{"status":"failed","session_id":child_id,"session_path":child_path}]
})).unwrap()
}
}),
);
append_test_event(
&session,
temp.path(),
"turn_status",
json!({"status":"failed"}),
);
append_test_event(
&session,
temp.path(),
"user_input",
json!({"text":"continue"}),
);
let replay = build_conversation_replay(Some(&session)).unwrap();
let material = replay_material(&replay);
assert!(material.contains("first second"));
assert!(!material.contains("first secondfirst second"));
}