use super::*;
pub(super) fn safe_durable_pre_failure_items(
items: &[ProviderConversationItem],
) -> Vec<ProviderConversationItem> {
let completed_call_ids = items
.iter()
.filter_map(|item| match item {
ProviderConversationItem::ToolResult(result) => Some(result.call_id.as_str()),
ProviderConversationItem::ResponseItem(value)
if value.get("type").and_then(Value::as_str) == Some("function_call_output") =>
{
value.get("call_id").and_then(Value::as_str)
}
_ => None,
})
.collect::<std::collections::HashSet<_>>();
items
.iter()
.filter_map(|item| match item {
ProviderConversationItem::ReasoningSelection { .. } => Some(item.clone()),
ProviderConversationItem::Message(message)
if matches!(
message.role,
crate::providers::MessageRole::User | crate::providers::MessageRole::Assistant
) =>
{
Some(ProviderConversationItem::Message(message.clone()))
}
ProviderConversationItem::ResponseItem(value)
if matches!(
value.get("type").and_then(Value::as_str),
Some("function_call" | "function_call_output")
) =>
{
let call_id = value.get("call_id").and_then(Value::as_str)?;
completed_call_ids
.contains(call_id)
.then(|| ProviderConversationItem::ResponseItem(value.clone()))
}
ProviderConversationItem::ToolResult(result)
if completed_call_ids.contains(result.call_id.as_str()) =>
{
Some(item.clone())
}
_ => None,
})
.collect()
}
pub(super) fn pending_user_only_items(
items: &[ProviderConversationItem],
) -> Vec<ProviderConversationItem> {
items
.iter()
.filter_map(|item| match item {
ProviderConversationItem::ReasoningSelection { .. } => Some(item.clone()),
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User =>
{
Some(ProviderConversationItem::Message(message.clone()))
}
_ => None,
})
.collect()
}
pub(super) fn abort_recovery_from_payload(payload: &Value) -> AbortRecovery {
AbortRecovery {
observed_reasoning_delta: payload
.get("observed_reasoning_delta")
.and_then(Value::as_bool)
.unwrap_or(false),
unsafe_tool_call_progress: payload
.get("unsafe_tool_call_progress")
.and_then(Value::as_bool)
.unwrap_or(false),
reasoning_preview: payload
.get("reasoning_preview")
.and_then(Value::as_str)
.map(|text| sanitize_recovery_text(text, FAILED_TURN_RECOVERY_FIELD_CHAR_LIMIT)),
partial_tool_call_summary: payload
.get("partial_tool_call_summary")
.and_then(Value::as_str)
.map(|text| sanitize_recovery_text(text, 512)),
}
}
pub(super) fn turn_needs_recovery_summary(turn: &ReplayTurn, sessions_root: Option<&Path>) -> bool {
!turn.abort_recoveries.is_empty()
|| failed_turn_tool_results(&turn.items)
.iter()
.any(|result| subagent_child_recovery_summary(&result.output, sessions_root).is_some())
|| !failed_turn_tool_calls_without_results(&turn.items).is_empty()
|| !failed_turn_legacy_notes(&turn.items).is_empty()
}
pub(super) fn turn_has_unsafe_tool_call_progress(turn: &ReplayTurn) -> bool {
turn.abort_recoveries
.iter()
.any(|recovery| recovery.unsafe_tool_call_progress)
}
pub(super) fn failed_turn_recovery_summary(
turn: &ReplayTurn,
oversized: bool,
sessions_root: Option<&Path>,
) -> Option<String> {
if turn.items.is_empty() {
return None;
}
let mut summary = String::new();
summary.push_str(
"Session recovery summary for failed historical turn; not exact structured replay.\n",
);
summary.push_str(
"Completed paired tool calls/results are replayed; unpaired or unsafe structured progress is summarized as best-effort context only.\n",
);
if !turn.abort_recoveries.is_empty() {
summary.push_str("Abort recovery facts:\n");
for recovery in &turn.abort_recoveries {
summary.push_str("- reasoning_delta_seen=");
summary.push_str(&recovery.observed_reasoning_delta.to_string());
summary.push_str(" unsafe_tool_call_progress=");
summary.push_str(&recovery.unsafe_tool_call_progress.to_string());
summary.push('\n');
if let Some(preview) = &recovery.reasoning_preview {
summary.push_str(" reasoning preview:\n");
summary.push_str(&indent_recovery_block(&sanitize_recovery_text(
preview,
FAILED_TURN_RECOVERY_FIELD_CHAR_LIMIT,
)));
}
if let Some(tool_summary) = &recovery.partial_tool_call_summary {
summary.push_str(" tool progress: ");
summary.push_str(&sanitize_recovery_text(tool_summary, 512));
summary.push('\n');
}
}
}
if let Some(last_category) = turn.provider_failure_categories.last() {
summary.push_str("Last diagnostic/error category: ");
summary.push_str(last_category);
summary.push('\n');
}
if !turn.provider_failure_categories.is_empty() {
summary.push_str("Diagnostic categories: ");
summary.push_str(&turn.provider_failure_categories.join(", "));
summary.push('\n');
}
if let Some(cwd) = &turn.cwd {
summary.push_str("CWD: ");
summary.push_str(&sanitize_recovery_text(
&cwd.display().to_string(),
FAILED_TURN_RECOVERY_FIELD_CHAR_LIMIT,
));
summary.push('\n');
}
let objectives = failed_turn_user_objectives(&turn.items);
if !objectives.is_empty() {
summary.push_str("User objective(s):\n");
for objective in objectives {
summary.push_str("- ");
summary.push_str(&sanitize_recovery_text(
&objective,
FAILED_TURN_RECOVERY_FIELD_CHAR_LIMIT,
));
summary.push('\n');
}
}
let tool_results = failed_turn_tool_results(&turn.items);
if !tool_results.is_empty() {
summary.push_str("Tool/subagent results:\n");
for result in tool_results {
summary.push_str("- tool=");
summary.push_str(&sanitize_recovery_text(&result.tool_name, 256));
summary.push_str(" success=");
summary.push_str(&result.success.to_string());
summary.push('\n');
let child_summary = subagent_child_recovery_summary(&result.output, sessions_root);
if oversized {
summary.push_str(
" output: [omitted because failed turn exceeded recovery safety limit]\n",
);
} else if child_summary.is_some() {
summary.push_str(
" output: [omitted because child session recovery snapshot is available]\n",
);
} else {
summary.push_str(" output:\n");
summary.push_str(&indent_recovery_block(&sanitize_recovery_text(
&result.output,
FAILED_TURN_RECOVERY_TOOL_OUTPUT_CHAR_LIMIT,
)));
}
if let Some(child_summary) = child_summary {
summary.push_str(&child_summary);
}
}
}
let tool_calls_without_results = failed_turn_tool_calls_without_results(&turn.items);
if !tool_calls_without_results.is_empty() {
summary.push_str("Tool calls without recovered result:\n");
for name in tool_calls_without_results {
summary.push_str("- tool=");
summary.push_str(&sanitize_recovery_text(&name, 256));
summary.push('\n');
}
}
let legacy_notes = failed_turn_legacy_notes(&turn.items);
if !legacy_notes.is_empty() {
summary.push_str("Downgraded/orphan session events:\n");
for note in legacy_notes {
summary.push_str("- ");
summary.push_str(&sanitize_recovery_text(
¬e,
FAILED_TURN_RECOVERY_FIELD_CHAR_LIMIT,
));
summary.push('\n');
}
}
summary.push_str(
"Pending next action: Continue from this recovered context and the later user prompt.\n",
);
let summary = sanitize_recovery_text(&summary, FAILED_TURN_RECOVERY_SUMMARY_CHAR_LIMIT);
if summary.trim().is_empty() {
None
} else {
Some(summary)
}
}
fn subagent_child_recovery_summary(output: &str, sessions_root: Option<&Path>) -> Option<String> {
let sessions_root = sessions_root?;
let value: Value = serde_json::from_str(output).ok()?;
let results = value.get("results")?.as_array()?;
let mut summary = String::new();
for result in results {
if result.get("status").and_then(Value::as_str) != Some("failed") {
continue;
}
let Some(child) = child_session_file_recovery_snapshot(result, sessions_root) else {
continue;
};
summary.push_str(" child session recovery snapshot:\n");
summary.push_str(&indent_recovery_block(&child));
}
(!summary.trim().is_empty()).then_some(summary)
}
fn child_session_file_recovery_snapshot(result: &Value, sessions_root: &Path) -> Option<String> {
let id = result
.get("session_id")
.and_then(Value::as_str)
.or_else(|| {
result
.get("session_path")
.and_then(Value::as_str)
.and_then(|path| Path::new(path).file_stem())
.and_then(std::ffi::OsStr::to_str)
})?;
crate::sessions::validate_session_id(id.to_string()).ok()?;
let path = sessions_root.join("subagents").join(format!("{id}.jsonl"));
let trusted_root = fs::canonicalize(sessions_root.join("subagents")).ok()?;
let trusted_path = fs::canonicalize(&path).ok()?;
if !trusted_path.starts_with(&trusted_root) {
return None;
}
let persisted_path = result.get("session_path").and_then(Value::as_str)?;
let persisted_path = fs::canonicalize(persisted_path).ok()?;
if persisted_path != trusted_path {
return None;
}
let lines = child_session_tail_lines(
&trusted_path,
CHILD_SESSION_RECOVERY_MAX_LINES,
CHILD_SESSION_RECOVERY_MAX_BYTES,
)?;
let mut assistant_chunks = String::new();
let mut authoritative_output = None;
let mut abort_recoveries = Vec::new();
for line in lines {
let Ok(event) = serde_json::from_slice::<SessionEvent>(&line) else {
continue;
};
if event.session_id != id {
continue;
}
match event.kind() {
Some(SessionEventKind::AssistantChunk) => {
if let Some(text) = event.payload.get("text").and_then(Value::as_str)
&& !text.trim().is_empty()
{
assistant_chunks.push_str(text);
}
}
Some(SessionEventKind::AssistantOutput) => {
if let Some(text) = event.payload.get("text").and_then(Value::as_str)
&& !text.trim().is_empty()
{
authoritative_output = Some(text.to_string());
}
}
Some(SessionEventKind::AbortRecovery) => {
abort_recoveries.push(event.payload.to_string());
}
_ => {}
}
}
let mut snapshot = authoritative_output.unwrap_or(assistant_chunks);
for recovery in abort_recoveries {
if !snapshot.is_empty() {
snapshot.push('\n');
}
snapshot.push_str("abort recovery: ");
snapshot.push_str(&recovery);
}
let snapshot = sanitize_recovery_text(&snapshot, FAILED_TURN_RECOVERY_FIELD_CHAR_LIMIT);
(!snapshot.trim().is_empty()).then_some(snapshot)
}
fn child_session_tail_lines(
path: &Path,
max_lines: usize,
max_bytes: usize,
) -> Option<Vec<Vec<u8>>> {
if max_lines == 0 || max_bytes == 0 {
return Some(Vec::new());
}
let mut file = fs::File::open(path).ok()?;
let file_len = file.metadata().ok()?.len();
let max_bytes_u64 = u64::try_from(max_bytes).ok()?;
let start = file_len.saturating_sub(max_bytes_u64);
file.seek(SeekFrom::Start(start)).ok()?;
let mut bytes = Vec::with_capacity(usize::try_from(file_len - start).ok()?);
file.take(max_bytes_u64).read_to_end(&mut bytes).ok()?;
let lines = if start > 0 {
bytes
.iter()
.position(|byte| *byte == b'\n')
.map_or(&[][..], |index| &bytes[index + 1..])
} else {
bytes.as_slice()
};
Some(
lines
.split_inclusive(|byte| *byte == b'\n')
.map(|line| line.strip_suffix(b"\n").unwrap_or(line).to_vec())
.rev()
.take(max_lines)
.collect::<Vec<_>>()
.into_iter()
.rev()
.collect(),
)
}
fn failed_turn_user_objectives(items: &[ProviderConversationItem]) -> Vec<String> {
items
.iter()
.filter_map(|item| match item {
ProviderConversationItem::Message(message)
if message.role == crate::providers::MessageRole::User =>
{
Some(message.content.clone())
}
_ => None,
})
.collect()
}
fn failed_turn_tool_results(items: &[ProviderConversationItem]) -> Vec<&ProviderToolResult> {
items
.iter()
.filter_map(|item| match item {
ProviderConversationItem::ToolResult(result) => Some(result),
_ => None,
})
.collect()
}
fn failed_turn_tool_calls_without_results(items: &[ProviderConversationItem]) -> Vec<String> {
let result_call_ids = items
.iter()
.filter_map(|item| match item {
ProviderConversationItem::ToolResult(result) => Some(result.call_id.as_str()),
_ => None,
})
.collect::<std::collections::HashSet<_>>();
items
.iter()
.filter_map(|item| match item {
ProviderConversationItem::ResponseItem(value)
if value.get("type").and_then(Value::as_str) == Some("function_call") =>
{
let call_id = value
.get("call_id")
.and_then(Value::as_str)
.unwrap_or_default();
if result_call_ids.contains(call_id) {
None
} else {
value
.get("name")
.and_then(Value::as_str)
.map(ToString::to_string)
}
}
_ => None,
})
.collect()
}
fn failed_turn_legacy_notes(items: &[ProviderConversationItem]) -> Vec<String> {
items
.iter()
.filter_map(|item| match item {
ProviderConversationItem::LegacyReplayNote {
event_type,
content,
} => Some(format!("{event_type}: {content}")),
_ => None,
})
.collect()
}
fn indent_recovery_block(text: &str) -> String {
let mut indented = text
.lines()
.map(|line| format!(" {line}"))
.collect::<Vec<_>>()
.join("\n");
indented.push('\n');
indented
}
pub(super) fn sanitize_recovery_text(text: &str, char_limit: usize) -> String {
let normalized = text
.chars()
.map(|ch| {
if ch.is_control() && !matches!(ch, '\n' | '\t') {
' '
} else {
ch
}
})
.collect::<String>();
let redacted = redact_sensitive_text(normalized.trim());
limit_recovery_chars(&redacted, char_limit)
}
fn limit_recovery_chars(text: &str, char_limit: usize) -> String {
let char_count = text.chars().count();
if char_count <= char_limit {
return text.to_string();
}
let head = char_limit.saturating_sub(160);
let prefix: String = text.chars().take(head).collect();
format!(
"{prefix}\n[recovery summary truncated: original_chars={char_count} limit={char_limit}]"
)
}