use super::*;
#[derive(Default)]
pub(super) struct ConversationReplayBuilder {
pub(super) items: Vec<ProviderConversationItem>,
pub(super) legacy_lossy_events: Vec<String>,
pub(super) replay_diagnostics: Vec<String>,
pub(super) warnings: Vec<String>,
pub(super) pending_assistant: String,
pub(super) assistant_segment_start: usize,
pub(super) exact_provider_call_ids: std::collections::HashSet<String>,
pub(super) pending_tool_call_ids: std::collections::HashSet<String>,
pub(super) active_turn: Option<ReplayTurn>,
pub(super) sessions_root: Option<PathBuf>,
}
#[derive(Default)]
pub(super) struct ReplayTurn {
pub(super) items: Vec<ProviderConversationItem>,
saw_completed_assistant: bool,
saw_provider_success_completion: bool,
saw_provider_failure_diagnostic: bool,
saw_turn_incomplete_marker: bool,
saw_turn_cancelled_marker: bool,
saw_turn_failed_marker: bool,
saw_turn_compaction_marker: bool,
saw_turn_complete_marker: bool,
pub(super) provider_failure_categories: Vec<&'static str>,
pub(super) abort_recoveries: Vec<AbortRecovery>,
provider_visible_chars: usize,
original_provider_visible_chars: usize,
provider_activity_after_user: bool,
pub(super) cwd: Option<PathBuf>,
tool_call_ids: std::collections::HashSet<String>,
code_mode: super::code_mode::CodeModeRecovery,
}
#[derive(Default)]
pub(super) struct AbortRecovery {
pub(super) observed_reasoning_delta: bool,
pub(super) unsafe_tool_call_progress: bool,
pub(super) reasoning_preview: Option<String>,
pub(super) partial_tool_call_summary: Option<String>,
}
impl ConversationReplayBuilder {
pub(super) fn push_compaction_summary(&mut self, summary: &str) {
self.items.push(ProviderConversationItem::Message(ChatMessage::user(
format!("Session compaction summary. Earlier session events are represented only by this summary:\n\n{summary}"),
)));
}
pub(super) fn push_event(&mut self, event: &SessionEvent, assessment: &Value) {
if let Some(turn) = self.active_turn.as_mut() {
turn.code_mode.observe(event);
}
match event.kind() {
Some(SessionEventKind::AssistantChunk) => {
if let Some(text) = event.payload.get("text").and_then(Value::as_str) {
if !text.trim().is_empty() {
self.pending_assistant.push_str(text);
}
} else {
self.push_legacy_note(event);
}
}
Some(SessionEventKind::AssistantOutput) => {
if let Some(text) = event.payload.get("text").and_then(Value::as_str) {
if let Some(turn) = self.active_turn.as_mut() {
turn.saw_completed_assistant = true;
}
self.reconcile_assistant_text(text, self.assistant_segment_start, false);
self.start_assistant_segment();
} else {
self.push_legacy_note(event);
}
}
Some(SessionEventKind::UserInput) => {
self.flush_assistant_chunks();
self.assistant_segment_start = 0;
if self.user_input_continues_active_turn(event) {
self.start_assistant_segment();
if let Some(text) = event.payload.get("text").and_then(Value::as_str) {
self.push_message(ChatMessage::user(text));
} else {
self.push_legacy_note(event);
}
return;
}
self.finalize_active_turn(false);
self.active_turn = Some(ReplayTurn {
cwd: Some(event.cwd.clone()),
..ReplayTurn::default()
});
if let Some(selection) =
ProviderConversationItem::reasoning_selection_from_payload(&event.payload)
{
self.push_replay_item(selection);
}
if let Some(text) = event.payload.get("text").and_then(Value::as_str) {
self.push_message(ChatMessage::user(text));
} else {
self.push_legacy_note(event);
}
}
Some(SessionEventKind::ProviderResponseItem) => {
self.flush_assistant_chunks();
if let Some(item) = event
.payload
.get("item")
.cloned()
.map(sanitize_provider_response_item)
{
let item_type = item.get("type").and_then(Value::as_str);
if item_type == Some("function_call")
&& let Some(call_id) = item.get("call_id").and_then(Value::as_str)
{
self.register_tool_call_id(call_id);
}
if item_type == Some("function_call_output") {
let Some(call_id) = item.get("call_id").and_then(Value::as_str) else {
self.push_legacy_note(event);
return;
};
if !self.pending_tool_call_ids.remove(call_id) {
self.replay_diagnostics.push(format!(
"downgraded_orphan_provider_tool_result call_id={}",
sanitize_recovery_text(call_id, 256)
));
self.push_legacy_note(event);
return;
}
self.start_assistant_segment();
}
if provider_response_item_is_success_completion(&item)
&& let Some(turn) = self.active_turn.as_mut()
{
turn.saw_provider_success_completion = true;
}
self.push_replay_item(ProviderConversationItem::ResponseItem(item));
} else {
self.push_legacy_note(event);
}
}
Some(SessionEventKind::ToolCall) => {
self.flush_assistant_chunks();
self.push_tool_call(event);
}
Some(SessionEventKind::ToolResult) => {
self.flush_assistant_chunks();
self.push_tool_result(event, assessment);
self.start_assistant_segment();
}
Some(SessionEventKind::ProviderContextItem) => {
self.flush_assistant_chunks();
self.push_provider_context_item(event);
}
Some(SessionEventKind::Rewind) => {
self.flush_assistant_chunks();
if let Some(content) = crate::checkpoints::replay_rewind_context(&event.payload) {
self.push_message(ChatMessage::user(content));
} else {
self.push_legacy_note(event);
}
}
Some(SessionEventKind::Compaction | SessionEventKind::ReasoningSummary) => {}
Some(SessionEventKind::AbortRecovery) => {
if let Some(turn) = self.active_turn.as_mut() {
turn.abort_recoveries
.push(abort_recovery_from_payload(&event.payload));
}
}
Some(SessionEventKind::SkillSuggestion) => {
if self.active_turn.is_some()
&& let Some(hint) = event.payload.get("hint").and_then(Value::as_str)
&& !hint.trim().is_empty()
{
self.push_message(ChatMessage::user(hint));
}
}
Some(SessionEventKind::TurnStatus) => {
if let Some(payload) = event.turn_status_payload()
&& let Some(turn) = self.active_turn.as_mut()
{
let restore_assistant_text = match payload.status {
TurnStatus::Incomplete => {
turn.saw_turn_incomplete_marker = true;
false
}
TurnStatus::Cancelled => {
turn.saw_turn_cancelled_marker = true;
true
}
TurnStatus::Failed => {
turn.saw_turn_failed_marker = true;
true
}
TurnStatus::CompactionRequired => {
turn.saw_turn_compaction_marker = true;
true
}
TurnStatus::Complete => {
turn.saw_turn_complete_marker = true;
true
}
};
if restore_assistant_text
&& let Some(text) = payload.assistant_text.as_deref()
&& !text.trim().is_empty()
{
self.reconcile_assistant_text(text, 0, true);
}
}
}
_ if is_local_only_session_event(event) => {
if let Some(category) = diagnostic_provider_failure_category(event)
&& let Some(turn) = self.active_turn.as_mut()
{
turn.saw_provider_failure_diagnostic = true;
if !turn.provider_failure_categories.contains(&category) {
turn.provider_failure_categories.push(category);
}
}
}
_ => {
self.flush_assistant_chunks();
self.push_legacy_note(event);
}
}
}
pub(super) fn finish(&mut self) {
self.flush_assistant_chunks();
self.finalize_active_turn(true);
}
fn user_input_continues_active_turn(&self, event: &SessionEvent) -> bool {
let Some(text) = event.payload.get("text").and_then(Value::as_str) else {
return false;
};
let is_steering = text.starts_with(STEERING_HEADER);
if !(is_steering || auto_recovery_continue(event)) {
return false;
}
self.active_turn.as_ref().is_some_and(|turn| {
turn.provider_activity_after_user
&& (is_steering || !turn.saw_completed_assistant)
&& !turn.saw_provider_success_completion
&& !turn.saw_provider_failure_diagnostic
&& !turn.saw_turn_incomplete_marker
&& !turn.saw_turn_cancelled_marker
&& !turn.saw_turn_failed_marker
&& !turn.saw_turn_compaction_marker
&& !turn.saw_turn_complete_marker
})
}
fn finalize_active_turn(&mut self, eof: bool) {
let Some(turn) = self.active_turn.take() else {
return;
};
for call_id in &turn.tool_call_ids {
self.pending_tool_call_ids.remove(call_id);
self.exact_provider_call_ids.remove(call_id);
}
let turn_incomplete = turn.saw_turn_incomplete_marker;
let turn_partial_terminal = turn.saw_turn_cancelled_marker
|| turn.saw_turn_failed_marker
|| turn.saw_turn_compaction_marker;
let committed = turn.saw_completed_assistant
|| turn.saw_provider_success_completion
|| turn.saw_turn_complete_marker;
let oversized =
turn.original_provider_visible_chars > UNCOMMITTED_TURN_PROVIDER_VISIBLE_CHAR_LIMIT;
let unsafe_tool_call_repair_failure = turn_has_unsafe_tool_call_progress(&turn);
let nested_recovery = turn
.code_mode
.pending_items(committed || turn_partial_terminal);
if committed {
self.items.extend(turn.items);
} else if turn.provider_activity_after_user {
if turn_partial_terminal {
self.items
.extend(safe_durable_pre_failure_items(&turn.items));
if unsafe_tool_call_repair_failure {
self.replay_diagnostics.push(format!(
"replayed_partial_historical_turn_without_recovery_summary unsafe_tool_call_progress=true cancelled={} failed={} chars={}",
turn.saw_turn_cancelled_marker,
turn.saw_turn_failed_marker,
turn.original_provider_visible_chars
));
}
if !unsafe_tool_call_repair_failure
&& turn_needs_recovery_summary(&turn, self.sessions_root.as_deref())
&& let Some(summary) = failed_turn_recovery_summary(
&turn,
oversized,
self.sessions_root.as_deref(),
)
{
let summary_chars = summary.chars().count();
self.items
.push(ProviderConversationItem::Message(ChatMessage::user(
summary,
)));
self.replay_diagnostics.push(format!(
"recovered_partial_historical_turn cancelled={} failed={} chars={} summary_chars={}",
turn.saw_turn_cancelled_marker,
turn.saw_turn_failed_marker,
turn.original_provider_visible_chars,
summary_chars
));
}
self.replay_diagnostics.push(format!(
"replayed_partial_historical_turn cancelled={} failed={} chars={}",
turn.saw_turn_cancelled_marker,
turn.saw_turn_failed_marker,
turn.original_provider_visible_chars
));
} else if turn.saw_provider_failure_diagnostic {
if let Some(summary) =
failed_turn_recovery_summary(&turn, oversized, self.sessions_root.as_deref())
{
let summary_chars = summary.chars().count();
self.items
.push(ProviderConversationItem::Message(ChatMessage::user(
summary,
)));
self.replay_diagnostics.push(format!(
"recovered_failed_historical_turn provider_failure=true oversized={} chars={} limit={} summary_chars={}",
oversized,
turn.original_provider_visible_chars,
UNCOMMITTED_TURN_PROVIDER_VISIBLE_CHAR_LIMIT,
summary_chars
));
} else {
self.replay_diagnostics.push(format!(
"skipped_failed_historical_turn provider_failure=true oversized={} chars={} limit={}",
oversized,
turn.original_provider_visible_chars,
UNCOMMITTED_TURN_PROVIDER_VISIBLE_CHAR_LIMIT
));
}
} else if oversized {
self.replay_diagnostics.push(format!(
"skipped_failed_historical_turn provider_failure=false oversized=true chars={} limit={}",
turn.original_provider_visible_chars,
UNCOMMITTED_TURN_PROVIDER_VISIBLE_CHAR_LIMIT
));
} else if eof {
self.replay_diagnostics.push(format!(
"dropped_uncommitted_eof_provider_activity chars={}",
turn.original_provider_visible_chars
));
}
if eof
&& !turn.saw_provider_failure_diagnostic
&& !oversized
&& !turn_incomplete
&& !turn_partial_terminal
{
self.items.extend(pending_user_only_items(&turn.items));
} else if turn_incomplete {
self.replay_diagnostics.push(format!(
"dropped_incomplete_historical_turn provider_activity_after_user=true chars={}",
turn.original_provider_visible_chars
));
}
} else if turn_incomplete {
self.replay_diagnostics
.push("dropped_incomplete_historical_user_only_turn".to_string());
} else if turn_partial_terminal && !turn.abort_recoveries.is_empty() {
if turn_has_unsafe_tool_call_progress(&turn) {
self.items.extend(pending_user_only_items(&turn.items));
self.replay_diagnostics.push(format!(
"replayed_partial_historical_user_only_turn_without_recovery_summary unsafe_tool_call_progress=true cancelled={} failed={}",
turn.saw_turn_cancelled_marker, turn.saw_turn_failed_marker
));
} else if let Some(summary) =
failed_turn_recovery_summary(&turn, oversized, self.sessions_root.as_deref())
{
self.items
.push(ProviderConversationItem::Message(ChatMessage::user(
summary,
)));
self.replay_diagnostics.push(format!(
"recovered_partial_historical_user_only_turn cancelled={} failed={}",
turn.saw_turn_cancelled_marker, turn.saw_turn_failed_marker
));
}
} else if turn.saw_provider_failure_diagnostic || oversized {
self.replay_diagnostics.push(format!(
"skipped_failed_historical_user_only_turn provider_failure={} oversized={} chars={} limit={}",
turn.saw_provider_failure_diagnostic,
oversized,
turn.original_provider_visible_chars,
UNCOMMITTED_TURN_PROVIDER_VISIBLE_CHAR_LIMIT
));
} else {
self.items.extend(turn.items);
}
self.items.extend(nested_recovery);
}
fn push_message(&mut self, message: ChatMessage) {
self.push_replay_item(ProviderConversationItem::Message(message));
}
fn push_replay_item(&mut self, item: ProviderConversationItem) {
if let Some(turn) = self.active_turn.as_mut() {
if !matches!(
item,
ProviderConversationItem::Message(ChatMessage {
role: crate::providers::MessageRole::User,
..
}) | ProviderConversationItem::ReasoningSelection { .. }
) {
turn.provider_activity_after_user = true;
}
let visible_chars = replay_item_visible_chars(&item);
turn.provider_visible_chars = turn.provider_visible_chars.saturating_add(visible_chars);
turn.original_provider_visible_chars = turn
.original_provider_visible_chars
.saturating_add(visible_chars);
turn.items.push(item);
} else {
self.items.push(item);
}
}
fn start_assistant_segment(&mut self) {
self.assistant_segment_start = self
.active_turn
.as_ref()
.map_or(self.items.len(), |turn| turn.items.len());
}
fn reconcile_assistant_text(&mut self, text: &str, start: usize, cumulative: bool) {
if text.trim().is_empty() {
return;
}
self.flush_assistant_chunks();
let items = self
.active_turn
.as_mut()
.map_or(&mut self.items, |turn| &mut turn.items);
let mut remaining = text;
let item_count = items.len();
for (index, item) in items.iter_mut().enumerate().skip(start) {
if let ProviderConversationItem::Message(message) = item
&& message.role == crate::providers::MessageRole::Assistant
{
let Some(suffix) = remaining.strip_prefix(&message.content) else {
return;
};
if index + 1 == item_count {
message.content = remaining.to_string();
remaining = "";
} else {
remaining = if cumulative {
suffix.strip_prefix("\n\n").unwrap_or(suffix)
} else {
suffix
};
}
}
}
if let Some(turn) = self.active_turn.as_mut() {
turn.provider_visible_chars = turn.items.iter().map(replay_item_visible_chars).sum();
turn.original_provider_visible_chars = turn.provider_visible_chars;
}
self.pending_assistant = remaining.to_string();
self.flush_assistant_chunks();
}
fn flush_assistant_chunks(&mut self) {
if self.pending_assistant.trim().is_empty() {
self.pending_assistant.clear();
return;
}
let text = std::mem::take(&mut self.pending_assistant);
self.push_message(ChatMessage::assistant(text));
}
fn register_tool_call_id(&mut self, call_id: &str) {
self.exact_provider_call_ids.insert(call_id.to_string());
self.pending_tool_call_ids.insert(call_id.to_string());
if let Some(turn) = self.active_turn.as_mut() {
turn.tool_call_ids.insert(call_id.to_string());
}
}
fn push_tool_call(&mut self, event: &SessionEvent) {
let Some(call_id) = event.payload.get("id").and_then(Value::as_str) else {
self.push_legacy_note(event);
return;
};
if call_id.trim().is_empty() || self.exact_provider_call_ids.contains(call_id) {
return;
}
let Some(name) = event.payload.get("name").and_then(Value::as_str) else {
self.push_legacy_note(event);
return;
};
let arguments = event
.payload
.get("arguments")
.cloned()
.unwrap_or(Value::Null);
self.register_tool_call_id(call_id);
self.push_replay_item(ProviderConversationItem::ResponseItem(json!({
"type": "function_call",
"call_id": call_id,
"name": name,
"arguments": arguments_json_text(&arguments),
"status": "completed",
})));
}
fn push_tool_result(&mut self, event: &SessionEvent, assessment: &Value) {
let Some(call_id) = event.payload.get("call_id").and_then(Value::as_str) else {
self.push_legacy_note(event);
return;
};
if call_id.trim().is_empty() {
self.push_legacy_note(event);
return;
}
if !self.pending_tool_call_ids.remove(call_id) {
self.replay_diagnostics.push(format!(
"downgraded_orphan_session_tool_result call_id={}",
sanitize_recovery_text(call_id, 256)
));
self.push_legacy_note(event);
return;
}
let result = &event.payload["result"];
let output = result
.get("content")
.and_then(Value::as_str)
.or_else(|| result.get("output").and_then(Value::as_str));
let Some(output) = output else {
self.push_legacy_note(event);
return;
};
let tool_name = result
.get("tool_name")
.and_then(Value::as_str)
.unwrap_or("unknown")
.to_string();
let success = result
.get("success")
.and_then(Value::as_bool)
.unwrap_or(false);
if assessment.is_null() {
const WARNING: &str = "History contains tool output without recorded prompt-injection policy; replay preserves it without claiming assessment. Start a new session if that history is untrusted.";
if !self.warnings.iter().any(|message| message == WARNING) {
self.warnings.push(WARNING.into());
}
}
let protected_output =
crate::protection::prompt_injection::replay_projection(output, assessment);
let projection_changed_output = protected_output
.as_deref()
.is_some_and(|protected| protected != output);
let output = protected_output.as_deref().unwrap_or(output);
let original_chars = output.chars().count();
let mut output = compact_historical_tool_output(output, &tool_name, success);
if tool_name == "code_mode"
&& (projection_changed_output || original_chars > REPLAY_TOOL_RESULT_OUTPUT_CHAR_LIMIT)
&& let Some(summary) = result.pointer("/metadata/code_mode_execution")
{
output.push_str("\nHost-owned Code Mode execution evidence and warnings:\n");
output.push_str(&summary.to_string());
}
let skill_reads = if original_chars <= REPLAY_TOOL_RESULT_OUTPUT_CHAR_LIMIT {
validated_skill_reads(
&tool_name,
success,
output.as_str(),
result.get("metadata").unwrap_or(&Value::Null),
)
} else {
Vec::new()
};
if original_chars > REPLAY_TOOL_RESULT_OUTPUT_CHAR_LIMIT {
self.replay_diagnostics.push(format!(
"compacted_historical_tool_result_output tool={tool_name} success={success} chars={original_chars} limit={REPLAY_TOOL_RESULT_OUTPUT_CHAR_LIMIT}"
));
}
self.push_tool_result_item(
ProviderToolResult {
call_id: call_id.to_string(),
tool_name,
success,
output,
skill_reads,
},
original_chars,
);
}
fn push_provider_context_item(&mut self, event: &SessionEvent) {
let role = event.payload.get("role").and_then(Value::as_str);
let Some(content) = event.payload.get("content").and_then(Value::as_str) else {
self.push_legacy_note(event);
return;
};
match role {
Some("user") => self.push_message(ChatMessage::user(content)),
_ => self.push_legacy_note(event),
}
}
fn push_tool_result_item(&mut self, result: ProviderToolResult, original_chars: usize) {
let item = ProviderConversationItem::ToolResult(result);
if let Some(turn) = self.active_turn.as_mut() {
turn.provider_activity_after_user = true;
turn.provider_visible_chars = turn
.provider_visible_chars
.saturating_add(replay_item_visible_chars(&item));
turn.original_provider_visible_chars = turn
.original_provider_visible_chars
.saturating_add(original_chars);
turn.items.push(item);
} else {
self.items.push(item);
}
}
fn push_legacy_note(&mut self, event: &SessionEvent) {
let content = if event.kind() == Some(SessionEventKind::ToolResult)
&& event.payload["result"]["metadata"]["prompt_injection_protection"]["enforced"]
.as_bool()
== Some(true)
{
"Unmatched protected tool result withheld; raw event remains in local session history."
.into()
} else {
event_line(event)
};
self.legacy_lossy_events.push(event.event_type.clone());
self.push_replay_item(ProviderConversationItem::LegacyReplayNote {
event_type: event.event_type.clone(),
content,
});
}
}