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() {
eprintln!("warning: session operation=auto_prune category=session_jsonl failed");
}
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;
}
for extension in ["metadata.json", "writer-owner"] {
if remove_session_sidecar(&root, &id, extension).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();
}
}
#[cfg(test)]
fn auto_prune_sessions_for_test(
&self,
now: SystemTime,
retention_days: u64,
active_session_id: Option<&str>,
) {
self.run_auto_prune(now, retention_days, active_session_id);
}
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}"),
);
}
for (extension, label) in
[("metadata.json", "metadata"), ("writer-owner", "ownership")]
{
if let Err(error) = remove_session_sidecar(&self.root, &id, extension) {
push_session_diagnostic(
&mut report.diagnostics,
Some(id.clone()),
format!("failed to prune session {label} 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,
})
}
#[cfg(test)]
fn prune_sessions_for_test(
&self,
now: SystemTime,
retention_days: u64,
active_session_id: Option<&str>,
remove_file: impl FnMut(&Path) -> std::io::Result<()>,
) -> anyhow::Result<PruneSessionsReport> {
self.prune_sessions_with(now, retention_days, active_session_id, remove_file)
}
}
#[cfg(test)]
mod tests {
use super::super::event::SessionEvent;
use super::super::metadata::metadata_path_for_session;
use super::*;
use chrono::{TimeZone, Utc};
use serde_json::json;
use tempfile::TempDir;
#[test]
fn prune_sessions_parser_accepts_default_and_positive_integer_only() {
assert_eq!(parse_prune_sessions_days(None).unwrap(), 30);
assert_eq!(parse_prune_sessions_days(Some(" ")).unwrap(), 30);
assert_eq!(parse_prune_sessions_days(Some("7")).unwrap(), 7);
assert_eq!(parse_prune_sessions_days(Some(" 07 ")).unwrap(), 7);
for invalid in [
"0",
"-1",
"+7",
"1.5",
"seven",
"7 now",
"18446744073709551616",
] {
assert_eq!(
parse_prune_sessions_days(Some(invalid)),
Err(PRUNE_SESSIONS_USAGE)
);
}
}
#[test]
fn prune_removes_owner_marker_only_after_successful_session_deletion() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let now = SystemTime::from(Utc.with_ymd_and_hms(2026, 1, 31, 0, 0, 0).unwrap());
let old = manager.open("old-session").unwrap();
let active = manager.open("active-session").unwrap();
let failed = manager.open("failed-session").unwrap();
for session in [&old, &active, &failed] {
append_event_at(session, temp.path(), 2025, 12, 1);
drop(session.try_daemon_writer().unwrap().unwrap());
assert!(session.path().with_extension("writer-owner").exists());
}
let active_lease = CrossProcessFileLock::try_acquire(&active_lease_target(active.path()))
.unwrap()
.unwrap();
let report = manager
.prune_sessions_for_test(now, 30, None, |path| {
if path == failed.path() {
return Err(io::Error::from(io::ErrorKind::PermissionDenied));
}
assert!(path.with_extension("writer-owner").exists());
fs::remove_file(path)
})
.unwrap();
assert_eq!(report.deleted_ids, vec![old.id()]);
assert_eq!(report.failed_ids, vec![failed.id()]);
assert_eq!(report.skipped_active_id.as_deref(), Some(active.id()));
assert!(!old.path().exists());
assert!(!old.path().with_extension("writer-owner").exists());
for session in [&active, &failed] {
assert!(session.path().exists());
assert!(session.path().with_extension("writer-owner").exists());
}
drop(active_lease);
}
#[test]
fn prune_sessions_deletes_only_old_inactive_top_level_jsonl() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
fs::create_dir_all(temp.path().join("sessions/nested")).unwrap();
crate::sessions::store::secure_test_session_root(temp.path().join("sessions").as_path());
let now = SystemTime::from(Utc.with_ymd_and_hms(2026, 1, 31, 0, 0, 0).unwrap());
let old = manager.open("old-session").unwrap();
let active = manager.open("active-session").unwrap();
let new = manager.open("new-session").unwrap();
let invalid = temp.path().join("sessions/bad.name.jsonl");
let wrong_ext = temp.path().join("sessions/old-note.txt");
let nested = temp.path().join("sessions/nested/nested-session.jsonl");
append_event_at(&old, temp.path(), 2025, 12, 1);
append_event_at(&active, temp.path(), 2025, 12, 1);
append_event_at(&new, temp.path(), 2026, 1, 15);
fs::write(&invalid, "{}").unwrap();
fs::write(&wrong_ext, "keep").unwrap();
fs::write(&nested, "keep").unwrap();
let report = manager
.prune_sessions_for_test(now, 30, Some(active.id()), |path| fs::remove_file(path))
.unwrap();
assert_eq!(report.deleted_ids, vec!["old-session"]);
assert_eq!(report.skipped_active_id.as_deref(), Some("active-session"));
assert!(report.failed_ids.is_empty());
assert_eq!(report.diagnostics.len(), 1);
assert!(report.diagnostics[0].message.contains("invalid session"));
assert!(!old.path().exists());
assert!(!metadata_path_for_session(&old).exists());
assert!(active.path().exists());
assert!(new.path().exists());
assert!(invalid.exists());
assert!(wrong_ext.exists());
assert!(nested.exists());
assert_eq!(
report.summary(),
"pruned 1 sessions older than 30 days; skipped active session"
);
}
#[test]
fn auto_prune_preserves_active_top_level_session() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let active = manager.open("active-auto").unwrap();
active
.append(&event_at(&active, temp.path(), 2025, 1, 1))
.unwrap();
drop(active.try_daemon_writer().unwrap().unwrap());
let owner_path = active.path().with_extension("writer-owner");
assert!(owner_path.exists());
let now = SystemTime::from(Utc.with_ymd_and_hms(2026, 1, 31, 0, 0, 0).unwrap());
manager.auto_prune_sessions_for_test(now, 30, None);
assert!(active.path().exists());
assert!(metadata_path_for_session(&active).exists());
assert!(owner_path.exists());
let path = active.path().to_path_buf();
drop(active);
manager.auto_prune_sessions_for_test(now, 30, None);
assert!(!path.exists());
assert!(!path.with_extension("metadata.json").exists());
assert!(!owner_path.exists());
}
#[test]
fn prune_sessions_retains_equal_cutoff_and_uses_latest_valid_event() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let now = SystemTime::from(Utc.with_ymd_and_hms(2026, 1, 31, 0, 0, 0).unwrap());
let equal = manager.open("equal-cutoff").unwrap();
let malformed_then_recent = manager.open("malformed-recent").unwrap();
append_event_at(&equal, temp.path(), 2026, 1, 1);
let old = event_at(&malformed_then_recent, temp.path(), 2025, 12, 1);
let recent = event_at(&malformed_then_recent, temp.path(), 2026, 1, 30);
fs::create_dir_all(malformed_then_recent.path().parent().unwrap()).unwrap();
fs::write(
malformed_then_recent.path(),
format!(
"{}\nnot json\n{}\n",
serde_json::to_string(&old).unwrap(),
serde_json::to_string(&recent).unwrap()
),
)
.unwrap();
crate::sessions::store::secure_test_session_root(
malformed_then_recent.path().parent().unwrap(),
);
let report = manager
.prune_sessions_for_test(now, 30, None, |path| fs::remove_file(path))
.unwrap();
assert!(report.deleted_ids.is_empty(), "{report:?}");
assert!(equal.path().exists());
assert!(malformed_then_recent.path().exists());
}
#[test]
fn prune_sessions_reports_partial_delete_failures_and_continues() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let now = SystemTime::from(Utc.with_ymd_and_hms(2026, 1, 31, 0, 0, 0).unwrap());
let fail = manager.open("fail-session").unwrap();
let delete = manager.open("delete-session").unwrap();
append_event_at(&fail, temp.path(), 2025, 12, 1);
append_event_at(&delete, temp.path(), 2025, 12, 1);
let report = manager
.prune_sessions_for_test(now, 30, None, |path| {
if path.file_stem().and_then(|stem| stem.to_str()) == Some("fail-session") {
Err(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"nope",
))
} else {
fs::remove_file(path)
}
})
.unwrap();
assert_eq!(report.failures[0].category, "permission_denied");
assert_eq!(report.failures[0].detail, "permission denied");
assert_eq!(report.deleted_ids, vec!["delete-session"]);
assert_eq!(report.failed_ids, vec!["fail-session"]);
assert!(fail.path().exists());
assert!(!delete.path().exists());
assert_eq!(
report.summary(),
"pruned 1 sessions older than 30 days; failed to delete 1 sessions: fail-session (permission_denied: permission denied)"
);
}
#[test]
fn prune_sessions_holds_cross_process_lock_while_deleting_jsonl() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let now = SystemTime::from(Utc.with_ymd_and_hms(2026, 1, 31, 0, 0, 0).unwrap());
let old = manager.open("locked-delete").unwrap();
append_event_at(&old, temp.path(), 2025, 12, 1);
let lock_path = lock_path_for_session(&old);
let report = manager
.prune_sessions_for_test(now, 30, None, |path| {
assert_eq!(path, old.path());
assert!(lock_path.exists());
fs::remove_file(path)
})
.unwrap();
assert_eq!(report.deleted_ids, vec!["locked-delete"]);
assert!(report.failed_ids.is_empty());
assert!(!old.path().exists());
}
#[test]
fn prune_sessions_skips_clone_shared_writer_lease() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let now = SystemTime::from(Utc.with_ymd_and_hms(2026, 1, 31, 0, 0, 0).unwrap());
let session = manager.open("clone-writer").unwrap();
session
.append(&event_at(&session, temp.path(), 2025, 12, 1))
.unwrap();
let append_session = session.clone();
let cwd = temp.path().to_path_buf();
std::thread::spawn(move || {
append_session
.append(&event_at(&append_session, &cwd, 2025, 12, 2))
.unwrap();
})
.join()
.unwrap();
let report = manager
.prune_sessions_for_test(now, 30, None, |path| fs::remove_file(path))
.unwrap();
assert!(report.deleted_ids.is_empty());
assert_eq!(report.skipped_active_id.as_deref(), Some("clone-writer"));
assert_eq!(session.read_events().unwrap().len(), 2);
}
#[test]
fn auto_prune_zero_retention_does_not_delete_sessions() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.open("retained").unwrap();
append_event_at(&session, temp.path(), 2025, 1, 1);
manager.auto_prune_sessions_for_test(SystemTime::now(), 0, None);
assert!(session.path().exists());
assert!(metadata_path_for_session(&session).exists());
}
#[test]
fn auto_prune_subagents_removes_old_artifacts_keeps_fresh_and_skips_locked() {
let temp = TempDir::new().unwrap();
let root = temp.path().join("sessions");
let manager = SessionManager::new(root.clone());
let subagents = SessionManager::new(root.join("subagents"));
let old = subagents.open("old-child").unwrap();
let fresh = subagents.open("fresh-child").unwrap();
let locked = subagents.open("locked-child").unwrap();
for session in [&old, &fresh, &locked] {
session
.append(&SessionEvent::new(
"tool_call",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"name":"read"}),
))
.unwrap();
drop(session.try_daemon_writer().unwrap().unwrap());
assert!(session.path().with_extension("writer-owner").exists());
session.release_active_lease_for_test();
}
crate::sessions::store::secure_test_session_root(&root);
let now = SystemTime::now();
let old_time = now - Duration::from_secs(2 * 86_400);
for session in [&old, &locked] {
set_modified(session.path(), old_time);
set_modified(&metadata_path_for_session(session), old_time);
}
let history = root.join("subagents/.history").join(old.id());
fs::create_dir_all(&history).unwrap();
fs::write(history.join("1.jsonl"), "archived").unwrap();
let _lock_guard = CrossProcessFileLock::acquire(locked.path()).unwrap();
manager.auto_prune_sessions_for_test(now, 1, None);
assert!(!old.path().exists());
assert!(!metadata_path_for_session(&old).exists());
assert!(!old.path().with_extension("writer-owner").exists());
assert!(!history.exists());
assert!(fresh.path().exists());
assert!(metadata_path_for_session(&fresh).exists());
assert!(fresh.path().with_extension("writer-owner").exists());
assert!(locked.path().exists());
assert!(metadata_path_for_session(&locked).exists());
assert!(locked.path().with_extension("writer-owner").exists());
}
#[test]
fn auto_prune_orphan_sidecars_uses_age_and_keeps_live_or_fresh_entries() {
let temp = TempDir::new().unwrap();
let root = temp.path().join("sessions");
let manager = SessionManager::new(root.clone());
let live = manager.open("live").unwrap();
live.append(&SessionEvent::new(
"user_input",
live.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"live"}),
))
.unwrap();
fs::create_dir_all(&root).unwrap();
let old_orphan = root.join("old-orphan.metadata.json");
let fresh_orphan = root.join("fresh-orphan.metadata.json");
fs::write(&old_orphan, "{} ").unwrap();
fs::write(&fresh_orphan, "{} ").unwrap();
crate::sessions::store::secure_test_session_root(&root);
let now = SystemTime::now();
set_modified(&old_orphan, now - Duration::from_secs(2 * 86_400));
manager.auto_prune_sessions_for_test(now, 1, None);
assert!(!old_orphan.exists());
assert!(fresh_orphan.exists());
assert!(live.path().exists());
assert!(metadata_path_for_session(&live).exists());
}
#[test]
fn auto_prune_retries_orphan_history_and_checkpoint_cleanup() {
let temp = TempDir::new().unwrap();
let root = temp.path().join("sessions");
let manager = SessionManager::new(root.clone());
fs::create_dir_all(&root).unwrap();
crate::sessions::store::secure_test_session_root(&root);
let history = root.join(".history/orphaned");
let subagent_history = root.join("subagents/.history/orphaned-child");
fs::create_dir_all(&history).unwrap();
fs::create_dir_all(&subagent_history).unwrap();
fs::write(history.join("1.jsonl"), "archived").unwrap();
fs::write(subagent_history.join("1.jsonl"), "archived").unwrap();
let ledger = temp.path().join("checkpoints/ledgers/orphaned.jsonl");
fs::create_dir_all(ledger.parent().unwrap()).unwrap();
fs::write(&ledger, "{}").unwrap();
let now = SystemTime::now();
let old = now - Duration::from_secs(2 * 86_400);
fs::File::open(&history).unwrap().set_modified(old).unwrap();
fs::File::open(&subagent_history)
.unwrap()
.set_modified(old)
.unwrap();
set_modified(&ledger, old);
manager.auto_prune_sessions_for_test(now, 1, None);
assert!(!history.exists());
assert!(!subagent_history.exists());
assert!(!ledger.exists());
let orphan_blob = temp.path().join("checkpoints/blobs/unreferenced");
fs::create_dir_all(orphan_blob.parent().unwrap()).unwrap();
fs::write(&orphan_blob, "orphan").unwrap();
manager.auto_prune_sessions_for_test(now, 1, None);
assert!(!orphan_blob.exists());
}
#[test]
fn auto_prune_preserves_checkpoint_ledger_for_active_subagent_without_primary() {
let temp = TempDir::new().unwrap();
let root = temp.path().join("sessions");
let manager = SessionManager::new(root.clone());
let active = SessionManager::new(root.join("subagents"))
.create()
.unwrap()
.activate()
.unwrap();
let ledger = temp
.path()
.join("checkpoints/ledgers")
.join(format!("{}.jsonl", active.id()));
fs::create_dir_all(ledger.parent().unwrap()).unwrap();
fs::write(&ledger, "{}").unwrap();
let now = SystemTime::now();
set_modified(&ledger, now - Duration::from_secs(2 * 86_400));
manager.auto_prune_sessions_for_test(now, 1, None);
assert!(ledger.exists());
drop(active);
manager.auto_prune_sessions_for_test(now, 1, None);
assert!(!ledger.exists());
}
fn set_modified(path: &Path, modified: SystemTime) {
fs::OpenOptions::new()
.write(true)
.open(path)
.unwrap()
.set_modified(modified)
.unwrap();
}
fn lock_path_for_session(session: &Session) -> std::path::PathBuf {
let file_name = session.path().file_name().unwrap().to_string_lossy();
session.path().with_file_name(format!(".{file_name}.lock"))
}
fn append_event_at(session: &Session, cwd: &Path, year: i32, month: u32, day: u32) {
session
.append(&event_at(session, cwd, year, month, day))
.unwrap();
session.release_active_lease_for_test();
}
fn event_at(session: &Session, cwd: &Path, year: i32, month: u32, day: u32) -> SessionEvent {
let mut event = SessionEvent::new(
"diagnostic",
session.id().to_string(),
cwd.to_path_buf(),
json!({}),
);
event.timestamp = Utc.with_ymd_and_hms(year, month, day, 0, 0, 0).unwrap();
event
}
}