use crate::artifacts::ArtifactRecord;
use crate::config::ApiProvider;
use crate::model_routing::AutoRouteReceipt;
use crate::models::{ContentBlock, Message, SystemPrompt};
use crate::session_tree::{SessionEntry, SessionImportContainer, SessionJournal};
use crate::tools::plan::PlanSnapshot;
use crate::tools::todo::TodoListSnapshot;
use crate::tui::file_mention::ContextReference;
use crate::utils::write_atomic;
use crate::work_graph::ReasoningEffortTier;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, BTreeSet};
use std::fs;
use std::io;
use std::path::{Component, Path, PathBuf};
use uuid::Uuid;
const MAX_SESSIONS: usize = 50;
pub const MAX_SESSION_TITLE_CHARS: usize = 100;
const WORK_GRAPH_IMPORT_ARCHIVE_DIR: &str = ".work-graph-import-archive";
const CURRENT_SESSION_SCHEMA_VERSION: u32 = 1;
const CURRENT_QUEUE_SCHEMA_VERSION: u32 = 1;
const fn default_session_schema_version() -> u32 {
CURRENT_SESSION_SCHEMA_VERSION
}
const fn default_queue_schema_version() -> u32 {
CURRENT_QUEUE_SCHEMA_VERSION
}
fn normalize_managed_dir(path: PathBuf) -> std::io::Result<PathBuf> {
if path.as_os_str().is_empty() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"managed directory path cannot be empty",
));
}
if path.components().any(|component| {
matches!(
component,
Component::ParentDir | Component::Prefix(_) | Component::RootDir
)
}) && path.is_relative()
{
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"managed directory path cannot contain traversal components",
));
}
if path.is_absolute() {
return Ok(path);
}
std::env::current_dir().map(|cwd| cwd.join(path))
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct QueuedSessionMessage {
pub display: String,
#[serde(default)]
pub skill_instruction: Option<String>,
#[serde(default)]
pub skill_provenance: Option<crate::plugins::types::PluginAuthority>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct OfflineQueueState {
#[serde(default = "default_queue_schema_version")]
pub schema_version: u32,
#[serde(default)]
pub session_id: Option<String>,
#[serde(default)]
pub messages: Vec<QueuedSessionMessage>,
#[serde(default)]
pub draft: Option<QueuedSessionMessage>,
}
#[derive(Debug, Clone)]
pub struct SessionRecovery {
pub session: SavedSession,
pub changed: bool,
pub repaired_call_count: usize,
pub duplicate_result_count: usize,
pub orphan_result_count: usize,
}
impl Default for OfflineQueueState {
fn default() -> Self {
Self {
schema_version: CURRENT_QUEUE_SCHEMA_VERSION,
session_id: None,
messages: Vec::new(),
draft: None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SessionContextReference {
pub message_index: usize,
pub reference: ContextReference,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SessionMetadata {
pub id: String,
pub title: String,
pub created_at: DateTime<Utc>,
pub updated_at: DateTime<Utc>,
pub message_count: usize,
pub total_tokens: u64,
pub model: String,
#[serde(default = "default_model_provider")]
pub model_provider: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model_provider_id: Option<String>,
pub workspace: PathBuf,
#[serde(default)]
pub mode: Option<String>,
#[serde(default)]
pub cost: SessionCostSnapshot,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent_session_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub forked_from_message_count: Option<usize>,
#[serde(default)]
pub cumulative_turn_secs: u64,
#[serde(default, skip_serializing_if = "is_not_archived")]
pub archived: bool,
#[serde(default)]
pub spawn_depth: u32,
}
fn is_not_archived(archived: &bool) -> bool {
!*archived
}
static LIVE_SESSIONS: std::sync::OnceLock<std::sync::RwLock<std::collections::HashSet<String>>> =
std::sync::OnceLock::new();
fn live_sessions() -> &'static std::sync::RwLock<std::collections::HashSet<String>> {
LIVE_SESSIONS.get_or_init(Default::default)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SessionMutator {
Owner,
External,
}
pub fn set_live_session(session_id: Option<&str>) {
if let Ok(mut live) = live_sessions().write() {
live.clear();
if let Some(id) = session_id.map(str::trim).filter(|id| !id.is_empty()) {
live.insert(id.to_string());
}
}
}
#[must_use]
pub fn is_live_session(session_id: &str) -> bool {
live_sessions()
.read()
.is_ok_and(|live| live.contains(session_id))
}
fn live_session_conflict(session_id: &str) -> std::io::Error {
std::io::Error::new(
std::io::ErrorKind::ResourceBusy,
format!(
"session '{session_id}' is open in an interactive Codewhale session; \
change it there instead — an external write would be reverted by its next autosave"
),
)
}
const SESSION_BOOT_OWNERS_STEM: &str = "session_boot_owners";
static SESSION_BOOT_ID: std::sync::OnceLock<String> = std::sync::OnceLock::new();
#[must_use]
pub fn current_session_boot_id() -> &'static str {
SESSION_BOOT_ID.get_or_init(|| format!("boot_{}", &Uuid::new_v4().to_string()[..12]))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum SessionListFilter {
#[default]
ActiveOnly,
IncludeArchived,
ArchivedOnly,
}
impl SessionListFilter {
#[must_use]
pub fn from_query(include_archived: Option<bool>, archived_only: Option<bool>) -> Self {
if archived_only.unwrap_or(false) {
Self::ArchivedOnly
} else if include_archived.unwrap_or(false) {
Self::IncludeArchived
} else {
Self::ActiveOnly
}
}
#[must_use]
pub fn admits(self, archived: bool) -> bool {
match self {
Self::ActiveOnly => !archived,
Self::IncludeArchived => true,
Self::ArchivedOnly => archived,
}
}
}
fn default_model_provider() -> String {
"deepseek".to_string()
}
impl SessionMetadata {
pub(crate) fn set_model_provider_route(&mut self, kind: &str, identity: Option<&str>) {
self.model_provider = kind.to_string();
self.model_provider_id = identity.map(str::to_string);
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct SessionCostSnapshot {
#[serde(default)]
pub session_cost_usd: f64,
#[serde(default)]
pub session_cost_cny: f64,
#[serde(default)]
pub subagent_cost_usd: f64,
#[serde(default)]
pub subagent_cost_cny: f64,
#[serde(default)]
pub displayed_cost_high_water_usd: f64,
#[serde(default)]
pub displayed_cost_high_water_cny: f64,
#[serde(default)]
pub priced_turns: u32,
#[serde(default)]
pub unpriced_turns: u32,
#[serde(default)]
pub cny_priced_turns: u32,
#[serde(default)]
pub cny_unpriced_turns: u32,
#[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
pub unpriced_reasons: BTreeSet<String>,
#[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
pub cny_unpriced_reasons: BTreeSet<String>,
#[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
pub unpriced_classes: BTreeSet<String>,
#[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
pub pricing_provenances: BTreeSet<String>,
#[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
pub live_pricing_defects: BTreeSet<String>,
#[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
pub live_pricing_unusable_defects: BTreeSet<String>,
#[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
pub route_receipts: BTreeSet<String>,
#[serde(default)]
pub coverage_recorded: bool,
}
impl SessionCostSnapshot {
#[must_use]
pub fn total_estimate(&self) -> crate::pricing::CostEstimate {
crate::pricing::CostEstimate {
usd: self.session_cost_usd,
cny: self.session_cost_cny,
}
.saturating_add(crate::pricing::CostEstimate {
usd: self.subagent_cost_usd,
cny: self.subagent_cost_cny,
})
}
pub fn total_usd(&self) -> f64 {
self.total_estimate()
.amount(crate::pricing::CostCurrency::Usd)
}
pub fn total_cny(&self) -> f64 {
self.total_estimate()
.amount(crate::pricing::CostCurrency::Cny)
}
#[must_use]
pub fn coverage_is_legacy_unknown(&self) -> bool {
!self.coverage_recorded
}
}
impl SessionMetadata {
#[allow(dead_code)]
pub fn copy_cost_from(&mut self, other: &SessionMetadata) {
self.cost = other.cost.clone();
}
pub fn mark_forked_from(&mut self, parent: &SessionMetadata) {
self.parent_session_id = Some(parent.id.clone());
self.forked_from_message_count = Some(parent.message_count);
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
pub struct SessionWorkState {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub graph: Option<crate::work_graph::WorkGraphSnapshot>,
#[serde(default, skip_serializing_if = "TodoListSnapshot::is_empty")]
pub todos: TodoListSnapshot,
#[serde(default, skip_serializing_if = "PlanSnapshot::is_empty")]
pub plan: PlanSnapshot,
}
impl SessionWorkState {
#[must_use]
pub fn is_empty(&self) -> bool {
self.graph
.as_ref()
.is_none_or(crate::work_graph::WorkGraphSnapshot::is_empty)
&& self.todos.is_empty()
&& self.plan.is_empty()
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub(crate) struct SavedAutoRouteReceipt {
pub(crate) provider: ApiProvider,
pub(crate) provider_identity: String,
pub(crate) model: String,
pub(crate) receipt: AutoRouteReceipt,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub(crate) effective_reasoning_effort: Option<ReasoningEffortTier>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SavedSession {
#[serde(default = "default_session_schema_version")]
pub schema_version: u32,
pub metadata: SessionMetadata,
pub messages: Vec<Message>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub journal: Option<SessionJournal>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub leaf_id: Option<String>,
pub system_prompt: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub context_references: Vec<SessionContextReference>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub artifacts: Vec<ArtifactRecord>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub work_state: Option<SessionWorkState>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub(crate) last_auto_route: Option<SavedAutoRouteReceipt>,
}
impl SavedSession {
pub(crate) fn compact_for_persistence_queue(&mut self) {
if self.journal.is_some() {
self.messages = Vec::new();
}
}
fn storage_compatible_copy(&self) -> Option<Self> {
let journal = self.journal.as_ref()?;
let active_messages = journal.to_messages();
if !self.messages.is_empty() && self.messages == active_messages {
return None;
}
let mut copy = self.clone();
if copy.messages.is_empty() {
copy.messages = active_messages;
} else if let Some(journal) = copy.journal.as_mut() {
journal.rebranch_active_messages(©.messages);
copy.leaf_id = journal.leaf_id.clone();
}
copy.metadata.message_count = copy.messages.len();
Some(copy)
}
pub fn ensure_journal(&mut self) {
if self.journal.is_some() {
if self.leaf_id.is_none() {
self.leaf_id = self.journal.as_ref().and_then(|j| j.leaf_id.clone());
}
let active = self
.journal
.as_ref()
.map(|j| j.to_messages())
.unwrap_or_default();
if !active.is_empty() {
self.messages = active;
self.metadata.message_count = self.messages.len();
}
return;
}
let journal =
SessionJournal::from_messages(self.messages.clone(), self.metadata.spawn_depth);
self.leaf_id = journal.leaf_id.clone();
self.journal = Some(journal);
}
pub fn journal_append_message(&mut self, message: Message) -> String {
self.ensure_journal();
let journal = self.journal.as_mut().expect("journal ensured");
let id = journal.append_message(message.clone());
self.leaf_id = journal.leaf_id.clone();
self.messages = journal.to_messages();
self.metadata.message_count = self.messages.len();
self.metadata.updated_at = Utc::now();
id
}
pub fn journal_branch_to(&mut self, entry_id: &str) -> Result<(), String> {
self.ensure_journal();
let journal = self.journal.as_mut().expect("journal ensured");
journal.branch_to(entry_id)?;
self.leaf_id = journal.leaf_id.clone();
self.messages = journal.to_messages();
self.metadata.updated_at = Utc::now();
Ok(())
}
pub fn active_entries(&self) -> Vec<SessionEntry> {
self.journal
.as_ref()
.map(|j| j.root_to_leaf().into_iter().cloned().collect())
.unwrap_or_default()
}
pub fn export_container(&self, source: &str) -> SessionImportContainer {
let journal = self.journal.clone().unwrap_or_else(|| {
SessionJournal::from_messages(self.messages.clone(), self.metadata.spawn_depth)
});
SessionImportContainer::new(
source.to_string(),
&journal,
serde_json::to_value(&self.metadata).ok(),
)
}
pub fn import_foreign(
container: SessionImportContainer,
workspace: PathBuf,
model: String,
) -> Result<Self, String> {
let journal = container.into_journal()?;
let leaf_id = journal.leaf_id.clone();
let messages = journal.to_messages();
let now = Utc::now();
let spawn_depth = journal.spawn_depth.saturating_add(1);
let title = messages
.iter()
.find(|m| m.role == "user")
.and_then(|m| {
m.content.iter().find_map(|b| match b {
ContentBlock::Text { text, .. } => Some(text.as_str()),
_ => None,
})
})
.map(|s| crate::session_manager::truncate_title(s, 50))
.unwrap_or_else(|| crate::session_manager::DEFAULT_SESSION_TITLE.to_string());
let metadata = SessionMetadata {
id: Uuid::new_v4().to_string(),
title,
created_at: now,
updated_at: now,
message_count: messages.len(),
total_tokens: 0,
model,
model_provider: default_model_provider(),
model_provider_id: None,
workspace,
mode: None,
cost: SessionCostSnapshot::default(),
parent_session_id: None,
forked_from_message_count: None,
cumulative_turn_secs: 0,
archived: false,
spawn_depth,
};
let mut journal = journal;
journal.spawn_depth = spawn_depth;
Ok(Self {
schema_version: CURRENT_SESSION_SCHEMA_VERSION,
metadata,
messages,
journal: Some(journal),
leaf_id,
system_prompt: None,
context_references: Vec::new(),
artifacts: Vec::new(),
work_state: None,
last_auto_route: None,
})
}
}
fn serialize_saved_session(session: &SavedSession) -> io::Result<String> {
let compatible = session.storage_compatible_copy();
serde_json::to_string_pretty(compatible.as_ref().unwrap_or(session))
.map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))
}
#[derive(Debug)]
pub struct SessionManager {
sessions_dir: PathBuf,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum CheckpointSource {
Session(String),
Legacy,
}
#[derive(Debug, Clone)]
pub struct CheckpointRef {
pub source: CheckpointSource,
pub path: PathBuf,
pub modified: std::time::SystemTime,
}
const LEGACY_CHECKPOINT_FILE: &str = "latest.json";
const OFFLINE_QUEUE_FILE: &str = "offline_queue.json";
impl SessionManager {
fn validated_session_id<'a>(&self, id: &'a str) -> std::io::Result<&'a str> {
let trimmed = id.trim();
if trimmed.is_empty() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"Session id cannot be empty",
));
}
if !trimmed
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_')
{
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!("Invalid session id '{id}'"),
));
}
if trimmed == SESSION_BOOT_OWNERS_STEM {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!("Session id '{trimmed}' collides with a reserved sessions file"),
));
}
Ok(trimmed)
}
fn validated_session_path(&self, id: &str) -> std::io::Result<PathBuf> {
let trimmed = self.validated_session_id(id)?;
Ok(self.sessions_dir.join(format!("{trimmed}.json")))
}
fn checkpoints_dir(&self) -> PathBuf {
self.sessions_dir.join("checkpoints")
}
fn validated_checkpoint_path(&self, session_id: &str) -> std::io::Result<PathBuf> {
let trimmed = self.validated_session_id(session_id)?;
if format!("{trimmed}.json") == LEGACY_CHECKPOINT_FILE
|| format!("{trimmed}.json") == OFFLINE_QUEUE_FILE
{
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!("Session id '{trimmed}' collides with a reserved checkpoint file"),
));
}
Ok(self.checkpoints_dir().join(format!("{trimmed}.json")))
}
pub fn new(sessions_dir: PathBuf) -> std::io::Result<Self> {
let sessions_dir = normalize_managed_dir(sessions_dir)?;
fs::create_dir_all(&sessions_dir)?;
Ok(Self { sessions_dir })
}
pub fn default_location() -> std::io::Result<Self> {
Self::new(default_sessions_dir()?)
}
pub fn sessions_dir(&self) -> &Path {
&self.sessions_dir
}
pub fn save_session(&self, session: &SavedSession) -> std::io::Result<PathBuf> {
let path = self.validated_session_path(&session.metadata.id)?;
let already_persisted = path.exists()
|| self
.validated_checkpoint_path(&session.metadata.id)
.is_ok_and(|checkpoint| checkpoint.exists());
self.archive_before_first_graph_write(session, &path)?;
let content = serialize_saved_session(session)?;
write_atomic(&path, content.as_bytes())?;
self.stamp_session_boot_owner_for_new_record(&session.metadata.id, already_persisted);
self.cleanup_old_sessions()?;
Ok(path)
}
pub fn save_checkpoint(&self, session: &SavedSession) -> std::io::Result<PathBuf> {
let path = self.validated_checkpoint_path(&session.metadata.id)?;
let session_path = self.validated_session_path(&session.metadata.id)?;
self.archive_before_first_graph_write(session, &session_path)?;
fs::create_dir_all(self.checkpoints_dir())?;
let already_persisted = path.exists() || session_path.exists();
let content = serialize_saved_session(session)?;
write_atomic(&path, content.as_bytes())?;
self.stamp_session_boot_owner_for_new_record(&session.metadata.id, already_persisted);
Ok(path)
}
fn session_boot_owners_path(&self) -> PathBuf {
self.sessions_dir
.join(format!("{SESSION_BOOT_OWNERS_STEM}.json"))
}
fn load_session_boot_owners(&self) -> BTreeMap<String, String> {
fs::read_to_string(self.session_boot_owners_path())
.ok()
.and_then(|content| serde_json::from_str(&content).ok())
.unwrap_or_default()
}
fn session_record_exists(&self, session_id: &str) -> bool {
self.validated_session_path(session_id)
.is_ok_and(|path| path.exists())
|| self
.validated_checkpoint_path(session_id)
.is_ok_and(|path| path.exists())
}
pub(crate) fn record_session_boot_owner(
&self,
session_id: &str,
boot_id: &str,
) -> std::io::Result<()> {
let id = self.validated_session_id(session_id)?.to_string();
let mut owners = self.load_session_boot_owners();
owners.retain(|owned, _| owned == &id || self.session_record_exists(owned));
owners.insert(id, boot_id.to_string());
let content = serde_json::to_string_pretty(&owners)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
write_atomic(&self.session_boot_owners_path(), content.as_bytes())
}
#[must_use]
pub fn session_boot_owner(&self, session_id: &str) -> Option<String> {
let id = self.validated_session_id(session_id).ok()?;
self.load_session_boot_owners().get(id).cloned()
}
#[must_use]
pub fn session_from_prior_instance(&self, session_id: &str) -> bool {
match self.session_boot_owner(session_id) {
Some(owner) => owner != current_session_boot_id(),
None => self.session_record_exists(session_id),
}
}
fn stamp_session_boot_owner_for_new_record(&self, session_id: &str, already_persisted: bool) {
if already_persisted || self.session_boot_owner(session_id).is_some() {
return;
}
if let Err(error) = self.record_session_boot_owner(session_id, current_session_boot_id()) {
tracing::warn!(session_id, %error, "could not stamp session boot owner");
}
}
fn clear_session_boot_owner(&self, session_id: &str) {
let Ok(id) = self.validated_session_id(session_id) else {
return;
};
let mut owners = self.load_session_boot_owners();
if owners.remove(id).is_none() {
return;
}
if let Ok(content) = serde_json::to_string_pretty(&owners) {
let _ = write_atomic(&self.session_boot_owners_path(), content.as_bytes());
}
}
fn archive_before_first_graph_write(
&self,
session: &SavedSession,
source: &Path,
) -> std::io::Result<()> {
let writes_graph = session
.work_state
.as_ref()
.and_then(|state| state.graph.as_ref())
.is_some_and(|graph| !graph.is_empty());
if !writes_graph || !source.exists() {
return Ok(());
}
let bytes = fs::read(source)?;
let already_graph_backed = serde_json::from_slice::<SavedSession>(&bytes)
.ok()
.and_then(|saved| saved.work_state)
.and_then(|state| state.graph)
.is_some_and(|graph| !graph.is_empty());
if already_graph_backed {
return Ok(());
}
let archive_dir = self.sessions_dir.join(WORK_GRAPH_IMPORT_ARCHIVE_DIR);
fs::create_dir_all(&archive_dir)?;
let archive =
archive_dir.join(source.file_name().ok_or_else(|| {
io::Error::new(io::ErrorKind::InvalidInput, "invalid session path")
})?);
if !archive.exists() {
write_atomic(&archive, &bytes)?;
}
Ok(())
}
fn read_checkpoint_file(&self, path: &Path) -> std::io::Result<Option<SavedSession>> {
if !path.exists() {
return Ok(None);
}
let content = fs::read_to_string(path)?;
let mut session: SavedSession = serde_json::from_str(&content)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
if session.schema_version > CURRENT_SESSION_SCHEMA_VERSION {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!(
"Checkpoint schema v{} is newer than supported v{}",
session.schema_version, CURRENT_SESSION_SCHEMA_VERSION
),
));
}
session.system_prompt = strip_legacy_truncation_note(session.system_prompt);
Ok(Some(session))
}
pub fn load_session_checkpoint(
&self,
session_id: &str,
) -> std::io::Result<Option<SavedSession>> {
let path = self.validated_checkpoint_path(session_id)?;
self.read_checkpoint_file(&path)
}
pub fn load_legacy_checkpoint(&self) -> std::io::Result<Option<SavedSession>> {
let path = self.checkpoints_dir().join(LEGACY_CHECKPOINT_FILE);
self.read_checkpoint_file(&path)
}
pub fn clear_session_checkpoint(&self, session_id: &str) -> std::io::Result<()> {
let path = self.validated_checkpoint_path(session_id)?;
if path.exists() {
fs::remove_file(path)?;
}
Ok(())
}
pub fn clear_legacy_checkpoint(&self) -> std::io::Result<()> {
let path = self.checkpoints_dir().join(LEGACY_CHECKPOINT_FILE);
if path.exists() {
fs::remove_file(path)?;
}
Ok(())
}
pub fn list_checkpoints(&self) -> std::io::Result<Vec<CheckpointRef>> {
let dir = self.checkpoints_dir();
let mut refs = Vec::new();
let entries = match fs::read_dir(&dir) {
Ok(entries) => entries,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(refs),
Err(err) => return Err(err),
};
for entry in entries {
let entry = entry?;
let path = entry.path();
if !path.is_file() || path.extension().is_none_or(|ext| ext != "json") {
continue;
}
let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
continue;
};
let source = if name == LEGACY_CHECKPOINT_FILE {
CheckpointSource::Legacy
} else if name == OFFLINE_QUEUE_FILE {
continue;
} else {
let session_id = name.trim_end_matches(".json").to_string();
if self.validated_checkpoint_path(&session_id).is_err() {
continue;
}
CheckpointSource::Session(session_id)
};
let Ok(modified) = entry.metadata().and_then(|m| m.modified()) else {
continue;
};
refs.push(CheckpointRef {
source,
path,
modified,
});
}
refs.sort_by_key(|r| std::cmp::Reverse(r.modified));
Ok(refs)
}
pub fn write_session_checkpoint_if_absent(
&self,
session: &SavedSession,
) -> std::io::Result<bool> {
let path = self.validated_checkpoint_path(&session.metadata.id)?;
if path.exists() {
return Ok(false);
}
self.save_checkpoint(session)?;
Ok(true)
}
pub fn save_offline_queue_state(
&self,
state: &OfflineQueueState,
session_id: Option<&str>,
) -> std::io::Result<PathBuf> {
let checkpoints = self.sessions_dir.join("checkpoints");
fs::create_dir_all(&checkpoints)?;
let path = checkpoints.join("offline_queue.json");
let mut state_with_id = state.clone();
state_with_id.session_id = session_id.map(|s| s.to_string());
let content = serde_json::to_string_pretty(&state_with_id)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
write_atomic(&path, content.as_bytes())?;
Ok(path)
}
pub fn load_offline_queue_state(&self) -> std::io::Result<Option<OfflineQueueState>> {
let path = self
.sessions_dir
.join("checkpoints")
.join("offline_queue.json");
if !path.exists() {
return Ok(None);
}
let content = fs::read_to_string(&path)?;
let state: OfflineQueueState = serde_json::from_str(&content)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
if state.schema_version > CURRENT_QUEUE_SCHEMA_VERSION {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!(
"Offline queue schema v{} is newer than supported v{}",
state.schema_version, CURRENT_QUEUE_SCHEMA_VERSION
),
));
}
Ok(Some(state))
}
pub fn clear_offline_queue_state(&self) -> std::io::Result<()> {
let path = self
.sessions_dir
.join("checkpoints")
.join("offline_queue.json");
if path.exists() {
fs::remove_file(path)?;
}
Ok(())
}
pub fn load_session_snapshot(&self, id: &str) -> std::io::Result<SavedSession> {
let path = self.validated_session_path(id)?;
let content = fs::read_to_string(&path)?;
let mut session: SavedSession = serde_json::from_str(&content)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
if session.schema_version > CURRENT_SESSION_SCHEMA_VERSION {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!(
"Session schema v{} is newer than supported v{}",
session.schema_version, CURRENT_SESSION_SCHEMA_VERSION
),
));
}
session.system_prompt = strip_legacy_truncation_note(session.system_prompt);
session.ensure_journal();
Ok(session)
}
pub fn recover_session_for_resume(&self, id: &str) -> std::io::Result<SessionRecovery> {
let mut session = self.load_session_snapshot(id)?;
let repair = crate::tool_history_repair::repair_tool_call_pairs(&mut session.messages);
let changed = !repair.is_empty();
if changed {
if let Some(journal) = session.journal.as_mut() {
journal.rebranch_active_messages(&session.messages);
session.leaf_id = journal.leaf_id.clone();
}
session.metadata.message_count = session.messages.len();
tracing::warn!(
session_id = %session.metadata.id,
repaired_call_ids = ?repair.repaired_call_ids,
duplicate_result_ids = ?repair.duplicate_result_ids,
orphan_result_ids = ?repair.orphan_result_ids,
"repaired persisted tool call/result history"
);
}
Ok(SessionRecovery {
session,
changed,
repaired_call_count: repair.repaired_call_ids.len(),
duplicate_result_count: repair.duplicate_result_ids.len(),
orphan_result_count: repair.orphan_result_ids.len(),
})
}
pub fn load_session(&self, id: &str) -> std::io::Result<SavedSession> {
self.recover_session_for_resume(id)
.map(|recovery| recovery.session)
}
pub fn load_session_by_prefix(&self, prefix: &str) -> std::io::Result<SavedSession> {
let sessions = self.list_sessions()?;
let matches: Vec<_> = sessions
.into_iter()
.filter(|s| s.id.starts_with(prefix))
.collect();
match matches.len() {
0 => Err(std::io::Error::new(
std::io::ErrorKind::NotFound,
format!("No session found with prefix: {prefix}"),
)),
1 => self.load_session(&matches[0].id),
_ => Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!(
"Ambiguous prefix '{}' matches {} sessions",
prefix,
matches.len()
),
)),
}
}
pub fn list_sessions(&self) -> std::io::Result<Vec<SessionMetadata>> {
let mut sessions = Vec::new();
for entry in fs::read_dir(&self.sessions_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().is_some_and(|ext| ext == "json")
&& let Ok(session) = Self::load_session_metadata(&path)
{
sessions.push(session);
}
}
sessions.sort_by_key(|s| std::cmp::Reverse(s.updated_at));
Ok(sessions)
}
pub fn set_session_archived(
&self,
id: &str,
archived: bool,
mutator: SessionMutator,
) -> std::io::Result<SessionMetadata> {
if mutator == SessionMutator::External && is_live_session(id) {
return Err(live_session_conflict(id));
}
let mut session = self.load_session(id)?;
if session.metadata.archived == archived {
return Ok(session.metadata);
}
session.metadata.archived = archived;
self.save_session(&session)?;
Ok(session.metadata)
}
pub fn merge_persisted_lifecycle(&self, metadata: &mut SessionMetadata) -> bool {
let Ok(path) = self.validated_session_path(&metadata.id) else {
return false;
};
let Ok(persisted) = Self::load_session_metadata(&path) else {
return false;
};
metadata.title = persisted.title;
metadata.archived = persisted.archived;
metadata.created_at = persisted.created_at;
metadata.parent_session_id = persisted.parent_session_id;
metadata.forked_from_message_count = persisted.forked_from_message_count;
true
}
pub fn rename_session(
&self,
id: &str,
title: &str,
mutator: SessionMutator,
) -> std::io::Result<SessionMetadata> {
let title = normalize_session_title(title)?;
if mutator == SessionMutator::External && is_live_session(id) {
return Err(live_session_conflict(id));
}
let mut session = self.load_session(id)?;
if session.metadata.title == title {
return Ok(session.metadata);
}
session.metadata.title = title;
self.save_session(&session)?;
Ok(session.metadata)
}
fn load_session_metadata(path: &Path) -> std::io::Result<SessionMetadata> {
use std::io::Read;
const PREFIX_BYTES: usize = 64 * 1024;
let mut file = fs::File::open(path)?;
let mut buf = Vec::with_capacity(PREFIX_BYTES);
file.by_ref()
.take(PREFIX_BYTES as u64)
.read_to_end(&mut buf)?;
if let Some(metadata) = extract_top_level_metadata(&buf) {
return Ok(metadata);
}
let mut rest = Vec::new();
file.read_to_end(&mut rest)?;
buf.extend_from_slice(&rest);
extract_top_level_metadata(&buf).ok_or_else(|| {
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"session file missing parseable `metadata` block",
)
})
}
pub fn delete_session(&self, id: &str) -> std::io::Result<()> {
let path = self.validated_session_path(id)?;
fs::remove_file(path)?;
self.clear_session_boot_owner(id);
let session_dir = self.sessions_dir.join(id.trim());
if session_dir.exists() {
fs::remove_dir_all(session_dir)?;
}
Ok(())
}
pub fn cleanup_old_sessions(&self) -> std::io::Result<()> {
self.cleanup_old_sessions_keeping(None)
}
pub fn cleanup_old_sessions_keeping(&self, keep: Option<&str>) -> std::io::Result<()> {
let sessions = self.list_sessions()?;
if sessions.len() > MAX_SESSIONS {
for session in sessions.iter().skip(MAX_SESSIONS) {
if keep.is_some_and(|id| id == session.id) {
continue;
}
let _ = self.delete_session(&session.id);
}
}
Ok(())
}
pub fn prune_sessions_older_than(
&self,
max_age: std::time::Duration,
) -> std::io::Result<usize> {
self.prune_sessions_older_than_keeping(max_age, None)
}
pub fn prune_sessions_older_than_keeping(
&self,
max_age: std::time::Duration,
keep: Option<&str>,
) -> std::io::Result<usize> {
let cutoff = Utc::now()
- chrono::Duration::from_std(max_age).unwrap_or(chrono::Duration::days(365 * 10));
let sessions = self.list_sessions()?;
let mut pruned = 0usize;
for session in sessions {
if keep.is_some_and(|id| id == session.id) {
continue;
}
if session.updated_at < cutoff {
if let Err(err) = self.delete_session(&session.id) {
tracing::warn!(
target: "session",
session = session.id,
?err,
"session prune skipped a record",
);
continue;
}
pruned += 1;
}
}
Ok(pruned)
}
pub fn get_latest_session_for_workspace(
&self,
workspace: &Path,
) -> std::io::Result<Option<SessionMetadata>> {
let sessions = self.list_sessions()?;
Ok(sessions.into_iter().find(|session| {
!session.archived
&& workspace_scope_matches(&session.workspace, workspace)
&& !is_empty_auto_created_session(session)
}))
}
pub fn search_sessions(&self, query: &str) -> std::io::Result<Vec<SessionMetadata>> {
let query_lower = query.to_lowercase();
let sessions = self.list_sessions()?;
Ok(sessions
.into_iter()
.filter(|s| s.title.to_lowercase().contains(&query_lower))
.collect())
}
}
pub(crate) fn is_title_format_char(ch: char) -> bool {
matches!(
ch,
'\u{00ad}'
| '\u{061c}'
| '\u{200b}'..='\u{200f}'
| '\u{2028}'..='\u{202e}'
| '\u{2060}'..='\u{2064}'
| '\u{2066}'..='\u{2069}'
| '\u{feff}'
)
}
pub fn sanitize_session_title(raw: &str) -> String {
raw.chars()
.filter(|ch| !ch.is_control() && !is_title_format_char(*ch))
.collect()
}
pub fn normalize_session_title(title: &str) -> std::io::Result<String> {
let sanitized = sanitize_session_title(title);
let trimmed = sanitized.trim();
if trimmed.is_empty() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"Session title cannot be empty",
));
}
if trimmed.chars().count() > MAX_SESSION_TITLE_CHARS {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!("Session title cannot exceed {MAX_SESSION_TITLE_CHARS} characters"),
));
}
Ok(trimmed.to_string())
}
pub(crate) fn workspace_scope_matches(saved_workspace: &Path, current_workspace: &Path) -> bool {
if paths_equivalent(saved_workspace, current_workspace) {
return true;
}
match (
find_git_root(saved_workspace),
find_git_root(current_workspace),
) {
(Some(saved_root), Some(current_root)) => paths_equivalent(&saved_root, ¤t_root),
_ => false,
}
}
fn is_empty_auto_created_session(session: &SessionMetadata) -> bool {
session.message_count == 0
&& session
.title
.trim()
.eq_ignore_ascii_case(DEFAULT_SESSION_TITLE)
}
fn paths_equivalent(lhs: &Path, rhs: &Path) -> bool {
let lhs_canonical = fs::canonicalize(lhs).ok();
let rhs_canonical = fs::canonicalize(rhs).ok();
match (lhs_canonical, rhs_canonical) {
(Some(lhs), Some(rhs)) => lhs == rhs,
_ => lhs == rhs,
}
}
fn find_git_root(path: &Path) -> Option<PathBuf> {
let mut current = fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf());
loop {
let git_entry = current.join(".git");
if git_entry.exists() {
return is_git_metadata_entry(&git_entry).then_some(current);
}
match current.parent() {
Some(parent) if parent != current => current = parent.to_path_buf(),
_ => return None,
}
}
}
fn is_git_metadata_entry(path: &Path) -> bool {
if path.is_dir() {
return path.join("HEAD").is_file();
}
fs::read_to_string(path)
.map(|content| content.trim_start().starts_with("gitdir:"))
.unwrap_or(false)
}
pub fn default_sessions_dir() -> std::io::Result<PathBuf> {
let dir = codewhale_config::ensure_state_dir("sessions")
.map_err(|e| std::io::Error::new(std::io::ErrorKind::NotFound, e.to_string()))?;
match merge_missing_legacy_session_entries(&dir) {
Ok(0) => {}
Ok(count) => {
tracing::info!(
target: "session::migration",
"Copied {count} missing legacy session entries into {}",
dir.display()
);
}
Err(err) => {
tracing::warn!(
target: "session::migration",
"Could not copy legacy sessions into {}: {err}",
dir.display()
);
}
}
Ok(dir)
}
fn merge_missing_legacy_session_entries(primary: &Path) -> io::Result<usize> {
if codewhale_paths::codewhale_home_is_explicit() {
return Ok(0);
}
let legacy = codewhale_config::legacy_deepseek_home()
.map_err(|e| io::Error::new(io::ErrorKind::NotFound, e.to_string()))?
.join("sessions");
if !legacy.is_dir() || paths_equivalent(primary, &legacy) {
return Ok(0);
}
copy_missing_dir_entries(&legacy, primary)
}
fn copy_missing_dir_entries(src: &Path, dst: &Path) -> io::Result<usize> {
fs::create_dir_all(dst)?;
let mut copied = 0;
for entry in fs::read_dir(src)? {
let entry = entry?;
let source = entry.path();
let target = dst.join(entry.file_name());
let file_type = entry.file_type()?;
if file_type.is_dir() {
if entry.file_name() == std::ffi::OsStr::new("checkpoints") || target.exists() {
continue;
}
copied += copy_missing_dir_entries(&source, &target)?;
} else if file_type.is_file() {
copied += usize::from(copy_file_create_new(&source, &target)?);
}
}
Ok(copied)
}
fn copy_file_create_new(src: &Path, dst: &Path) -> io::Result<bool> {
let mut source = fs::File::open(src)?;
let mut target = match fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(dst)
{
Ok(file) => file,
Err(err) if err.kind() == io::ErrorKind::AlreadyExists => return Ok(false),
Err(err) => return Err(err),
};
if let Err(err) = io::copy(&mut source, &mut target) {
let _ = fs::remove_file(dst);
return Err(err);
}
Ok(true)
}
pub fn prune_workspace_snapshots(workspace: &Path, max_age: std::time::Duration) {
match crate::snapshot::prune_older_than(workspace, max_age) {
Ok(0) => {}
Ok(n) => {
tracing::debug!(target: "snapshot", "boot prune removed {n} snapshot(s)");
}
Err(e) => {
tracing::warn!(target: "snapshot", "boot prune failed: {e}");
}
}
}
pub fn create_saved_session(
messages: &[Message],
model: &str,
workspace: &Path,
total_tokens: u64,
system_prompt: Option<&SystemPrompt>,
) -> SavedSession {
create_saved_session_with_mode(
messages,
model,
workspace,
total_tokens,
system_prompt,
None,
)
}
pub(crate) const DEFAULT_SESSION_TITLE: &str = "New Session";
pub fn create_saved_session_with_mode(
messages: &[Message],
model: &str,
workspace: &Path,
total_tokens: u64,
system_prompt: Option<&SystemPrompt>,
mode: Option<&str>,
) -> SavedSession {
create_saved_session_with_id_and_mode(
Uuid::new_v4().to_string(),
messages,
model,
workspace,
total_tokens,
system_prompt,
mode,
)
}
pub fn create_saved_session_with_id_and_mode(
id: String,
messages: &[Message],
model: &str,
workspace: &Path,
total_tokens: u64,
system_prompt: Option<&SystemPrompt>,
mode: Option<&str>,
) -> SavedSession {
let now = Utc::now();
let title = messages
.iter()
.find(|m| m.role == "user")
.and_then(|m| {
m.content.iter().find_map(|block| match block {
ContentBlock::Text { text, .. } => {
let prompt = extract_user_prompt(text);
if prompt.is_empty() {
None
} else {
Some(truncate_title(prompt, 50))
}
}
_ => None,
})
})
.unwrap_or_else(|| DEFAULT_SESSION_TITLE.to_string());
let journal = SessionJournal::from_messages(messages.to_vec(), 0);
let leaf_id = journal.leaf_id.clone();
SavedSession {
schema_version: CURRENT_SESSION_SCHEMA_VERSION,
metadata: SessionMetadata {
id,
title,
created_at: now,
updated_at: now,
message_count: messages.len(),
total_tokens,
model: model.to_string(),
model_provider: default_model_provider(),
model_provider_id: None,
workspace: workspace.to_path_buf(),
mode: mode.map(str::to_string),
cost: SessionCostSnapshot::default(),
parent_session_id: None,
forked_from_message_count: None,
cumulative_turn_secs: 0,
archived: false,
spawn_depth: 0,
},
messages: messages.to_vec(),
journal: Some(journal),
leaf_id,
system_prompt: system_prompt_to_string(system_prompt),
context_references: Vec::new(),
artifacts: Vec::new(),
work_state: None,
last_auto_route: None,
}
}
pub fn update_session(
mut session: SavedSession,
messages: &[Message],
total_tokens: u64,
system_prompt: Option<&SystemPrompt>,
) -> SavedSession {
session.schema_version = CURRENT_SESSION_SCHEMA_VERSION;
session.ensure_journal();
let old_len = session.messages.len();
let new_len = messages.len();
if new_len >= old_len && messages[..old_len] == session.messages[..] {
if let Some(journal) = session.journal.as_mut() {
for msg in &messages[old_len..] {
journal.append_message(msg.clone());
}
session.leaf_id = journal.leaf_id.clone();
}
} else if (new_len != old_len || messages != session.messages.as_slice())
&& let Some(journal) = session.journal.as_mut()
{
let common = messages
.iter()
.zip(session.messages.iter())
.take_while(|(a, b)| a == b)
.count();
if common > 0 && common <= journal.entries.len() {
let target_id = journal
.root_to_leaf()
.get(common - 1)
.map(|entry| entry.id.clone());
if let Some(target_id) = target_id {
let _ = journal.branch_to(&target_id);
} else {
journal.leaf_id = None;
}
} else if common == 0 {
journal.leaf_id = journal.entries.first().and_then(|e| e.parent_id.clone());
if journal.leaf_id.is_none() && !journal.entries.is_empty() {
journal.leaf_id = None;
}
}
for msg in messages.iter().skip(common) {
journal.append_message(msg.clone());
}
session.leaf_id = journal.leaf_id.clone();
}
session.messages.clear();
session.messages.extend_from_slice(messages);
session.metadata.updated_at = Utc::now();
session.metadata.message_count = messages.len();
session.metadata.total_tokens = total_tokens;
session.system_prompt = system_prompt_to_string(system_prompt);
session
}
fn strip_legacy_truncation_note(system_prompt: Option<String>) -> Option<String> {
let sp = system_prompt?;
let Some(trimmed) = sp.strip_prefix("[Session note]\n") else {
return Some(sp);
};
if !trimmed.contains("older messages were dropped") {
return Some(sp);
}
trimmed
.find("\n\n---\n\n")
.map(|pos| trimmed[pos + 7..].to_string())
}
fn extract_top_level_metadata(buf: &[u8]) -> Option<SessionMetadata> {
let s = std::str::from_utf8(buf).ok()?;
let bytes = s.as_bytes();
let key_pat = b"\"metadata\"";
let mut idx = 0usize;
let mut in_string = false;
let mut escape = false;
let key_offset = loop {
if idx >= bytes.len() {
return None;
}
let c = bytes[idx];
if escape {
escape = false;
idx += 1;
continue;
}
if c == b'\\' {
escape = true;
idx += 1;
continue;
}
if c == b'"' {
if !in_string && bytes[idx..].starts_with(key_pat) {
break idx;
}
in_string = !in_string;
idx += 1;
continue;
}
idx += 1;
};
let after_key = key_offset + key_pat.len();
let mut after_colon = after_key;
while after_colon < bytes.len() && (bytes[after_colon] as char).is_whitespace() {
after_colon += 1;
}
if after_colon >= bytes.len() || bytes[after_colon] != b':' {
return None;
}
after_colon += 1;
while after_colon < bytes.len() && (bytes[after_colon] as char).is_whitespace() {
after_colon += 1;
}
if after_colon >= bytes.len() || bytes[after_colon] != b'{' {
return None;
}
let mut depth = 0i32;
let mut in_string = false;
let mut escape = false;
let mut end = None;
for (i, &c) in bytes[after_colon..].iter().enumerate() {
let abs = after_colon + i;
if escape {
escape = false;
continue;
}
if c == b'\\' {
escape = true;
continue;
}
if c == b'"' {
in_string = !in_string;
continue;
}
if in_string {
continue;
}
match c {
b'{' => depth += 1,
b'}' => {
depth -= 1;
if depth == 0 {
end = Some(abs + 1);
break;
}
}
_ => {}
}
}
let end = end?;
serde_json::from_str::<SessionMetadata>(&s[after_colon..end]).ok()
}
fn system_prompt_to_string(system_prompt: Option<&SystemPrompt>) -> Option<String> {
match system_prompt {
Some(SystemPrompt::Text(text)) => Some(text.clone()),
Some(SystemPrompt::Blocks(blocks)) => Some(
blocks
.iter()
.map(|b| b.text.clone())
.collect::<Vec<_>>()
.join("\n\n---\n\n"),
),
None => None,
}
}
pub fn truncate_id(id: &str) -> &str {
id.get(..8).unwrap_or(id)
}
pub(crate) fn extract_user_prompt(raw: &str) -> &str {
let trimmed = raw.trim_start();
let Some(after_open) = trimmed.strip_prefix("<turn_meta>") else {
return trimmed;
};
if let Some(close_pos) = after_open.find("</turn_meta>") {
return after_open[close_pos + "</turn_meta>".len()..].trim_start();
}
after_open.trim_start()
}
pub(crate) fn extract_title(raw: &str) -> &str {
let title = extract_user_prompt(raw);
if title.is_empty() { "Session" } else { title }
}
pub(crate) fn strip_thinking_tags(text: &str) -> String {
if !text.contains("<think") && !text.contains("<thinking") && !text.contains("<reasoning") {
return text.to_string();
}
let tags = ["think", "thinking", "reasoning"];
let mut result = text.to_string();
for tag in tags {
let open = format!("<{tag}>");
let close = format!("</{tag}>");
while let Some(start) = result.find(&open) {
let Some(end) = result[start..].find(&close) else {
break;
};
let end_abs = start + end + close.len();
result.replace_range(start..end_abs, "");
}
}
result
}
fn truncate_title(s: &str, max_len: usize) -> String {
let s = s.trim();
let first_line = sanitize_session_title(s.lines().next().unwrap_or(s));
let first_line = first_line.trim();
let char_count = first_line.chars().count();
if char_count <= max_len {
first_line.to_string()
} else {
let truncated: String = first_line.chars().take(max_len - 3).collect();
format!("{truncated}...")
}
}
pub fn format_session_line(meta: &SessionMetadata) -> String {
let age = format_age(&meta.updated_at);
let updated = format_session_updated_at(&meta.updated_at, &age);
let truncated_title = truncate_title(extract_title(&meta.title), 40);
let fork_label = if meta.parent_session_id.is_some() {
" | fork"
} else {
""
};
format!(
"{} | {} | {} msgs{} | {}",
truncate_id(&meta.id),
truncated_title,
meta.message_count,
fork_label,
updated
)
}
pub(crate) fn format_session_updated_at(dt: &DateTime<Utc>, age: &str) -> String {
format!("{} ({age})", dt.format("%Y-%m-%d %H:%M UTC"))
}
fn format_age(dt: &DateTime<Utc>) -> String {
let now = Utc::now();
let duration = now.signed_duration_since(*dt);
if duration.num_minutes() < 1 {
"just now".to_string()
} else if duration.num_hours() < 1 {
format!("{}m ago", duration.num_minutes())
} else if duration.num_days() < 1 {
format!("{}h ago", duration.num_hours())
} else if duration.num_weeks() < 1 {
format!("{}d ago", duration.num_days())
} else {
format!("{}w ago", duration.num_weeks())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::models::ContentBlock;
use crate::tools::plan::StepStatus;
use crate::tui::history::{HistoryCell, ToolCell, history_cells_from_message};
use std::fs;
use tempfile::tempdir;
fn make_test_message(role: &str, text: &str) -> Message {
Message {
role: role.to_string(),
content: vec![ContentBlock::Text {
text: text.to_string(),
cache_control: None,
}],
}
}
#[test]
fn cost_snapshot_round_trips_coverage_and_detects_legacy_unknown() {
let legacy: SessionCostSnapshot = serde_json::from_value(serde_json::json!({
"session_cost_usd": 1.25,
"session_cost_cny": 0.0,
"subagent_cost_usd": 0.0,
"subagent_cost_cny": 0.0,
"displayed_cost_high_water_usd": 1.25,
"displayed_cost_high_water_cny": 0.0
}))
.expect("legacy cost snapshot stays readable");
assert_eq!(legacy.priced_turns, 0);
assert_eq!(legacy.unpriced_turns, 0);
assert!(!legacy.coverage_recorded);
assert!(
legacy.coverage_is_legacy_unknown(),
"a non-zero total with no coverage evidence must not read as complete"
);
let empty = SessionCostSnapshot::default();
assert!(empty.coverage_is_legacy_unknown());
let recorded_zero = SessionCostSnapshot {
session_cost_usd: 1.25,
coverage_recorded: true,
..SessionCostSnapshot::default()
};
assert!(!recorded_zero.coverage_is_legacy_unknown());
let full = SessionCostSnapshot {
session_cost_usd: 2.5,
session_cost_cny: 3.0,
subagent_cost_usd: 0.5,
subagent_cost_cny: 0.25,
displayed_cost_high_water_usd: 3.0,
displayed_cost_high_water_cny: 3.25,
priced_turns: 7,
unpriced_turns: 2,
cny_priced_turns: 1,
cny_unpriced_turns: 8,
unpriced_reasons: ["missing_class_price".to_string()].into(),
cny_unpriced_reasons: ["currency_not_published".to_string()].into(),
unpriced_classes: ["cache_write".to_string()].into(),
pricing_provenances: ["models_dev_bundled".to_string()].into(),
live_pricing_defects: ["live_pricing_stale".to_string()].into(),
live_pricing_unusable_defects: ["live_pricing_scope_mismatch".to_string()].into(),
route_receipts: ["provider=anthropic identity=- model=claude-haiku-4-5 \
surface=first-party-payg endpoint_fp=abc123 currency=usd"
.to_string()]
.into(),
coverage_recorded: true,
};
let json = serde_json::to_string(&full).expect("serialize");
let back: SessionCostSnapshot = serde_json::from_str(&json).expect("round-trip");
assert_eq!(back.priced_turns, 7);
assert_eq!(back.unpriced_turns, 2);
assert_eq!(back.cny_priced_turns, 1);
assert_eq!(back.cny_unpriced_turns, 8);
assert_eq!(back.unpriced_reasons, full.unpriced_reasons);
assert_eq!(back.cny_unpriced_reasons, full.cny_unpriced_reasons);
assert_eq!(back.unpriced_classes, full.unpriced_classes);
assert_eq!(back.pricing_provenances, full.pricing_provenances);
assert_eq!(back.live_pricing_defects, full.live_pricing_defects);
assert_eq!(
back.live_pricing_unusable_defects,
full.live_pricing_unusable_defects
);
assert_eq!(back.route_receipts, full.route_receipts);
assert!(back.coverage_recorded);
assert!(!back.coverage_is_legacy_unknown());
let lower = json.to_lowercase();
for needle in ["http", "api_key", "authorization", "bearer", "sk-"] {
assert!(!lower.contains(needle), "{needle} leaked into {json}");
}
}
#[test]
fn cost_snapshot_currency_totals_are_projections_of_one_accumulator() {
use crate::pricing::CostEstimate;
let turn_sequences: &[&[CostEstimate]] = &[
&[
CostEstimate {
usd: 0.01,
cny: 0.07,
},
CostEstimate {
usd: 0.02,
cny: 0.14,
},
],
&[
CostEstimate {
usd: 0.25,
cny: 0.0,
},
CostEstimate { usd: 1.5, cny: 0.0 },
],
&[
CostEstimate { usd: 0.5, cny: 0.0 },
CostEstimate { usd: 0.0, cny: 3.5 },
CostEstimate {
usd: 0.125,
cny: 0.875,
},
],
&[
CostEstimate {
usd: f64::NAN,
cny: 0.25,
},
CostEstimate {
usd: 0.75,
cny: -1.0,
},
CostEstimate {
usd: f64::INFINITY,
cny: 0.25,
},
],
];
for turns in turn_sequences {
let joint = turns.iter().fold(CostEstimate::default(), |acc, turn| {
acc.saturating_add(*turn)
});
let usd_alone = turns.iter().fold(CostEstimate::default(), |acc, turn| {
acc.saturating_add(CostEstimate {
usd: turn.usd,
cny: 0.0,
})
});
let cny_alone = turns.iter().fold(CostEstimate::default(), |acc, turn| {
acc.saturating_add(CostEstimate {
usd: 0.0,
cny: turn.cny,
})
});
let snapshot = SessionCostSnapshot {
session_cost_usd: joint.usd,
session_cost_cny: joint.cny,
..SessionCostSnapshot::default()
};
assert_eq!(
snapshot.total_usd(),
usd_alone.usd,
"USD projection drifted from independent accumulation for {turns:?}"
);
assert_eq!(
snapshot.total_cny(),
cny_alone.cny,
"CNY projection drifted from independent accumulation for {turns:?}"
);
assert_eq!(snapshot.total_estimate().usd, snapshot.total_usd());
assert_eq!(snapshot.total_estimate().cny, snapshot.total_cny());
}
let usd_only = SessionCostSnapshot {
session_cost_usd: 2.5,
subagent_cost_usd: 0.5,
..SessionCostSnapshot::default()
};
assert_eq!(usd_only.total_usd(), 3.0);
assert_eq!(usd_only.total_cny(), 0.0);
}
fn write_session_record(
manager: &SessionManager,
id: &str,
workspace: &Path,
updated_at: DateTime<Utc>,
) {
let session = SavedSession {
schema_version: CURRENT_SESSION_SCHEMA_VERSION,
messages: vec![make_test_message("user", "hi")],
metadata: SessionMetadata {
id: id.to_string(),
title: format!("session-{id}"),
created_at: updated_at,
updated_at,
message_count: 1,
total_tokens: 0,
model: "deepseek-v4-flash".to_string(),
model_provider: "deepseek".to_string(),
model_provider_id: None,
workspace: workspace.to_path_buf(),
mode: None,
cost: SessionCostSnapshot::default(),
parent_session_id: None,
forked_from_message_count: None,
cumulative_turn_secs: 0,
archived: false,
spawn_depth: 0,
},
journal: None,
leaf_id: None,
system_prompt: None,
context_references: Vec::new(),
artifacts: Vec::new(),
work_state: None,
last_auto_route: None,
};
manager.save_session(&session).expect("save");
}
fn write_empty_session_record(
manager: &SessionManager,
id: &str,
workspace: &Path,
updated_at: DateTime<Utc>,
) {
let session = SavedSession {
schema_version: CURRENT_SESSION_SCHEMA_VERSION,
messages: Vec::new(),
metadata: SessionMetadata {
id: id.to_string(),
title: DEFAULT_SESSION_TITLE.to_string(),
created_at: updated_at,
updated_at,
message_count: 0,
total_tokens: 0,
model: "deepseek-v4-pro".to_string(),
model_provider: "deepseek".to_string(),
model_provider_id: None,
workspace: workspace.to_path_buf(),
mode: Some("yolo".to_string()),
cost: SessionCostSnapshot::default(),
parent_session_id: None,
forked_from_message_count: None,
cumulative_turn_secs: 0,
archived: false,
spawn_depth: 0,
},
journal: None,
leaf_id: None,
system_prompt: None,
context_references: Vec::new(),
artifacts: Vec::new(),
work_state: None,
last_auto_route: None,
};
manager.save_session(&session).expect("save empty");
}
#[test]
fn session_boot_owner_stamps_only_the_creating_instance() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().to_path_buf()).expect("manager");
let workspace = tmp.path().join("ws");
write_session_record(&manager, "mine", &workspace, Utc::now());
assert_eq!(
manager.session_boot_owner("mine").as_deref(),
Some(current_session_boot_id())
);
assert!(!manager.session_from_prior_instance("mine"));
assert!(!manager.session_from_prior_instance("unsaved"));
manager
.record_session_boot_owner("theirs", "boot_other_instance")
.expect("stamp");
write_session_record(&manager, "theirs", &workspace, Utc::now());
assert_eq!(
manager.session_boot_owner("theirs").as_deref(),
Some("boot_other_instance")
);
assert!(manager.session_from_prior_instance("theirs"));
write_session_record(&manager, "legacy", &workspace, Utc::now());
manager.clear_session_boot_owner("legacy");
assert!(manager.session_from_prior_instance("legacy"));
write_session_record(&manager, "legacy", &workspace, Utc::now());
assert!(manager.session_from_prior_instance("legacy"));
manager.delete_session("theirs").expect("delete");
assert_eq!(manager.session_boot_owner("theirs"), None);
}
#[test]
fn session_boot_owner_sidecar_never_lists_as_a_session() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().to_path_buf()).expect("manager");
write_session_record(&manager, "real", &tmp.path().join("ws"), Utc::now());
assert!(manager.session_boot_owners_path().exists());
let listed = manager.list_sessions().expect("list");
assert_eq!(listed.len(), 1);
assert_eq!(listed[0].id, "real");
assert!(manager.load_session("session_boot_owners").is_err());
}
#[test]
fn test_session_manager_new() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
assert!(tmp.path().join("sessions").exists());
let _ = manager;
}
#[test]
fn test_save_and_load_session() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let messages = vec![
make_test_message("user", "Hello!"),
make_test_message("assistant", "Hi there!"),
];
let session = create_saved_session(&messages, "test-model", tmp.path(), 100, None);
let session_id = session.metadata.id.clone();
manager.save_session(&session).expect("save");
let loaded = manager.load_session(&session_id).expect("load");
assert_eq!(loaded.metadata.id, session_id);
assert_eq!(loaded.messages.len(), 2);
}
#[test]
fn rehydrated_turn_meta_blocks_never_render_in_history_cells() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let turn_meta = "<turn_meta>\nCurrent local date: 2026-08-01\n</turn_meta>";
let trailing_shape = Message {
role: "user".to_string(),
content: vec![
ContentBlock::Text {
text: "Fix the flaky test".to_string(),
cache_control: None,
},
ContentBlock::Text {
text: turn_meta.to_string(),
cache_control: None,
},
],
};
let legacy_leading_shape = Message {
role: "user".to_string(),
content: vec![
ContentBlock::Text {
text: turn_meta.to_string(),
cache_control: None,
},
ContentBlock::Text {
text: "Now add the docs".to_string(),
cache_control: None,
},
],
};
let messages = vec![
trailing_shape,
make_test_message("assistant", "Done."),
legacy_leading_shape,
];
let session = create_saved_session(&messages, "test-model", tmp.path(), 100, None);
let session_id = session.metadata.id.clone();
manager.save_session(&session).expect("save");
let loaded = manager.load_session(&session_id).expect("load");
let rendered: Vec<HistoryCell> = loaded
.messages
.iter()
.flat_map(history_cells_from_message)
.collect();
let user_texts: Vec<&str> = rendered
.iter()
.filter_map(|cell| match cell {
HistoryCell::User { content } => Some(content.as_str()),
_ => None,
})
.collect();
assert_eq!(user_texts, vec!["Fix the flaky test", "Now add the docs"]);
assert!(
!user_texts.iter().any(|text| text.contains("<turn_meta")),
"rendered cells must not contain turn_meta markup: {user_texts:?}"
);
let replayed_envelopes = loaded
.messages
.iter()
.flat_map(|message| &message.content)
.filter(|block| {
matches!(block, ContentBlock::Text { text, .. } if text.contains("<turn_meta>"))
})
.count();
assert_eq!(replayed_envelopes, 2);
}
#[test]
fn runtime_snapshot_load_preserves_in_flight_tool_call() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let messages = vec![Message {
role: "assistant".to_string(),
content: vec![ContentBlock::ToolUse {
id: "call-in-flight".to_string(),
name: "read_file".to_string(),
input: serde_json::json!({"path": "README.md"}),
caller: None,
thought_signature: None,
}],
}];
let session = create_saved_session(&messages, "test-model", tmp.path(), 0, None);
let session_id = session.metadata.id.clone();
manager.save_session(&session).expect("save");
let loaded = manager
.load_session_snapshot(&session_id)
.expect("snapshot load");
assert_eq!(loaded.messages, messages);
assert_eq!(loaded.metadata.message_count, 1);
assert!(!loaded.messages.iter().any(|message| {
message.content.iter().any(|block| {
matches!(
block,
ContentBlock::ToolResult { content, .. }
if content.contains("crashed_and_repaired")
)
})
}));
}
#[test]
fn explicit_session_recovery_is_reported_and_idempotent_after_save() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let messages = vec![Message {
role: "assistant".to_string(),
content: vec![ContentBlock::ToolUse {
id: "call-crashed".to_string(),
name: "read_file".to_string(),
input: serde_json::json!({"path": "README.md"}),
caller: None,
thought_signature: None,
}],
}];
let session = create_saved_session(&messages, "test-model", tmp.path(), 0, None);
let session_id = session.metadata.id.clone();
manager.save_session(&session).expect("save");
let recovered = manager
.recover_session_for_resume(&session_id)
.expect("recover");
assert!(recovered.changed);
assert_eq!(recovered.repaired_call_count, 1);
assert_eq!(recovered.duplicate_result_count, 0);
assert_eq!(recovered.orphan_result_count, 0);
manager
.save_session(&recovered.session)
.expect("persist recovery");
let second = manager
.recover_session_for_resume(&session_id)
.expect("recover twice");
assert!(!second.changed);
assert_eq!(second.repaired_call_count, 0);
assert_eq!(second.session.messages, recovered.session.messages);
}
#[test]
fn load_session_repairs_dangling_tool_call_with_visible_receipt() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let messages = vec![Message {
role: "assistant".to_string(),
content: vec![ContentBlock::ToolUse {
id: "call-crashed".to_string(),
name: "read_file".to_string(),
input: serde_json::json!({"path": "README.md"}),
caller: None,
thought_signature: None,
}],
}];
let session = create_saved_session(&messages, "test-model", tmp.path(), 0, None);
let session_id = session.metadata.id.clone();
manager.save_session(&session).expect("save");
let loaded = manager.load_session(&session_id).expect("load");
assert_eq!(loaded.metadata.message_count, loaded.messages.len());
assert!(loaded.messages.iter().any(|message| {
message.content.iter().any(|block| {
matches!(
block,
ContentBlock::ToolResult {
tool_use_id,
content,
is_error: Some(true),
..
} if tool_use_id == "call-crashed" && content.contains("crashed_and_repaired")
)
})
}));
assert_eq!(
loaded.journal.as_ref().map(SessionJournal::to_messages),
Some(loaded.messages.clone()),
"the append-only journal must follow the repaired active branch"
);
assert!(loaded.messages.iter().any(|message| {
(message.role == "assistant"
|| message.role == crate::models::INTERRUPTED_ASSISTANT_ROLE)
&& message.content.iter().any(|block| {
matches!(
block,
ContentBlock::Text { text, .. }
if text.contains("[tool_history_repair]")
)
})
}));
}
#[test]
fn save_and_load_session_preserves_rich_update_plan_tool_payload() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let messages = vec![
make_test_message("user", "plan this carefully"),
Message {
role: "assistant".to_string(),
content: vec![ContentBlock::ToolUse {
id: "plan-1".to_string(),
name: "update_plan".to_string(),
input: serde_json::json!({
"objective": "Make Plan mode reviewable",
"sources_used": ["gh issue view 2691"],
"critical_files": ["crates/tui/src/tools/plan.rs"],
"constraints": ["Preserve legacy update_plan payloads"],
"verification_plan": "Run focused plan tests",
"handoff_packet": "Next agent should inspect replay",
"plan": [
{ "step": "render replay card", "status": "completed" }
]
}),
caller: None,
thought_signature: None,
}],
},
Message {
role: "user".to_string(),
content: vec![ContentBlock::ToolResult {
tool_use_id: "plan-1".to_string(),
content: "Plan updated".to_string(),
is_error: None,
content_blocks: None,
}],
},
];
let session = create_saved_session(&messages, "deepseek-v4-flash", tmp.path(), 42, None);
let session_id = session.metadata.id.clone();
manager.save_session(&session).expect("save");
let loaded = manager.load_session(&session_id).expect("load");
assert_eq!(loaded.messages.len(), 3);
let cells = history_cells_from_message(&loaded.messages[1]);
let Some(HistoryCell::Tool(ToolCell::PlanUpdate(cell))) = cells.first() else {
panic!("expected loaded update_plan to replay as a PlanUpdate cell");
};
assert_eq!(
cell.snapshot.objective.as_deref(),
Some("Make Plan mode reviewable")
);
assert_eq!(
cell.snapshot.critical_files,
vec!["crates/tui/src/tools/plan.rs"]
);
assert_eq!(cell.snapshot.items[0].status, StepStatus::Completed);
}
#[test]
fn save_session_preserves_large_tool_outputs_for_cache_fidelity() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let raw = "RAW_SESSION_SENTINEL\n".repeat(2_000);
let messages = vec![
Message {
role: "assistant".to_string(),
content: vec![ContentBlock::ToolUse {
id: "call-big".to_string(),
name: "exec_shell".to_string(),
input: serde_json::json!({"command": "cargo test -p codewhale-tui"}),
caller: None,
thought_signature: None,
}],
},
Message {
role: "user".to_string(),
content: vec![ContentBlock::ToolResult {
tool_use_id: "call-big".to_string(),
content: raw.clone(),
is_error: None,
content_blocks: None,
}],
},
];
let mut session = create_saved_session(&messages, "test-model", tmp.path(), 100, None);
session.artifacts.push(crate::artifacts::ArtifactRecord {
id: "art_call-big".to_string(),
kind: crate::artifacts::ArtifactKind::ToolOutput,
session_id: session.metadata.id.clone(),
tool_call_id: "call-big".to_string(),
tool_name: "exec_shell".to_string(),
created_at: Utc::now(),
byte_size: raw.len() as u64,
preview: "checking crate ... error[E0425]".to_string(),
storage_path: PathBuf::from("artifacts/art_call-big.txt"),
});
let path = manager.save_session(&session).expect("save");
let persisted_json = fs::read_to_string(path).expect("read persisted session");
assert!(persisted_json.contains("RAW_SESSION_SENTINEL"));
let loaded = manager.load_session(&session.metadata.id).expect("load");
let ContentBlock::ToolResult { content, .. } = &loaded.messages[1].content[0] else {
panic!("expected loaded tool result");
};
assert!(content.contains("RAW_SESSION_SENTINEL"));
assert!(!content.contains("[TOOL_OUTPUT_RECEIPT]"));
}
#[test]
fn load_session_preserves_legacy_large_tool_outputs_for_cache_fidelity() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let raw = "RAW_LEGACY_RESUME_SENTINEL\n".repeat(2_000);
let messages = vec![
Message {
role: "assistant".to_string(),
content: vec![ContentBlock::ToolUse {
id: "call-legacy".to_string(),
name: "exec_shell".to_string(),
input: serde_json::json!({"command": "cargo check"}),
caller: None,
thought_signature: None,
}],
},
Message {
role: "user".to_string(),
content: vec![ContentBlock::ToolResult {
tool_use_id: "call-legacy".to_string(),
content: raw.clone(),
is_error: None,
content_blocks: None,
}],
},
];
let mut session = create_saved_session(&messages, "test-model", tmp.path(), 100, None);
session.artifacts.push(crate::artifacts::ArtifactRecord {
id: "art_call-legacy".to_string(),
kind: crate::artifacts::ArtifactKind::ToolOutput,
session_id: session.metadata.id.clone(),
tool_call_id: "call-legacy".to_string(),
tool_name: "exec_shell".to_string(),
created_at: Utc::now(),
byte_size: raw.len() as u64,
preview: "cargo check output".to_string(),
storage_path: PathBuf::from("artifacts/art_call-legacy.txt"),
});
let path = manager
.validated_session_path(&session.metadata.id)
.expect("path");
fs::write(
&path,
serde_json::to_string_pretty(&session).expect("serialize legacy session"),
)
.expect("write legacy raw session");
assert!(
fs::read_to_string(&path)
.expect("read legacy raw")
.contains("RAW_LEGACY_RESUME_SENTINEL")
);
let loaded = manager.load_session(&session.metadata.id).expect("load");
let ContentBlock::ToolResult { content, .. } = &loaded.messages[1].content[0] else {
panic!("expected loaded tool result");
};
assert!(content.contains("RAW_LEGACY_RESUME_SENTINEL"));
assert!(!content.contains("[TOOL_OUTPUT_RECEIPT]"));
}
#[test]
fn test_list_sessions() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
for i in 0..3 {
let messages = vec![make_test_message("user", &format!("Session {i}"))];
let session = create_saved_session(&messages, "test-model", tmp.path(), 100, None);
manager.save_session(&session).expect("save");
}
let sessions = manager.list_sessions().expect("list");
assert_eq!(sessions.len(), 3);
}
#[test]
fn default_manager_copies_legacy_sessions_when_primary_already_exists() {
let _lock = crate::test_support::lock_test_env();
let tmp = tempdir().expect("tempdir");
let home = tmp.path().join("home");
let _home = crate::test_support::EnvVarGuard::set("HOME", &home);
let _codewhale_home = crate::test_support::EnvVarGuard::remove("CODEWHALE_HOME");
let primary_sessions = home.join(".codewhale").join("sessions");
let legacy_sessions = home.join(".deepseek").join("sessions");
fs::create_dir_all(&primary_sessions).expect("primary sessions");
fs::create_dir_all(&legacy_sessions).expect("legacy sessions");
fs::create_dir_all(legacy_sessions.join("checkpoints")).expect("legacy checkpoints");
fs::write(
legacy_sessions.join("checkpoints").join("latest.json"),
"{}",
)
.expect("legacy checkpoint");
let mut legacy_session = create_saved_session(
&[make_test_message("user", "find my old session")],
"test-model",
tmp.path(),
100,
None,
);
legacy_session.metadata.id = "legacy-visible".to_string();
legacy_session.metadata.title = "session from legacy home".to_string();
fs::write(
legacy_sessions.join("legacy-visible.json"),
serde_json::to_string_pretty(&legacy_session).expect("serialize legacy session"),
)
.expect("write legacy session");
let manager = SessionManager::default_location().expect("default manager");
assert_eq!(manager.sessions_dir(), primary_sessions.as_path());
assert!(primary_sessions.join("legacy-visible.json").exists());
assert!(!primary_sessions.join("checkpoints").exists());
assert!(legacy_sessions.join("legacy-visible.json").exists());
let sessions = manager.list_sessions().expect("list");
assert_eq!(sessions.len(), 1);
assert_eq!(sessions[0].id, "legacy-visible");
}
#[test]
fn legacy_session_copy_never_overwrites_primary_session() {
let _lock = crate::test_support::lock_test_env();
let tmp = tempdir().expect("tempdir");
let home = tmp.path().join("home");
let _home = crate::test_support::EnvVarGuard::set("HOME", &home);
let _codewhale_home = crate::test_support::EnvVarGuard::remove("CODEWHALE_HOME");
let primary_sessions = home.join(".codewhale").join("sessions");
let legacy_sessions = home.join(".deepseek").join("sessions");
fs::create_dir_all(&primary_sessions).expect("primary sessions");
fs::create_dir_all(&legacy_sessions).expect("legacy sessions");
let primary_path = primary_sessions.join("same-id.json");
fs::write(&primary_path, "primary data wins").expect("write primary session");
fs::write(
legacy_sessions.join("same-id.json"),
"legacy data must not overwrite",
)
.expect("write legacy session");
let dir = default_sessions_dir().expect("default session dir");
assert_eq!(dir, primary_sessions);
assert_eq!(
fs::read_to_string(primary_path).expect("read primary session"),
"primary data wins"
);
}
#[test]
fn explicit_codewhale_home_disables_legacy_session_copy() {
let _lock = crate::test_support::lock_test_env();
let tmp = tempdir().expect("tempdir");
let home = tmp.path().join("home");
let explicit_home = tmp.path().join("explicit-codewhale");
let _home = crate::test_support::EnvVarGuard::set("HOME", &home);
let _codewhale_home =
crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", &explicit_home);
let legacy_sessions = home.join(".deepseek").join("sessions");
fs::create_dir_all(&legacy_sessions).expect("legacy sessions");
fs::write(legacy_sessions.join("legacy-visible.json"), "{}").expect("write legacy session");
let dir = default_sessions_dir().expect("default session dir");
assert_eq!(dir, explicit_home.join("sessions"));
assert!(!dir.join("legacy-visible.json").exists());
}
#[cfg(unix)]
#[test]
fn non_unicode_codewhale_home_is_still_an_explicit_session_boundary() {
use std::os::unix::ffi::OsStringExt;
let _lock = crate::test_support::lock_test_env();
let tmp = tempdir().expect("tempdir");
let home = tmp.path().join("home");
let explicit_home = tmp.path().join(std::ffi::OsString::from_vec(
b"codewhale-\xff-home".to_vec(),
));
let _home = crate::test_support::EnvVarGuard::set("HOME", &home);
let _codewhale_home =
crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", &explicit_home);
let legacy_sessions = home.join(".deepseek").join("sessions");
fs::create_dir_all(&legacy_sessions).expect("legacy sessions");
fs::write(legacy_sessions.join("ambient.json"), "ambient").expect("ambient legacy session");
let safe_primary = tmp.path().join("safe-primary");
fs::create_dir_all(&safe_primary).expect("safe primary");
assert_eq!(
merge_missing_legacy_session_entries(&safe_primary).expect("merge decision"),
0
);
assert!(!safe_primary.join("ambient.json").exists());
}
#[test]
fn latest_session_for_workspace_ignores_newer_other_directory() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let workspace_a = tmp.path().join("aa").join("aaa");
let workspace_b = tmp.path().join("bb").join("bbb");
fs::create_dir_all(&workspace_a).expect("mkdir workspace a");
fs::create_dir_all(&workspace_b).expect("mkdir workspace b");
fs::create_dir_all(tmp.path().join(".git")).expect("mkdir invalid git boundary");
write_session_record(
&manager,
"current-workspace",
&workspace_a,
Utc::now() - chrono::Duration::minutes(10),
);
write_session_record(&manager, "other-workspace", &workspace_b, Utc::now());
let global = manager
.list_sessions()
.expect("list")
.into_iter()
.next()
.expect("global latest");
assert_eq!(global.id, "other-workspace");
let scoped = manager
.get_latest_session_for_workspace(&workspace_a)
.expect("latest for workspace")
.expect("scoped latest");
assert_eq!(scoped.id, "current-workspace");
}
#[test]
fn latest_session_for_workspace_ignores_invalid_parent_git_marker() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let workspace_a = tmp.path().join("aa").join("aaa");
let workspace_b = tmp.path().join("bb").join("bbb");
fs::create_dir_all(&workspace_a).expect("mkdir workspace a");
fs::create_dir_all(&workspace_b).expect("mkdir workspace b");
fs::create_dir_all(tmp.path().join(".git")).expect("mkdir invalid git marker");
write_session_record(
&manager,
"current-workspace",
&workspace_a,
Utc::now() - chrono::Duration::minutes(10),
);
write_session_record(&manager, "other-workspace", &workspace_b, Utc::now());
let scoped = manager
.get_latest_session_for_workspace(&workspace_a)
.expect("latest for workspace")
.expect("scoped latest");
assert_eq!(scoped.id, "current-workspace");
}
#[test]
fn latest_session_for_workspace_matches_same_git_repository() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let repo = tmp.path().join("repo");
let repo_app = repo.join("apps").join("client");
let repo_crate = repo.join("crates").join("server");
let other_repo = tmp.path().join("other").join("project");
fs::create_dir_all(repo.join(".git")).expect("mkdir .git");
fs::write(repo.join(".git").join("HEAD"), "ref: refs/heads/main\n").expect("write HEAD");
fs::create_dir_all(&repo_app).expect("mkdir repo app");
fs::create_dir_all(&repo_crate).expect("mkdir repo crate");
fs::create_dir_all(&other_repo).expect("mkdir other repo");
write_session_record(
&manager,
"same-repo",
&repo_app,
Utc::now() - chrono::Duration::minutes(5),
);
write_session_record(&manager, "other-repo", &other_repo, Utc::now());
let scoped = manager
.get_latest_session_for_workspace(&repo_crate)
.expect("latest for workspace")
.expect("same repo latest");
assert_eq!(scoped.id, "same-repo");
}
#[test]
fn latest_session_for_workspace_skips_empty_auto_created_session() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let workspace = tmp.path().join("repo");
fs::create_dir_all(&workspace).expect("mkdir workspace");
write_session_record(
&manager,
"interrupted-user-turn",
&workspace,
Utc::now() - chrono::Duration::minutes(5),
);
write_empty_session_record(&manager, "empty-auto-shell", &workspace, Utc::now());
let global = manager
.list_sessions()
.expect("list")
.into_iter()
.next()
.expect("global latest");
assert_eq!(global.id, "empty-auto-shell");
let scoped = manager
.get_latest_session_for_workspace(&workspace)
.expect("latest for workspace")
.expect("scoped latest");
assert_eq!(scoped.id, "interrupted-user-turn");
}
#[test]
fn test_load_by_prefix() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let messages = vec![make_test_message("user", "Test session")];
let session = create_saved_session(&messages, "test-model", tmp.path(), 100, None);
let prefix = truncate_id(&session.metadata.id).to_string();
manager.save_session(&session).expect("save");
let loaded = manager.load_session_by_prefix(&prefix).expect("load");
assert_eq!(loaded.messages.len(), 1);
}
#[test]
fn test_delete_session() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let messages = vec![make_test_message("user", "To be deleted")];
let session = create_saved_session(&messages, "test-model", tmp.path(), 100, None);
let session_id = session.metadata.id.clone();
manager.save_session(&session).expect("save");
assert!(manager.load_session(&session_id).is_ok());
manager.delete_session(&session_id).expect("delete");
assert!(manager.load_session(&session_id).is_err());
}
#[test]
fn delete_session_removes_artifact_directory() {
let tmp = tempdir().expect("tempdir");
let sessions_dir = tmp.path().join("sessions");
let manager = SessionManager::new(sessions_dir.clone()).expect("new");
let session = create_saved_session(
&[make_test_message("user", "artifact session")],
"test-model",
tmp.path(),
100,
None,
);
let session_id = session.metadata.id.clone();
let artifact_dir = sessions_dir.join(&session_id).join("artifacts");
fs::create_dir_all(&artifact_dir).expect("artifact dir");
fs::write(artifact_dir.join("art_call.txt"), "raw output").expect("artifact file");
manager.save_session(&session).expect("save");
manager.delete_session(&session_id).expect("delete");
assert!(!sessions_dir.join(format!("{session_id}.json")).exists());
assert!(!sessions_dir.join(&session_id).exists());
}
#[test]
fn test_session_id_rejects_invalid_characters() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let err = manager
.load_session("../outside")
.expect_err("invalid id should fail");
assert_eq!(err.kind(), std::io::ErrorKind::InvalidInput);
let err = manager
.delete_session("sess bad")
.expect_err("invalid id should fail");
assert_eq!(err.kind(), std::io::ErrorKind::InvalidInput);
}
#[test]
fn test_session_manager_rejects_relative_traversal_dir() {
let err = SessionManager::new(PathBuf::from("../sessions"))
.expect_err("relative traversal directory should fail");
assert_eq!(err.kind(), std::io::ErrorKind::InvalidInput);
}
#[test]
fn test_truncate_title() {
assert_eq!(truncate_title("Short", 50), "Short");
assert_eq!(
truncate_title("This is a very long title that should be truncated", 20),
"This is a very lo..."
);
assert_eq!(truncate_title("Line 1\nLine 2", 50), "Line 1");
}
#[test]
fn extract_user_prompt_strips_turn_meta_prefix() {
assert_eq!(
extract_user_prompt("<turn_meta>{\"cache\":\"x\"}</turn_meta>\nReal prompt"),
"Real prompt"
);
assert_eq!(extract_user_prompt(" Real prompt"), "Real prompt");
assert_eq!(
extract_user_prompt("<turn_meta>{\"unterminated\":true}\nReal prompt"),
"{\"unterminated\":true}\nReal prompt"
);
}
#[test]
fn create_saved_session_uses_prompt_after_turn_meta_for_title() {
let tmp = tempdir().expect("tempdir");
let messages = vec![make_test_message(
"user",
"<turn_meta>{\"cache\":\"x\"}</turn_meta>\nFix the session picker history pane",
)];
let session = create_saved_session(&messages, "test-model", tmp.path(), 100, None);
assert_eq!(
session.metadata.title,
"Fix the session picker history pane"
);
}
#[test]
fn strip_thinking_tags_removes_common_inline_blocks() {
let text = "Before <think>private</think> middle <reasoning>hidden</reasoning> after";
let cleaned = strip_thinking_tags(text);
assert_eq!(cleaned, "Before middle after");
assert_eq!(strip_thinking_tags("plain answer"), "plain answer");
}
#[test]
fn test_format_age() {
let now = Utc::now();
assert_eq!(format_age(&now), "just now");
let hour_ago = now - chrono::Duration::hours(2);
assert_eq!(format_age(&hour_ago), "2h ago");
let day_ago = now - chrono::Duration::days(3);
assert_eq!(format_age(&day_ago), "3d ago");
}
#[test]
fn session_titles_never_keep_terminal_controls_or_bidi_format_chars() {
let raw = "Ev\u{1b}]0;PWNED\u{7}il\u{202e}R\u{200b}Z\u{9d}0;X\u{9c}After\u{2066}B\u{2069} 会議 🐳";
assert_eq!(
sanitize_session_title(raw),
"Ev]0;PWNEDilRZ0;XAfterB 会議 🐳"
);
assert_eq!(
normalize_session_title(raw).unwrap(),
"Ev]0;PWNEDilRZ0;XAfterB 会議 🐳"
);
assert!(normalize_session_title("\u{1b}\u{7}\u{200b}").is_err());
assert_eq!(truncate_title(raw, 40), "Ev]0;PWNEDilRZ0;XAfterB 会議 🐳");
}
#[test]
fn format_session_line_includes_absolute_updated_timestamp() {
let mut session = create_saved_session(
&[make_test_message("user", "Find Friday work")],
"test-model",
Path::new("/tmp/project"),
100,
None,
);
session.metadata.updated_at = DateTime::parse_from_rfc3339("2026-06-01T12:34:00Z")
.expect("timestamp")
.with_timezone(&Utc);
let line = format_session_line(&session.metadata);
assert!(
line.contains("2026-06-01 12:34 UTC"),
"session list should include an absolute timestamp, got {line:?}"
);
}
#[test]
fn test_update_session() {
let tmp = tempdir().expect("tempdir");
let messages = vec![make_test_message("user", "Hello")];
let session = create_saved_session(&messages, "test-model", tmp.path(), 50, None);
let new_messages = vec![
make_test_message("user", "Hello"),
make_test_message("assistant", "Hi!"),
];
let updated = update_session(session, &new_messages, 100, None);
assert_eq!(updated.messages.len(), 2);
assert_eq!(updated.metadata.total_tokens, 100);
}
#[test]
fn save_load_round_trip_preserves_all_messages_for_cache_fidelity() {
#[derive(serde::Deserialize)]
struct LegacySession {
messages: Vec<Message>,
}
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
for count in [0, 1, 500, 501, 600, 1000] {
let original: Vec<_> = (0..count)
.map(|i| {
make_test_message(
if i % 2 == 0 { "user" } else { "assistant" },
&format!("round-trip message {i}"),
)
})
.collect();
let mut session = create_saved_session(&original, "test-model", tmp.path(), 0, None);
let expected_journal = session.journal.clone();
session.compact_for_persistence_queue();
let path = manager.save_session(&session).expect("save");
let legacy: LegacySession =
serde_json::from_slice(&fs::read(path).expect("read")).expect("legacy reader");
let loaded = manager.load_session(&session.metadata.id).expect("load");
assert_eq!(
legacy.messages, original,
"legacy messages for count={count}"
);
assert_eq!(
loaded.journal, expected_journal,
"journal for count={count}"
);
assert_eq!(
loaded.messages.len(),
count,
"count preserved for count={count}"
);
assert_eq!(
loaded.messages, original,
"every message byte-identical after round-trip for count={count}"
);
}
}
#[test]
fn test_checkpoint_round_trip_and_clear() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let messages = vec![make_test_message("user", "checkpoint me")];
let mut session = create_saved_session(&messages, "test-model", tmp.path(), 12, None);
session.work_state = Some(SessionWorkState {
todos: crate::tools::todo::TodoListSnapshot {
items: vec![crate::tools::todo::TodoItem {
id: 1,
content: "verify checkpoint durability".to_string(),
status: crate::tools::todo::TodoStatus::InProgress,
}],
completion_pct: 0,
in_progress_id: Some(1),
},
..SessionWorkState::default()
});
let expected_messages = session.messages.clone();
let expected_journal = session.journal.clone();
session.compact_for_persistence_queue();
let path = manager.save_checkpoint(&session).expect("save checkpoint");
assert_eq!(
path.file_name().and_then(|n| n.to_str()),
Some(format!("{}.json", session.metadata.id).as_str()),
"checkpoint file must be keyed by session id"
);
let loaded = manager
.load_session_checkpoint(&session.metadata.id)
.expect("load checkpoint")
.expect("checkpoint exists");
assert_eq!(loaded.metadata.id, session.metadata.id);
assert_eq!(loaded.messages, expected_messages);
assert_eq!(loaded.journal, expected_journal);
assert_eq!(
loaded.work_state, session.work_state,
"work state must survive the checkpoint round trip"
);
manager
.clear_session_checkpoint(&session.metadata.id)
.expect("clear checkpoint");
assert!(
manager
.load_session_checkpoint(&session.metadata.id)
.expect("load checkpoint")
.is_none()
);
}
#[test]
fn graph_backed_work_state_remains_readable_by_legacy_shape() {
#[derive(serde::Deserialize)]
struct LegacyWorkState {
#[serde(default)]
todos: crate::tools::todo::TodoListSnapshot,
#[serde(default)]
plan: crate::tools::plan::PlanSnapshot,
}
let fixture = include_bytes!("../tests/fixtures/work_graph_session_v1_reader.json");
let current: SavedSession = serde_json::from_slice(fixture).expect("current reader");
let state = current.work_state.expect("fixture Work state");
let legacy: LegacyWorkState = serde_json::from_value(
serde_json::from_slice::<serde_json::Value>(fixture)
.expect("fixture JSON")["work_state"]
.clone(),
)
.expect("v1 reader ignores graph");
assert_eq!(legacy.todos, state.todos);
assert_eq!(legacy.plan, state.plan);
let graph = state.graph.expect("fixture graph");
crate::work_graph::validate(&graph).expect("valid fixture graph");
assert_eq!(crate::work_graph::project_todos(&graph), state.todos);
assert_eq!(crate::work_graph::project_plan(&graph), state.plan);
}
#[test]
fn first_graph_write_archives_exact_legacy_session_once() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let mut session = create_saved_session(
&[make_test_message("user", "archive before import")],
"test-model",
tmp.path(),
0,
None,
);
let plan = crate::tools::plan::PlanSnapshot {
items: vec![crate::tools::plan::PlanItemArg {
step: "Import".to_string(),
status: crate::tools::plan::StepStatus::Pending,
}],
..crate::tools::plan::PlanSnapshot::default()
};
let todos = crate::tools::todo::TodoListSnapshot::default();
session.work_state = Some(SessionWorkState {
graph: None,
todos: todos.clone(),
plan: plan.clone(),
});
let path = manager.save_session(&session).expect("save legacy session");
let legacy_bytes = fs::read(&path).expect("read legacy bytes");
let graph = crate::work_graph::import_legacy(&session.metadata.id, &plan, &todos)
.expect("import graph");
session.work_state = Some(SessionWorkState {
graph: Some(graph),
todos,
plan,
});
manager.save_session(&session).expect("first graph write");
let archive = manager
.sessions_dir
.join(WORK_GRAPH_IMPORT_ARCHIVE_DIR)
.join(path.file_name().expect("session filename"));
assert_eq!(fs::read(&archive).expect("archive exists"), legacy_bytes);
session.metadata.title = "later graph write".to_string();
manager.save_session(&session).expect("second graph write");
assert_eq!(
fs::read(&archive).expect("archive still exists"),
legacy_bytes,
"later graph writes must not replace the pre-import receipt"
);
}
#[test]
fn checkpoints_are_independent_per_session() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let first = create_saved_session(
&[make_test_message("user", "session one")],
"test-model",
tmp.path(),
0,
None,
);
let second = create_saved_session(
&[make_test_message("user", "session two")],
"test-model",
tmp.path(),
0,
None,
);
manager.save_checkpoint(&first).expect("save first");
manager.save_checkpoint(&second).expect("save second");
manager
.clear_session_checkpoint(&first.metadata.id)
.expect("clear first");
assert!(
manager
.load_session_checkpoint(&first.metadata.id)
.expect("load first")
.is_none(),
"clearing one session must remove only that session's file"
);
let survivor = manager
.load_session_checkpoint(&second.metadata.id)
.expect("load second")
.expect("second checkpoint survives");
assert_eq!(survivor.metadata.id, second.metadata.id);
}
#[test]
fn list_checkpoints_includes_legacy_slot_and_skips_offline_queue() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let session = create_saved_session(
&[make_test_message("user", "list me")],
"test-model",
tmp.path(),
0,
None,
);
manager.save_checkpoint(&session).expect("save checkpoint");
let checkpoints = tmp.path().join("sessions").join("checkpoints");
fs::write(checkpoints.join("latest.json"), "{}").expect("write legacy slot");
fs::write(checkpoints.join("offline_queue.json"), "{}").expect("write offline queue");
let refs = manager.list_checkpoints().expect("list checkpoints");
assert_eq!(refs.len(), 2, "offline queue must not be a candidate");
assert!(
refs.iter()
.any(|r| r.source == CheckpointSource::Session(session.metadata.id.clone()))
);
assert!(refs.iter().any(|r| r.source == CheckpointSource::Legacy));
}
#[test]
fn legacy_migration_never_overwrites_existing_per_session_checkpoint() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let mut session = create_saved_session(
&[make_test_message("user", "original")],
"test-model",
tmp.path(),
0,
None,
);
manager.save_checkpoint(&session).expect("save checkpoint");
session.messages = vec![make_test_message("user", "stale legacy copy")];
let written = manager
.write_session_checkpoint_if_absent(&session)
.expect("migration attempt");
assert!(!written, "migration must not overwrite an existing file");
let loaded = manager
.load_session_checkpoint(&session.metadata.id)
.expect("load")
.expect("checkpoint exists");
assert_eq!(
loaded.messages,
vec![make_test_message("user", "original")],
"existing per-session checkpoint content must be preserved"
);
}
#[test]
fn workspace_scope_matches_subdirectories_in_same_git_checkout() {
let tmp = tempdir().expect("tempdir");
let repo = tmp.path().join("repo");
let nested = repo.join("crates").join("tui");
fs::create_dir_all(&nested).expect("mkdir nested");
fs::write(repo.join(".git"), "gitdir: .git/worktrees/repo").expect("write git marker");
assert!(workspace_scope_matches(&repo, &nested));
}
#[test]
fn workspace_scope_rejects_sibling_git_checkouts() {
let tmp = tempdir().expect("tempdir");
let first = tmp.path().join("repo-a");
let second = tmp.path().join("repo-b");
fs::create_dir_all(&first).expect("mkdir first");
fs::create_dir_all(&second).expect("mkdir second");
fs::write(first.join(".git"), "gitdir: .git/worktrees/a").expect("write first marker");
fs::write(second.join(".git"), "gitdir: .git/worktrees/b").expect("write second marker");
assert!(!workspace_scope_matches(&first, &second));
}
#[test]
fn test_offline_queue_round_trip_and_clear() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let state = OfflineQueueState {
messages: vec![QueuedSessionMessage {
display: "queued message".to_string(),
skill_instruction: Some("Use skill".to_string()),
skill_provenance: None,
}],
draft: Some(QueuedSessionMessage {
display: "draft message".to_string(),
skill_instruction: None,
skill_provenance: None,
}),
..OfflineQueueState::default()
};
manager
.save_offline_queue_state(&state, Some("test-session"))
.expect("save queue state");
let loaded = manager
.load_offline_queue_state()
.expect("load queue state")
.expect("queue state exists");
assert_eq!(loaded.messages.len(), 1);
assert_eq!(loaded.messages[0].display, "queued message");
assert!(loaded.draft.is_some());
manager
.clear_offline_queue_state()
.expect("clear queue state");
assert!(
manager
.load_offline_queue_state()
.expect("load queue state")
.is_none()
);
}
#[test]
fn test_offline_queue_stamps_session_id_on_save() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let state = OfflineQueueState {
messages: vec![QueuedSessionMessage {
display: "first parked".to_string(),
skill_instruction: None,
skill_provenance: None,
}],
..OfflineQueueState::default()
};
manager
.save_offline_queue_state(&state, Some("session-A"))
.expect("save with session id");
let loaded = manager
.load_offline_queue_state()
.expect("ok")
.expect("present");
assert_eq!(loaded.session_id.as_deref(), Some("session-A"));
manager
.save_offline_queue_state(&state, Some("session-B"))
.expect("re-save");
let reloaded = manager
.load_offline_queue_state()
.expect("ok")
.expect("present");
assert_eq!(reloaded.session_id.as_deref(), Some("session-B"));
manager
.save_offline_queue_state(&state, None)
.expect("save without session id");
let unscoped = manager
.load_offline_queue_state()
.expect("ok")
.expect("present");
assert!(
unscoped.session_id.is_none(),
"save with None must persist a missing session_id"
);
}
#[test]
fn test_session_context_references_round_trip() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let mut session = create_saved_session(
&[make_test_message("user", "read @src/main.rs")],
"deepseek-v4-pro",
tmp.path(),
0,
None,
);
session.context_references.push(SessionContextReference {
message_index: 0,
reference: ContextReference {
kind: crate::tui::file_mention::ContextReferenceKind::File,
source: crate::tui::file_mention::ContextReferenceSource::AtMention,
badge: "file".to_string(),
label: "src/main.rs".to_string(),
target: tmp.path().join("src/main.rs").display().to_string(),
included: true,
expanded: true,
detail: Some("included".to_string()),
},
});
let path = manager.save_session(&session).expect("save session");
let loaded = manager
.load_session(&session.metadata.id)
.expect("load session");
assert!(path.exists());
assert_eq!(loaded.context_references, session.context_references);
}
#[test]
fn test_checkpoint_rejects_newer_schema() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let checkpoints = tmp.path().join("sessions").join("checkpoints");
fs::create_dir_all(&checkpoints).expect("create checkpoints dir");
let path = checkpoints.join("latest.json");
fs::write(
&path,
r#"{
"schema_version": 999,
"metadata": {
"id": "sid",
"title": "bad",
"created_at": "2026-01-01T00:00:00Z",
"updated_at": "2026-01-01T00:00:00Z",
"message_count": 0,
"total_tokens": 0,
"model": "m",
"workspace": "/tmp",
"mode": null
},
"messages": [],
"system_prompt": null
}"#,
)
.expect("write checkpoint");
let err = manager
.load_legacy_checkpoint()
.expect_err("should reject schema");
assert!(err.to_string().contains("newer than supported"));
fs::rename(&path, checkpoints.join("sid.json")).expect("rename to per-session file");
let err = manager
.load_session_checkpoint("sid")
.expect_err("should reject schema");
assert!(err.to_string().contains("newer than supported"));
}
#[test]
fn test_load_session_rejects_newer_schema() {
let tmp = tempdir().expect("tempdir");
let sessions_dir = tmp.path().join("sessions");
let manager = SessionManager::new(sessions_dir.clone()).expect("new");
let id = "future-session";
let path = sessions_dir.join(format!("{id}.json"));
fs::write(
&path,
r#"{
"schema_version": 999,
"metadata": {
"id": "future-session",
"title": "future",
"created_at": "2026-01-01T00:00:00Z",
"updated_at": "2026-01-01T00:00:00Z",
"message_count": 0,
"total_tokens": 0,
"model": "m",
"workspace": "/tmp",
"mode": null
},
"messages": [],
"system_prompt": null
}"#,
)
.expect("write session");
let err = manager.load_session(id).expect_err("should reject schema");
assert!(
err.to_string().contains("newer than supported"),
"unexpected error: {err}"
);
}
#[test]
fn extract_top_level_metadata_skips_huge_messages_array() {
let big_text = format!(
r#"this message references "metadata" inside it, repeated:{}"#,
"x".repeat(20_000)
);
let json = format!(
r#"{{
"schema_version": 1,
"metadata": {{
"id": "abc-123",
"title": "Real Session",
"created_at": "2026-01-01T00:00:00Z",
"updated_at": "2026-01-02T00:00:00Z",
"message_count": 12,
"total_tokens": 4096,
"model": "deepseek-v4-flash",
"workspace": "/tmp"
}},
"messages": [
{{ "role": "user", "content": [ {{ "Text": {{ "text": {big_text:?} }} }} ] }}
]
}}"#
);
let extracted =
extract_top_level_metadata(json.as_bytes()).expect("metadata extractable from prefix");
assert_eq!(extracted.id, "abc-123");
assert_eq!(extracted.title, "Real Session");
assert_eq!(extracted.message_count, 12);
assert_eq!(extracted.total_tokens, 4096);
}
#[test]
fn extract_top_level_metadata_handles_braces_inside_strings() {
let json = r#"{
"metadata": {
"id": "x",
"title": "weird { title } with braces",
"created_at": "2026-01-01T00:00:00Z",
"updated_at": "2026-01-01T00:00:00Z",
"message_count": 0,
"total_tokens": 0,
"model": "m",
"workspace": "/tmp"
},
"messages": []
}"#;
let extracted = extract_top_level_metadata(json.as_bytes())
.expect("brace-in-string survives the scanner");
assert_eq!(extracted.title, "weird { title } with braces");
}
#[test]
fn saved_session_deserializes_without_artifacts_as_empty_registry() {
let json = r#"{
"schema_version": 1,
"metadata": {
"id": "legacy-session",
"title": "legacy",
"created_at": "2026-05-08T00:00:00Z",
"updated_at": "2026-05-08T00:00:00Z",
"message_count": 0,
"total_tokens": 0,
"model": "deepseek-v4-pro",
"workspace": "/tmp"
},
"messages": [],
"system_prompt": null
}"#;
let session: SavedSession = serde_json::from_str(json).expect("legacy session loads");
assert!(session.artifacts.is_empty());
assert!(session.last_auto_route.is_none());
assert!(session.metadata.parent_session_id.is_none());
assert!(session.metadata.forked_from_message_count.is_none());
}
#[test]
fn fork_lineage_metadata_round_trips_and_formats() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let parent = create_saved_session(
&[
make_test_message("user", "try approach A"),
make_test_message("assistant", "A looks viable"),
],
"deepseek-v4-pro",
Path::new("/tmp"),
42,
None,
);
let mut forked = create_saved_session(
&parent.messages,
&parent.metadata.model,
&parent.metadata.workspace,
parent.metadata.total_tokens,
None,
);
forked.metadata.mark_forked_from(&parent.metadata);
manager.save_session(&forked).expect("save fork");
let loaded = manager
.load_session(&forked.metadata.id)
.expect("load fork");
assert_eq!(
loaded.metadata.parent_session_id.as_deref(),
Some(parent.metadata.id.as_str())
);
assert_eq!(loaded.metadata.forked_from_message_count, Some(2));
let line = format_session_line(&loaded.metadata);
assert!(line.contains("fork"));
assert!(!line.contains(parent.metadata.id.as_str()));
}
#[test]
fn save_and_load_session_preserves_artifact_metadata() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let mut session = create_saved_session(
&[make_test_message("user", "run tests")],
"deepseek-v4-pro",
Path::new("/tmp"),
0,
None,
);
session.artifacts.push(crate::artifacts::ArtifactRecord {
id: "art_call_big".to_string(),
kind: crate::artifacts::ArtifactKind::ToolOutput,
session_id: session.metadata.id.clone(),
tool_call_id: "call-big".to_string(),
tool_name: "exec_shell".to_string(),
created_at: Utc::now(),
byte_size: 512_000,
preview: "cargo test output".to_string(),
storage_path: PathBuf::from("/tmp/tool_outputs/call-big.txt"),
});
manager.save_session(&session).expect("save");
let loaded = manager.load_session(&session.metadata.id).expect("load");
assert_eq!(loaded.artifacts, session.artifacts);
}
fn write_session_with_updated_at(
manager: &SessionManager,
id: &str,
updated_at: DateTime<Utc>,
) {
write_session_record(manager, id, Path::new("/tmp"), updated_at);
}
#[test]
fn prune_sessions_older_than_returns_zero_for_empty_dir() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
let pruned = manager
.prune_sessions_older_than(std::time::Duration::from_secs(3600))
.expect("prune");
assert_eq!(pruned, 0);
}
#[test]
fn prune_sessions_older_than_keeps_fresh_records() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
write_session_with_updated_at(
&manager,
"fresh-1",
Utc::now() - chrono::Duration::minutes(30),
);
write_session_with_updated_at(
&manager,
"fresh-2",
Utc::now() - chrono::Duration::minutes(5),
);
let pruned = manager
.prune_sessions_older_than(std::time::Duration::from_secs(3600))
.expect("prune");
assert_eq!(pruned, 0);
assert_eq!(manager.list_sessions().expect("list").len(), 2);
}
#[test]
fn prune_sessions_older_than_removes_stale_records() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
write_session_with_updated_at(&manager, "stale-1", Utc::now() - chrono::Duration::days(8));
write_session_with_updated_at(&manager, "stale-2", Utc::now() - chrono::Duration::days(30));
let pruned = manager
.prune_sessions_older_than(std::time::Duration::from_secs(7 * 24 * 3600))
.expect("prune");
assert_eq!(pruned, 2);
assert_eq!(manager.list_sessions().expect("list").len(), 0);
}
#[test]
fn prune_sessions_older_than_only_removes_stale_records_in_mixed_dir() {
let tmp = tempdir().expect("tempdir");
let manager = SessionManager::new(tmp.path().join("sessions")).expect("new");
write_session_with_updated_at(&manager, "fresh", Utc::now() - chrono::Duration::hours(1));
write_session_with_updated_at(&manager, "stale", Utc::now() - chrono::Duration::days(60));
let pruned = manager
.prune_sessions_older_than(std::time::Duration::from_secs(7 * 24 * 3600))
.expect("prune");
assert_eq!(pruned, 1);
let remaining = manager.list_sessions().expect("list");
assert_eq!(remaining.len(), 1);
assert_eq!(remaining[0].id, "fresh");
}
#[test]
fn prune_sessions_older_than_skips_checkpoint_directory() {
let tmp = tempdir().expect("tempdir");
let sessions_dir = tmp.path().join("sessions");
let manager = SessionManager::new(sessions_dir.clone()).expect("new");
let checkpoint_dir = sessions_dir.join("checkpoints");
fs::create_dir_all(&checkpoint_dir).expect("mkdir checkpoints");
let checkpoint_file = checkpoint_dir.join("latest.json");
fs::write(&checkpoint_file, "{}").expect("write checkpoint");
write_session_with_updated_at(&manager, "stale", Utc::now() - chrono::Duration::days(60));
let pruned = manager
.prune_sessions_older_than(std::time::Duration::from_secs(7 * 24 * 3600))
.expect("prune");
assert_eq!(pruned, 1, "the top-level stale session should be removed");
assert!(
checkpoint_file.exists(),
"checkpoint file should be untouched"
);
}
#[test]
fn test_load_offline_queue_rejects_newer_schema() {
let tmp = tempdir().expect("tempdir");
let sessions_dir = tmp.path().join("sessions");
let manager = SessionManager::new(sessions_dir.clone()).expect("new");
let checkpoints = sessions_dir.join("checkpoints");
fs::create_dir_all(&checkpoints).expect("create checkpoints dir");
let path = checkpoints.join("offline_queue.json");
fs::write(
&path,
r#"{
"schema_version": 999,
"messages": [],
"draft": null
}"#,
)
.expect("write queue");
let err = manager
.load_offline_queue_state()
.expect_err("should reject schema");
assert!(
err.to_string().contains("newer than supported"),
"unexpected error: {err}"
);
}
}