mod builder;
mod recovery;
use builder::AbortRecovery;
use builder::ConversationReplayBuilder;
use builder::ReplayTurn;
use recovery::abort_recovery_from_payload;
use recovery::failed_turn_recovery_summary;
use recovery::pending_user_only_items;
use recovery::safe_durable_pre_failure_items;
use recovery::turn_has_unsafe_tool_call_progress;
use recovery::turn_needs_recovery_summary;
use recovery::sanitize_recovery_text;
use crate::{
agent::{
recovery::{DANGLING_TOOL_INTENT_PROMPT, DANGLING_TOOL_INTENT_REASON},
steering::STEERING_HEADER,
},
output::redact_sensitive_text,
providers::{ChatMessage, ProviderConversationItem, ProviderToolResult},
sessions::{
BoundedReadError, Session, SessionEvent, SessionEventKind, SessionReadDiagnostic,
SessionReadStats, TurnStatus, latest_valid_compaction_checkpoint_for_replay,
},
tools::skill_provenance::validated_skill_reads,
};
use serde_json::{Value, json};
use std::{
fs,
io::{Read, Seek, SeekFrom},
path::{Path, PathBuf},
};
use super::{canonical_json, code_mode, event_line, is_local_only_session_event};
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ConversationReplay {
pub items: Vec<ProviderConversationItem>,
pub legacy_lossy_events: Vec<String>,
pub session_read_diagnostics: Vec<SessionReadDiagnostic>,
pub replay_diagnostics: Vec<String>,
pub warnings: Vec<String>,
pub(crate) task_scope: Result<Option<crate::sessions::task_scope::TaskScope>, String>,
}
pub(super) struct ConversationReplayLoad {
pub replay: ConversationReplay,
pub read: SessionReadStats,
pub events_parsed: usize,
pub events_before_cutoff: usize,
pub events_after_cutoff: usize,
}
pub(crate) const REPLAY_TOOL_RESULT_OUTPUT_CHAR_LIMIT: usize = 24_000;
pub(crate) const REPLAY_JSONL_MAX_LINES: usize = 100_000;
pub(crate) const REPLAY_JSONL_MAX_BYTES: usize = 64 * 1024 * 1024;
const UNCOMMITTED_TURN_PROVIDER_VISIBLE_CHAR_LIMIT: usize = 128_000;
const FAILED_TURN_RECOVERY_SUMMARY_CHAR_LIMIT: usize = 12_000;
const FAILED_TURN_RECOVERY_FIELD_CHAR_LIMIT: usize = 4_000;
const FAILED_TURN_RECOVERY_TOOL_OUTPUT_CHAR_LIMIT: usize = 3_000;
const CHILD_SESSION_RECOVERY_MAX_LINES: usize = 200;
const CHILD_SESSION_RECOVERY_MAX_BYTES: usize = 128 * 1024;
const MAX_SESSION_ID_MISMATCH_DIAGNOSTICS: usize = 64;
fn sanitize_provider_response_item(mut item: Value) -> Value {
fn remove_encrypted_content(value: &mut Value) {
match value {
Value::Object(object) => {
object.remove("encrypted_content");
for child in object.values_mut() {
remove_encrypted_content(child);
}
}
Value::Array(items) => {
for item in items {
remove_encrypted_content(item);
}
}
_ => {}
}
}
remove_encrypted_content(&mut item);
item
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct ConversationReplayLimits {
pub(crate) max_lines: usize,
pub(crate) max_bytes: usize,
}
impl Default for ConversationReplayLimits {
fn default() -> Self {
Self {
max_lines: REPLAY_JSONL_MAX_LINES,
max_bytes: REPLAY_JSONL_MAX_BYTES,
}
}
}
pub(crate) fn build_conversation_replay(
session: Option<&Session>,
) -> anyhow::Result<ConversationReplay> {
build_conversation_replay_with_limits(session, ConversationReplayLimits::default())
}
pub(crate) fn build_conversation_replay_with_limits(
session: Option<&Session>,
limits: ConversationReplayLimits,
) -> anyhow::Result<ConversationReplay> {
let Some(session) = session else {
return Ok(empty_conversation_replay());
};
Ok(build_conversation_replay_with_limits_and_stats(session, limits)?.replay)
}
fn empty_conversation_replay() -> ConversationReplay {
ConversationReplay {
items: Vec::new(),
legacy_lossy_events: Vec::new(),
session_read_diagnostics: Vec::new(),
replay_diagnostics: Vec::new(),
warnings: Vec::new(),
task_scope: Ok(None),
}
}
pub(super) fn build_conversation_replay_with_read_stats(
session: &Session,
) -> anyhow::Result<ConversationReplayLoad> {
build_conversation_replay_with_limits_and_stats(session, ConversationReplayLimits::default())
}
fn build_conversation_replay_with_limits_and_stats(
session: &Session,
limits: ConversationReplayLimits,
) -> anyhow::Result<ConversationReplayLoad> {
let (tolerant, read_stats) = session
.read_events_tolerant_bounded_with_stats(limits.max_lines, limits.max_bytes)
.map_err(map_replay_read_error)?;
let events_parsed = tolerant.events.len();
let replay_diagnostics = session_id_mismatch_diagnostics(session.id(), &tolerant.events);
let start_index = latest_valid_compaction_checkpoint_for_replay(session.id(), &tolerant.events)
.0
.map(|checkpoint| checkpoint.cutoff_event_count)
.unwrap_or(0)
.min(events_parsed);
let sessions_root = session.path().parent().map(Path::to_path_buf);
let mut replay = build_conversation_replay_from_events_with_diagnostics(
session.id(),
&tolerant.events,
replay_diagnostics,
sessions_root.as_deref(),
);
replay.session_read_diagnostics = tolerant.diagnostics;
if !replay.session_read_diagnostics.is_empty()
&& replay.task_scope.as_ref().is_ok_and(Option::is_some)
{
replay.task_scope =
Err("Unreadable delegated session control history; execution blocked".to_string());
}
Ok(ConversationReplayLoad {
replay,
read: read_stats,
events_parsed,
events_before_cutoff: start_index,
events_after_cutoff: events_parsed.saturating_sub(start_index),
})
}
pub(super) fn map_replay_read_error(error: anyhow::Error) -> anyhow::Error {
if error.downcast_ref::<BoundedReadError>().is_some() {
anyhow::anyhow!(
"session replay exceeded bounded history budget; durable session history was preserved; start /new and inspect or back up session data"
)
} else {
error
}
}
pub(crate) fn build_conversation_replay_from_events(
session_id: &str,
events: &[SessionEvent],
) -> ConversationReplay {
build_conversation_replay_from_events_with_sessions_root(session_id, events, None)
}
fn build_conversation_replay_from_events_with_sessions_root(
session_id: &str,
events: &[SessionEvent],
sessions_root: Option<&Path>,
) -> ConversationReplay {
let replay_diagnostics = session_id_mismatch_diagnostics(session_id, events);
build_conversation_replay_from_events_with_diagnostics(
session_id,
events,
replay_diagnostics,
sessions_root,
)
}
fn session_id_mismatch_diagnostics(session_id: &str, events: &[SessionEvent]) -> Vec<String> {
let mut diagnostics = Vec::new();
let mut omitted_diagnostics = 0usize;
for (event_index, event) in events.iter().enumerate() {
if event.session_id == session_id {
continue;
}
if diagnostics.len() < MAX_SESSION_ID_MISMATCH_DIAGNOSTICS {
diagnostics.push(format!(
"skipped session replay event at event_index={event_index}: session_id_mismatch"
));
} else {
omitted_diagnostics = omitted_diagnostics.saturating_add(1);
}
}
if omitted_diagnostics != 0 {
if diagnostics.len() == MAX_SESSION_ID_MISMATCH_DIAGNOSTICS {
diagnostics.pop();
omitted_diagnostics = omitted_diagnostics.saturating_add(1);
}
diagnostics.push(format!(
"omitted {omitted_diagnostics} additional session_id_mismatch diagnostics after cap of {MAX_SESSION_ID_MISMATCH_DIAGNOSTICS}"
));
}
diagnostics
}
fn build_conversation_replay_from_events_with_diagnostics(
session_id: &str,
events: &[SessionEvent],
mut replay_diagnostics: Vec<String>,
sessions_root: Option<&Path>,
) -> ConversationReplay {
let mut builder = ConversationReplayBuilder {
sessions_root: sessions_root.map(Path::to_path_buf),
..ConversationReplayBuilder::default()
};
let (checkpoint, checkpoint_diagnostics) =
latest_valid_compaction_checkpoint_for_replay(session_id, events);
replay_diagnostics.extend(checkpoint_diagnostics);
let start_index = checkpoint
.as_ref()
.map(|checkpoint| checkpoint.cutoff_event_count)
.unwrap_or(0);
if let Some(checkpoint) = checkpoint {
builder.push_compaction_summary(&checkpoint.summary);
if let Some(note) = crate::sessions::effects::ExecutionEffects::from_checkpoint(
events[checkpoint.event_index]
.payload
.get("execution_effects"),
)
.provider_note()
{
builder
.items
.push(ProviderConversationItem::Message(ChatMessage::user(note)));
}
}
let assessments = crate::protection::prompt_injection::RecordedAssessments::new(events);
for (index, event) in events.iter().enumerate().skip(start_index) {
if event.session_id == session_id {
builder.push_event(event, assessments.for_result(index));
}
}
builder.finish();
replay_diagnostics.extend(builder.replay_diagnostics);
ConversationReplay {
items: builder.items,
legacy_lossy_events: builder.legacy_lossy_events,
session_read_diagnostics: Vec::new(),
replay_diagnostics,
warnings: builder.warnings,
task_scope: crate::sessions::task_scope::from_events(session_id, events),
}
}
fn auto_recovery_continue(event: &SessionEvent) -> bool {
if event.payload.get("auto_recovery").and_then(Value::as_bool) != Some(true) {
return false;
}
match event.payload.get("reason").and_then(Value::as_str) {
Some("incomplete_semantic_progress_timeout") => {
event.payload.get("text").and_then(Value::as_str) == Some("Continue")
}
Some(DANGLING_TOOL_INTENT_REASON) => {
event.payload.get("text").and_then(Value::as_str) == Some(DANGLING_TOOL_INTENT_PROMPT)
}
_ => false,
}
}
fn arguments_json_text(arguments: &Value) -> String {
match arguments {
Value::String(text) => text.clone(),
value => value.to_string(),
}
}
fn provider_response_item_is_success_completion(item: &Value) -> bool {
let item_type = item.get("type").and_then(Value::as_str);
let status = item.get("status").and_then(Value::as_str);
matches!(
status,
Some("completed" | "complete" | "succeeded" | "success")
) && matches!(
item_type,
Some("message" | "output_text" | "assistant_message")
)
}
fn diagnostic_provider_failure_category(event: &SessionEvent) -> Option<&'static str> {
if event.kind() != Some(SessionEventKind::Diagnostic) {
return None;
}
let mut text = event.payload.to_string();
text.make_ascii_lowercase();
if [
"request or response body error",
"request body",
"response body",
"failed to read provider error body",
]
.iter()
.any(|needle| text.contains(needle))
{
return Some("provider_body_error");
}
if text.contains("response.failed") || text.contains("failed or incomplete response") {
return Some("provider_response_failed");
}
if text.contains("response.incomplete") || text.contains("incomplete response") {
return Some("provider_response_incomplete");
}
if text.contains("provider stream no semantic progress before timeout")
|| text.contains("provider stream idle timeout")
{
return Some("provider_stream_timeout");
}
if text.contains("provider stream ended") {
return Some("provider_stream_ended");
}
None
}
fn replay_item_visible_chars(item: &ProviderConversationItem) -> usize {
match item {
ProviderConversationItem::ReasoningSelection { .. } => 0,
ProviderConversationItem::Message(message) => message.content.chars().count(),
ProviderConversationItem::ResponseItem(value) => canonical_json(value).chars().count(),
ProviderConversationItem::ToolResult(result) => result.output.chars().count(),
ProviderConversationItem::LegacyReplayNote { content, .. } => content.chars().count(),
}
}
pub(crate) fn compact_historical_tool_output(
output: &str,
tool_name: &str,
success: bool,
) -> String {
let original_chars = output.chars().count();
if original_chars <= REPLAY_TOOL_RESULT_OUTPUT_CHAR_LIMIT {
return output.to_string();
}
let edge_chars = REPLAY_TOOL_RESULT_OUTPUT_CHAR_LIMIT / 4;
let prefix: String = output.chars().take(edge_chars).collect();
let suffix_start = output
.char_indices()
.nth_back(edge_chars.saturating_sub(1))
.map(|(index, _)| index)
.unwrap_or(output.len());
let suffix = &output[suffix_start..];
format!(
"[historical replay compacted tool result: tool={tool_name} success={success} original_chars={original_chars} limit={REPLAY_TOOL_RESULT_OUTPUT_CHAR_LIMIT}]\n{prefix}\n[... historical replay output omitted ...]\n{suffix}"
)
}