use super::{AgentOutputSink, util::safe_error_message};
use crate::{
output::OutputEvent,
sessions::{
Session, SessionAppendOutcome, SessionEventKind, SessionReadDiagnostic, TurnStatus,
TurnStatusPayload, record_session_event,
},
};
use serde_json::Value;
use std::path::Path;
pub(super) struct SessionPersistence<'a> {
session: Option<&'a Session>,
cwd: &'a Path,
warning_emitted: bool,
metadata_warning_emitted: bool,
incomplete_marker_attempted: bool,
terminal_status_recorded: bool,
task_scope: Option<std::sync::Arc<std::sync::Mutex<crate::sessions::task_scope::TaskScope>>>,
pub(super) code_mode_parent: Option<String>,
pub(super) code_mode_sequence: u64,
pub(super) code_mode_audit: std::cell::RefCell<crate::code_mode::audit::CodeModeAudit>,
code_mode_history_budget: std::cell::RefCell<crate::sessions::SessionAppendBudget>,
}
impl<'a> SessionPersistence<'a> {
pub(super) fn new(session: Option<&'a Session>, cwd: &'a Path) -> Self {
Self {
session,
cwd,
warning_emitted: false,
metadata_warning_emitted: false,
incomplete_marker_attempted: false,
terminal_status_recorded: false,
task_scope: None,
code_mode_parent: None,
code_mode_sequence: 0,
code_mode_audit: Default::default(),
code_mode_history_budget: std::cell::RefCell::new(
crate::sessions::SessionAppendBudget {
remaining_bytes: crate::code_mode::limits::MAX_BATCH_HISTORY_BYTES,
max_session_bytes: crate::code_mode::limits::MAX_NESTED_SESSION_BYTES as u64,
},
),
}
}
pub(super) fn with_task_scope(mut self, tools: Option<&crate::tools::ToolRuntime>) -> Self {
self.task_scope = tools.and_then(|tools| tools.task_scope.clone());
self
}
pub(super) fn nested(&self, parent: &str) -> Self {
let mut nested = Self::new(self.session, self.cwd);
nested.code_mode_parent = Some(parent.to_owned());
nested
}
pub(super) fn is_degraded(&self) -> bool {
self.warning_emitted
}
pub(super) fn warn_once(
&mut self,
output_sink: &mut Option<&mut dyn AgentOutputSink>,
error: &anyhow::Error,
) -> anyhow::Result<()> {
if !self.incomplete_marker_attempted {
self.incomplete_marker_attempted = true;
let _ = record_session_event(
self.session,
self.cwd,
SessionEventKind::TurnStatus,
TurnStatusPayload::new(TurnStatus::Incomplete, None).into_value()?,
);
}
if self.warning_emitted {
return Ok(());
}
self.warning_emitted = true;
if let Some(sink) = output_sink.as_deref_mut() {
sink.persistence_degraded();
sink.output_event(OutputEvent::Diagnostic {
level: "warning".to_string(),
message: format!(
"session persistence failed; future resume may be incomplete: {}",
safe_error_message(error)
),
})?;
}
Ok(())
}
pub(super) fn try_record(
&mut self,
kind: SessionEventKind,
payload: Value,
output_sink: &mut Option<&mut dyn AgentOutputSink>,
) -> anyhow::Result<()> {
if self.code_mode_parent.is_some() {
return self.record_required(kind, payload);
}
match record_session_event(self.session, self.cwd, kind, payload) {
Ok(SessionAppendOutcome::Durable) => {}
Ok(SessionAppendOutcome::DurableJsonlMetadataUpdateFailed(error)) => {
if !self.metadata_warning_emitted {
self.metadata_warning_emitted = true;
if let Some(sink) = output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: "error".to_string(),
message: format!(
"Session metadata update failed: {error}. Session events remain saved. Start a new session; if this repeats, check session-directory permissions and free disk space."
),
})?;
}
}
}
Err(error) => self.warn_once(output_sink, &error)?,
}
Ok(())
}
pub(super) fn record_terminal_status(
&mut self,
status: TurnStatus,
assistant_text: &str,
output_sink: &mut Option<&mut dyn AgentOutputSink>,
) -> anyhow::Result<()> {
self.record_terminal_status_with_error(status, assistant_text, None, output_sink)
}
pub(super) fn record_terminal_status_with_error(
&mut self,
status: TurnStatus,
assistant_text: &str,
error: Option<&anyhow::Error>,
output_sink: &mut Option<&mut dyn AgentOutputSink>,
) -> anyhow::Result<()> {
if self.terminal_status_recorded {
return Ok(());
}
self.terminal_status_recorded = true;
let error_summary = (status == TurnStatus::Failed)
.then(|| error.map(safe_error_message))
.flatten();
let _ = self.try_record(
SessionEventKind::TurnStatus,
TurnStatusPayload::new_with_error_summary(
status,
Some(assistant_text),
error_summary.as_deref(),
)
.into_value()?,
output_sink,
);
Ok(())
}
pub(super) fn record_abort_recovery(
&mut self,
payload: Value,
output_sink: &mut Option<&mut dyn AgentOutputSink>,
) -> anyhow::Result<()> {
self.try_record(SessionEventKind::AbortRecovery, payload, output_sink)
}
pub(super) fn record_provider_stream_trace(
&mut self,
payload: Value,
output_sink: &mut Option<&mut dyn AgentOutputSink>,
) -> anyhow::Result<()> {
self.try_record(SessionEventKind::ProviderStreamTrace, payload, output_sink)
}
pub(super) fn code_mode_budget(&self) -> anyhow::Result<Value> {
let Some(session) = self.session else {
return Ok(
serde_json::json!({"batch_history_bytes": null, "session_history_bytes": null}),
);
};
let bytes = session.history_bytes()?;
let budget = self.code_mode_history_budget.borrow();
Ok(serde_json::json!({
"batch_history_bytes": crate::code_mode::limits::budget_snapshot(
(crate::code_mode::limits::MAX_BATCH_HISTORY_BYTES - budget.remaining_bytes) as u64,
crate::code_mode::limits::MAX_BATCH_HISTORY_BYTES as u64,
),
"session_history_bytes": crate::code_mode::limits::budget_snapshot(bytes, budget.max_session_bytes),
}))
}
pub(super) fn record_required(
&self,
kind: SessionEventKind,
mut payload: Value,
) -> anyhow::Result<()> {
if let Some(parent) = &self.code_mode_parent {
let kind = match kind {
SessionEventKind::ToolCall => SessionEventKind::CodeModeToolCall,
SessionEventKind::ToolResult | SessionEventKind::ToolDisplayResult => {
SessionEventKind::CodeModeToolResult
}
SessionEventKind::ProviderContextItem => SessionEventKind::CodeModeContext,
other => other,
};
payload["parent_call_id"] = Value::String(parent.clone());
payload["sequence"] = self.code_mode_sequence.into();
if kind == SessionEventKind::CodeModeContext {
payload["code_mode_context_id"] = Value::String(format!(
"{parent}:context:{}",
self.code_mode_audit.borrow().contexts.len()
));
}
if let Some(session) = self.session {
let event = crate::sessions::SessionEvent::new_kind(
kind,
session.id().to_owned(),
self.cwd.to_path_buf(),
payload.clone(),
);
session
.append_with_budget(event, &mut self.code_mode_history_budget.borrow_mut())
.map_err(|error| {
let limit = match error
.downcast_ref::<crate::sessions::SessionAppendError>()
{
Some(crate::sessions::SessionAppendError::BatchByteLimitExceeded {
..
}) => Some("batch persisted history bytes"),
Some(
crate::sessions::SessionAppendError::SessionByteLimitExceeded {
..
},
) => Some("session persisted history bytes"),
_ => None,
};
match limit {
Some(limit) => anyhow::Error::from(
crate::code_mode::engine::ResourceLimitExceeded(limit),
)
.context(error.to_string()),
None => error,
}
})?;
}
self.code_mode_audit.borrow_mut().observe(kind, &payload);
return Ok(());
}
if kind == SessionEventKind::UserInput
&& let Some(scope) = &self.task_scope
{
let mut scope = scope
.lock()
.map_err(|_| anyhow::anyhow!("task scope lock poisoned"))?;
let updated = scope.with_input(&mut payload)?;
record_session_event(self.session, self.cwd, kind, payload)?;
*scope = updated;
return Ok(());
}
record_session_event(self.session, self.cwd, kind, payload).map(|_| ())
}
}
pub(super) fn emit_replay_diagnostics(
output_sink: &mut Option<&mut dyn AgentOutputSink>,
diagnostics: &[SessionReadDiagnostic],
warnings: &[String],
) -> anyhow::Result<()> {
if let Some(sink) = output_sink.as_deref_mut() {
if !diagnostics.is_empty() {
sink.output_event(OutputEvent::Diagnostic {
level: "warning".to_string(),
message: format!(
"session replay skipped {} unreadable JSONL line(s); continuing with valid session events only",
diagnostics.len()
),
})?;
}
for message in warnings {
sink.output_event(OutputEvent::Diagnostic {
level: "warning".to_string(),
message: message.clone(),
})?;
}
}
Ok(())
}