use super::manager::active_lease_target;
use super::manager::{Session, SessionInternalDiagnostic, SessionManager};
use super::metadata::{push_session_diagnostic, session_from_entry, session_metadata_for_listing};
use super::read::validate_session_id;
use super::store::{remove_primary, validate_existing_file, validate_path_file};
use crate::persistence::CrossProcessFileLock;
use std::{
fs, io,
path::{Path, PathBuf},
thread,
time::{Duration, SystemTime},
};
pub(crate) const PRUNE_SESSIONS_DEFAULT_DAYS: u64 = crate::config::DEFAULT_SESSION_RETENTION_DAYS;
pub(crate) const PRUNE_SESSIONS_USAGE: &str =
"usage: /prune-sessions [days]; days must be a positive integer";
const MAX_PRUNE_SUMMARY_CHARS: usize = 4096;
const MAX_PRUNE_FAILURES: usize = 64;
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct PruneDeletionFailure {
pub(crate) session_id: String,
pub(crate) category: &'static str,
pub(crate) detail: &'static str,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct PruneSessionsReport {
pub(crate) retention_days: u64,
pub(crate) deleted_ids: Vec<String>,
pub(crate) skipped_active_id: Option<String>,
pub(crate) failed_ids: Vec<String>,
pub(crate) failures: Vec<PruneDeletionFailure>,
pub(crate) omitted_failures: usize,
pub(crate) diagnostics: Vec<SessionInternalDiagnostic>,
}
impl PruneSessionsReport {
pub(crate) fn summary(&self) -> String {
let mut parts = vec![format!(
"pruned {} sessions older than {} days",
self.deleted_ids.len(),
self.retention_days
)];
if self.skipped_active_id.is_some() {
parts.push("skipped active session".to_string());
}
if !self.failed_ids.is_empty() {
let mut failures = self
.failures
.iter()
.map(|failure| {
format!(
"{} ({}: {})",
failure.session_id, failure.category, failure.detail
)
})
.collect::<Vec<_>>();
if self.omitted_failures > 0 {
failures.push(format!(
"omitted {} additional failures",
self.omitted_failures
));
}
parts.push(format!(
"failed to delete {} sessions: {}",
self.failed_ids.len(),
failures.join(", ")
));
}
let summary = parts.join("; ");
if summary.chars().count() <= MAX_PRUNE_SUMMARY_CHARS {
return summary;
}
let mut bounded = summary
.chars()
.take(MAX_PRUNE_SUMMARY_CHARS - 32)
.collect::<String>();
bounded.push_str("; output truncated");
bounded
}
}
pub(crate) fn parse_prune_sessions_days(arg: Option<&str>) -> Result<u64, &'static str> {
let tokens = arg
.unwrap_or("")
.split_whitespace()
.filter(|token| !token.is_empty())
.collect::<Vec<_>>();
match tokens.as_slice() {
[] => Ok(PRUNE_SESSIONS_DEFAULT_DAYS),
[token] if token.chars().all(|ch| ch.is_ascii_digit()) && !token.is_empty() => {
let days = token.parse::<u64>().map_err(|_| PRUNE_SESSIONS_USAGE)?;
if days == 0 || days.checked_mul(86_400).is_none() {
return Err(PRUNE_SESSIONS_USAGE);
}
Ok(days)
}
_ => Err(PRUNE_SESSIONS_USAGE),
}
}
fn prune_delete_error_category(kind: std::io::ErrorKind) -> (&'static str, &'static str) {
match kind {
std::io::ErrorKind::PermissionDenied => ("permission_denied", "permission denied"),
std::io::ErrorKind::NotFound => ("not_found", "session file was not found"),
std::io::ErrorKind::IsADirectory
| std::io::ErrorKind::NotADirectory
| std::io::ErrorKind::InvalidInput => ("wrong_type", "session path has wrong type"),
_ => ("filesystem", "filesystem deletion failed"),
}
}
fn record_prune_failure(
report: &mut PruneSessionsReport,
session_id: String,
category: &'static str,
detail: &'static str,
) {
report.failed_ids.push(session_id.clone());
if report.failures.len() < MAX_PRUNE_FAILURES {
report.failures.push(PruneDeletionFailure {
session_id,
category,
detail,
});
}
}
#[derive(Debug, Clone)]
struct PruneCandidate {
session: Session,
activity_time: SystemTime,
}
struct PruneCandidatesReport {
candidates: Vec<PruneCandidate>,
diagnostics: Vec<SessionInternalDiagnostic>,
failed_ids: Vec<String>,
failures: Vec<PruneDeletionFailure>,
}
fn validated_history_path(root: &Path, id: &str) -> Option<PathBuf> {
validate_session_id(id.to_string()).ok()?;
let root_metadata = fs::symlink_metadata(root).ok()?;
if root_metadata.file_type().is_symlink() || !root_metadata.file_type().is_dir() {
return None;
}
let history_parent = root.join(".history");
let history_metadata = fs::symlink_metadata(&history_parent).ok()?;
if history_metadata.file_type().is_symlink() || !history_metadata.file_type().is_dir() {
return None;
}
let history = history_parent.join(id);
let metadata = fs::symlink_metadata(&history).ok()?;
if metadata.file_type().is_symlink() || !metadata.file_type().is_dir() {
return None;
}
if !history.starts_with(root) {
return None;
}
Some(history)
}
fn remove_session_sidecar(root: &Path, id: &str, extension: &str) -> io::Result<()> {
validate_session_id(id.to_string())
.map_err(|_| io::Error::from(io::ErrorKind::InvalidInput))?;
let path = root.join(format!("{id}.{extension}"));
match validate_existing_file(&path) {
Ok(Some(_)) => fs::remove_file(path),
Ok(None) => Ok(()),
Err(error) => Err(io::Error::other(error.to_string())),
}
}
fn retention_cutoff(now: SystemTime, retention_days: u64) -> Option<SystemTime> {
let seconds = retention_days.checked_mul(86_400)?;
Some(
now.checked_sub(Duration::from_secs(seconds))
.unwrap_or(SystemTime::UNIX_EPOCH),
)
}
fn log_auto_prune_failure() {
crate::output::emit_terminal_warning(
"warning: session operation=auto_prune category=session_jsonl failed".to_string(),
);
}
impl SessionManager {
pub(crate) fn prune_sessions(
&self,
retention_days: u64,
active_session_id: Option<&str>,
) -> anyhow::Result<PruneSessionsReport> {
self.prune_sessions_at(SystemTime::now(), retention_days, active_session_id)
}
fn prune_sessions_at(
&self,
now: SystemTime,
retention_days: u64,
active_session_id: Option<&str>,
) -> anyhow::Result<PruneSessionsReport> {
self.prune_sessions_with(now, retention_days, active_session_id, |path| {
let id = path
.file_name()
.and_then(|name| name.to_str())
.and_then(|name| name.strip_suffix(".jsonl"))
.unwrap_or_default();
remove_primary(&self.root, id)
})
}
pub(crate) fn spawn_auto_prune(&self, retention_days: u64, active_session_id: Option<&str>) {
if retention_days == 0 {
return;
}
let manager = self.clone();
let active_session_id = active_session_id.map(str::to_string);
if thread::Builder::new()
.name("magi-session-auto-prune".to_string())
.spawn(move || {
manager.run_auto_prune(
SystemTime::now(),
retention_days,
active_session_id.as_deref(),
);
})
.is_err()
{
log_auto_prune_failure();
}
}
fn run_auto_prune(
&self,
now: SystemTime,
retention_days: u64,
active_session_id: Option<&str>,
) {
if retention_days == 0 {
return;
}
let Some(cutoff) = retention_cutoff(now, retention_days) else {
log_auto_prune_failure();
return;
};
match self.prune_sessions_at(now, retention_days, active_session_id) {
Ok(report) if report.failed_ids.is_empty() && report.diagnostics.is_empty() => {}
Ok(_) | Err(_) => log_auto_prune_failure(),
}
self.prune_subagent_sessions(cutoff, active_session_id);
Self::prune_orphan_sidecars(&self.root, cutoff, active_session_id);
Self::prune_orphan_sidecars(&self.root.join("subagents"), cutoff, None);
Self::prune_orphan_history(&self.root, cutoff);
Self::prune_orphan_history(&self.root.join("subagents"), cutoff);
self.prune_orphan_checkpoints(cutoff);
}
fn prune_subagent_sessions(&self, cutoff: SystemTime, active_session_id: Option<&str>) {
let root = self.root.join("subagents");
let Ok(root_metadata) = fs::symlink_metadata(&root) else {
return;
};
if root_metadata.file_type().is_symlink() || !root_metadata.file_type().is_dir() {
log_auto_prune_failure();
return;
}
let Ok(entries) = fs::read_dir(&root) else {
log_auto_prune_failure();
return;
};
for entry in entries.flatten() {
let Ok(file_type) = entry.file_type() else {
log_auto_prune_failure();
continue;
};
if !file_type.is_file() {
continue;
}
let path = entry.path();
let Some(id) = path
.file_name()
.and_then(|name| name.to_str())
.and_then(|name| name.strip_suffix(".jsonl"))
.map(str::to_string)
else {
continue;
};
if validate_session_id(id.clone()).is_err() || active_session_id == Some(id.as_str()) {
continue;
}
let Ok(metadata) = fs::symlink_metadata(&path) else {
continue;
};
if metadata
.modified()
.ok()
.is_none_or(|modified| modified >= cutoff)
{
continue;
}
let Some(_active_lease) =
(match CrossProcessFileLock::try_acquire(&active_lease_target(&path)) {
Ok(lock) => lock,
Err(_) => {
log_auto_prune_failure();
continue;
}
})
else {
continue;
};
let Some(_lock_guard) = (match CrossProcessFileLock::try_acquire(&path) {
Ok(lock) => lock,
Err(_) => {
log_auto_prune_failure();
continue;
}
}) else {
continue;
};
let Ok(current_metadata) = fs::symlink_metadata(&path) else {
continue;
};
if current_metadata.file_type().is_symlink()
|| !current_metadata.file_type().is_file()
|| current_metadata
.modified()
.ok()
.is_none_or(|modified| modified >= cutoff)
{
continue;
}
if remove_primary(&root, &id).is_err() {
log_auto_prune_failure();
continue;
}
if remove_session_sidecar(&root, &id, "metadata.json").is_err() {
log_auto_prune_failure();
}
if let Some(history) = validated_history_path(&root, &id)
&& fs::remove_dir_all(history).is_err()
{
log_auto_prune_failure();
}
}
}
fn prune_orphan_sidecars(root: &Path, cutoff: SystemTime, active_session_id: Option<&str>) {
let Ok(root_metadata) = fs::symlink_metadata(root) else {
return;
};
if root_metadata.file_type().is_symlink() || !root_metadata.file_type().is_dir() {
log_auto_prune_failure();
return;
}
let Ok(entries) = fs::read_dir(root) else {
log_auto_prune_failure();
return;
};
for entry in entries.flatten() {
let Ok(file_type) = entry.file_type() else {
log_auto_prune_failure();
continue;
};
if !file_type.is_file() {
continue;
}
let path = entry.path();
let Some(id) = path
.file_name()
.and_then(|name| name.to_str())
.and_then(|name| name.strip_suffix(".metadata.json"))
.map(str::to_string)
else {
continue;
};
if validate_session_id(id.clone()).is_err() || active_session_id == Some(id.as_str()) {
continue;
}
let primary = root.join(format!("{id}.jsonl"));
if fs::symlink_metadata(primary).is_ok() {
continue;
}
let Ok(metadata) = fs::symlink_metadata(&path) else {
continue;
};
if metadata
.modified()
.ok()
.is_none_or(|modified| modified >= cutoff)
{
continue;
}
let Ok(Some(_active_lease)) = CrossProcessFileLock::try_acquire(&active_lease_target(
&root.join(format!("{id}.jsonl")),
)) else {
continue;
};
if remove_session_sidecar(root, &id, "metadata.json").is_err() {
log_auto_prune_failure();
}
}
}
fn prune_orphan_history(root: &Path, cutoff: SystemTime) {
let Ok(root_metadata) = fs::symlink_metadata(root) else {
return;
};
if root_metadata.file_type().is_symlink() || !root_metadata.file_type().is_dir() {
log_auto_prune_failure();
return;
}
let history_root = root.join(".history");
let Ok(metadata) = fs::symlink_metadata(&history_root) else {
return;
};
if metadata.file_type().is_symlink() || !metadata.file_type().is_dir() {
log_auto_prune_failure();
return;
}
let Ok(entries) = fs::read_dir(&history_root) else {
log_auto_prune_failure();
return;
};
for entry in entries.flatten() {
let Ok(file_type) = entry.file_type() else {
log_auto_prune_failure();
continue;
};
if !file_type.is_dir() || file_type.is_symlink() {
continue;
}
let Some(id) = entry.file_name().to_str().map(str::to_string) else {
continue;
};
if validate_session_id(id.clone()).is_err()
|| fs::symlink_metadata(root.join(format!("{id}.jsonl"))).is_ok()
{
continue;
}
let path = entry.path();
let Ok(metadata) = fs::symlink_metadata(&path) else {
continue;
};
if metadata
.modified()
.ok()
.is_none_or(|modified| modified >= cutoff)
{
continue;
}
let Ok(Some(_active_lease)) = CrossProcessFileLock::try_acquire(&active_lease_target(
&root.join(format!("{id}.jsonl")),
)) else {
continue;
};
if fs::remove_dir_all(path).is_err() {
log_auto_prune_failure();
}
}
}
fn prune_orphan_checkpoints(&self, cutoff: SystemTime) {
let Some(parent) = self.root.parent() else {
return;
};
let checkpoint_root = parent.join("checkpoints");
let ledgers = checkpoint_root.join("ledgers");
let subagent_root = self.root.join("subagents");
let store = crate::checkpoints::CheckpointStore::new(checkpoint_root);
match fs::symlink_metadata(&ledgers) {
Ok(metadata) if !metadata.file_type().is_symlink() && metadata.file_type().is_dir() => {
let Ok(entries) = fs::read_dir(&ledgers) else {
log_auto_prune_failure();
return;
};
for entry in entries.flatten() {
let Ok(file_type) = entry.file_type() else {
log_auto_prune_failure();
continue;
};
if !file_type.is_file() || file_type.is_symlink() {
continue;
}
let path = entry.path();
let Some(id) = path
.file_name()
.and_then(|name| name.to_str())
.and_then(|name| name.strip_suffix(".jsonl"))
.map(str::to_string)
else {
continue;
};
if validate_session_id(id.clone()).is_err()
|| fs::symlink_metadata(self.root.join(format!("{id}.jsonl"))).is_ok()
|| fs::symlink_metadata(subagent_root.join(format!("{id}.jsonl"))).is_ok()
{
continue;
}
let Ok(metadata) = fs::symlink_metadata(&path) else {
continue;
};
if metadata
.modified()
.ok()
.is_none_or(|modified| modified >= cutoff)
{
continue;
}
let Ok(Some(_primary_lease)) = CrossProcessFileLock::try_acquire(
&active_lease_target(&self.root.join(format!("{id}.jsonl"))),
) else {
continue;
};
let Ok(Some(_subagent_lease)) = CrossProcessFileLock::try_acquire(
&active_lease_target(&subagent_root.join(format!("{id}.jsonl"))),
) else {
continue;
};
if store.prune_session(&id).is_err() {
log_auto_prune_failure();
}
}
}
Ok(_) => {
log_auto_prune_failure();
}
Err(error) if error.kind() == io::ErrorKind::NotFound => {}
Err(_) => log_auto_prune_failure(),
}
if store.prune_unreferenced_blobs().is_err() {
log_auto_prune_failure();
}
}
fn prune_sessions_with(
&self,
now: SystemTime,
retention_days: u64,
active_session_id: Option<&str>,
mut remove_file: impl FnMut(&Path) -> std::io::Result<()>,
) -> anyhow::Result<PruneSessionsReport> {
let retention_secs = retention_days
.checked_mul(86_400)
.ok_or_else(|| anyhow::anyhow!(PRUNE_SESSIONS_USAGE))?;
let cutoff = now
.checked_sub(Duration::from_secs(retention_secs))
.unwrap_or(SystemTime::UNIX_EPOCH);
let mut report = PruneSessionsReport {
retention_days,
deleted_ids: Vec::new(),
skipped_active_id: None,
failed_ids: Vec::new(),
failures: Vec::new(),
omitted_failures: 0,
diagnostics: Vec::new(),
};
let candidates_report = self.prune_candidates()?;
report.diagnostics = candidates_report.diagnostics;
report.failed_ids.extend(candidates_report.failed_ids);
report.failures = candidates_report.failures;
report.omitted_failures = report
.failed_ids
.len()
.saturating_sub(report.failures.len());
for candidate in candidates_report.candidates {
let id = candidate.session.id().to_string();
if candidate.activity_time >= cutoff {
continue;
}
if active_session_id == Some(&id) {
report.skipped_active_id = Some(id);
continue;
}
let _active_lease = match CrossProcessFileLock::try_acquire(&active_lease_target(
candidate.session.path(),
)) {
Ok(Some(guard)) => guard,
Ok(None) => {
report.skipped_active_id = Some(id);
continue;
}
Err(_) => {
record_prune_failure(
&mut report,
id,
"lock",
"could not lock session for deletion",
);
continue;
}
};
let _lock_guard = match CrossProcessFileLock::try_acquire(candidate.session.path()) {
Ok(Some(guard)) => guard,
Ok(None) | Err(_) => {
record_prune_failure(
&mut report,
id,
"lock",
"could not lock session for deletion",
);
continue;
}
};
if validate_path_file(&self.root, &id, candidate.session.path()).is_err() {
record_prune_failure(
&mut report,
id,
"changed",
"session changed before deletion",
);
continue;
}
let current_metadata = match session_metadata_for_listing(&candidate.session) {
Ok(metadata) => metadata,
Err(_) => {
record_prune_failure(
&mut report,
id,
"metadata",
"session metadata could not be rechecked",
);
continue;
}
};
if current_metadata.activity_time(&candidate.session) >= cutoff {
continue;
}
match remove_file(candidate.session.path()) {
Ok(()) => {
if let Some(history) = validated_history_path(&self.root, &id)
&& let Err(error) = fs::remove_dir_all(history)
{
push_session_diagnostic(
&mut report.diagnostics,
Some(id.clone()),
format!("failed to prune archived session history: {error}"),
);
}
if let Err(error) = remove_session_sidecar(&self.root, &id, "metadata.json") {
push_session_diagnostic(
&mut report.diagnostics,
Some(id.clone()),
format!("failed to prune session metadata sidecar: {error}"),
);
}
if let Some(parent) = self.root.parent() {
let store =
crate::checkpoints::CheckpointStore::new(parent.join("checkpoints"));
if store.prune_session(&id).is_err() {
push_session_diagnostic(
&mut report.diagnostics,
Some(id.clone()),
"failed to prune checkpoint storage for deleted session"
.to_string(),
);
}
}
report.deleted_ids.push(id);
}
Err(error) => {
let (category, detail) = prune_delete_error_category(error.kind());
record_prune_failure(&mut report, id, category, detail);
}
}
}
report.omitted_failures = report
.failed_ids
.len()
.saturating_sub(report.failures.len());
Ok(report)
}
fn prune_candidates(&self) -> anyhow::Result<PruneCandidatesReport> {
if !self.root.exists() {
return Ok(PruneCandidatesReport {
candidates: Vec::new(),
diagnostics: Vec::new(),
failed_ids: Vec::new(),
failures: Vec::new(),
});
}
let mut candidates = Vec::new();
let mut diagnostics = Vec::new();
let mut failed_ids = Vec::new();
let mut failures = Vec::new();
let mut entries = fs::read_dir(&self.root)?.collect::<Result<Vec<_>, _>>()?;
entries.sort_by_key(|entry| entry.file_name());
for entry in entries {
let Ok(file_type) = entry.file_type() else {
push_session_diagnostic(
&mut diagnostics,
None,
"failed to inspect session file type".to_string(),
);
continue;
};
if !file_type.is_file() {
continue;
}
let path = entry.path();
let Some((id, path)) = session_from_entry(path) else {
continue;
};
let id = match validate_session_id(id) {
Ok(id) => id,
Err(error) => {
push_session_diagnostic(
&mut diagnostics,
None,
format!("ignored invalid session file: {error}"),
);
continue;
}
};
let expected_path = self.path_for_valid_id(&id)?;
if path != expected_path {
push_session_diagnostic(
&mut diagnostics,
Some(id),
"ignored session file with unexpected path".to_string(),
);
continue;
}
let session = Session::new(id, path);
let metadata = match session_metadata_for_listing(&session) {
Ok(metadata) => metadata,
Err(_) => {
failed_ids.push(session.id.clone());
if failures.len() < MAX_PRUNE_FAILURES {
failures.push(PruneDeletionFailure {
session_id: session.id.clone(),
category: "metadata",
detail: "session metadata could not be loaded",
});
}
continue;
}
};
candidates.push(PruneCandidate {
activity_time: metadata.activity_time(&session),
session,
});
}
Ok(PruneCandidatesReport {
candidates,
diagnostics,
failed_ids,
failures,
})
}
}