| 1 | //! Session management for resuming conversations. |
| 2 | //! |
| 3 | //! This module provides functionality for: |
| 4 | //! - Saving sessions to disk |
| 5 | //! - Listing previous sessions |
| 6 | //! - Resuming sessions by ID |
| 7 | //! - Managing session lifecycle |
| 8 | |
| 9 | use crate::approval_log::{ApprovalReceipt, ApprovalReceiptStore, ApprovalReplay}; |
| 10 | use crate::artifacts::ArtifactRecord; |
| 11 | use crate::config::ProviderKind; |
| 12 | use crate::model_routing::AutoRouteReceipt; |
| 13 | use crate::project_context::find_git_root; |
| 14 | use crate::session_tree::{SessionEntry, SessionImportContainer, SessionJournal}; |
| 15 | use crate::tools::goal::{GoalPauseReason, GoalSnapshot}; |
| 16 | use crate::tools::plan::PlanSnapshot; |
| 17 | use crate::tools::todo::TodoListSnapshot; |
| 18 | use crate::utils::write_atomic; |
| 19 | use crate::work_graph::ReasoningEffortTier; |
| 20 | use chrono::{DateTime, Utc}; |
| 21 | use codewhale_core::ContextReference; |
| 22 | #[cfg(test)] |
| 23 | use codewhale_core::{ContextReferenceKind, ContextReferenceSource}; |
| 24 | use codewhale_models::{ContentBlock, Message, SystemPrompt}; |
| 25 | use serde::{Deserialize, Serialize}; |
| 26 | use std::collections::{BTreeMap, BTreeSet}; |
| 27 | use std::fs::{self, OpenOptions}; |
| 28 | use std::io; |
| 29 | use std::path::{Component, Path, PathBuf}; |
| 30 | use std::sync::atomic::AtomicBool; |
| 31 | use uuid::Uuid; |
| 32 | |
| 33 | /// Maximum number of active (non-archived) transcripts to retain. |
| 34 | /// |
| 35 | /// A transcript that falls out of this window is archived, never unlinked |
| 36 | /// (#6136); archived records sit outside the cap until the user prunes them, |
| 37 | /// and empty auto-created stubs are capped separately (#6137). |
| 38 | const MAX_SESSIONS: usize = 50; |
| 39 | /// Maximum empty auto-created stubs ("New Session", zero messages) to keep. |
| 40 | /// |
| 41 | /// The product writes one per boot; they are junk that must never occupy a |
| 42 | /// transcript's slot in the cap (#6137). |
| 43 | const MAX_EMPTY_SESSION_STUBS: usize = 10; |
| 44 | /// Maximum session title length, in `char`s. Matches the bound the session |
| 45 | /// picker's rename prompt has always enforced. |
| 46 | pub const MAX_SESSION_TITLE_CHARS: usize = 100; |
| 47 | pub(crate) const WORK_GRAPH_IMPORT_ARCHIVE_DIR: &str = ".work-graph-import-archive"; |
| 48 | /// Per-session JSONL sidecar holding journal entries a bounded save moved out |
| 49 | /// of the session document (#6842): `<sessions>/.journal-archive/<id>.jsonl`. |
| 50 | pub(crate) const JOURNAL_ARCHIVE_DIR: &str = ".journal-archive"; |
| 51 | /// Off-branch entries a session document may carry before an autosave |
| 52 | /// archives the excess. Together with [`JOURNAL_RETAINED_DEAD_ENTRIES`] this |
| 53 | /// keeps a compacted journal far below `MAX_CANONICAL_HISTORY_ENTRIES`. |
| 54 | const JOURNAL_PRUNE_DEAD_THRESHOLD: usize = 4_096; |
| 55 | /// Recent off-branch entries kept in the document after archiving, so `/tree` |
| 56 | /// and `/branch` stay cheap for recent history. |
| 57 | const JOURNAL_RETAINED_DEAD_ENTRIES: usize = 1_024; |
| 58 | const SESSION_GOALS_DIR: &str = ".goals"; |
| 59 | const CURRENT_SESSION_GOAL_SCHEMA_VERSION: u32 = 2; |
| 60 | const MAX_SESSION_GOAL_OBJECTIVE_CHARS: usize = 8_192; |
| 61 | const MAX_SESSION_GOAL_FILE_BYTES: u64 = 64 * 1_024; |
| 62 | pub(crate) const CURRENT_SESSION_SCHEMA_VERSION: u32 = 1; |
| 63 | const CURRENT_QUEUE_SCHEMA_VERSION: u32 = 1; |
| 64 | const LATE_USAGE_DIR: &str = ".late-usage"; |
| 65 | const CURRENT_LATE_USAGE_SCHEMA_VERSION: u32 = 1; |
| 66 | const MAX_LATE_USAGE_UNRESOLVED_RECORDS_PER_SESSION: usize = 64; |
| 67 | const MAX_LATE_USAGE_LEDGER_BYTES: u64 = 1024 * 1024; |
| 68 | const LATE_USAGE_DELETED: &[u8] = b"codewhale-session-deleted-v1\n"; |
| 69 | fn is_false(value: &bool) -> bool { |
| 70 | !*value |
| 71 | } |
| 72 | |
| 73 | const LATE_USAGE_UNAVAILABLE_REASON: &str = "late_usage_ledger_unavailable"; |
| 74 | |
| 75 | #[derive(Clone, Copy)] |
| 76 | enum SessionRemoval { |
| 77 | Explicit, |
| 78 | Retention, |
| 79 | } |
| 80 | |
| 81 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 82 | struct LateUsageRecord { |
| 83 | source_fingerprint: String, |
| 84 | turn_fingerprint: String, |
| 85 | route: crate::cost_status::EffectiveRouteEnvelope, |
| 86 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 87 | usage: Option<codewhale_models::Usage>, |
| 88 | #[serde( |
| 89 | default, |
| 90 | skip_serializing_if = "crate::cost_status::RuntimeUsageMissingReason::is_success" |
| 91 | )] |
| 92 | reason: crate::cost_status::RuntimeUsageMissingReason, |
| 93 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 94 | decision: Option<crate::cost_status::RuntimeDecisionReceipt>, |
| 95 | } |
| 96 | |
| 97 | impl LateUsageRecord { |
| 98 | fn captured( |
| 99 | turn_id: &str, |
| 100 | source_id: &str, |
| 101 | route: &crate::cost_status::EffectiveRouteEnvelope, |
| 102 | usage: Option<&codewhale_models::Usage>, |
| 103 | decision: Option<&crate::cost_status::RuntimeDecisionReceipt>, |
| 104 | reason: crate::cost_status::RuntimeUsageMissingReason, |
| 105 | ) -> Self { |
| 106 | Self { |
| 107 | source_fingerprint: crate::cost_status::usage_source_fingerprint(source_id), |
| 108 | turn_fingerprint: crate::cost_status::usage_source_fingerprint(turn_id), |
| 109 | route: route.sanitized_for_persistence(), |
| 110 | usage: usage |
| 111 | .filter(|usage| **usage != codewhale_models::Usage::default()) |
| 112 | .cloned(), |
| 113 | reason, |
| 114 | decision: decision.map(crate::cost_status::RuntimeDecisionReceipt::sanitized), |
| 115 | } |
| 116 | } |
| 117 | } |
| 118 | |
| 119 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 120 | struct LateUsageLedger { |
| 121 | schema_version: u32, |
| 122 | #[serde(default)] |
| 123 | records: Vec<LateUsageRecord>, |
| 124 | #[serde(default)] |
| 125 | overflowed: bool, |
| 126 | } |
| 127 | |
| 128 | impl Default for LateUsageLedger { |
| 129 | fn default() -> Self { |
| 130 | Self { |
| 131 | schema_version: CURRENT_LATE_USAGE_SCHEMA_VERSION, |
| 132 | records: Vec::new(), |
| 133 | overflowed: false, |
| 134 | } |
| 135 | } |
| 136 | } |
| 137 | |
| 138 | fn is_sha256_fingerprint(value: &str) -> bool { |
| 139 | value.len() == 64 |
| 140 | && value |
| 141 | .bytes() |
| 142 | .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) |
| 143 | } |
| 144 | |
| 145 | const fn default_session_schema_version() -> u32 { |
| 146 | CURRENT_SESSION_SCHEMA_VERSION |
| 147 | } |
| 148 | |
| 149 | const fn default_queue_schema_version() -> u32 { |
| 150 | CURRENT_QUEUE_SCHEMA_VERSION |
| 151 | } |
| 152 | |
| 153 | fn normalize_managed_dir(path: PathBuf) -> std::io::Result<PathBuf> { |
| 154 | if path.as_os_str().is_empty() { |
| 155 | return Err(std::io::Error::new( |
| 156 | std::io::ErrorKind::InvalidInput, |
| 157 | "managed directory path cannot be empty", |
| 158 | )); |
| 159 | } |
| 160 | if path.components().any(|component| { |
| 161 | matches!( |
| 162 | component, |
| 163 | Component::ParentDir | Component::Prefix(_) | Component::RootDir |
| 164 | ) |
| 165 | }) && path.is_relative() |
| 166 | { |
| 167 | return Err(std::io::Error::new( |
| 168 | std::io::ErrorKind::InvalidInput, |
| 169 | "managed directory path cannot contain traversal components", |
| 170 | )); |
| 171 | } |
| 172 | if path.is_absolute() { |
| 173 | return Ok(path); |
| 174 | } |
| 175 | std::env::current_dir().map(|cwd| cwd.join(path)) |
| 176 | } |
| 177 | |
| 178 | /// Open (creating if needed) an advisory-lock sidecar that is owner-only |
| 179 | /// (0600), one regular link, and never followed through a symlink. |
| 180 | pub(crate) fn open_private_lock_file(path: &Path) -> io::Result<fs::File> { |
| 181 | let mut options = OpenOptions::new(); |
| 182 | options.create(true).read(true).write(true); |
| 183 | #[cfg(unix)] |
| 184 | { |
| 185 | use std::os::unix::fs::{OpenOptionsExt, PermissionsExt}; |
| 186 | options |
| 187 | .mode(0o600) |
| 188 | .custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC | libc::O_NONBLOCK); |
| 189 | let file = options.open(path)?; |
| 190 | validate_private_regular_file(&file, path)?; |
| 191 | file.set_permissions(fs::Permissions::from_mode(0o600))?; |
| 192 | Ok(file) |
| 193 | } |
| 194 | #[cfg(windows)] |
| 195 | { |
| 196 | use std::os::windows::fs::OpenOptionsExt as _; |
| 197 | use windows_sys::Win32::Storage::FileSystem::FILE_FLAG_OPEN_REPARSE_POINT; |
| 198 | options.custom_flags(FILE_FLAG_OPEN_REPARSE_POINT); |
| 199 | let file = options.open(path)?; |
| 200 | validate_private_regular_file(&file, path)?; |
| 201 | Ok(file) |
| 202 | } |
| 203 | #[cfg(all(not(unix), not(windows)))] |
| 204 | { |
| 205 | let file = options.open(path)?; |
| 206 | validate_private_regular_file(&file, path)?; |
| 207 | Ok(file) |
| 208 | } |
| 209 | } |
| 210 | |
| 211 | fn open_private_read_file(path: &Path) -> io::Result<fs::File> { |
| 212 | let mut options = OpenOptions::new(); |
| 213 | options.read(true); |
| 214 | #[cfg(unix)] |
| 215 | { |
| 216 | use std::os::unix::fs::OpenOptionsExt as _; |
| 217 | options.custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC | libc::O_NONBLOCK); |
| 218 | } |
| 219 | #[cfg(windows)] |
| 220 | { |
| 221 | use std::os::windows::fs::OpenOptionsExt as _; |
| 222 | use windows_sys::Win32::Storage::FileSystem::FILE_FLAG_OPEN_REPARSE_POINT; |
| 223 | options.custom_flags(FILE_FLAG_OPEN_REPARSE_POINT); |
| 224 | } |
| 225 | let file = options.open(path)?; |
| 226 | validate_private_regular_file(&file, path)?; |
| 227 | Ok(file) |
| 228 | } |
| 229 | |
| 230 | #[cfg(unix)] |
| 231 | fn validate_private_regular_file(file: &fs::File, path: &Path) -> io::Result<()> { |
| 232 | use std::os::unix::fs::MetadataExt as _; |
| 233 | |
| 234 | let metadata = file.metadata()?; |
| 235 | if !metadata.is_file() || metadata.nlink() != 1 { |
| 236 | return Err(io::Error::new( |
| 237 | io::ErrorKind::InvalidData, |
| 238 | format!( |
| 239 | "private sidecar file {} must be one regular filesystem link", |
| 240 | path.display() |
| 241 | ), |
| 242 | )); |
| 243 | } |
| 244 | Ok(()) |
| 245 | } |
| 246 | |
| 247 | #[cfg(windows)] |
| 248 | fn validate_private_regular_file(file: &fs::File, path: &Path) -> io::Result<()> { |
| 249 | use std::os::windows::fs::MetadataExt as _; |
| 250 | use std::os::windows::io::AsRawHandle as _; |
| 251 | use windows_sys::Win32::Storage::FileSystem::{ |
| 252 | BY_HANDLE_FILE_INFORMATION, FILE_ATTRIBUTE_REPARSE_POINT, GetFileInformationByHandle, |
| 253 | }; |
| 254 | |
| 255 | let metadata = file.metadata()?; |
| 256 | if !metadata.is_file() || metadata.file_attributes() & FILE_ATTRIBUTE_REPARSE_POINT != 0 { |
| 257 | return Err(io::Error::new( |
| 258 | io::ErrorKind::InvalidData, |
| 259 | format!( |
| 260 | "private sidecar file {} must be a non-reparse regular file", |
| 261 | path.display() |
| 262 | ), |
| 263 | )); |
| 264 | } |
| 265 | let mut info = BY_HANDLE_FILE_INFORMATION::default(); |
| 266 | // SAFETY: `file` keeps the handle valid and `info` is writable for the call. |
| 267 | if unsafe { GetFileInformationByHandle(file.as_raw_handle(), &mut info) } == 0 { |
| 268 | return Err(io::Error::last_os_error()); |
| 269 | } |
| 270 | if info.nNumberOfLinks != 1 { |
| 271 | return Err(io::Error::new( |
| 272 | io::ErrorKind::InvalidData, |
| 273 | format!( |
| 274 | "private sidecar file {} must have exactly one filesystem link", |
| 275 | path.display() |
| 276 | ), |
| 277 | )); |
| 278 | } |
| 279 | Ok(()) |
| 280 | } |
| 281 | |
| 282 | #[cfg(all(not(unix), not(windows)))] |
| 283 | fn validate_private_regular_file(file: &fs::File, path: &Path) -> io::Result<()> { |
| 284 | if !file.metadata()?.is_file() { |
| 285 | return Err(io::Error::new( |
| 286 | io::ErrorKind::InvalidData, |
| 287 | format!("private sidecar file {} must be regular", path.display()), |
| 288 | )); |
| 289 | } |
| 290 | Ok(()) |
| 291 | } |
| 292 | |
| 293 | /// Persisted queued message for offline/degraded mode. |
| 294 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 295 | pub struct QueuedSessionMessage { |
| 296 | pub display: String, |
| 297 | #[serde(default)] |
| 298 | pub skill_instruction: Option<String>, |
| 299 | #[serde(default)] |
| 300 | pub skill_provenance: Option<crate::skills::SkillProvenance>, |
| 301 | } |
| 302 | |
| 303 | /// Persisted queue state for recovery after restart/crash. |
| 304 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 305 | pub struct OfflineQueueState { |
| 306 | #[serde(default = "default_queue_schema_version")] |
| 307 | pub schema_version: u32, |
| 308 | /// Session ID this queue belongs to. Redundant with the per-session file |
| 309 | /// name it is stored under; the UI's restore path still compares it |
| 310 | /// against the live session before adopting the messages. |
| 311 | #[serde(default)] |
| 312 | pub session_id: Option<String>, |
| 313 | #[serde(default)] |
| 314 | pub messages: Vec<QueuedSessionMessage>, |
| 315 | #[serde(default)] |
| 316 | pub draft: Option<QueuedSessionMessage>, |
| 317 | } |
| 318 | |
| 319 | /// Result of explicitly repairing a persisted session for process resume. |
| 320 | /// |
| 321 | /// Normal snapshot reads must not infer that an unmatched tool call crashed: |
| 322 | /// an embedding host can persist and inspect a session while that tool is |
| 323 | /// still running. Hosts should use [`SessionManager::load_session_snapshot`] |
| 324 | /// during normal operation and reserve this recovery path for a known process |
| 325 | /// or engine restart. |
| 326 | #[derive(Debug, Clone)] |
| 327 | pub struct SessionRecovery { |
| 328 | pub session: SavedSession, |
| 329 | pub changed: bool, |
| 330 | #[cfg_attr(not(test), expect(dead_code))] |
| 331 | pub repaired_call_count: usize, |
| 332 | #[cfg_attr(not(test), expect(dead_code))] |
| 333 | pub duplicate_result_count: usize, |
| 334 | #[cfg_attr(not(test), expect(dead_code))] |
| 335 | pub orphan_result_count: usize, |
| 336 | } |
| 337 | |
| 338 | impl Default for OfflineQueueState { |
| 339 | fn default() -> Self { |
| 340 | Self { |
| 341 | schema_version: CURRENT_QUEUE_SCHEMA_VERSION, |
| 342 | session_id: None, |
| 343 | messages: Vec::new(), |
| 344 | draft: None, |
| 345 | } |
| 346 | } |
| 347 | } |
| 348 | |
| 349 | /// Durable context-reference metadata attached to a user message. |
| 350 | #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] |
| 351 | pub struct SessionContextReference { |
| 352 | pub message_index: usize, |
| 353 | pub reference: ContextReference, |
| 354 | } |
| 355 | |
| 356 | /// Session metadata stored with each saved session |
| 357 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 358 | pub struct SessionMetadata { |
| 359 | /// Unique session identifier |
| 360 | pub id: String, |
| 361 | /// Actual host Runtime authority; independent of the conversation id. |
| 362 | /// Legacy/imported conversations have no binding until saved by a host. |
| 363 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 364 | pub runtime_store: Option<crate::runtime_threads::RuntimeStoreBinding>, |
| 365 | /// Human-readable title (derived from first message) |
| 366 | pub title: String, |
| 367 | /// When the session was created |
| 368 | pub created_at: DateTime<Utc>, |
| 369 | /// When the session was last updated |
| 370 | pub updated_at: DateTime<Utc>, |
| 371 | /// Number of messages in the session |
| 372 | pub message_count: usize, |
| 373 | /// Total tokens used |
| 374 | pub total_tokens: u64, |
| 375 | /// Model used for the session |
| 376 | pub model: String, |
| 377 | /// Provider used for the session model. Defaults for legacy saved sessions. |
| 378 | #[serde(default = "default_model_provider")] |
| 379 | pub model_provider: String, |
| 380 | /// Exact configured provider key. This is separate from `model_provider` |
| 381 | /// so old consumers can keep treating that field as the built-in provider |
| 382 | /// kind (`custom` for every named custom route). |
| 383 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 384 | pub model_provider_id: Option<String>, |
| 385 | /// Workspace directory |
| 386 | pub workspace: PathBuf, |
| 387 | /// Optional mode label (agent/plan/etc.) |
| 388 | #[serde(default)] |
| 389 | pub mode: Option<String>, |
| 390 | /// Accumulated cost data for persisted billing and high-water mark. |
| 391 | #[serde(default)] |
| 392 | pub cost: SessionCostSnapshot, |
| 393 | /// Source session id when this session was created with `deepseek fork`. |
| 394 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 395 | pub parent_session_id: Option<String>, |
| 396 | /// Source message count at fork time. This is intentionally coarse: |
| 397 | /// current saved sessions are linear JSON files, not per-entry trees. |
| 398 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 399 | pub forked_from_message_count: Option<usize>, |
| 400 | /// Cumulative turn duration in seconds (sum of completed turn elapsed |
| 401 | /// times). Persisted so the footer "worked" chip survives restarts |
| 402 | /// (#2038). |
| 403 | #[serde(default)] |
| 404 | pub cumulative_turn_secs: u64, |
| 405 | /// Durable archive flag (#2934 / #4397). Archived sessions stay on disk |
| 406 | /// and stay loadable; they are hidden from the default browse surfaces |
| 407 | /// and are never chosen by auto-resume. |
| 408 | /// |
| 409 | /// This mirrors `ThreadRecord::archived` in [`crate::runtime_threads`] so |
| 410 | /// the TUI session surfaces and the Runtime API/web dashboard project the |
| 411 | /// same lifecycle field instead of two divergent notions of "put away". |
| 412 | /// Additive and `skip_serializing_if`-guarded: sessions written before |
| 413 | /// v0.9.2 load as `archived = false` and round-trip byte-identically |
| 414 | /// until the flag is actually set. |
| 415 | #[serde(default, skip_serializing_if = "is_not_archived")] |
| 416 | pub archived: bool, |
| 417 | #[serde(default)] |
| 418 | pub spawn_depth: u32, |
| 419 | } |
| 420 | |
| 421 | fn is_not_archived(archived: &bool) -> bool { |
| 422 | !*archived |
| 423 | } |
| 424 | |
| 425 | /// Sessions currently owned by an in-process interactive surface (the TUI). |
| 426 | /// |
| 427 | /// A saved session is a file, and a running TUI holds the authoritative copy |
| 428 | /// in memory: it autosaves the whole document from `App` state. That makes an |
| 429 | /// out-of-band write to the *same* session unsafe — the next autosave would |
| 430 | /// silently revert it. Rather than let that happen quietly, the owner claims |
| 431 | /// the id here and any external writer is refused. |
| 432 | /// |
| 433 | /// A static registry rather than a field on `RuntimeApiState` because the |
| 434 | /// embedded Runtime API runs inside the TUI process. A standalone |
| 435 | /// `codewhale web` has an empty registry, so external writers consult |
| 436 | /// [`SessionManager::is_session_live_anywhere`], which also sees the |
| 437 | /// cross-process lease a TUI in another process holds. |
| 438 | static LIVE_SESSIONS: std::sync::OnceLock<std::sync::RwLock<std::collections::HashSet<String>>> = |
| 439 | std::sync::OnceLock::new(); |
| 440 | |
| 441 | fn live_sessions() -> &'static std::sync::RwLock<std::collections::HashSet<String>> { |
| 442 | LIVE_SESSIONS.get_or_init(Default::default) |
| 443 | } |
| 444 | |
| 445 | /// Who is asking to mutate a saved session. |
| 446 | /// |
| 447 | /// This is an authority distinction, not a convenience one: the owner may |
| 448 | /// write because it will update its in-memory copy in the same step; anyone |
| 449 | /// else may not, because it cannot. |
| 450 | #[derive(Debug, Clone, Copy, PartialEq, Eq)] |
| 451 | pub enum SessionMutator { |
| 452 | /// The in-process surface that currently owns the session (the TUI). It |
| 453 | /// is responsible for updating its cached metadata atomically with the |
| 454 | /// write — see `App::apply_session_mutation`. |
| 455 | Owner, |
| 456 | /// Any other writer: the Runtime API, the web dashboard, a second |
| 457 | /// process. Refused while the session is claimed. |
| 458 | External, |
| 459 | } |
| 460 | |
| 461 | /// Set the claimed session to exactly `session_id` (or nothing). |
| 462 | /// |
| 463 | /// The TUI owns at most one session at a time, so switching sessions must |
| 464 | /// release the previous claim in the same step — otherwise a `/new` would |
| 465 | /// leave the old id permanently locked against the dashboard. |
| 466 | pub fn set_live_session(session_id: Option<&str>) { |
| 467 | let session_id = session_id.map(str::trim).filter(|id| !id.is_empty()); |
| 468 | if let Ok(mut live) = live_sessions().write() { |
| 469 | live.clear(); |
| 470 | if let Some(id) = session_id { |
| 471 | live.insert(id.to_string()); |
| 472 | } |
| 473 | } |
| 474 | // A lease for any other session is released with the claim it backed. |
| 475 | if let Ok(mut lease) = LIVE_SESSION_LEASE.lock() |
| 476 | && lease.as_ref().map(|(id, _)| id.as_str()) != session_id |
| 477 | { |
| 478 | *lease = None; |
| 479 | } |
| 480 | } |
| 481 | |
| 482 | /// The cross-process half of the live claim (#6144): an exclusive lock on |
| 483 | /// `.late-usage/<id>.live`, held for as long as this process owns the |
| 484 | /// session. The registry above only protects against writers in this |
| 485 | /// process; a standalone `codewhale serve` or a second TUI could still |
| 486 | /// rewrite or delete the document an interactive session is about to |
| 487 | /// autosave over. This lock is separate from the per-write `<id>.lock`, so the |
| 488 | /// owner's own saves never contend with it. It is a liveness signal only — |
| 489 | /// the file carries no data. `None` for the file records a lease that could |
| 490 | /// not be taken (another process holds it), so it is not retried per save. |
| 491 | static LIVE_SESSION_LEASE: std::sync::Mutex<Option<(String, Option<fs::File>)>> = |
| 492 | std::sync::Mutex::new(None); |
| 493 | |
| 494 | /// Lock attempts [`SessionManager::reserve_session_for_attach`] makes |
| 495 | /// before it treats contention as a running owner. |
| 496 | const LIVE_LEASE_ATTACH_ATTEMPTS: u32 = 3; |
| 497 | |
| 498 | /// The refusal for attaching to a session another process has open. |
| 499 | pub(crate) fn session_open_elsewhere(session_id: &str) -> std::io::Error { |
| 500 | std::io::Error::new( |
| 501 | std::io::ErrorKind::ResourceBusy, |
| 502 | format!( |
| 503 | "session {session_id} is open in another Codewhale window. \ |
| 504 | Continue it there, or run `codewhale fork {session_id}` to work on a copy" |
| 505 | ), |
| 506 | ) |
| 507 | } |
| 508 | |
| 509 | /// A reservation of a session's cross-process live lease, taken by |
| 510 | /// [`SessionManager::reserve_session_for_attach`] before a surface attaches. |
| 511 | /// Dropping it releases the reservation; [`Self::commit`] makes the session |
| 512 | /// this process's live session and releases the previous one's lease. |
| 513 | #[must_use = "an uncommitted reservation is released when dropped"] |
| 514 | #[derive(Debug)] |
| 515 | pub struct SessionLease { |
| 516 | id: String, |
| 517 | /// `None` when this process already holds the session's lease. |
| 518 | file: Option<fs::File>, |
| 519 | } |
| 520 | |
| 521 | impl SessionLease { |
| 522 | /// Make the reserved session this process's live session. Called after |
| 523 | /// the attach has been validated and applied, so a failed attach never |
| 524 | /// costs the session this process already owns its lease. |
| 525 | pub fn commit(self) { |
| 526 | set_live_session(Some(&self.id)); |
| 527 | if let Some(file) = self.file |
| 528 | && let Ok(mut lease) = LIVE_SESSION_LEASE.lock() |
| 529 | { |
| 530 | *lease = Some((self.id, Some(file))); |
| 531 | } |
| 532 | } |
| 533 | } |
| 534 | |
| 535 | /// Is this session currently owned by **this process's** interactive surface? |
| 536 | /// |
| 537 | /// The registry is process-local. Reclamation must not treat a missing entry |
| 538 | /// here as proof that no other Codewhale process still owns the directory. |
| 539 | #[must_use] |
| 540 | pub fn is_live_session(session_id: &str) -> bool { |
| 541 | live_sessions() |
| 542 | .read() |
| 543 | .is_ok_and(|live| live.contains(session_id)) |
| 544 | } |
| 545 | |
| 546 | /// The error an external writer gets when the session is live. |
| 547 | /// |
| 548 | /// `ResourceBusy` so callers can map it to a typed conflict rather than |
| 549 | /// pattern-matching on a message. |
| 550 | pub(crate) fn live_session_conflict(session_id: &str) -> std::io::Error { |
| 551 | std::io::Error::new( |
| 552 | std::io::ErrorKind::ResourceBusy, |
| 553 | format!( |
| 554 | "session '{session_id}' is open in an interactive Codewhale session; \ |
| 555 | change it there instead — an external write would be reverted by its next autosave" |
| 556 | ), |
| 557 | ) |
| 558 | } |
| 559 | |
| 560 | /// File-name stem of the sidecar mapping session ids to the session |
| 561 | /// instance (process boot) that created their persisted record. Lives in |
| 562 | /// the sessions directory next to the `<id>.json` records it describes. |
| 563 | const SESSION_BOOT_OWNERS_STEM: &str = "session_boot_owners"; |
| 564 | |
| 565 | static SESSION_BOOT_ID: std::sync::OnceLock<String> = std::sync::OnceLock::new(); |
| 566 | |
| 567 | /// Identity of this running session instance (one per process boot). |
| 568 | /// |
| 569 | /// Mirrors the `SubAgentManager` boot id from #405: persisted records are |
| 570 | /// stamped with the instance that created them, so a later Codewhale |
| 571 | /// instance in the same workspace can tell restored rows from its own live |
| 572 | /// work (#4416). |
| 573 | #[must_use] |
| 574 | pub fn current_session_boot_id() -> &'static str { |
| 575 | SESSION_BOOT_ID.get_or_init(|| format!("boot_{}", &Uuid::new_v4().to_string()[..12])) |
| 576 | } |
| 577 | |
| 578 | /// Which archive states a session listing includes. |
| 579 | /// |
| 580 | /// Deliberately the same three-way shape as |
| 581 | /// [`crate::runtime_threads::ThreadListFilter`] so `/v1/sessions` and |
| 582 | /// `/v1/threads` answer the same `include_archived` / `archived_only` query |
| 583 | /// pair with the same semantics. |
| 584 | #[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] |
| 585 | pub enum SessionListFilter { |
| 586 | /// Only `archived = false` sessions. The browse default. |
| 587 | #[default] |
| 588 | ActiveOnly, |
| 589 | /// Active and archived sessions, newest first. |
| 590 | IncludeArchived, |
| 591 | /// Only `archived = true` sessions. |
| 592 | ArchivedOnly, |
| 593 | } |
| 594 | |
| 595 | impl SessionListFilter { |
| 596 | /// Resolve the `include_archived` / `archived_only` query pair the same |
| 597 | /// way the threads routes do. |
| 598 | #[must_use] |
| 599 | pub fn from_query(include_archived: Option<bool>, archived_only: Option<bool>) -> Self { |
| 600 | if archived_only.unwrap_or(false) { |
| 601 | Self::ArchivedOnly |
| 602 | } else if include_archived.unwrap_or(false) { |
| 603 | Self::IncludeArchived |
| 604 | } else { |
| 605 | Self::ActiveOnly |
| 606 | } |
| 607 | } |
| 608 | |
| 609 | #[must_use] |
| 610 | pub fn admits(self, archived: bool) -> bool { |
| 611 | match self { |
| 612 | Self::ActiveOnly => !archived, |
| 613 | Self::IncludeArchived => true, |
| 614 | Self::ArchivedOnly => archived, |
| 615 | } |
| 616 | } |
| 617 | } |
| 618 | |
| 619 | fn default_model_provider() -> String { |
| 620 | "deepseek".to_string() |
| 621 | } |
| 622 | |
| 623 | impl SessionMetadata { |
| 624 | pub(crate) fn set_model_provider_route(&mut self, kind: &str, identity: Option<&str>) { |
| 625 | self.model_provider = kind.to_string(); |
| 626 | self.model_provider_id = identity.map(str::to_string); |
| 627 | } |
| 628 | } |
| 629 | |
| 630 | /// Cost and high-water-mark fields persisted with each session. |
| 631 | /// |
| 632 | /// The coverage fields below are persisted **alongside** the money so a restored |
| 633 | /// session can still say what its total covers. Without them a reload produced a |
| 634 | /// dollar figure with no completeness information, which then rendered as "0 of 0 |
| 635 | /// turns priced" — a fabricated claim of a complete total. Sessions written |
| 636 | /// before these fields existed deserialize them from `Default`, which is |
| 637 | /// indistinguishable from that same false reading, so the load path detects the |
| 638 | /// legacy shape explicitly (see [`Self::coverage_is_legacy_unknown`]) rather than |
| 639 | /// trusting the defaults (#4318). |
| 640 | #[derive(Debug, Clone, Default, Serialize, Deserialize)] |
| 641 | pub struct SessionCostSnapshot { |
| 642 | /// Accumulated parent-turn session cost in USD. |
| 643 | #[serde(default)] |
| 644 | pub session_cost_usd: f64, |
| 645 | /// Accumulated parent-turn session cost in CNY. |
| 646 | #[serde(default)] |
| 647 | pub session_cost_cny: f64, |
| 648 | /// Accumulated sub-agent/background LLM cost in USD. |
| 649 | #[serde(default)] |
| 650 | pub subagent_cost_usd: f64, |
| 651 | /// Accumulated sub-agent/background LLM cost in CNY. |
| 652 | #[serde(default)] |
| 653 | pub subagent_cost_cny: f64, |
| 654 | /// Max-ever displayed session+subagent cost in USD (preserves #244 |
| 655 | /// monotonic guarantee across session restarts). |
| 656 | #[serde(default)] |
| 657 | pub displayed_cost_high_water_usd: f64, |
| 658 | /// Max-ever displayed session+subagent cost in CNY. |
| 659 | #[serde(default)] |
| 660 | pub displayed_cost_high_water_cny: f64, |
| 661 | /// Turns whose route was money-metered and produced an authoritative price. |
| 662 | /// These are exactly the turns the persisted totals contain. |
| 663 | #[serde(default)] |
| 664 | pub priced_turns: u32, |
| 665 | /// Money-metered (or unknown-basis) turns that produced no authoritative |
| 666 | /// price, so their spend is missing from the persisted totals. |
| 667 | #[serde(default)] |
| 668 | pub unpriced_turns: u32, |
| 669 | /// CNY-specific coverage. USD-only routes are unpriced in CNY rather than |
| 670 | /// silently contributing a fabricated zero. |
| 671 | #[serde(default)] |
| 672 | pub cny_priced_turns: u32, |
| 673 | #[serde(default)] |
| 674 | pub cny_unpriced_turns: u32, |
| 675 | /// Stable reason labels for the unpriced turns. |
| 676 | #[serde(default, skip_serializing_if = "BTreeSet::is_empty")] |
| 677 | pub unpriced_reasons: BTreeSet<String>, |
| 678 | #[serde(default, skip_serializing_if = "BTreeSet::is_empty")] |
| 679 | pub cny_unpriced_reasons: BTreeSet<String>, |
| 680 | /// Token classes used on some route that carry no published price. |
| 681 | #[serde(default, skip_serializing_if = "BTreeSet::is_empty")] |
| 682 | pub unpriced_classes: BTreeSet<String>, |
| 683 | /// Provenance labels of the pricing rows the totals were built from. |
| 684 | #[serde(default, skip_serializing_if = "BTreeSet::is_empty")] |
| 685 | pub pricing_provenances: BTreeSet<String>, |
| 686 | /// Live-pricing downgrade receipts recorded while building the totals. |
| 687 | #[serde(default, skip_serializing_if = "BTreeSet::is_empty")] |
| 688 | pub live_pricing_defects: BTreeSet<String>, |
| 689 | /// Live rows that failed validation and had no usable bundled fallback. |
| 690 | #[serde(default, skip_serializing_if = "BTreeSet::is_empty")] |
| 691 | pub live_pricing_unusable_defects: BTreeSet<String>, |
| 692 | /// Redacted per-route receipts: provider, configured identity, wire model, |
| 693 | /// billing surface, endpoint fingerprint, billing mode, currency. Never a URL, a |
| 694 | /// credential, or a filesystem path. |
| 695 | #[serde(default, skip_serializing_if = "BTreeSet::is_empty")] |
| 696 | pub route_receipts: BTreeSet<String>, |
| 697 | /// Redacted provider-response identities already included in the live and |
| 698 | /// durable sub-agent totals. Worker records persist the same fingerprints. |
| 699 | /// |
| 700 | /// Known limit (#6842): unbounded, ~70 bytes per background call. It is |
| 701 | /// the replay-idempotency set for money (late-usage overlay and restored |
| 702 | /// dedupe), so it is never capped; `load_session_metadata` grows its read |
| 703 | /// to fit a large block instead of falling back to a whole-file read. |
| 704 | #[serde(default, skip_serializing_if = "BTreeSet::is_empty")] |
| 705 | pub usage_source_fingerprints: BTreeSet<String>, |
| 706 | #[serde( |
| 707 | default, |
| 708 | skip_serializing_if = "BTreeMap::is_empty", |
| 709 | deserialize_with = "crate::cost_status::deserialize_missing_usage_sources" |
| 710 | )] |
| 711 | pub missing_usage_sources: BTreeMap<String, crate::cost_status::MissingUsageCoverage>, |
| 712 | #[serde(default, skip_serializing_if = "is_false")] |
| 713 | pub missing_usage_overflowed: bool, |
| 714 | /// Written by builds that track coverage, so a reader can tell "this session |
| 715 | /// genuinely had zero money-metered turns" apart from "this session predates |
| 716 | /// coverage tracking". Absent on legacy rows. |
| 717 | #[serde(default)] |
| 718 | pub coverage_recorded: bool, |
| 719 | } |
| 720 | |
| 721 | impl SessionCostSnapshot { |
| 722 | fn absorb_late_background_cost(&mut self, pool: &crate::cost_status::PendingBackgroundCost) { |
| 723 | let pool = crate::cost_status::project_missing_usage_ledger( |
| 724 | &mut self.missing_usage_sources, |
| 725 | &mut self.missing_usage_overflowed, |
| 726 | &mut self.unpriced_turns, |
| 727 | &mut self.cny_unpriced_turns, |
| 728 | pool, |
| 729 | ); |
| 730 | let estimate = crate::pricing::CostEstimate { |
| 731 | usd: self.subagent_cost_usd, |
| 732 | cny: self.subagent_cost_cny, |
| 733 | } |
| 734 | .saturating_add(pool.estimate); |
| 735 | self.subagent_cost_usd = estimate.usd; |
| 736 | self.subagent_cost_cny = estimate.cny; |
| 737 | self.priced_turns = self.priced_turns.saturating_add(pool.priced_turns); |
| 738 | self.unpriced_turns = self.unpriced_turns.saturating_add(pool.unpriced_turns); |
| 739 | self.cny_priced_turns = self.cny_priced_turns.saturating_add(pool.cny_priced_turns); |
| 740 | self.cny_unpriced_turns = self |
| 741 | .cny_unpriced_turns |
| 742 | .saturating_add(pool.cny_unpriced_turns); |
| 743 | self.unpriced_reasons |
| 744 | .extend(pool.unpriced_reasons.iter().map(ToString::to_string)); |
| 745 | self.cny_unpriced_reasons |
| 746 | .extend(pool.cny_unpriced_reasons.iter().map(ToString::to_string)); |
| 747 | self.unpriced_classes |
| 748 | .extend(pool.unpriced_classes.iter().map(ToString::to_string)); |
| 749 | self.pricing_provenances |
| 750 | .extend(pool.pricing_provenances.iter().map(ToString::to_string)); |
| 751 | self.live_pricing_defects |
| 752 | .extend(pool.live_pricing_defects.iter().map(ToString::to_string)); |
| 753 | self.live_pricing_unusable_defects.extend( |
| 754 | pool.live_pricing_unusable_defects |
| 755 | .iter() |
| 756 | .map(ToString::to_string), |
| 757 | ); |
| 758 | self.route_receipts |
| 759 | .extend(pool.route_receipts.iter().cloned()); |
| 760 | self.usage_source_fingerprints |
| 761 | .extend(pool.usage_source_fingerprints.iter().cloned()); |
| 762 | self.coverage_recorded = true; |
| 763 | let total = self.total_estimate(); |
| 764 | self.displayed_cost_high_water_usd = self.displayed_cost_high_water_usd.max(total.usd); |
| 765 | self.displayed_cost_high_water_cny = self.displayed_cost_high_water_cny.max(total.cny); |
| 766 | } |
| 767 | |
| 768 | /// Session + subagent spend as **one** dual-currency accumulator. |
| 769 | /// |
| 770 | /// The persisted USD and CNY columns are projections of per-turn |
| 771 | /// [`crate::pricing::CostEstimate`]s that were accumulated jointly; every |
| 772 | /// display total is derived from this single fold so the two currencies |
| 773 | /// cannot be re-summed by separate code paths that then drift (#4939). |
| 774 | /// CNY is *not* an FX multiple of USD: a turn carries CNY only when its |
| 775 | /// route published an authoritative CNY row (provider-published |
| 776 | /// dual-currency pricing, e.g. DeepSeek's CNY table), and a USD-only turn |
| 777 | /// contributes exactly zero CNY while `cny_unpriced_turns` records the gap. |
| 778 | #[must_use] |
| 779 | pub fn total_estimate(&self) -> crate::pricing::CostEstimate { |
| 780 | crate::pricing::CostEstimate { |
| 781 | usd: self.session_cost_usd, |
| 782 | cny: self.session_cost_cny, |
| 783 | } |
| 784 | .saturating_add(crate::pricing::CostEstimate { |
| 785 | usd: self.subagent_cost_usd, |
| 786 | cny: self.subagent_cost_cny, |
| 787 | }) |
| 788 | } |
| 789 | |
| 790 | /// Session + subagent cost in USD. |
| 791 | pub fn total_usd(&self) -> f64 { |
| 792 | self.total_estimate() |
| 793 | .amount(crate::pricing::CostCurrency::Usd) |
| 794 | } |
| 795 | |
| 796 | /// Session + subagent cost in CNY. |
| 797 | pub fn total_cny(&self) -> f64 { |
| 798 | self.total_estimate() |
| 799 | .amount(crate::pricing::CostCurrency::Cny) |
| 800 | } |
| 801 | |
| 802 | /// Whether this snapshot's coverage state must be shown as unknown. |
| 803 | /// |
| 804 | /// True when the snapshot has no coverage evidence — the signature of a |
| 805 | /// session written before coverage was persisted. Reporting any such |
| 806 | /// session as "0 of 0 priced" would claim completeness without evidence, |
| 807 | /// including when the saved amount is zero. |
| 808 | #[must_use] |
| 809 | pub fn coverage_is_legacy_unknown(&self) -> bool { |
| 810 | !self.coverage_recorded |
| 811 | } |
| 812 | } |
| 813 | |
| 814 | impl SessionMetadata { |
| 815 | /// Copy cost fields from another metadata (used when forking a session). |
| 816 | pub fn copy_cost_from(&mut self, other: &SessionMetadata) { |
| 817 | self.cost = other.cost.clone(); |
| 818 | } |
| 819 | |
| 820 | /// Record additive lineage metadata for a forked saved session. |
| 821 | pub fn mark_forked_from(&mut self, parent: &SessionMetadata) { |
| 822 | self.parent_session_id = Some(parent.id.clone()); |
| 823 | self.forked_from_message_count = Some(parent.message_count); |
| 824 | } |
| 825 | } |
| 826 | |
| 827 | /// Durable Work-panel state. Optional on [`SavedSession`] so every session |
| 828 | /// written before v0.8.68 remains loadable without migration. |
| 829 | #[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)] |
| 830 | pub struct SessionWorkState { |
| 831 | /// Authoritative Work Graph. Optional so pre-Work-Graph sessions and old |
| 832 | /// binaries continue to exchange fully populated Plan/To-do views. |
| 833 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 834 | pub graph: Option<crate::work_graph::WorkGraphSnapshot>, |
| 835 | #[serde(default, skip_serializing_if = "TodoListSnapshot::is_empty")] |
| 836 | pub todos: TodoListSnapshot, |
| 837 | #[serde(default, skip_serializing_if = "PlanSnapshot::is_empty")] |
| 838 | pub plan: PlanSnapshot, |
| 839 | } |
| 840 | |
| 841 | /// Bounded goal projection persisted beside the owning saved session. |
| 842 | /// |
| 843 | /// This intentionally excludes completion prose, verifier output, transcripts, |
| 844 | /// and filesystem evidence. The saved session already owns conversation |
| 845 | /// history; restart only needs the typed control state that makes the next turn |
| 846 | /// continue the same objective without trusting text reconstructed from it. |
| 847 | #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] |
| 848 | #[serde(deny_unknown_fields)] |
| 849 | pub struct SessionGoalState { |
| 850 | #[serde(default = "current_session_goal_schema_version")] |
| 851 | pub schema_version: u32, |
| 852 | pub objective: String, |
| 853 | pub status: SessionGoalStatus, |
| 854 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 855 | pub token_budget: Option<u32>, |
| 856 | #[serde(default)] |
| 857 | pub tokens_used: u64, |
| 858 | #[serde(default)] |
| 859 | pub time_used_seconds: u64, |
| 860 | #[serde(default)] |
| 861 | pub continuation_count: u32, |
| 862 | #[serde(default)] |
| 863 | pub elapsed_seconds: u64, |
| 864 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 865 | pub pause_reason: Option<GoalPauseReason>, |
| 866 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 867 | pub goal_id: Option<String>, |
| 868 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 869 | pub last_gap_fingerprint: Option<String>, |
| 870 | #[serde(default)] |
| 871 | pub repeated_gap_count: u32, |
| 872 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 873 | pub last_gap_pass: Option<u32>, |
| 874 | } |
| 875 | |
| 876 | #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] |
| 877 | #[serde(rename_all = "snake_case")] |
| 878 | pub enum SessionGoalStatus { |
| 879 | Active, |
| 880 | Paused, |
| 881 | Complete, |
| 882 | Blocked, |
| 883 | } |
| 884 | |
| 885 | const fn current_session_goal_schema_version() -> u32 { |
| 886 | CURRENT_SESSION_GOAL_SCHEMA_VERSION |
| 887 | } |
| 888 | |
| 889 | impl SessionGoalState { |
| 890 | /// Convert a runtime update into the durable, bounded session contract. |
| 891 | /// The canonical empty runtime snapshot removes the sidecar. |
| 892 | pub fn from_runtime(snapshot: &GoalSnapshot) -> io::Result<Option<Self>> { |
| 893 | if snapshot.objective.is_none() && snapshot.status.trim() == "none" { |
| 894 | return Ok(None); |
| 895 | } |
| 896 | let objective = snapshot |
| 897 | .objective |
| 898 | .as_deref() |
| 899 | .map(str::trim) |
| 900 | .filter(|objective| !objective.is_empty()) |
| 901 | .ok_or_else(|| { |
| 902 | io::Error::new(io::ErrorKind::InvalidData, "goal snapshot has no objective") |
| 903 | })?; |
| 904 | let status = match snapshot.status.trim() { |
| 905 | "active" => SessionGoalStatus::Active, |
| 906 | "paused" => SessionGoalStatus::Paused, |
| 907 | "complete" => SessionGoalStatus::Complete, |
| 908 | "blocked" => SessionGoalStatus::Blocked, |
| 909 | other => { |
| 910 | return Err(io::Error::new( |
| 911 | io::ErrorKind::InvalidData, |
| 912 | format!("goal snapshot has unsupported status '{other}'"), |
| 913 | )); |
| 914 | } |
| 915 | }; |
| 916 | let state = Self { |
| 917 | schema_version: CURRENT_SESSION_GOAL_SCHEMA_VERSION, |
| 918 | objective: objective.to_string(), |
| 919 | status, |
| 920 | token_budget: snapshot.token_budget, |
| 921 | tokens_used: snapshot.tokens_used, |
| 922 | time_used_seconds: snapshot.time_used_seconds, |
| 923 | continuation_count: snapshot.continuation_count, |
| 924 | elapsed_seconds: snapshot.elapsed_seconds.unwrap_or_default(), |
| 925 | pause_reason: snapshot.pause_reason, |
| 926 | goal_id: snapshot.goal_id.clone(), |
| 927 | last_gap_fingerprint: snapshot.last_gap_fingerprint.clone(), |
| 928 | repeated_gap_count: snapshot.repeated_gap_count, |
| 929 | last_gap_pass: snapshot.last_gap_pass, |
| 930 | }; |
| 931 | state.validate()?; |
| 932 | Ok(Some(state)) |
| 933 | } |
| 934 | |
| 935 | pub fn validate(&self) -> io::Result<()> { |
| 936 | if self.schema_version > CURRENT_SESSION_GOAL_SCHEMA_VERSION { |
| 937 | return Err(io::Error::new( |
| 938 | io::ErrorKind::InvalidData, |
| 939 | format!( |
| 940 | "Session goal schema v{} is newer than supported v{}", |
| 941 | self.schema_version, CURRENT_SESSION_GOAL_SCHEMA_VERSION |
| 942 | ), |
| 943 | )); |
| 944 | } |
| 945 | let objective = self.objective.trim(); |
| 946 | if objective.is_empty() || objective.chars().count() > MAX_SESSION_GOAL_OBJECTIVE_CHARS { |
| 947 | return Err(io::Error::new( |
| 948 | io::ErrorKind::InvalidData, |
| 949 | format!( |
| 950 | "Session goal objective must contain 1..={MAX_SESSION_GOAL_OBJECTIVE_CHARS} characters" |
| 951 | ), |
| 952 | )); |
| 953 | } |
| 954 | if self |
| 955 | .goal_id |
| 956 | .as_ref() |
| 957 | .is_some_and(|id| id.is_empty() || id.len() > 128) |
| 958 | { |
| 959 | return Err(io::Error::new( |
| 960 | io::ErrorKind::InvalidData, |
| 961 | "invalid session goal revision", |
| 962 | )); |
| 963 | } |
| 964 | codewhale_protocol::validate_goal_stall_state( |
| 965 | self.last_gap_fingerprint.as_deref(), |
| 966 | self.repeated_gap_count, |
| 967 | self.last_gap_pass, |
| 968 | self.continuation_count, |
| 969 | ) |
| 970 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error)) |
| 971 | } |
| 972 | |
| 973 | #[must_use] |
| 974 | pub fn to_runtime_snapshot(&self) -> GoalSnapshot { |
| 975 | GoalSnapshot { |
| 976 | goal_id: self.goal_id.clone(), |
| 977 | objective: Some(self.objective.clone()), |
| 978 | status: match self.status { |
| 979 | SessionGoalStatus::Active => "active", |
| 980 | SessionGoalStatus::Paused => "paused", |
| 981 | SessionGoalStatus::Complete => "complete", |
| 982 | SessionGoalStatus::Blocked => "blocked", |
| 983 | } |
| 984 | .to_string(), |
| 985 | token_budget: self.token_budget, |
| 986 | tokens_used: self.tokens_used, |
| 987 | time_used_seconds: self.time_used_seconds, |
| 988 | continuation_count: self.continuation_count, |
| 989 | elapsed_seconds: Some(self.elapsed_seconds), |
| 990 | evidence: None, |
| 991 | blocker: None, |
| 992 | pause_reason: self.pause_reason, |
| 993 | completion_verification: None, |
| 994 | advisories: Vec::new(), |
| 995 | last_gap_fingerprint: self.last_gap_fingerprint.clone(), |
| 996 | repeated_gap_count: self.repeated_gap_count, |
| 997 | last_gap_pass: self.last_gap_pass, |
| 998 | progress: None, |
| 999 | } |
| 1000 | } |
| 1001 | } |
| 1002 | |
| 1003 | impl SessionWorkState { |
| 1004 | #[must_use] |
| 1005 | pub fn is_empty(&self) -> bool { |
| 1006 | self.graph |
| 1007 | .as_ref() |
| 1008 | .is_none_or(crate::work_graph::WorkGraphSnapshot::is_empty) |
| 1009 | && self.todos.is_empty() |
| 1010 | && self.plan.is_empty() |
| 1011 | } |
| 1012 | } |
| 1013 | |
| 1014 | /// Latest concrete Auto route and the decision receipt that produced it. |
| 1015 | /// |
| 1016 | /// This is additive, optional session metadata: sessions written before |
| 1017 | /// v0.9.1 deserialize with no receipt and keep their legacy restore behavior. |
| 1018 | #[derive(Debug, Clone, PartialEq, Eq)] |
| 1019 | pub(crate) struct SavedAutoRouteReceipt { |
| 1020 | pub(crate) provider: ProviderKind, |
| 1021 | pub(crate) provider_identity: String, |
| 1022 | pub(crate) model: String, |
| 1023 | pub(crate) receipt: AutoRouteReceipt, |
| 1024 | /// Canonical effective reasoning receipt for the selected route, including |
| 1025 | /// routes where a concrete tier cannot be proven. Optional so older |
| 1026 | /// sessions remain loadable. |
| 1027 | pub(crate) effective_reasoning_effort: Option<ReasoningEffortTier>, |
| 1028 | } |
| 1029 | |
| 1030 | #[derive(Serialize, Deserialize)] |
| 1031 | struct SavedAutoRouteReceiptWire { |
| 1032 | provider: String, |
| 1033 | provider_identity: String, |
| 1034 | model: String, |
| 1035 | receipt: AutoRouteReceipt, |
| 1036 | /// Canonical effective reasoning receipt for the selected route, including |
| 1037 | /// routes where a concrete tier cannot be proven. Optional so older |
| 1038 | /// sessions remain loadable. |
| 1039 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 1040 | effective_reasoning_effort: Option<ReasoningEffortTier>, |
| 1041 | } |
| 1042 | |
| 1043 | impl Serialize for SavedAutoRouteReceipt { |
| 1044 | fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error> |
| 1045 | where |
| 1046 | S: serde::Serializer, |
| 1047 | { |
| 1048 | let provider = codewhale_config::descriptors::tui_wire_tag_for_route( |
| 1049 | self.provider, |
| 1050 | &self.provider_identity, |
| 1051 | ) |
| 1052 | .ok_or_else(|| serde::ser::Error::custom("contradictory saved provider identity"))?; |
| 1053 | SavedAutoRouteReceiptWire { |
| 1054 | provider: provider.into(), |
| 1055 | provider_identity: self.provider_identity.clone(), |
| 1056 | model: self.model.clone(), |
| 1057 | receipt: self.receipt.clone(), |
| 1058 | effective_reasoning_effort: self.effective_reasoning_effort, |
| 1059 | } |
| 1060 | .serialize(serializer) |
| 1061 | } |
| 1062 | } |
| 1063 | impl<'de> Deserialize<'de> for SavedAutoRouteReceipt { |
| 1064 | fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> |
| 1065 | where |
| 1066 | D: serde::Deserializer<'de>, |
| 1067 | { |
| 1068 | let wire = SavedAutoRouteReceiptWire::deserialize(deserializer)?; |
| 1069 | let provider = codewhale_config::descriptors::kind_from_tui_wire_tag( |
| 1070 | &wire.provider, |
| 1071 | &wire.provider_identity, |
| 1072 | ) |
| 1073 | .ok_or_else(|| serde::de::Error::custom("contradictory saved provider identity"))?; |
| 1074 | Ok(Self { |
| 1075 | provider, |
| 1076 | provider_identity: wire.provider_identity, |
| 1077 | model: wire.model, |
| 1078 | receipt: wire.receipt, |
| 1079 | effective_reasoning_effort: wire.effective_reasoning_effort, |
| 1080 | }) |
| 1081 | } |
| 1082 | } |
| 1083 | |
| 1084 | /// Most turn outcomes one session record keeps; the oldest drop first. |
| 1085 | pub(crate) const MAX_SAVED_TURN_OUTCOMES: usize = 64; |
| 1086 | /// Longest error text one saved outcome keeps, in `char`s. |
| 1087 | const MAX_SAVED_TURN_OUTCOME_ERROR_CHARS: usize = 4_000; |
| 1088 | |
| 1089 | /// A turn that ended `Failed`, as the person saw it end. |
| 1090 | /// |
| 1091 | /// The transcript only holds messages, so before this record a failed turn |
| 1092 | /// left nothing but the user's prompt behind: once the TUI closed, resume, |
| 1093 | /// export, and the app had no way to say why the turn stopped. The error is |
| 1094 | /// the text the live transcript showed, passed through the shared secret |
| 1095 | /// redactor before it is stored. |
| 1096 | #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] |
| 1097 | pub struct SavedTurnOutcome { |
| 1098 | pub status: crate::core::events::TurnOutcomeStatus, |
| 1099 | /// User-facing error text, secrets redacted. |
| 1100 | pub error: String, |
| 1101 | pub ended_at: DateTime<Utc>, |
| 1102 | /// Transcript messages that existed when the turn ended. Resume places |
| 1103 | /// the notice after that many messages (clamped to the transcript). |
| 1104 | pub after_message_count: usize, |
| 1105 | } |
| 1106 | |
| 1107 | impl SavedTurnOutcome { |
| 1108 | /// Build the persisted record for a failed turn. Redacts before bounding |
| 1109 | /// so a cut can never leave half a secret behind. |
| 1110 | pub(crate) fn failed(error: &str, after_message_count: usize) -> Self { |
| 1111 | let redacted = codewhale_secrets::redact::redact_secrets(error.trim()); |
| 1112 | let error = if redacted.chars().count() > MAX_SAVED_TURN_OUTCOME_ERROR_CHARS { |
| 1113 | let mut cut: String = redacted |
| 1114 | .chars() |
| 1115 | .take(MAX_SAVED_TURN_OUTCOME_ERROR_CHARS) |
| 1116 | .collect(); |
| 1117 | cut.push('…'); |
| 1118 | cut |
| 1119 | } else { |
| 1120 | redacted |
| 1121 | }; |
| 1122 | Self { |
| 1123 | status: crate::core::events::TurnOutcomeStatus::Failed, |
| 1124 | error, |
| 1125 | ended_at: Utc::now(), |
| 1126 | after_message_count, |
| 1127 | } |
| 1128 | } |
| 1129 | } |
| 1130 | |
| 1131 | /// Append `outcome`, keeping only the newest [`MAX_SAVED_TURN_OUTCOMES`]. |
| 1132 | pub(crate) fn push_turn_outcome(outcomes: &mut Vec<SavedTurnOutcome>, outcome: SavedTurnOutcome) { |
| 1133 | outcomes.push(outcome); |
| 1134 | if outcomes.len() > MAX_SAVED_TURN_OUTCOMES { |
| 1135 | let excess = outcomes.len() - MAX_SAVED_TURN_OUTCOMES; |
| 1136 | outcomes.drain(..excess); |
| 1137 | } |
| 1138 | } |
| 1139 | |
| 1140 | /// A saved session containing full conversation history |
| 1141 | /// Starting with v0.9.5 (#5262) the canonical history is the append-only entry journal (`journal` / `leaf_id`). |
| 1142 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 1143 | pub struct SavedSession { |
| 1144 | /// Schema version for migration compatibility |
| 1145 | #[serde(default = "default_session_schema_version")] |
| 1146 | pub schema_version: u32, |
| 1147 | /// Session metadata |
| 1148 | pub metadata: SessionMetadata, |
| 1149 | /// Conversation messages — derived from the journal's active branch (kept for compat). |
| 1150 | pub messages: Vec<Message>, |
| 1151 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 1152 | pub journal: Option<SessionJournal>, |
| 1153 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 1154 | pub leaf_id: Option<String>, |
| 1155 | /// System prompt if any |
| 1156 | pub system_prompt: Option<String>, |
| 1157 | /// Compact linked context references for user-visible `@path` and |
| 1158 | /// `/attach` mentions. Optional for backward-compatible session loads. |
| 1159 | #[serde(default, skip_serializing_if = "Vec::is_empty")] |
| 1160 | pub context_references: Vec<SessionContextReference>, |
| 1161 | /// Metadata registry of large outputs produced during this session. |
| 1162 | /// Artifact contents are stored in the session-owned artifact directory. |
| 1163 | #[serde(default, skip_serializing_if = "Vec::is_empty")] |
| 1164 | pub artifacts: Vec<ArtifactRecord>, |
| 1165 | /// Session-owned approval evidence. The append-only sidecar is canonical |
| 1166 | /// during a live turn; this projection makes saved snapshots self- |
| 1167 | /// describing without putting receipts in the model transcript. |
| 1168 | #[serde(default, skip_serializing_if = "Vec::is_empty")] |
| 1169 | pub(crate) approval_receipts: Vec<ApprovalReceipt>, |
| 1170 | /// To-do and plan state shown in the Work sidebar. |
| 1171 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 1172 | pub work_state: Option<SessionWorkState>, |
| 1173 | /// User-configured tab/window title for this session (`/title`), shown as |
| 1174 | /// `[title] …` in front of the terminal window title. Optional for |
| 1175 | /// backward-compatible session loads; absent sessions use the `title` |
| 1176 | /// config default instead. |
| 1177 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 1178 | pub window_title: Option<String>, |
| 1179 | /// Most recent accepted/completed Auto decision, when the saved model mode |
| 1180 | /// is `auto`. Optional for backward-compatible session loads. |
| 1181 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 1182 | pub(crate) last_auto_route: Option<SavedAutoRouteReceipt>, |
| 1183 | /// Turns that ended `Failed`, oldest first, bounded to |
| 1184 | /// [`MAX_SAVED_TURN_OUTCOMES`]. Not model context: the terminal-outcome |
| 1185 | /// record resume, export, and the Runtime API read back. |
| 1186 | #[serde(default, skip_serializing_if = "Vec::is_empty")] |
| 1187 | pub(crate) turn_outcomes: Vec<SavedTurnOutcome>, |
| 1188 | } |
| 1189 | impl SavedSession { |
| 1190 | /// Drop the journal-derived compatibility projection before an async |
| 1191 | /// persistence request takes ownership. Disk serialization restores it. |
| 1192 | pub(crate) fn compact_for_persistence_queue(&mut self) { |
| 1193 | if self.journal.is_some() { |
| 1194 | self.messages = Vec::new(); |
| 1195 | } |
| 1196 | } |
| 1197 | |
| 1198 | /// Bring `messages` and the journal into the shape the on-disk schema |
| 1199 | /// expects, in place. |
| 1200 | /// |
| 1201 | /// This used to be `storage_compatible_copy`, which cloned the whole |
| 1202 | /// session to do it. On the debounced persistence path the caller already |
| 1203 | /// owns the value and `compact_for_persistence_queue` has already emptied |
| 1204 | /// `messages`, so the clone was pure waste — two full deep copies of the |
| 1205 | /// history per write (#6214 T3). |
| 1206 | /// |
| 1207 | /// The no-op cases are load-bearing and must stay no-ops: with no journal, |
| 1208 | /// or with `messages` already equal to the journal's active branch, the |
| 1209 | /// session serializes exactly as it arrived — including a |
| 1210 | /// `metadata.message_count` that disagrees with `messages.len()`. Rewriting |
| 1211 | /// that count here would silently edit live data on every save. |
| 1212 | pub(crate) fn make_storage_compatible(&mut self) { |
| 1213 | let Some(journal) = self.journal.as_ref() else { |
| 1214 | return; |
| 1215 | }; |
| 1216 | if self.messages.is_empty() { |
| 1217 | self.messages = journal.to_messages(); |
| 1218 | } else { |
| 1219 | if self.messages == journal.to_messages() { |
| 1220 | return; |
| 1221 | } |
| 1222 | // Split the `journal` / `messages` borrows; the take is returned |
| 1223 | // before this function ends, so the session is never left short. |
| 1224 | let messages = std::mem::take(&mut self.messages); |
| 1225 | if let Some(journal) = self.journal.as_mut() { |
| 1226 | journal.rebranch_active_messages(&messages); |
| 1227 | self.leaf_id = journal.leaf_id.clone(); |
| 1228 | } |
| 1229 | self.messages = messages; |
| 1230 | } |
| 1231 | self.metadata.message_count = self.messages.len(); |
| 1232 | } |
| 1233 | |
| 1234 | pub fn ensure_journal(&mut self) { |
| 1235 | if self.journal.is_some() { |
| 1236 | if self.leaf_id.is_none() { |
| 1237 | self.leaf_id = self.journal.as_ref().and_then(|j| j.leaf_id.clone()); |
| 1238 | } |
| 1239 | let active = self |
| 1240 | .journal |
| 1241 | .as_ref() |
| 1242 | .map(|j| j.to_messages()) |
| 1243 | .unwrap_or_default(); |
| 1244 | if !active.is_empty() { |
| 1245 | self.messages = active; |
| 1246 | self.metadata.message_count = self.messages.len(); |
| 1247 | } |
| 1248 | return; |
| 1249 | } |
| 1250 | let journal = |
| 1251 | SessionJournal::from_messages(self.messages.clone(), self.metadata.spawn_depth); |
| 1252 | self.leaf_id = journal.leaf_id.clone(); |
| 1253 | self.journal = Some(journal); |
| 1254 | } |
| 1255 | #[expect(dead_code)] |
| 1256 | pub fn journal_append_message(&mut self, message: Message) -> String { |
| 1257 | self.ensure_journal(); |
| 1258 | let journal = self.journal.as_mut().expect("journal ensured"); |
| 1259 | let id = journal.append_message(message.clone()); |
| 1260 | self.leaf_id = journal.leaf_id.clone(); |
| 1261 | self.messages = journal.to_messages(); |
| 1262 | self.metadata.message_count = self.messages.len(); |
| 1263 | self.metadata.updated_at = Utc::now(); |
| 1264 | id |
| 1265 | } |
| 1266 | pub fn journal_branch_to(&mut self, entry_id: &str) -> Result<(), String> { |
| 1267 | self.ensure_journal(); |
| 1268 | let journal = self.journal.as_mut().expect("journal ensured"); |
| 1269 | journal.branch_to(entry_id)?; |
| 1270 | self.leaf_id = journal.leaf_id.clone(); |
| 1271 | self.messages = journal.to_messages(); |
| 1272 | self.metadata.message_count = self.messages.len(); |
| 1273 | self.metadata.updated_at = Utc::now(); |
| 1274 | Ok(()) |
| 1275 | } |
| 1276 | #[expect(dead_code)] |
| 1277 | pub fn active_entries(&self) -> Vec<SessionEntry> { |
| 1278 | self.journal |
| 1279 | .as_ref() |
| 1280 | .map(|j| j.root_to_leaf().into_iter().cloned().collect()) |
| 1281 | .unwrap_or_default() |
| 1282 | } |
| 1283 | /// `created_at` of the active branch's message entries, in order — the |
| 1284 | /// stamps a resumed session hands back to the live message log so the |
| 1285 | /// next save preserves append times instead of rewriting them to resume |
| 1286 | /// time. |
| 1287 | pub fn journal_message_stamps(&self) -> Vec<DateTime<Utc>> { |
| 1288 | self.journal |
| 1289 | .as_ref() |
| 1290 | .map(|journal| { |
| 1291 | journal |
| 1292 | .root_to_leaf() |
| 1293 | .iter() |
| 1294 | .filter(|entry| { |
| 1295 | matches!( |
| 1296 | entry.kind, |
| 1297 | crate::session_tree::SessionEntryKind::Message { .. } |
| 1298 | ) |
| 1299 | }) |
| 1300 | .map(|entry| entry.created_at) |
| 1301 | .collect() |
| 1302 | }) |
| 1303 | .unwrap_or_default() |
| 1304 | } |
| 1305 | pub fn export_container(&self, source: &str) -> SessionImportContainer { |
| 1306 | let journal = self.journal.clone().unwrap_or_else(|| { |
| 1307 | SessionJournal::from_messages(self.messages.clone(), self.metadata.spawn_depth) |
| 1308 | }); |
| 1309 | SessionImportContainer::new( |
| 1310 | source.to_string(), |
| 1311 | &journal, |
| 1312 | serde_json::to_value(&self.metadata).ok(), |
| 1313 | ) |
| 1314 | } |
| 1315 | pub fn import_foreign( |
| 1316 | container: SessionImportContainer, |
| 1317 | workspace: PathBuf, |
| 1318 | model: String, |
| 1319 | ) -> Result<Self, String> { |
| 1320 | let journal = container.into_journal()?; |
| 1321 | validate_saved_journal(&journal).map_err(|error| error.to_string())?; |
| 1322 | let leaf_id = journal.leaf_id.clone(); |
| 1323 | let messages = journal.to_messages(); |
| 1324 | let now = Utc::now(); |
| 1325 | let spawn_depth = journal.spawn_depth.saturating_add(1); |
| 1326 | // Reuse the conversation-derived title so an imported session that |
| 1327 | // opens with runtime-owned control traffic (Operate contract, restore |
| 1328 | // checkpoint) is named after the real prompt, not the envelope. |
| 1329 | let title = conversation_derived_title(&messages) |
| 1330 | .unwrap_or_else(|| crate::session_manager::DEFAULT_SESSION_TITLE.to_string()); |
| 1331 | let metadata = SessionMetadata { |
| 1332 | id: Uuid::new_v4().to_string(), |
| 1333 | title, |
| 1334 | created_at: now, |
| 1335 | updated_at: now, |
| 1336 | message_count: messages.len(), |
| 1337 | total_tokens: 0, |
| 1338 | model, |
| 1339 | model_provider: default_model_provider(), |
| 1340 | model_provider_id: None, |
| 1341 | workspace, |
| 1342 | mode: None, |
| 1343 | cost: SessionCostSnapshot::default(), |
| 1344 | parent_session_id: None, |
| 1345 | forked_from_message_count: None, |
| 1346 | runtime_store: None, |
| 1347 | cumulative_turn_secs: 0, |
| 1348 | archived: false, |
| 1349 | spawn_depth, |
| 1350 | }; |
| 1351 | let mut journal = journal; |
| 1352 | journal.spawn_depth = spawn_depth; |
| 1353 | Ok(Self { |
| 1354 | schema_version: CURRENT_SESSION_SCHEMA_VERSION, |
| 1355 | metadata, |
| 1356 | messages, |
| 1357 | journal: Some(journal), |
| 1358 | leaf_id, |
| 1359 | system_prompt: None, |
| 1360 | context_references: Vec::new(), |
| 1361 | artifacts: Vec::new(), |
| 1362 | approval_receipts: Vec::new(), |
| 1363 | work_state: None, |
| 1364 | window_title: None, |
| 1365 | last_auto_route: None, |
| 1366 | turn_outcomes: Vec::new(), |
| 1367 | }) |
| 1368 | } |
| 1369 | } |
| 1370 | |
| 1371 | /// Validate every stored branch before deriving a projection. The journal's |
| 1372 | /// traversal safely stops on malformed edges, but persistence must refuse such |
| 1373 | /// a document instead of treating a truncated path as the complete conversation. |
| 1374 | fn validate_saved_journal(journal: &SessionJournal) -> io::Result<()> { |
| 1375 | let invalid = |message| io::Error::new(io::ErrorKind::InvalidData, message); |
| 1376 | if journal.schema_version > crate::session_tree::CURRENT_JOURNAL_SCHEMA_VERSION { |
| 1377 | return Err(invalid("saved journal schema is newer than supported")); |
| 1378 | } |
| 1379 | journal |
| 1380 | .validate() |
| 1381 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?; |
| 1382 | let mut parents = std::collections::HashMap::with_capacity(journal.entries.len()); |
| 1383 | for entry in &journal.entries { |
| 1384 | if entry.id.is_empty() |
| 1385 | || parents |
| 1386 | .insert(entry.id.as_str(), entry.parent_id.as_deref()) |
| 1387 | .is_some() |
| 1388 | { |
| 1389 | return Err(invalid( |
| 1390 | "saved journal contains an empty or duplicate entry id", |
| 1391 | )); |
| 1392 | } |
| 1393 | } |
| 1394 | let mut colors = std::collections::HashMap::with_capacity(parents.len()); |
| 1395 | for id in parents.keys().copied() { |
| 1396 | let mut cursor = Some(id); |
| 1397 | let mut path = Vec::new(); |
| 1398 | while let Some(id) = cursor { |
| 1399 | match colors.get(id) { |
| 1400 | Some(1) => return Err(invalid("saved journal contains a cycle")), |
| 1401 | Some(2) => break, |
| 1402 | _ => { |
| 1403 | colors.insert(id, 1); |
| 1404 | path.push(id); |
| 1405 | cursor = parents[id]; |
| 1406 | } |
| 1407 | } |
| 1408 | } |
| 1409 | for id in path { |
| 1410 | colors.insert(id, 2); |
| 1411 | } |
| 1412 | } |
| 1413 | Ok(()) |
| 1414 | } |
| 1415 | |
| 1416 | /// Archived ids per session that the live app has not yet dropped from its |
| 1417 | /// in-memory journal. A bounded save records them only after its document |
| 1418 | /// write lands; the UI takes them on its next snapshot (no I/O) so RAM |
| 1419 | /// matches disk, and a racing save skips re-appending ids still listed here. |
| 1420 | fn archived_journal_ids() |
| 1421 | -> &'static std::sync::Mutex<std::collections::HashMap<String, std::collections::HashSet<String>>> { |
| 1422 | static IDS: std::sync::OnceLock< |
| 1423 | std::sync::Mutex<std::collections::HashMap<String, std::collections::HashSet<String>>>, |
| 1424 | > = std::sync::OnceLock::new(); |
| 1425 | IDS.get_or_init(Default::default) |
| 1426 | } |
| 1427 | |
| 1428 | /// Take the ids a bounded save archived for `session_id`, for the live |
| 1429 | /// journal to drop. Never touches the disk. |
| 1430 | pub(crate) fn take_archived_journal_ids(session_id: &str) -> std::collections::HashSet<String> { |
| 1431 | archived_journal_ids() |
| 1432 | .lock() |
| 1433 | .unwrap_or_else(std::sync::PoisonError::into_inner) |
| 1434 | .remove(session_id) |
| 1435 | .unwrap_or_default() |
| 1436 | } |
| 1437 | |
| 1438 | fn journal_archive_path(sessions_dir: &Path, session_id: &str) -> io::Result<PathBuf> { |
| 1439 | let id = session_id.trim(); |
| 1440 | let mut components = Path::new(id).components(); |
| 1441 | if id.is_empty() |
| 1442 | || !matches!(components.next(), Some(Component::Normal(_))) |
| 1443 | || components.next().is_some() |
| 1444 | { |
| 1445 | return Err(io::Error::new( |
| 1446 | io::ErrorKind::InvalidInput, |
| 1447 | "invalid session id for journal archive", |
| 1448 | )); |
| 1449 | } |
| 1450 | Ok(sessions_dir |
| 1451 | .join(JOURNAL_ARCHIVE_DIR) |
| 1452 | .join(format!("{id}.jsonl"))) |
| 1453 | } |
| 1454 | |
| 1455 | /// Every entry archived for `session_id`, oldest first, deduplicated by id |
| 1456 | /// (a save that crashed before its document landed re-archives the same |
| 1457 | /// entries). A torn final line — a crash mid-append, whose entries are still |
| 1458 | /// in the document — is skipped; damage anywhere else is an error. |
| 1459 | pub(crate) fn load_journal_archive( |
| 1460 | sessions_dir: &Path, |
| 1461 | session_id: &str, |
| 1462 | ) -> io::Result<Vec<SessionEntry>> { |
| 1463 | use std::io::Read as _; |
| 1464 | let path = journal_archive_path(sessions_dir, session_id)?; |
| 1465 | let mut file = match open_private_read_file(&path) { |
| 1466 | Ok(file) => file, |
| 1467 | Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(Vec::new()), |
| 1468 | Err(error) => return Err(error), |
| 1469 | }; |
| 1470 | let mut raw = String::new(); |
| 1471 | file.read_to_string(&mut raw)?; |
| 1472 | let lines: Vec<&str> = raw.lines().filter(|line| !line.trim().is_empty()).collect(); |
| 1473 | let mut seen = std::collections::HashSet::new(); |
| 1474 | let mut entries = Vec::with_capacity(lines.len()); |
| 1475 | for (index, line) in lines.iter().enumerate() { |
| 1476 | match serde_json::from_str::<SessionEntry>(line) { |
| 1477 | Ok(entry) => { |
| 1478 | if seen.insert(entry.id.clone()) { |
| 1479 | entries.push(entry); |
| 1480 | } |
| 1481 | } |
| 1482 | Err(error) if index + 1 == lines.len() && !raw.ends_with('\n') => { |
| 1483 | tracing::warn!( |
| 1484 | path = %path.display(), |
| 1485 | %error, |
| 1486 | "skipping torn final journal-archive line" |
| 1487 | ); |
| 1488 | } |
| 1489 | Err(error) => { |
| 1490 | return Err(io::Error::new( |
| 1491 | io::ErrorKind::InvalidData, |
| 1492 | format!( |
| 1493 | "journal archive {} line {}: {error}", |
| 1494 | path.display(), |
| 1495 | index + 1 |
| 1496 | ), |
| 1497 | )); |
| 1498 | } |
| 1499 | } |
| 1500 | } |
| 1501 | Ok(entries) |
| 1502 | } |
| 1503 | |
| 1504 | fn serialize_saved_session(mut session: SavedSession) -> io::Result<String> { |
| 1505 | if let Some(journal) = session.journal.as_ref() { |
| 1506 | validate_saved_journal(journal)?; |
| 1507 | } |
| 1508 | session.make_storage_compatible(); |
| 1509 | serde_json::to_string_pretty(&session) |
| 1510 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error)) |
| 1511 | } |
| 1512 | |
| 1513 | /// Repair dangling tool-call/result pairs in an already-loaded session and |
| 1514 | /// rebranch its journal to the repaired messages. Returns the repair receipt; |
| 1515 | /// callers decide whether the result gets persisted (`SessionManager::resume_*` |
| 1516 | /// does, foreign `/load` files do not). |
| 1517 | pub(crate) fn repair_recovered_session( |
| 1518 | session: &mut SavedSession, |
| 1519 | ) -> crate::tool_history_repair::ToolRepairReceipt { |
| 1520 | let repair = crate::tool_history_repair::repair_tool_call_pairs(&mut session.messages); |
| 1521 | if !repair.is_empty() { |
| 1522 | if let Some(journal) = session.journal.as_mut() { |
| 1523 | journal.rebranch_active_messages(&session.messages); |
| 1524 | session.leaf_id = journal.leaf_id.clone(); |
| 1525 | } |
| 1526 | session.metadata.message_count = session.messages.len(); |
| 1527 | tracing::warn!( |
| 1528 | session_id = %session.metadata.id, |
| 1529 | repaired_call_ids = ?repair.repaired_call_ids, |
| 1530 | duplicate_result_ids = ?repair.duplicate_result_ids, |
| 1531 | orphan_result_ids = ?repair.orphan_result_ids, |
| 1532 | "repaired persisted tool call/result history" |
| 1533 | ); |
| 1534 | } |
| 1535 | repair |
| 1536 | } |
| 1537 | |
| 1538 | /// Manager for session persistence operations |
| 1539 | #[derive(Debug)] |
| 1540 | pub struct SessionManager { |
| 1541 | /// Directory where sessions are stored |
| 1542 | sessions_dir: PathBuf, |
| 1543 | /// Re-entrancy guard: archiving a record saves it, and every save runs |
| 1544 | /// retention. Without this, a backlog past the cap would nest one |
| 1545 | /// cleanup per archived transcript instead of draining in one pass. |
| 1546 | retention_in_progress: AtomicBool, |
| 1547 | } |
| 1548 | |
| 1549 | /// One interactive editor owns a session's unsent text until its last queued |
| 1550 | /// write finishes. The stable lock file is never unlinked: replacing it would |
| 1551 | /// let two processes lock different files for the same session. |
| 1552 | #[derive(Debug)] |
| 1553 | pub struct OfflineQueueLease { |
| 1554 | session_id: String, |
| 1555 | _file: fs::File, |
| 1556 | } |
| 1557 | |
| 1558 | impl OfflineQueueLease { |
| 1559 | pub fn session_id(&self) -> &str { |
| 1560 | &self.session_id |
| 1561 | } |
| 1562 | } |
| 1563 | |
| 1564 | impl Drop for OfflineQueueLease { |
| 1565 | fn drop(&mut self) { |
| 1566 | // A forked child can briefly retain the same open-file description. |
| 1567 | // Release the editor's lock now, rather than waiting for every inherited |
| 1568 | // descriptor to close, as RuntimeProcessOwnerLock does on shutdown. |
| 1569 | #[cfg(all(unix, not(target_os = "solaris")))] |
| 1570 | { |
| 1571 | use std::os::fd::AsRawFd as _; |
| 1572 | // SAFETY: the lease still owns this descriptor throughout Drop. |
| 1573 | unsafe { |
| 1574 | libc::flock(self._file.as_raw_fd(), libc::LOCK_UN); |
| 1575 | } |
| 1576 | } |
| 1577 | #[cfg(windows)] |
| 1578 | { |
| 1579 | use std::os::windows::io::AsRawHandle as _; |
| 1580 | use windows_sys::Win32::Storage::FileSystem::UnlockFile; |
| 1581 | // SAFETY: the lease owns the handle; fd-lock locks byte 0 only. |
| 1582 | unsafe { |
| 1583 | UnlockFile(self._file.as_raw_handle() as _, 0, 0, 1, 0); |
| 1584 | } |
| 1585 | } |
| 1586 | // fd-lock uses process-associated fcntl locks on Solaris. They are not |
| 1587 | // inherited by fork and closing this descriptor releases the lock. |
| 1588 | } |
| 1589 | } |
| 1590 | |
| 1591 | /// Origin of a crash-recovery checkpoint file. |
| 1592 | #[derive(Debug, Clone, PartialEq, Eq)] |
| 1593 | pub enum CheckpointSource { |
| 1594 | /// Per-session checkpoint file `checkpoints/<session_id>.json`. |
| 1595 | Session(String), |
| 1596 | /// Legacy single-slot checkpoint file `checkpoints/latest.json`. |
| 1597 | Legacy, |
| 1598 | } |
| 1599 | |
| 1600 | /// A crash-recovery checkpoint file discovered on disk (metadata only — |
| 1601 | /// callers load the session content separately). |
| 1602 | #[derive(Debug, Clone)] |
| 1603 | pub struct CheckpointRef { |
| 1604 | pub source: CheckpointSource, |
| 1605 | #[cfg_attr(not(test), expect(dead_code))] |
| 1606 | pub path: PathBuf, |
| 1607 | pub modified: std::time::SystemTime, |
| 1608 | } |
| 1609 | |
| 1610 | /// File names in `checkpoints/` that are never per-session checkpoints. |
| 1611 | const LEGACY_CHECKPOINT_FILE: &str = "latest.json"; |
| 1612 | /// Pre-per-session global offline queue, still read once for migration. |
| 1613 | const OFFLINE_QUEUE_FILE: &str = "offline_queue.json"; |
| 1614 | /// Per-session offline queue file: `checkpoints/<session_id>.offline_queue.json`. |
| 1615 | const OFFLINE_QUEUE_SUFFIX: &str = ".offline_queue.json"; |
| 1616 | |
| 1617 | pub(crate) fn is_offline_queue_file(name: &str) -> bool { |
| 1618 | name == OFFLINE_QUEUE_FILE || name.ends_with(OFFLINE_QUEUE_SUFFIX) |
| 1619 | } |
| 1620 | |
| 1621 | impl SessionManager { |
| 1622 | fn approval_receipt_store(&self) -> ApprovalReceiptStore { |
| 1623 | ApprovalReceiptStore::new(self.sessions_dir.clone()) |
| 1624 | } |
| 1625 | |
| 1626 | fn hydrate_approval_receipts(&self, session: &mut SavedSession) -> io::Result<()> { |
| 1627 | if let Some(durable) = self |
| 1628 | .approval_receipt_store() |
| 1629 | .load_if_present(&session.metadata.id)? |
| 1630 | { |
| 1631 | // Only a missing log permits legacy embedded evidence to stand. |
| 1632 | // An empty or torn-first log must not resurrect an old approval. |
| 1633 | session.approval_receipts = durable; |
| 1634 | } |
| 1635 | ApprovalReplay::from_receipts(&session.approval_receipts) |
| 1636 | .map_err(|err| io::Error::new(io::ErrorKind::InvalidData, err))?; |
| 1637 | Ok(()) |
| 1638 | } |
| 1639 | |
| 1640 | /// Reconstruct completed approvals and interrupted unmatched asks for one |
| 1641 | /// session without consulting the model transcript. |
| 1642 | #[cfg_attr(not(test), expect(dead_code))] |
| 1643 | pub(crate) fn replay_approvals(&self, session_id: &str) -> io::Result<ApprovalReplay> { |
| 1644 | self.approval_receipt_store().replay(session_id) |
| 1645 | } |
| 1646 | |
| 1647 | fn validated_session_id<'a>(&self, id: &'a str) -> std::io::Result<&'a str> { |
| 1648 | let trimmed = id.trim(); |
| 1649 | if trimmed.is_empty() { |
| 1650 | return Err(std::io::Error::new( |
| 1651 | std::io::ErrorKind::InvalidInput, |
| 1652 | "Session id cannot be empty", |
| 1653 | )); |
| 1654 | } |
| 1655 | if !trimmed |
| 1656 | .chars() |
| 1657 | .all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_') |
| 1658 | { |
| 1659 | return Err(std::io::Error::new( |
| 1660 | std::io::ErrorKind::InvalidInput, |
| 1661 | format!("Invalid session id '{id}'"), |
| 1662 | )); |
| 1663 | } |
| 1664 | if trimmed == SESSION_BOOT_OWNERS_STEM { |
| 1665 | return Err(std::io::Error::new( |
| 1666 | std::io::ErrorKind::InvalidInput, |
| 1667 | format!("Session id '{trimmed}' collides with a reserved sessions file"), |
| 1668 | )); |
| 1669 | } |
| 1670 | Ok(trimmed) |
| 1671 | } |
| 1672 | |
| 1673 | /// Metadata of saved session `id`, read without loading its transcript. |
| 1674 | pub fn load_session_metadata_by_id(&self, id: &str) -> std::io::Result<SessionMetadata> { |
| 1675 | Self::load_session_metadata(&self.validated_session_path(id)?) |
| 1676 | } |
| 1677 | |
| 1678 | fn validated_session_path(&self, id: &str) -> std::io::Result<PathBuf> { |
| 1679 | let trimmed = self.validated_session_id(id)?; |
| 1680 | Ok(self.sessions_dir.join(format!("{trimmed}.json"))) |
| 1681 | } |
| 1682 | |
| 1683 | fn checkpoints_dir(&self) -> PathBuf { |
| 1684 | self.sessions_dir.join("checkpoints") |
| 1685 | } |
| 1686 | |
| 1687 | fn session_goals_dir(&self) -> PathBuf { |
| 1688 | self.sessions_dir.join(SESSION_GOALS_DIR) |
| 1689 | } |
| 1690 | |
| 1691 | fn checked_existing_session_goals_dir(&self) -> std::io::Result<Option<PathBuf>> { |
| 1692 | let dir = self.session_goals_dir(); |
| 1693 | let metadata = match fs::symlink_metadata(&dir) { |
| 1694 | Ok(metadata) => metadata, |
| 1695 | Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None), |
| 1696 | Err(error) => return Err(error), |
| 1697 | }; |
| 1698 | if metadata.file_type().is_symlink() || !metadata.is_dir() { |
| 1699 | return Err(io::Error::new( |
| 1700 | io::ErrorKind::InvalidData, |
| 1701 | format!( |
| 1702 | "Session goal store {} must be a real directory", |
| 1703 | dir.display() |
| 1704 | ), |
| 1705 | )); |
| 1706 | } |
| 1707 | Ok(Some(dir)) |
| 1708 | } |
| 1709 | |
| 1710 | fn ensure_session_goals_dir(&self) -> std::io::Result<PathBuf> { |
| 1711 | if let Some(dir) = self.checked_existing_session_goals_dir()? { |
| 1712 | return Ok(dir); |
| 1713 | } |
| 1714 | let dir = self.session_goals_dir(); |
| 1715 | match fs::create_dir(&dir) { |
| 1716 | Ok(()) => {} |
| 1717 | Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {} |
| 1718 | Err(error) => return Err(error), |
| 1719 | } |
| 1720 | self.checked_existing_session_goals_dir()?.ok_or_else(|| { |
| 1721 | io::Error::new( |
| 1722 | io::ErrorKind::NotFound, |
| 1723 | format!("Session goal store {} was not created", dir.display()), |
| 1724 | ) |
| 1725 | }) |
| 1726 | } |
| 1727 | |
| 1728 | fn validated_session_goal_path(&self, session_id: &str) -> std::io::Result<PathBuf> { |
| 1729 | let id = self.validated_session_id(session_id)?; |
| 1730 | Ok(self.session_goals_dir().join(format!("{id}.json"))) |
| 1731 | } |
| 1732 | |
| 1733 | fn checked_existing_session_goal_file(path: &Path) -> std::io::Result<bool> { |
| 1734 | let metadata = match fs::symlink_metadata(path) { |
| 1735 | Ok(metadata) => metadata, |
| 1736 | Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(false), |
| 1737 | Err(error) => return Err(error), |
| 1738 | }; |
| 1739 | if metadata.file_type().is_symlink() || !metadata.is_file() { |
| 1740 | return Err(io::Error::new( |
| 1741 | io::ErrorKind::InvalidData, |
| 1742 | format!("Session goal {} must be a regular file", path.display()), |
| 1743 | )); |
| 1744 | } |
| 1745 | Ok(true) |
| 1746 | } |
| 1747 | |
| 1748 | fn validated_checkpoint_path(&self, session_id: &str) -> std::io::Result<PathBuf> { |
| 1749 | let trimmed = self.validated_session_id(session_id)?; |
| 1750 | // Reserved file names inside `checkpoints/` must never collide with a |
| 1751 | // per-session checkpoint file. |
| 1752 | if format!("{trimmed}.json") == LEGACY_CHECKPOINT_FILE |
| 1753 | || format!("{trimmed}.json") == OFFLINE_QUEUE_FILE |
| 1754 | { |
| 1755 | return Err(std::io::Error::new( |
| 1756 | std::io::ErrorKind::InvalidInput, |
| 1757 | format!("Session id '{trimmed}' collides with a reserved checkpoint file"), |
| 1758 | )); |
| 1759 | } |
| 1760 | Ok(self.checkpoints_dir().join(format!("{trimmed}.json"))) |
| 1761 | } |
| 1762 | |
| 1763 | /// Create a new `SessionManager` with the specified sessions directory |
| 1764 | pub fn new(sessions_dir: PathBuf) -> std::io::Result<Self> { |
| 1765 | let sessions_dir = normalize_managed_dir(sessions_dir)?; |
| 1766 | // Ensure the sessions directory exists |
| 1767 | fs::create_dir_all(&sessions_dir)?; |
| 1768 | Ok(Self { |
| 1769 | sessions_dir, |
| 1770 | retention_in_progress: AtomicBool::new(false), |
| 1771 | }) |
| 1772 | } |
| 1773 | |
| 1774 | /// Create a `SessionManager` using the default location. |
| 1775 | pub fn default_location() -> std::io::Result<Self> { |
| 1776 | Self::new(default_sessions_dir()?) |
| 1777 | } |
| 1778 | |
| 1779 | /// Return the resolved sessions directory path. |
| 1780 | pub fn sessions_dir(&self) -> &Path { |
| 1781 | &self.sessions_dir |
| 1782 | } |
| 1783 | |
| 1784 | /// The live-lease file for `session_id`; see [`set_live_session`]. |
| 1785 | fn live_lease_path(&self, session_id: &str, create_dir: bool) -> io::Result<PathBuf> { |
| 1786 | let (late_path, _) = if create_dir { |
| 1787 | self.ensure_late_usage_paths(session_id)? |
| 1788 | } else { |
| 1789 | self.late_usage_paths(session_id)? |
| 1790 | }; |
| 1791 | Ok(late_path.with_extension("live")) |
| 1792 | } |
| 1793 | |
| 1794 | /// Claim `session_id` for this process's interactive surface: the |
| 1795 | /// in-process registry ([`set_live_session`]) plus the cross-process |
| 1796 | /// lease in this store, so writers in other processes see it too. |
| 1797 | /// |
| 1798 | /// Known limitation: this claim is best-effort. Attach paths reserve the |
| 1799 | /// lease up front ([`Self::reserve_session_for_attach`]), but a session |
| 1800 | /// that starts fresh is claimed only at its first snapshot, and a claim |
| 1801 | /// that loses the lock to another process's probe or external write, or |
| 1802 | /// cannot open the lease file, runs unleased until a later snapshot |
| 1803 | /// retries it (an open failure is logged only at debug). In that window |
| 1804 | /// another process's external writer can take the lease and write, and |
| 1805 | /// this session's next autosave reverts that write. |
| 1806 | pub fn claim_live_session(&self, session_id: &str) { |
| 1807 | set_live_session(Some(session_id)); |
| 1808 | let id = session_id.trim(); |
| 1809 | let Ok(mut lease) = LIVE_SESSION_LEASE.lock() else { |
| 1810 | return; |
| 1811 | }; |
| 1812 | // A lease already held is kept. One that could not be taken is tried |
| 1813 | // again: another process's liveness probe briefly holds the same lock, |
| 1814 | // and a claim that lost that race once must not run unleased for the |
| 1815 | // rest of the session. |
| 1816 | let retry = match lease.as_ref() { |
| 1817 | Some((held, Some(_))) if held == id => return, |
| 1818 | Some((held, None)) => held == id, |
| 1819 | _ => false, |
| 1820 | }; |
| 1821 | let file = self.live_lease_path(id, true).and_then(|path| { |
| 1822 | let file = open_private_lock_file(&path)?; |
| 1823 | Ok(crate::runtime_threads::try_lock_file_exclusive(&file)?.then_some(file)) |
| 1824 | }); |
| 1825 | match &file { |
| 1826 | Ok(Some(_)) => {} |
| 1827 | Ok(None) if !retry => tracing::warn!( |
| 1828 | session_id = id, |
| 1829 | "another Codewhale process already holds this session open" |
| 1830 | ), |
| 1831 | Ok(None) => {} |
| 1832 | Err(error) => { |
| 1833 | tracing::debug!(session_id = id, %error, "session live lease unavailable"); |
| 1834 | } |
| 1835 | } |
| 1836 | *lease = Some((id.to_string(), file.ok().flatten())); |
| 1837 | } |
| 1838 | |
| 1839 | /// Reserve `session_id` before attaching to it (`resume`, `--continue`, |
| 1840 | /// `exec --resume`, the session picker, `/load`). Unlike |
| 1841 | /// [`Self::claim_live_session`], which records whatever it gets, this |
| 1842 | /// refuses with `ResourceBusy` when another process holds the session |
| 1843 | /// open: attaching anyway gives the document two autosaving writers, and |
| 1844 | /// the last save silently drops the other's turns. Taking the lock is the |
| 1845 | /// check, so nothing can claim the session between a check and the attach. |
| 1846 | /// |
| 1847 | /// The reservation does not touch the session this process owns now: that |
| 1848 | /// claim and its lease stay until [`SessionLease::commit`], which the |
| 1849 | /// caller runs only once the new session is loaded and applied. Dropping |
| 1850 | /// the reservation (a load or apply that failed) releases it and leaves |
| 1851 | /// the current session leased. |
| 1852 | /// |
| 1853 | /// It fails closed: a lease that cannot be opened or locked is an error, |
| 1854 | /// not an unguarded attach. |
| 1855 | pub fn reserve_session_for_attach(&self, session_id: &str) -> io::Result<SessionLease> { |
| 1856 | let id = self.validated_session_id(session_id)?; |
| 1857 | if LIVE_SESSION_LEASE |
| 1858 | .lock() |
| 1859 | .is_ok_and(|lease| matches!(lease.as_ref(), Some((held, Some(_))) if held == id)) |
| 1860 | { |
| 1861 | return Ok(SessionLease { |
| 1862 | id: id.to_string(), |
| 1863 | file: None, |
| 1864 | }); |
| 1865 | } |
| 1866 | let lease_unavailable = |error: io::Error| { |
| 1867 | io::Error::new( |
| 1868 | error.kind(), |
| 1869 | format!("could not take session {id}'s live lease: {error}"), |
| 1870 | ) |
| 1871 | }; |
| 1872 | let file = self |
| 1873 | .live_lease_path(id, true) |
| 1874 | .and_then(|path| open_private_lock_file(&path)) |
| 1875 | .map_err(lease_unavailable)?; |
| 1876 | // A liveness probe from another process holds this lock for a moment; |
| 1877 | // a few short retries tell that apart from a running owner. |
| 1878 | for attempt in 0..LIVE_LEASE_ATTACH_ATTEMPTS { |
| 1879 | if crate::runtime_threads::try_lock_file_exclusive(&file).map_err(lease_unavailable)? { |
| 1880 | return Ok(SessionLease { |
| 1881 | id: id.to_string(), |
| 1882 | file: Some(file), |
| 1883 | }); |
| 1884 | } |
| 1885 | if attempt + 1 < LIVE_LEASE_ATTACH_ATTEMPTS { |
| 1886 | std::thread::sleep(std::time::Duration::from_millis(10)); |
| 1887 | } |
| 1888 | } |
| 1889 | Err(session_open_elsewhere(id)) |
| 1890 | } |
| 1891 | |
| 1892 | /// Hold the existing live lease throughout an external mutation. A |
| 1893 | /// liveness probe releases its lock before returning and cannot protect |
| 1894 | /// the subsequent read/write from another surface attaching meanwhile. |
| 1895 | /// |
| 1896 | /// Refuses with `ResourceBusy` when an interactive session in this or |
| 1897 | /// another process holds the session, and with `InvalidInput` for a |
| 1898 | /// malformed id. It may sleep briefly between lock attempts (see |
| 1899 | /// [`Self::reserve_session_for_attach`]), so async callers run it under |
| 1900 | /// `spawn_blocking`. Keep the returned lease alive until the write lands. |
| 1901 | pub(crate) fn reserve_session_for_external_write(&self, id: &str) -> io::Result<SessionLease> { |
| 1902 | let live = live_sessions() |
| 1903 | .read() |
| 1904 | .map_err(|_| io::Error::other("session ownership registry unavailable"))?; |
| 1905 | if live.contains(id.trim()) { |
| 1906 | return Err(live_session_conflict(id)); |
| 1907 | } |
| 1908 | drop(live); |
| 1909 | let lease = self.reserve_session_for_attach(id)?; |
| 1910 | // A same-process attach may have committed after the registry read. |
| 1911 | // Its already-held lease is borrowed, never ours to mutate through. |
| 1912 | if lease.file.is_none() { |
| 1913 | return Err(live_session_conflict(id)); |
| 1914 | } |
| 1915 | Ok(lease) |
| 1916 | } |
| 1917 | |
| 1918 | /// Is `session_id` open in an interactive session in this process *or any |
| 1919 | /// other*? A read-only hint for listings and recovery candidates (the |
| 1920 | /// session picker, interrupted-work discovery, retention's candidate |
| 1921 | /// scan). It never authorizes a write: the probe releases its lock before |
| 1922 | /// returning, so every mutation — rename, archive, delete, the Runtime |
| 1923 | /// API's export and save, `scrub-secrets` — holds |
| 1924 | /// [`Self::reserve_session_for_external_write`] across its load and save |
| 1925 | /// instead (#6144). |
| 1926 | /// |
| 1927 | /// Fails closed: a malformed id, or a lease that cannot be opened, reads |
| 1928 | /// as live, so a caller skips it rather than acting on it. Callers that |
| 1929 | /// must tell a malformed id apart (a 400, not a 409) validate first or |
| 1930 | /// reserve the lease, which reports `InvalidInput`. |
| 1931 | /// |
| 1932 | /// Known limitation: startup's stale-checkpoint pruning |
| 1933 | /// (`load_recent_checkpoints` in `lib.rs`) still clears a checkpoint |
| 1934 | /// older than a day after this released probe. A session that attached |
| 1935 | /// in that gap and is mid-turn refreshes its checkpoint, so only a |
| 1936 | /// day-old checkpoint of a session attached in the same instant is |
| 1937 | /// exposed. |
| 1938 | #[must_use] |
| 1939 | pub fn is_session_live_anywhere(&self, session_id: &str) -> bool { |
| 1940 | if is_live_session(session_id) { |
| 1941 | return true; |
| 1942 | } |
| 1943 | let Ok(path) = self.live_lease_path(session_id, false) else { |
| 1944 | return true; |
| 1945 | }; |
| 1946 | let file = match open_private_read_file(&path) { |
| 1947 | Ok(file) => file, |
| 1948 | Err(error) if error.kind() == io::ErrorKind::NotFound => return false, |
| 1949 | Err(_) => return true, |
| 1950 | }; |
| 1951 | // Contention means a live holder; acquiring proves none, and the |
| 1952 | // probe's lock is released when `file` drops here. |
| 1953 | !matches!( |
| 1954 | crate::runtime_threads::try_lock_file_exclusive(&file), |
| 1955 | Ok(true) |
| 1956 | ) |
| 1957 | } |
| 1958 | |
| 1959 | /// Hold `session_id`'s live lease the way another running Codewhale |
| 1960 | /// process does: an OS lock on its own open file description. |
| 1961 | #[cfg(test)] |
| 1962 | pub(crate) fn hold_live_lease_elsewhere(&self, session_id: &str) -> fs::File { |
| 1963 | let path = self.live_lease_path(session_id, true).expect("lease path"); |
| 1964 | let lease = open_private_lock_file(&path).expect("open lease"); |
| 1965 | assert!(crate::runtime_threads::try_lock_file_exclusive(&lease).expect("lock lease")); |
| 1966 | lease |
| 1967 | } |
| 1968 | |
| 1969 | /// Whether a saved document exists for `session_id`. |
| 1970 | #[must_use] |
| 1971 | pub fn session_document_exists(&self, session_id: &str) -> bool { |
| 1972 | self.validated_session_path(session_id) |
| 1973 | .is_ok_and(|path| path.is_file()) |
| 1974 | } |
| 1975 | |
| 1976 | fn late_usage_paths(&self, session_id: &str) -> io::Result<(PathBuf, PathBuf)> { |
| 1977 | let session_id = self.validated_session_id(session_id)?; |
| 1978 | let dir = self.sessions_dir.join(LATE_USAGE_DIR); |
| 1979 | match fs::symlink_metadata(&dir) { |
| 1980 | Ok(metadata) => { |
| 1981 | #[cfg(windows)] |
| 1982 | let linked = { |
| 1983 | use std::os::windows::fs::MetadataExt as _; |
| 1984 | metadata.file_attributes() |
| 1985 | & windows_sys::Win32::Storage::FileSystem::FILE_ATTRIBUTE_REPARSE_POINT |
| 1986 | != 0 |
| 1987 | }; |
| 1988 | #[cfg(not(windows))] |
| 1989 | let linked = metadata.file_type().is_symlink(); |
| 1990 | if linked || !metadata.is_dir() { |
| 1991 | return Err(io::Error::new( |
| 1992 | io::ErrorKind::InvalidData, |
| 1993 | "late usage store must be a real directory", |
| 1994 | )); |
| 1995 | } |
| 1996 | } |
| 1997 | Err(error) if error.kind() == io::ErrorKind::NotFound => {} |
| 1998 | Err(error) => return Err(error), |
| 1999 | } |
| 2000 | Ok(( |
| 2001 | dir.join(format!("{session_id}.json")), |
| 2002 | dir.join(format!("{session_id}.lock")), |
| 2003 | )) |
| 2004 | } |
| 2005 | |
| 2006 | /// Only mutations create accounting storage. Snapshot/list reads must work |
| 2007 | /// for a healthy transcript even when no sidecar has ever been written. |
| 2008 | fn ensure_late_usage_paths(&self, session_id: &str) -> io::Result<(PathBuf, PathBuf)> { |
| 2009 | self.late_usage_paths(session_id)?; |
| 2010 | let dir = self.sessions_dir.join(LATE_USAGE_DIR); |
| 2011 | match fs::create_dir(&dir) { |
| 2012 | Ok(()) => {} |
| 2013 | Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {} |
| 2014 | Err(error) => return Err(error), |
| 2015 | } |
| 2016 | let paths = self.late_usage_paths(session_id)?; |
| 2017 | #[cfg(unix)] |
| 2018 | { |
| 2019 | use std::os::unix::fs::{OpenOptionsExt as _, PermissionsExt as _}; |
| 2020 | OpenOptions::new() |
| 2021 | .read(true) |
| 2022 | .custom_flags(libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC) |
| 2023 | .open(&dir)? |
| 2024 | .set_permissions(fs::Permissions::from_mode(0o700))?; |
| 2025 | } |
| 2026 | Ok(paths) |
| 2027 | } |
| 2028 | |
| 2029 | /// A deletion marker and its stable lock survive deletion, without any |
| 2030 | /// route or usage data. A captured callback must never recreate the ledger. |
| 2031 | fn late_usage_is_deleted(path: &Path) -> io::Result<bool> { |
| 2032 | use std::io::Read as _; |
| 2033 | let tombstone = match open_private_read_file(&path.with_extension("deleted")) { |
| 2034 | Ok(file) => file, |
| 2035 | Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(false), |
| 2036 | Err(error) => return Err(error), |
| 2037 | }; |
| 2038 | let mut marker = Vec::with_capacity(LATE_USAGE_DELETED.len()); |
| 2039 | tombstone |
| 2040 | .take(u64::try_from(LATE_USAGE_DELETED.len()).unwrap_or(u64::MAX) + 1) |
| 2041 | .read_to_end(&mut marker)?; |
| 2042 | if marker != LATE_USAGE_DELETED { |
| 2043 | return Err(io::Error::new( |
| 2044 | io::ErrorKind::InvalidData, |
| 2045 | "invalid late usage deletion marker", |
| 2046 | )); |
| 2047 | } |
| 2048 | Ok(true) |
| 2049 | } |
| 2050 | |
| 2051 | fn write_late_usage_ledger(path: &Path, ledger: &LateUsageLedger) -> io::Result<()> { |
| 2052 | let bytes = serde_json::to_vec(ledger) |
| 2053 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?; |
| 2054 | if u64::try_from(bytes.len()).unwrap_or(u64::MAX) > MAX_LATE_USAGE_LEDGER_BYTES { |
| 2055 | return Err(io::Error::new( |
| 2056 | io::ErrorKind::InvalidData, |
| 2057 | "late usage ledger exceeds its size bound", |
| 2058 | )); |
| 2059 | } |
| 2060 | write_atomic(path, &bytes) |
| 2061 | } |
| 2062 | |
| 2063 | fn load_late_usage_unlocked(path: &Path) -> io::Result<LateUsageLedger> { |
| 2064 | let file = match open_private_read_file(path) { |
| 2065 | Ok(file) => file, |
| 2066 | Err(error) if error.kind() == io::ErrorKind::NotFound => { |
| 2067 | return Ok(LateUsageLedger::default()); |
| 2068 | } |
| 2069 | Err(error) => return Err(error), |
| 2070 | }; |
| 2071 | let metadata = file.metadata()?; |
| 2072 | if metadata.len() > MAX_LATE_USAGE_LEDGER_BYTES { |
| 2073 | return Err(io::Error::new( |
| 2074 | io::ErrorKind::InvalidData, |
| 2075 | format!( |
| 2076 | "late usage ledger {} exceeds its size bound", |
| 2077 | path.display() |
| 2078 | ), |
| 2079 | )); |
| 2080 | } |
| 2081 | use std::io::Read as _; |
| 2082 | let mut raw = Vec::with_capacity( |
| 2083 | usize::try_from(metadata.len().min(MAX_LATE_USAGE_LEDGER_BYTES)).unwrap_or(0), |
| 2084 | ); |
| 2085 | file.take(MAX_LATE_USAGE_LEDGER_BYTES.saturating_add(1)) |
| 2086 | .read_to_end(&mut raw)?; |
| 2087 | if u64::try_from(raw.len()).unwrap_or(u64::MAX) > MAX_LATE_USAGE_LEDGER_BYTES { |
| 2088 | return Err(io::Error::new( |
| 2089 | io::ErrorKind::InvalidData, |
| 2090 | format!( |
| 2091 | "late usage ledger {} exceeds its size bound", |
| 2092 | path.display() |
| 2093 | ), |
| 2094 | )); |
| 2095 | } |
| 2096 | let ledger: LateUsageLedger = serde_json::from_slice(&raw) |
| 2097 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?; |
| 2098 | if ledger.schema_version != CURRENT_LATE_USAGE_SCHEMA_VERSION |
| 2099 | || ledger |
| 2100 | .records |
| 2101 | .iter() |
| 2102 | .filter(|record| { |
| 2103 | record |
| 2104 | .usage |
| 2105 | .as_ref() |
| 2106 | .is_none_or(|usage| *usage == codewhale_models::Usage::default()) |
| 2107 | }) |
| 2108 | .count() |
| 2109 | > MAX_LATE_USAGE_UNRESOLVED_RECORDS_PER_SESSION |
| 2110 | || ledger.records.iter().any(|record| { |
| 2111 | !is_sha256_fingerprint(&record.source_fingerprint) |
| 2112 | || !is_sha256_fingerprint(&record.turn_fingerprint) |
| 2113 | || record.decision.as_ref().is_some_and(|r| !r.is_bounded()) |
| 2114 | }) |
| 2115 | { |
| 2116 | return Err(io::Error::new( |
| 2117 | io::ErrorKind::InvalidData, |
| 2118 | "late usage ledger has an unsupported or unbounded shape", |
| 2119 | )); |
| 2120 | } |
| 2121 | Ok(ledger) |
| 2122 | } |
| 2123 | |
| 2124 | fn with_session_write_admission<T>( |
| 2125 | &self, |
| 2126 | session_id: &str, |
| 2127 | write: impl FnOnce() -> io::Result<T>, |
| 2128 | ) -> io::Result<Option<T>> { |
| 2129 | let (path, lock_path) = self.ensure_late_usage_paths(session_id)?; |
| 2130 | let lock_file = open_private_lock_file(&lock_path)?; |
| 2131 | let mut lock = fd_lock::RwLock::new(lock_file); |
| 2132 | let _guard = lock.write()?; |
| 2133 | if Self::late_usage_is_deleted(&path)? { |
| 2134 | return Ok(None); |
| 2135 | } |
| 2136 | write().map(Some) |
| 2137 | } |
| 2138 | |
| 2139 | /// Run an out-of-band rewrite of a session's files (`scrub-secrets`) |
| 2140 | /// under the same per-session lock every save takes, so it cannot |
| 2141 | /// interleave with a save, and while holding the session's live lease, so |
| 2142 | /// no interactive session holds the conversation in memory to put the old |
| 2143 | /// content back on its next autosave. `None` when the session was |
| 2144 | /// deleted; an invalid id is an `InvalidInput` error, and a session open |
| 2145 | /// in an interactive surface is `ResourceBusy`, left untouched. |
| 2146 | /// |
| 2147 | /// The lease is taken inside the save lock with non-blocking attempts, so |
| 2148 | /// the reverse order elsewhere (lease, then a save) cannot deadlock: one |
| 2149 | /// side reports busy instead. |
| 2150 | pub(crate) fn with_session_file_lock<T>( |
| 2151 | &self, |
| 2152 | session_id: &str, |
| 2153 | rewrite: impl FnOnce() -> io::Result<T>, |
| 2154 | ) -> io::Result<Option<T>> { |
| 2155 | let session_id = self.validated_session_id(session_id)?; |
| 2156 | self.with_session_write_admission(session_id, || { |
| 2157 | let _lease = self.reserve_session_for_external_write(session_id)?; |
| 2158 | rewrite() |
| 2159 | }) |
| 2160 | } |
| 2161 | |
| 2162 | /// Serialize active accounting admission with deletion of its origin. |
| 2163 | /// A retired origin is handled without running the callback. Callers must |
| 2164 | /// release this boundary before attempting a late-usage append, which |
| 2165 | /// independently checks retirement under the same stable lock. |
| 2166 | pub(crate) fn with_live_session_origin( |
| 2167 | &self, |
| 2168 | session_id: &str, |
| 2169 | accept: impl FnOnce() -> bool, |
| 2170 | ) -> io::Result<Option<bool>> { |
| 2171 | self.with_session_write_admission(session_id, || Ok(accept())) |
| 2172 | } |
| 2173 | |
| 2174 | fn retired_session_write_error() -> io::Error { |
| 2175 | io::Error::new(io::ErrorKind::NotFound, "session was deleted") |
| 2176 | } |
| 2177 | |
| 2178 | fn persist_late_usage_record( |
| 2179 | &self, |
| 2180 | session_id: &str, |
| 2181 | incoming: LateUsageRecord, |
| 2182 | ) -> io::Result<bool> { |
| 2183 | let (path, lock_path) = self.ensure_late_usage_paths(session_id)?; |
| 2184 | let lock_file = open_private_lock_file(&lock_path)?; |
| 2185 | let mut lock = fd_lock::RwLock::new(lock_file); |
| 2186 | let _guard = lock.write()?; |
| 2187 | if Self::late_usage_is_deleted(&path)? { |
| 2188 | // Handled, rather than a failed sink that should queue a retry. |
| 2189 | return Ok(true); |
| 2190 | } |
| 2191 | let mut ledger = Self::load_late_usage_unlocked(&path)?; |
| 2192 | let original = ledger.clone(); |
| 2193 | let source_fingerprint = incoming.source_fingerprint.clone(); |
| 2194 | if let Some(record) = ledger |
| 2195 | .records |
| 2196 | .iter_mut() |
| 2197 | .find(|record| record.source_fingerprint == source_fingerprint) |
| 2198 | { |
| 2199 | if record.turn_fingerprint != incoming.turn_fingerprint |
| 2200 | || record.route != incoming.route |
| 2201 | { |
| 2202 | return Err(io::Error::new( |
| 2203 | io::ErrorKind::InvalidData, |
| 2204 | "late usage source changed its captured origin or route", |
| 2205 | )); |
| 2206 | } |
| 2207 | let mut changed = false; |
| 2208 | if record |
| 2209 | .usage |
| 2210 | .as_ref() |
| 2211 | .is_none_or(|usage| *usage == codewhale_models::Usage::default()) |
| 2212 | && let Some(usage) = incoming.usage.as_ref() |
| 2213 | { |
| 2214 | record.usage = Some(usage.clone()); |
| 2215 | changed = true; |
| 2216 | } |
| 2217 | if let Some(decision) = incoming.decision.as_ref() { |
| 2218 | if !decision.is_bounded() { |
| 2219 | return Err(io::Error::new( |
| 2220 | io::ErrorKind::InvalidData, |
| 2221 | "decision receipt exceeds its bound", |
| 2222 | )); |
| 2223 | } |
| 2224 | record.decision = Some(decision.sanitized()); |
| 2225 | changed = true; |
| 2226 | } |
| 2227 | if changed { |
| 2228 | Self::write_late_usage_candidate(&path, &original, &ledger)?; |
| 2229 | } |
| 2230 | return Ok(true); |
| 2231 | } |
| 2232 | let known_usage = incoming.usage.as_ref(); |
| 2233 | if known_usage.is_none() |
| 2234 | && ledger |
| 2235 | .records |
| 2236 | .iter() |
| 2237 | .filter(|record| { |
| 2238 | record |
| 2239 | .usage |
| 2240 | .as_ref() |
| 2241 | .is_none_or(|usage| *usage == codewhale_models::Usage::default()) |
| 2242 | }) |
| 2243 | .count() |
| 2244 | == MAX_LATE_USAGE_UNRESOLVED_RECORDS_PER_SESSION |
| 2245 | { |
| 2246 | if !ledger.overflowed { |
| 2247 | ledger.overflowed = true; |
| 2248 | Self::write_late_usage_ledger(&path, &ledger)?; |
| 2249 | } |
| 2250 | return Ok(true); |
| 2251 | } |
| 2252 | ledger.records.push(incoming); |
| 2253 | Self::write_late_usage_candidate(&path, &original, &ledger)?; |
| 2254 | Ok(true) |
| 2255 | } |
| 2256 | |
| 2257 | /// Preserve the last complete ledger on a failed known-receipt append. |
| 2258 | /// Its existing overflow flag records settlement incompleteness; the error |
| 2259 | /// is returned, never treated as handled or a reason to repeat a provider. |
| 2260 | fn write_late_usage_candidate( |
| 2261 | path: &Path, |
| 2262 | original: &LateUsageLedger, |
| 2263 | candidate: &LateUsageLedger, |
| 2264 | ) -> io::Result<()> { |
| 2265 | match Self::write_late_usage_ledger(path, candidate) { |
| 2266 | Ok(()) => Ok(()), |
| 2267 | Err(error) => { |
| 2268 | if !original.overflowed { |
| 2269 | let mut retained = original.clone(); |
| 2270 | retained.overflowed = true; |
| 2271 | if let Err(gap_error) = Self::write_late_usage_ledger(path, &retained) { |
| 2272 | tracing::warn!(%gap_error, "late usage settlement gap could not be persisted"); |
| 2273 | } |
| 2274 | } |
| 2275 | Err(error) |
| 2276 | } |
| 2277 | } |
| 2278 | } |
| 2279 | |
| 2280 | pub(crate) fn persist_late_decision_receipt( |
| 2281 | &self, |
| 2282 | session_id: &str, |
| 2283 | turn_id: &str, |
| 2284 | receipt: &crate::cost_status::RuntimeDecisionReceipt, |
| 2285 | ) -> io::Result<bool> { |
| 2286 | if !receipt.is_bounded() { |
| 2287 | return Err(io::Error::new( |
| 2288 | io::ErrorKind::InvalidData, |
| 2289 | "decision receipt exceeds its bound", |
| 2290 | )); |
| 2291 | } |
| 2292 | self.persist_late_usage_record( |
| 2293 | session_id, |
| 2294 | LateUsageRecord::captured( |
| 2295 | turn_id, |
| 2296 | &receipt.source_id, |
| 2297 | &receipt.route, |
| 2298 | receipt.usage.as_ref().filter(|_| receipt.usage_complete), |
| 2299 | Some(receipt), |
| 2300 | crate::cost_status::RuntimeUsageMissingReason::SuccessWithoutUsage, |
| 2301 | ), |
| 2302 | ) |
| 2303 | } |
| 2304 | |
| 2305 | /// Read the bounded provider decision evidence from the existing origin ledger. |
| 2306 | #[cfg(test)] |
| 2307 | pub(crate) fn decision_receipts_for_session( |
| 2308 | &self, |
| 2309 | session_id: &str, |
| 2310 | ) -> io::Result<Vec<crate::cost_status::RuntimeDecisionReceipt>> { |
| 2311 | Ok(self |
| 2312 | .load_late_usage(session_id)? |
| 2313 | .records |
| 2314 | .into_iter() |
| 2315 | .filter_map(|record| record.decision) |
| 2316 | .collect()) |
| 2317 | } |
| 2318 | |
| 2319 | pub(crate) fn persist_late_runtime_usage( |
| 2320 | &self, |
| 2321 | session_id: &str, |
| 2322 | turn_id: &str, |
| 2323 | record: &crate::cost_status::RuntimeUsageRecord, |
| 2324 | ) -> io::Result<bool> { |
| 2325 | self.persist_late_usage_record( |
| 2326 | session_id, |
| 2327 | LateUsageRecord::captured( |
| 2328 | turn_id, |
| 2329 | &record.source_id, |
| 2330 | &record.usage.route, |
| 2331 | Some(&record.usage.usage), |
| 2332 | None, |
| 2333 | crate::cost_status::RuntimeUsageMissingReason::SuccessWithoutUsage, |
| 2334 | ), |
| 2335 | ) |
| 2336 | } |
| 2337 | |
| 2338 | pub(crate) fn persist_late_runtime_drop( |
| 2339 | &self, |
| 2340 | session_id: &str, |
| 2341 | turn_id: &str, |
| 2342 | record: &crate::cost_status::RuntimeUsageDropRecord, |
| 2343 | ) -> io::Result<bool> { |
| 2344 | self.persist_late_usage_record( |
| 2345 | session_id, |
| 2346 | LateUsageRecord::captured( |
| 2347 | turn_id, |
| 2348 | &record.source_id, |
| 2349 | &record.route, |
| 2350 | None, |
| 2351 | None, |
| 2352 | record.reason, |
| 2353 | ), |
| 2354 | ) |
| 2355 | } |
| 2356 | |
| 2357 | fn with_session_read_lock<T>( |
| 2358 | &self, |
| 2359 | session_id: &str, |
| 2360 | read: impl FnOnce(&Path) -> io::Result<T>, |
| 2361 | ) -> io::Result<T> { |
| 2362 | let (path, lock_path) = self.late_usage_paths(session_id)?; |
| 2363 | let lock_file = match open_private_read_file(&lock_path) { |
| 2364 | Ok(file) => file, |
| 2365 | Err(error) if error.kind() == io::ErrorKind::NotFound => { |
| 2366 | // Atomic replacement makes a copied ledger readable without |
| 2367 | // creating a lock. No writer can have published a tombstone |
| 2368 | // without first creating the stable lock. |
| 2369 | return read(&path); |
| 2370 | } |
| 2371 | Err(error) => return Err(error), |
| 2372 | }; |
| 2373 | let lock = fd_lock::RwLock::new(lock_file); |
| 2374 | let _guard = lock.read()?; |
| 2375 | read(&path) |
| 2376 | } |
| 2377 | |
| 2378 | fn load_late_usage(&self, session_id: &str) -> io::Result<LateUsageLedger> { |
| 2379 | self.with_session_read_lock(session_id, |path| { |
| 2380 | if Self::late_usage_is_deleted(path)? { |
| 2381 | return Err(io::Error::new( |
| 2382 | io::ErrorKind::NotFound, |
| 2383 | "session accounting was deleted", |
| 2384 | )); |
| 2385 | } |
| 2386 | Self::load_late_usage_unlocked(path) |
| 2387 | }) |
| 2388 | } |
| 2389 | |
| 2390 | fn apply_late_usage_to_metadata(&self, metadata: &mut SessionMetadata) { |
| 2391 | let ledger = match self.load_late_usage(&metadata.id) { |
| 2392 | Ok(ledger) => ledger, |
| 2393 | Err(_) => { |
| 2394 | // The transcript is independent of optional accounting data. |
| 2395 | // Keep a stable gap receipt even when this projection is later |
| 2396 | // saved; loading it again must not invent another missing call. |
| 2397 | let fingerprint = crate::cost_status::usage_source_fingerprint(&format!( |
| 2398 | "late-usage-unavailable:{}", |
| 2399 | crate::cost_status::usage_source_fingerprint(&metadata.id) |
| 2400 | )); |
| 2401 | if metadata.cost.usage_source_fingerprints.insert(fingerprint) { |
| 2402 | metadata.cost.unpriced_turns = metadata.cost.unpriced_turns.saturating_add(1); |
| 2403 | metadata.cost.cny_unpriced_turns = |
| 2404 | metadata.cost.cny_unpriced_turns.saturating_add(1); |
| 2405 | } |
| 2406 | metadata |
| 2407 | .cost |
| 2408 | .unpriced_reasons |
| 2409 | .insert(LATE_USAGE_UNAVAILABLE_REASON.to_string()); |
| 2410 | metadata |
| 2411 | .cost |
| 2412 | .cny_unpriced_reasons |
| 2413 | .insert(LATE_USAGE_UNAVAILABLE_REASON.to_string()); |
| 2414 | metadata.cost.coverage_recorded = true; |
| 2415 | return; |
| 2416 | } |
| 2417 | }; |
| 2418 | for mut record in ledger.records { |
| 2419 | if record.usage.as_ref() == Some(&codewhale_models::Usage::default()) { |
| 2420 | record.usage = None; |
| 2421 | } |
| 2422 | if let Some(receipt) = &record.decision |
| 2423 | && receipt.is_bounded() |
| 2424 | { |
| 2425 | metadata |
| 2426 | .cost |
| 2427 | .route_receipts |
| 2428 | .insert(receipt.diagnostic_receipt()); |
| 2429 | } |
| 2430 | let coverage = |
| 2431 | crate::cost_status::MissingUsageCoverage::for_route(&record.route, record.reason); |
| 2432 | let source_fingerprint = record.source_fingerprint.clone(); |
| 2433 | let source_id = format!("late:{}", record.source_fingerprint); |
| 2434 | let mut pending = if let Some(usage) = record.usage.as_ref() { |
| 2435 | crate::cost_status::background_cost_for_runtime_usage( |
| 2436 | &crate::cost_status::RuntimeUsageRecord { |
| 2437 | source_id, |
| 2438 | usage: crate::cost_status::EffectiveRouteUsage { |
| 2439 | route: record.route, |
| 2440 | usage: usage.clone(), |
| 2441 | }, |
| 2442 | }, |
| 2443 | ) |
| 2444 | } else { |
| 2445 | crate::cost_status::background_cost_for_runtime_drop( |
| 2446 | &crate::cost_status::RuntimeUsageDropRecord { |
| 2447 | reason: record.reason, |
| 2448 | source_id, |
| 2449 | route: record.route, |
| 2450 | }, |
| 2451 | ) |
| 2452 | }; |
| 2453 | // The sidecar already stores the canonical SHA-256 identity. Do |
| 2454 | // not hash it again while projecting the receipt into the saved |
| 2455 | // session, or a concurrent main-snapshot writer that already |
| 2456 | // contains the response would not dedupe against this overlay. |
| 2457 | pending.usage_source_fingerprints.clear(); |
| 2458 | pending |
| 2459 | .usage_source_fingerprints |
| 2460 | .insert(source_fingerprint.clone()); |
| 2461 | let was_missing = metadata |
| 2462 | .cost |
| 2463 | .missing_usage_sources |
| 2464 | .contains_key(&source_fingerprint); |
| 2465 | if metadata |
| 2466 | .cost |
| 2467 | .usage_source_fingerprints |
| 2468 | .contains(&source_fingerprint) |
| 2469 | && (!was_missing || record.usage.is_none()) |
| 2470 | { |
| 2471 | continue; |
| 2472 | } |
| 2473 | pending.missing_usage_sources.clear(); |
| 2474 | pending.resolved_missing_usage_sources.clear(); |
| 2475 | if record.usage.is_some() { |
| 2476 | pending |
| 2477 | .resolved_missing_usage_sources |
| 2478 | .insert(source_fingerprint.clone()); |
| 2479 | } else { |
| 2480 | pending |
| 2481 | .missing_usage_sources |
| 2482 | .insert(source_fingerprint.clone(), coverage); |
| 2483 | } |
| 2484 | if let Some(usage) = record.usage { |
| 2485 | metadata.total_tokens = metadata |
| 2486 | .total_tokens |
| 2487 | .saturating_add(u64::from(usage.input_tokens)) |
| 2488 | .saturating_add(u64::from(usage.output_tokens)); |
| 2489 | } |
| 2490 | metadata.cost.absorb_late_background_cost(&pending); |
| 2491 | } |
| 2492 | if ledger.overflowed { |
| 2493 | let fingerprint = crate::cost_status::usage_source_fingerprint(&format!( |
| 2494 | "late-usage-overflow:{}", |
| 2495 | crate::cost_status::usage_source_fingerprint(&metadata.id) |
| 2496 | )); |
| 2497 | if metadata.cost.usage_source_fingerprints.insert(fingerprint) { |
| 2498 | metadata.cost.unpriced_turns = metadata.cost.unpriced_turns.saturating_add(1); |
| 2499 | metadata.cost.cny_unpriced_turns = |
| 2500 | metadata.cost.cny_unpriced_turns.saturating_add(1); |
| 2501 | metadata |
| 2502 | .cost |
| 2503 | .unpriced_reasons |
| 2504 | .insert("late_usage_ledger_overflow".to_string()); |
| 2505 | metadata |
| 2506 | .cost |
| 2507 | .cny_unpriced_reasons |
| 2508 | .insert("late_usage_ledger_overflow".to_string()); |
| 2509 | metadata.cost.coverage_recorded = true; |
| 2510 | } |
| 2511 | } |
| 2512 | } |
| 2513 | |
| 2514 | /// Persist the bounded goal control state for one saved session. |
| 2515 | /// `None` is the canonical clear operation and is idempotent. |
| 2516 | pub fn save_session_goal( |
| 2517 | &self, |
| 2518 | session_id: &str, |
| 2519 | goal: Option<&SessionGoalState>, |
| 2520 | ) -> std::io::Result<()> { |
| 2521 | let path = self.validated_session_goal_path(session_id)?; |
| 2522 | let Some(goal) = goal else { |
| 2523 | if self.checked_existing_session_goals_dir()?.is_some() && path.exists() { |
| 2524 | fs::remove_file(path)?; |
| 2525 | } |
| 2526 | return Ok(()); |
| 2527 | }; |
| 2528 | goal.validate()?; |
| 2529 | self.ensure_session_goals_dir()?; |
| 2530 | Self::checked_existing_session_goal_file(&path)?; |
| 2531 | let content = serde_json::to_string_pretty(goal) |
| 2532 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?; |
| 2533 | write_atomic(&path, content.as_bytes()) |
| 2534 | } |
| 2535 | |
| 2536 | /// Load a saved session's durable goal, rejecting malformed or future |
| 2537 | /// records instead of silently starting a different objective. |
| 2538 | pub fn load_session_goal(&self, session_id: &str) -> std::io::Result<Option<SessionGoalState>> { |
| 2539 | let path = self.validated_session_goal_path(session_id)?; |
| 2540 | if self.checked_existing_session_goals_dir()?.is_none() |
| 2541 | || !Self::checked_existing_session_goal_file(&path)? |
| 2542 | { |
| 2543 | return Ok(None); |
| 2544 | } |
| 2545 | let file_len = fs::metadata(&path)?.len(); |
| 2546 | if file_len > MAX_SESSION_GOAL_FILE_BYTES { |
| 2547 | return Err(io::Error::new( |
| 2548 | io::ErrorKind::InvalidData, |
| 2549 | format!( |
| 2550 | "Session goal {} is {file_len} bytes; maximum is {MAX_SESSION_GOAL_FILE_BYTES}", |
| 2551 | path.display() |
| 2552 | ), |
| 2553 | )); |
| 2554 | } |
| 2555 | let raw = fs::read_to_string(path)?; |
| 2556 | let goal: SessionGoalState = serde_json::from_str(&raw) |
| 2557 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?; |
| 2558 | goal.validate()?; |
| 2559 | Ok(Some(goal)) |
| 2560 | } |
| 2561 | |
| 2562 | fn hydrate_recovered_runtime_binding(&self, session: &mut SavedSession) -> std::io::Result<()> { |
| 2563 | // Compare under the session write lock. A stale process may neither |
| 2564 | // resurrect a missing binding nor replace a different recovered owner. |
| 2565 | // An adoptable empty store (#6207) counts as abandonable on either |
| 2566 | // side, exactly like a missing one: there is no durable work to lose |
| 2567 | // in either direction. |
| 2568 | if let Some(incoming) = session.metadata.runtime_store.as_ref() |
| 2569 | && let Ok(persisted) = |
| 2570 | Self::load_session_metadata(&self.validated_session_path(&session.metadata.id)?) |
| 2571 | && let Some(binding) = persisted.runtime_store |
| 2572 | && incoming != &binding |
| 2573 | { |
| 2574 | let incoming_abandonable = incoming.is_missing_session_store().unwrap_or(false) |
| 2575 | || incoming.is_adoptable_empty_store().unwrap_or(false); |
| 2576 | if incoming_abandonable && binding.validate_existing_store().is_ok() { |
| 2577 | session.metadata.runtime_store = Some(binding); |
| 2578 | } else { |
| 2579 | let persisted_abandonable = binding.is_missing_session_store().unwrap_or(false) |
| 2580 | || binding.is_adoptable_empty_store().unwrap_or(false); |
| 2581 | if !persisted_abandonable { |
| 2582 | return Err(io::Error::new( |
| 2583 | io::ErrorKind::PermissionDenied, |
| 2584 | "Session Runtime ownership changed; reopen the session before saving", |
| 2585 | )); |
| 2586 | } |
| 2587 | } |
| 2588 | } |
| 2589 | Ok(()) |
| 2590 | } |
| 2591 | |
| 2592 | /// Save a session to disk using atomic write (temp file + fsync + rename). |
| 2593 | /// |
| 2594 | /// Borrowing form: clones once so the ~150 existing `&session` call sites |
| 2595 | /// keep working. The debounced persistence path already owns its value and |
| 2596 | /// calls [`Self::save_session_owned`] instead (#6214 T3). |
| 2597 | pub fn save_session(&self, session: &SavedSession) -> std::io::Result<PathBuf> { |
| 2598 | self.save_session_owned(session.clone()) |
| 2599 | } |
| 2600 | |
| 2601 | /// Save a session to disk, consuming it. |
| 2602 | pub(crate) fn save_session_owned(&self, session: SavedSession) -> std::io::Result<PathBuf> { |
| 2603 | self.save_session_inner(session, false) |
| 2604 | } |
| 2605 | |
| 2606 | /// The autosave path (#6842): like [`Self::save_session_owned`], but when |
| 2607 | /// the journal carries more than `JOURNAL_PRUNE_DEAD_THRESHOLD` |
| 2608 | /// off-branch entries, first append the excess to the session's journal |
| 2609 | /// archive (fsynced), then write the document without them. If the |
| 2610 | /// archive cannot be written nothing is pruned. Runs on the persistence |
| 2611 | /// actor (`tui/persistence_actor.rs` `flush_inner`), never the UI loop. |
| 2612 | pub(crate) fn save_session_bounded(&self, session: SavedSession) -> std::io::Result<PathBuf> { |
| 2613 | self.save_session_inner(session, true) |
| 2614 | } |
| 2615 | |
| 2616 | fn save_session_inner(&self, session: SavedSession, bound: bool) -> std::io::Result<PathBuf> { |
| 2617 | let session_id = session.metadata.id.clone(); |
| 2618 | let path = self.validated_session_path(&session_id)?; |
| 2619 | // Not a `move` closure: `session` is consumed inside, so inference |
| 2620 | // captures it by value while `path` and `session_id` stay borrowed for |
| 2621 | // the caller to use after the write. |
| 2622 | self.with_session_write_admission(&session_id, || { |
| 2623 | let already_persisted = path.exists() |
| 2624 | || self |
| 2625 | .validated_checkpoint_path(&session_id) |
| 2626 | .is_ok_and(|checkpoint| checkpoint.exists()); |
| 2627 | |
| 2628 | // Still the pre-hydration value, and still before write_atomic. |
| 2629 | self.archive_before_first_graph_write(&session, &path)?; |
| 2630 | |
| 2631 | let mut durable_session = session; |
| 2632 | let archived = if bound { |
| 2633 | self.archive_dead_journal_entries(&mut durable_session) |
| 2634 | } else { |
| 2635 | Vec::new() |
| 2636 | }; |
| 2637 | self.hydrate_recovered_runtime_binding(&mut durable_session)?; |
| 2638 | self.hydrate_approval_receipts(&mut durable_session)?; |
| 2639 | let content = serialize_saved_session(durable_session)?; |
| 2640 | |
| 2641 | // Atomic write via write_atomic (NamedTempFile + fsync + persist) |
| 2642 | write_atomic(&path, content.as_bytes())?; |
| 2643 | self.stamp_session_boot_owner_for_new_record(&session_id, already_persisted); |
| 2644 | if !archived.is_empty() { |
| 2645 | archived_journal_ids() |
| 2646 | .lock() |
| 2647 | .unwrap_or_else(std::sync::PoisonError::into_inner) |
| 2648 | .entry(session_id.clone()) |
| 2649 | .or_default() |
| 2650 | .extend(archived); |
| 2651 | } |
| 2652 | Ok(()) |
| 2653 | })? |
| 2654 | .ok_or_else(Self::retired_session_write_error)?; |
| 2655 | |
| 2656 | // Cleanup may delete sessions, so release this session's lifecycle |
| 2657 | // lock first instead of recursively acquiring it during cleanup. |
| 2658 | self.cleanup_old_sessions()?; |
| 2659 | |
| 2660 | Ok(path) |
| 2661 | } |
| 2662 | |
| 2663 | /// Save a crash-recovery checkpoint for in-flight turns. |
| 2664 | /// |
| 2665 | /// Checkpoints are keyed per session (`checkpoints/<session_id>.json`) so |
| 2666 | /// concurrent sessions never overwrite each other's crash-recovery state. |
| 2667 | pub fn save_checkpoint(&self, session: &SavedSession) -> std::io::Result<PathBuf> { |
| 2668 | self.save_checkpoint_owned(session.clone()) |
| 2669 | } |
| 2670 | |
| 2671 | /// Save a crash-recovery checkpoint, consuming the session. |
| 2672 | pub(crate) fn save_checkpoint_owned(&self, session: SavedSession) -> std::io::Result<PathBuf> { |
| 2673 | let session_id = session.metadata.id.clone(); |
| 2674 | let path = self.validated_checkpoint_path(&session_id)?; |
| 2675 | self.with_session_write_admission(&session_id, || { |
| 2676 | let session_path = self.validated_session_path(&session_id)?; |
| 2677 | self.archive_before_first_graph_write(&session, &session_path)?; |
| 2678 | fs::create_dir_all(self.checkpoints_dir())?; |
| 2679 | let already_persisted = path.exists() || session_path.exists(); |
| 2680 | let mut durable_session = session; |
| 2681 | self.hydrate_recovered_runtime_binding(&mut durable_session)?; |
| 2682 | self.hydrate_approval_receipts(&mut durable_session)?; |
| 2683 | let content = serialize_saved_session(durable_session)?; |
| 2684 | write_atomic(&path, content.as_bytes())?; |
| 2685 | self.stamp_session_boot_owner_for_new_record(&session_id, already_persisted); |
| 2686 | Ok(()) |
| 2687 | })? |
| 2688 | .ok_or_else(Self::retired_session_write_error)?; |
| 2689 | Ok(path) |
| 2690 | } |
| 2691 | |
| 2692 | fn session_boot_owners_path(&self) -> PathBuf { |
| 2693 | self.sessions_dir |
| 2694 | .join(format!("{SESSION_BOOT_OWNERS_STEM}.json")) |
| 2695 | } |
| 2696 | |
| 2697 | fn load_session_boot_owners(&self) -> BTreeMap<String, String> { |
| 2698 | fs::read_to_string(self.session_boot_owners_path()) |
| 2699 | .ok() |
| 2700 | .and_then(|content| serde_json::from_str(&content).ok()) |
| 2701 | .unwrap_or_default() |
| 2702 | } |
| 2703 | |
| 2704 | /// Does any durable record (session file or crash checkpoint) exist for |
| 2705 | /// this session id? |
| 2706 | fn session_record_exists(&self, session_id: &str) -> bool { |
| 2707 | self.validated_session_path(session_id) |
| 2708 | .is_ok_and(|path| path.exists()) |
| 2709 | || self |
| 2710 | .validated_checkpoint_path(session_id) |
| 2711 | .is_ok_and(|path| path.exists()) |
| 2712 | } |
| 2713 | |
| 2714 | /// Record which session instance owns `session_id`'s persisted record. |
| 2715 | /// |
| 2716 | /// Entries whose durable record no longer exists are pruned on the same |
| 2717 | /// write, so the sidecar cannot grow without bound. |
| 2718 | pub(crate) fn record_session_boot_owner( |
| 2719 | &self, |
| 2720 | session_id: &str, |
| 2721 | boot_id: &str, |
| 2722 | ) -> std::io::Result<()> { |
| 2723 | let id = self.validated_session_id(session_id)?.to_string(); |
| 2724 | self.with_boot_owners_lock(|| { |
| 2725 | let mut owners = self.load_session_boot_owners(); |
| 2726 | owners.retain(|owned, _| owned == &id || self.session_record_exists(owned)); |
| 2727 | owners.insert(id, boot_id.to_string()); |
| 2728 | let content = serde_json::to_string_pretty(&owners) |
| 2729 | .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?; |
| 2730 | write_atomic(&self.session_boot_owners_path(), content.as_bytes()) |
| 2731 | }) |
| 2732 | } |
| 2733 | |
| 2734 | /// Serialize the sidecar's read-modify-write across processes. Two |
| 2735 | /// processes stamping at once each read the old map and the second rename |
| 2736 | /// dropped the first one's entry (#6144 P8). |
| 2737 | fn with_boot_owners_lock<T>(&self, update: impl FnOnce() -> io::Result<T>) -> io::Result<T> { |
| 2738 | let lock_file = open_private_lock_file( |
| 2739 | &self |
| 2740 | .sessions_dir |
| 2741 | .join(format!("{SESSION_BOOT_OWNERS_STEM}.lock")), |
| 2742 | )?; |
| 2743 | let mut lock = fd_lock::RwLock::new(lock_file); |
| 2744 | let _guard = lock.write()?; |
| 2745 | update() |
| 2746 | } |
| 2747 | |
| 2748 | /// The session-instance boot id stamped on this session's persisted |
| 2749 | /// record, when one was recorded. |
| 2750 | #[must_use] |
| 2751 | pub fn session_boot_owner(&self, session_id: &str) -> Option<String> { |
| 2752 | let id = self.validated_session_id(session_id).ok()?; |
| 2753 | self.load_session_boot_owners().get(id).cloned() |
| 2754 | } |
| 2755 | |
| 2756 | /// Was this session's persisted record created by a different session |
| 2757 | /// instance (an earlier or sibling Codewhale process)? |
| 2758 | /// |
| 2759 | /// Mirrors `SubAgentManager::is_from_prior_session` (#405): a durable |
| 2760 | /// record with no stamped owner predates the marker and is classified as |
| 2761 | /// prior-instance work, while an id with no durable record at all is |
| 2762 | /// this instance's own not-yet-persisted session. |
| 2763 | #[must_use] |
| 2764 | pub fn session_from_prior_instance(&self, session_id: &str) -> bool { |
| 2765 | match self.session_boot_owner(session_id) { |
| 2766 | Some(owner) => owner != current_session_boot_id(), |
| 2767 | None => self.session_record_exists(session_id), |
| 2768 | } |
| 2769 | } |
| 2770 | |
| 2771 | /// Stamp this instance as creator when a save writes the first durable |
| 2772 | /// record for `session_id`. A record that already existed keeps its |
| 2773 | /// original owner: re-serializing another instance's work (crash |
| 2774 | /// recovery, external mutation) must not re-badge it as ours. |
| 2775 | fn stamp_session_boot_owner_for_new_record(&self, session_id: &str, already_persisted: bool) { |
| 2776 | if already_persisted || self.session_boot_owner(session_id).is_some() { |
| 2777 | return; |
| 2778 | } |
| 2779 | if let Err(error) = self.record_session_boot_owner(session_id, current_session_boot_id()) { |
| 2780 | tracing::warn!(session_id, %error, "could not stamp session boot owner"); |
| 2781 | } |
| 2782 | } |
| 2783 | |
| 2784 | fn clear_session_boot_owner(&self, session_id: &str) { |
| 2785 | let Ok(id) = self.validated_session_id(session_id) else { |
| 2786 | return; |
| 2787 | }; |
| 2788 | let _ = self.with_boot_owners_lock(|| { |
| 2789 | let mut owners = self.load_session_boot_owners(); |
| 2790 | if owners.remove(id).is_none() { |
| 2791 | return Ok(()); |
| 2792 | } |
| 2793 | if let Ok(content) = serde_json::to_string_pretty(&owners) { |
| 2794 | write_atomic(&self.session_boot_owners_path(), content.as_bytes())?; |
| 2795 | } |
| 2796 | Ok(()) |
| 2797 | }); |
| 2798 | } |
| 2799 | |
| 2800 | /// Move the journal's excess off-branch entries into the archive sidecar. |
| 2801 | /// Order is the data-loss guard: plan on the unchanged journal, append and |
| 2802 | /// fsync the planned entries, and only then remove them from the document |
| 2803 | /// about to be written. Any failure leaves the journal whole. Returns the |
| 2804 | /// removed ids. |
| 2805 | fn archive_dead_journal_entries(&self, session: &mut SavedSession) -> Vec<String> { |
| 2806 | let session_id = session.metadata.id.clone(); |
| 2807 | let Some(journal) = session.journal.as_mut() else { |
| 2808 | return Vec::new(); |
| 2809 | }; |
| 2810 | let plan = journal.prune_plan(JOURNAL_PRUNE_DEAD_THRESHOLD, JOURNAL_RETAINED_DEAD_ENTRIES); |
| 2811 | if plan.is_empty() { |
| 2812 | return Vec::new(); |
| 2813 | } |
| 2814 | let ids: std::collections::HashSet<String> = plan.into_iter().collect(); |
| 2815 | // Ids a previous save archived but the live journal still carries |
| 2816 | // are already durable; append only the rest. |
| 2817 | let already = archived_journal_ids() |
| 2818 | .lock() |
| 2819 | .unwrap_or_else(std::sync::PoisonError::into_inner) |
| 2820 | .get(&session_id) |
| 2821 | .cloned() |
| 2822 | .unwrap_or_default(); |
| 2823 | let fresh: Vec<&SessionEntry> = journal |
| 2824 | .entries |
| 2825 | .iter() |
| 2826 | .filter(|entry| ids.contains(&entry.id) && !already.contains(&entry.id)) |
| 2827 | .collect(); |
| 2828 | if let Err(error) = self.append_journal_archive(&session_id, &fresh) { |
| 2829 | tracing::warn!( |
| 2830 | %session_id, |
| 2831 | %error, |
| 2832 | "journal archive write failed; keeping the full journal" |
| 2833 | ); |
| 2834 | return Vec::new(); |
| 2835 | } |
| 2836 | match journal.remove_entries(&ids) { |
| 2837 | Ok(_) => { |
| 2838 | session.leaf_id = journal.leaf_id.clone(); |
| 2839 | ids.into_iter().collect() |
| 2840 | } |
| 2841 | Err(error) => { |
| 2842 | tracing::warn!(%session_id, %error, "journal prune refused; keeping the full journal"); |
| 2843 | Vec::new() |
| 2844 | } |
| 2845 | } |
| 2846 | } |
| 2847 | |
| 2848 | /// Append `entries` to the session's journal archive, one JSON object per |
| 2849 | /// line, and fsync before returning. The sidecar is owner-only and never |
| 2850 | /// followed through a link. |
| 2851 | fn append_journal_archive( |
| 2852 | &self, |
| 2853 | session_id: &str, |
| 2854 | entries: &[&SessionEntry], |
| 2855 | ) -> io::Result<()> { |
| 2856 | use std::io::{Seek as _, SeekFrom, Write as _}; |
| 2857 | if entries.is_empty() { |
| 2858 | return Ok(()); |
| 2859 | } |
| 2860 | let path = journal_archive_path(&self.sessions_dir, session_id)?; |
| 2861 | let dir = self.sessions_dir.join(JOURNAL_ARCHIVE_DIR); |
| 2862 | if let Ok(metadata) = fs::symlink_metadata(&dir) |
| 2863 | && crate::plugins::metadata_is_link_or_reparse(&metadata) |
| 2864 | { |
| 2865 | return Err(io::Error::new( |
| 2866 | io::ErrorKind::InvalidInput, |
| 2867 | "journal archive directory is a link", |
| 2868 | )); |
| 2869 | } |
| 2870 | let created_dir = !dir.exists(); |
| 2871 | fs::create_dir_all(&dir)?; |
| 2872 | let mut buf = Vec::new(); |
| 2873 | for entry in entries { |
| 2874 | serde_json::to_writer(&mut buf, entry) |
| 2875 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?; |
| 2876 | buf.push(b'\n'); |
| 2877 | } |
| 2878 | let created_file = !path.exists(); |
| 2879 | let mut file = open_private_lock_file(&path)?; |
| 2880 | file.seek(SeekFrom::End(0))?; |
| 2881 | file.write_all(&buf)?; |
| 2882 | file.sync_all()?; |
| 2883 | #[cfg(unix)] |
| 2884 | if created_file || created_dir { |
| 2885 | // Make the new directory entries durable, not only the bytes. |
| 2886 | fs::File::open(&dir)?.sync_all()?; |
| 2887 | if created_dir { |
| 2888 | fs::File::open(&self.sessions_dir)?.sync_all()?; |
| 2889 | } |
| 2890 | } |
| 2891 | #[cfg(not(unix))] |
| 2892 | let _ = (created_file, created_dir); |
| 2893 | Ok(()) |
| 2894 | } |
| 2895 | |
| 2896 | /// Entries archived out of `session_id`'s journal by bounded saves. |
| 2897 | pub(crate) fn load_journal_archive(&self, session_id: &str) -> io::Result<Vec<SessionEntry>> { |
| 2898 | load_journal_archive(&self.sessions_dir, self.validated_session_id(session_id)?) |
| 2899 | } |
| 2900 | |
| 2901 | /// Bring an archived entry, and every ancestor the journal no longer |
| 2902 | /// holds, back into `session`'s journal so `/branch` can select it. |
| 2903 | /// `Ok(false)` when `entry_id` is in neither the journal nor the archive. |
| 2904 | pub(crate) fn restore_archived_journal_chain( |
| 2905 | &self, |
| 2906 | session: &mut SavedSession, |
| 2907 | entry_id: &str, |
| 2908 | ) -> io::Result<bool> { |
| 2909 | session.ensure_journal(); |
| 2910 | let session_id = session.metadata.id.clone(); |
| 2911 | let journal = session.journal.as_mut().expect("journal ensured"); |
| 2912 | if journal.contains(entry_id) { |
| 2913 | return Ok(true); |
| 2914 | } |
| 2915 | let archive = self.load_journal_archive(&session_id)?; |
| 2916 | let by_id: std::collections::HashMap<&str, &SessionEntry> = archive |
| 2917 | .iter() |
| 2918 | .map(|entry| (entry.id.as_str(), entry)) |
| 2919 | .collect(); |
| 2920 | let mut chain = Vec::new(); |
| 2921 | let mut cursor = Some(entry_id); |
| 2922 | while let Some(id) = cursor { |
| 2923 | if journal.contains(id) { |
| 2924 | break; |
| 2925 | } |
| 2926 | let Some(entry) = by_id.get(id) else { |
| 2927 | if chain.is_empty() { |
| 2928 | return Ok(false); |
| 2929 | } |
| 2930 | return Err(io::Error::new( |
| 2931 | io::ErrorKind::InvalidData, |
| 2932 | format!("archived entry {entry_id} is missing ancestor {id}"), |
| 2933 | )); |
| 2934 | }; |
| 2935 | if chain.len() > by_id.len() { |
| 2936 | return Err(io::Error::new( |
| 2937 | io::ErrorKind::InvalidData, |
| 2938 | "journal archive contains a cycle", |
| 2939 | )); |
| 2940 | } |
| 2941 | chain.push((*entry).clone()); |
| 2942 | cursor = entry.parent_id.as_deref(); |
| 2943 | } |
| 2944 | let restored: Vec<String> = chain.iter().map(|entry| entry.id.clone()).collect(); |
| 2945 | journal.entries.extend(chain.into_iter().rev()); |
| 2946 | validate_saved_journal(journal)?; |
| 2947 | // The live app must not later drop what was just brought back. |
| 2948 | if let Some(pending) = archived_journal_ids() |
| 2949 | .lock() |
| 2950 | .unwrap_or_else(std::sync::PoisonError::into_inner) |
| 2951 | .get_mut(&session_id) |
| 2952 | { |
| 2953 | for id in &restored { |
| 2954 | pending.remove(id); |
| 2955 | } |
| 2956 | } |
| 2957 | Ok(true) |
| 2958 | } |
| 2959 | |
| 2960 | /// Remove the session's journal archive with the session itself. A |
| 2961 | /// linked archive directory is not followed. |
| 2962 | fn remove_journal_archive(&self, id: &str) -> io::Result<()> { |
| 2963 | let dir = self.sessions_dir.join(JOURNAL_ARCHIVE_DIR); |
| 2964 | match fs::symlink_metadata(&dir) { |
| 2965 | Ok(metadata) if crate::plugins::metadata_is_link_or_reparse(&metadata) => { |
| 2966 | tracing::warn!( |
| 2967 | path = %dir.display(), |
| 2968 | "journal archive directory is a link; its copies were not removed" |
| 2969 | ); |
| 2970 | return Ok(()); |
| 2971 | } |
| 2972 | Ok(_) => {} |
| 2973 | Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()), |
| 2974 | Err(error) => return Err(error), |
| 2975 | } |
| 2976 | match fs::remove_file(journal_archive_path(&self.sessions_dir, id)?) { |
| 2977 | Ok(()) => Ok(()), |
| 2978 | Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()), |
| 2979 | Err(error) => Err(error), |
| 2980 | } |
| 2981 | } |
| 2982 | |
| 2983 | /// Preserve the exact pre-import session once, before the first graph- |
| 2984 | /// bearing session or checkpoint write can replace it. |
| 2985 | fn archive_before_first_graph_write( |
| 2986 | &self, |
| 2987 | session: &SavedSession, |
| 2988 | source: &Path, |
| 2989 | ) -> std::io::Result<()> { |
| 2990 | let writes_graph = session |
| 2991 | .work_state |
| 2992 | .as_ref() |
| 2993 | .and_then(|state| state.graph.as_ref()) |
| 2994 | .is_some_and(|graph| !graph.is_empty()); |
| 2995 | if !writes_graph || !source.exists() { |
| 2996 | return Ok(()); |
| 2997 | } |
| 2998 | let bytes = fs::read(source)?; |
| 2999 | let already_graph_backed = serde_json::from_slice::<SavedSession>(&bytes) |
| 3000 | .ok() |
| 3001 | .and_then(|saved| saved.work_state) |
| 3002 | .and_then(|state| state.graph) |
| 3003 | .is_some_and(|graph| !graph.is_empty()); |
| 3004 | if already_graph_backed { |
| 3005 | return Ok(()); |
| 3006 | } |
| 3007 | let archive_dir = self.sessions_dir.join(WORK_GRAPH_IMPORT_ARCHIVE_DIR); |
| 3008 | fs::create_dir_all(&archive_dir)?; |
| 3009 | let archive = |
| 3010 | archive_dir.join(source.file_name().ok_or_else(|| { |
| 3011 | io::Error::new(io::ErrorKind::InvalidInput, "invalid session path") |
| 3012 | })?); |
| 3013 | if !archive.exists() { |
| 3014 | write_atomic(&archive, &bytes)?; |
| 3015 | } |
| 3016 | Ok(()) |
| 3017 | } |
| 3018 | |
| 3019 | fn read_checkpoint_file(&self, path: &Path) -> std::io::Result<Option<SavedSession>> { |
| 3020 | if !path.exists() { |
| 3021 | return Ok(None); |
| 3022 | } |
| 3023 | let content = fs::read_to_string(path)?; |
| 3024 | let mut session: SavedSession = serde_json::from_str(&content) |
| 3025 | .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?; |
| 3026 | if session.schema_version > CURRENT_SESSION_SCHEMA_VERSION { |
| 3027 | return Err(std::io::Error::new( |
| 3028 | std::io::ErrorKind::InvalidData, |
| 3029 | format!( |
| 3030 | "Checkpoint schema v{} is newer than supported v{}", |
| 3031 | session.schema_version, CURRENT_SESSION_SCHEMA_VERSION |
| 3032 | ), |
| 3033 | )); |
| 3034 | } |
| 3035 | // A crash after retirement but before checkpoint removal must not |
| 3036 | // offer the deleted origin for recovery. Optional accounting damage |
| 3037 | // still permits recovery and is projected as incomplete below. |
| 3038 | if self |
| 3039 | .with_session_read_lock(&session.metadata.id, Self::late_usage_is_deleted) |
| 3040 | .unwrap_or(false) |
| 3041 | { |
| 3042 | return Ok(None); |
| 3043 | } |
| 3044 | session.system_prompt = strip_legacy_truncation_note(session.system_prompt); |
| 3045 | self.hydrate_approval_receipts(&mut session)?; |
| 3046 | self.apply_late_usage_to_metadata(&mut session.metadata); |
| 3047 | Ok(Some(session)) |
| 3048 | } |
| 3049 | |
| 3050 | /// Load a specific session's crash-recovery checkpoint if present. |
| 3051 | pub fn load_session_checkpoint( |
| 3052 | &self, |
| 3053 | session_id: &str, |
| 3054 | ) -> std::io::Result<Option<SavedSession>> { |
| 3055 | let path = self.validated_checkpoint_path(session_id)?; |
| 3056 | self.read_checkpoint_file(&path) |
| 3057 | } |
| 3058 | |
| 3059 | /// Load the legacy single-slot checkpoint (`checkpoints/latest.json`) if |
| 3060 | /// present. Compatibility read only — this release no longer writes it. |
| 3061 | pub fn load_legacy_checkpoint(&self) -> std::io::Result<Option<SavedSession>> { |
| 3062 | let path = self.checkpoints_dir().join(LEGACY_CHECKPOINT_FILE); |
| 3063 | self.read_checkpoint_file(&path) |
| 3064 | } |
| 3065 | |
| 3066 | pub(crate) fn legacy_checkpoint_origin(&self) -> io::Result<Option<String>> { |
| 3067 | use std::io::Read as _; |
| 3068 | |
| 3069 | let path = self.checkpoints_dir().join(LEGACY_CHECKPOINT_FILE); |
| 3070 | let file = match open_private_read_file(&path) { |
| 3071 | Ok(file) => file, |
| 3072 | Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None), |
| 3073 | Err(error) => return Err(error), |
| 3074 | }; |
| 3075 | // Lifecycle cleanup only needs the leading metadata. Never follow |
| 3076 | // links or read an unbounded legacy transcript to identify its owner. |
| 3077 | let mut prefix = Vec::new(); |
| 3078 | file.take(1024 * 1024).read_to_end(&mut prefix)?; |
| 3079 | extract_top_level_metadata(&prefix) |
| 3080 | .map(|metadata| Some(metadata.id)) |
| 3081 | .ok_or_else(|| { |
| 3082 | io::Error::new( |
| 3083 | io::ErrorKind::InvalidData, |
| 3084 | "unknown legacy checkpoint origin", |
| 3085 | ) |
| 3086 | }) |
| 3087 | } |
| 3088 | |
| 3089 | /// Clear one session's crash-recovery checkpoint. Scoped: this can never |
| 3090 | /// remove another session's checkpoint file or the legacy slot. |
| 3091 | pub fn clear_session_checkpoint(&self, session_id: &str) -> std::io::Result<()> { |
| 3092 | let path = self.validated_checkpoint_path(session_id)?; |
| 3093 | if path.exists() { |
| 3094 | fs::remove_file(path)?; |
| 3095 | } |
| 3096 | Ok(()) |
| 3097 | } |
| 3098 | |
| 3099 | /// Remove the legacy single-slot checkpoint file. |
| 3100 | pub fn clear_legacy_checkpoint(&self) -> std::io::Result<()> { |
| 3101 | let path = self.checkpoints_dir().join(LEGACY_CHECKPOINT_FILE); |
| 3102 | if path.exists() { |
| 3103 | fs::remove_file(path)?; |
| 3104 | } |
| 3105 | Ok(()) |
| 3106 | } |
| 3107 | |
| 3108 | /// Enumerate all crash-recovery checkpoint files (per-session files plus |
| 3109 | /// the legacy single slot), sorted most recently modified first. Only |
| 3110 | /// file metadata is read here; callers load content per candidate. |
| 3111 | pub fn list_checkpoints(&self) -> std::io::Result<Vec<CheckpointRef>> { |
| 3112 | let dir = self.checkpoints_dir(); |
| 3113 | let mut refs = Vec::new(); |
| 3114 | let entries = match fs::read_dir(&dir) { |
| 3115 | Ok(entries) => entries, |
| 3116 | Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(refs), |
| 3117 | Err(err) => return Err(err), |
| 3118 | }; |
| 3119 | for entry in entries { |
| 3120 | let entry = entry?; |
| 3121 | let path = entry.path(); |
| 3122 | if !path.is_file() || path.extension().is_none_or(|ext| ext != "json") { |
| 3123 | continue; |
| 3124 | } |
| 3125 | let Some(name) = path.file_name().and_then(|n| n.to_str()) else { |
| 3126 | continue; |
| 3127 | }; |
| 3128 | let source = if name == LEGACY_CHECKPOINT_FILE { |
| 3129 | CheckpointSource::Legacy |
| 3130 | } else if is_offline_queue_file(name) { |
| 3131 | // Parked offline queues live in this directory but are not |
| 3132 | // crash-recovery checkpoints. |
| 3133 | continue; |
| 3134 | } else { |
| 3135 | let session_id = name.trim_end_matches(".json").to_string(); |
| 3136 | if self.validated_checkpoint_path(&session_id).is_err() { |
| 3137 | continue; |
| 3138 | } |
| 3139 | CheckpointSource::Session(session_id) |
| 3140 | }; |
| 3141 | let Ok(modified) = entry.metadata().and_then(|m| m.modified()) else { |
| 3142 | continue; |
| 3143 | }; |
| 3144 | refs.push(CheckpointRef { |
| 3145 | source, |
| 3146 | path, |
| 3147 | modified, |
| 3148 | }); |
| 3149 | } |
| 3150 | refs.sort_by_key(|r| std::cmp::Reverse(r.modified)); |
| 3151 | Ok(refs) |
| 3152 | } |
| 3153 | |
| 3154 | /// Does `session_id` still hold a crash-recovery checkpoint — the |
| 3155 | /// durable sign the session ended mid-turn (#5715)? |
| 3156 | #[must_use] |
| 3157 | pub fn session_has_checkpoint(&self, session_id: &str) -> bool { |
| 3158 | self.validated_checkpoint_path(session_id) |
| 3159 | .is_ok_and(|path| path.exists()) |
| 3160 | } |
| 3161 | |
| 3162 | /// The most recent workspace-scoped session that still holds a |
| 3163 | /// crash-recovery checkpoint — durable evidence a prior session in this |
| 3164 | /// workspace ended mid-turn (#5715). Metadata only; the transcript is |
| 3165 | /// never read. `exclude` is the live session's own id: its in-flight |
| 3166 | /// checkpoint is current work, not prior work, and engine respawns |
| 3167 | /// inside one session must not report the session's own checkpoint. |
| 3168 | /// Sessions this process instance created are likewise excluded. |
| 3169 | pub fn interrupted_workspace_session( |
| 3170 | &self, |
| 3171 | workspace: &Path, |
| 3172 | exclude: Option<&str>, |
| 3173 | ) -> Option<SessionMetadata> { |
| 3174 | // Newest-first already; a checkpoint file only survives a session |
| 3175 | // that never reached a settled save. |
| 3176 | for checkpoint in self.list_checkpoints().ok()? { |
| 3177 | let CheckpointSource::Session(id) = checkpoint.source else { |
| 3178 | continue; |
| 3179 | }; |
| 3180 | // A checkpoint another running session is still refreshing is |
| 3181 | // that session's in-flight work, not a prior crash. |
| 3182 | if Some(id.as_str()) == exclude |
| 3183 | || !self.session_from_prior_instance(&id) |
| 3184 | || self.is_session_live_anywhere(&id) |
| 3185 | { |
| 3186 | continue; |
| 3187 | } |
| 3188 | // One malformed id or unreadable record must not hide a later |
| 3189 | // valid checkpoint — skip and keep scanning. |
| 3190 | let Ok(path) = self.validated_session_path(&id) else { |
| 3191 | continue; |
| 3192 | }; |
| 3193 | if let Ok(meta) = Self::load_session_metadata(&path) |
| 3194 | && workspace_scope_matches(&meta.workspace, workspace) |
| 3195 | { |
| 3196 | return Some(meta); |
| 3197 | } |
| 3198 | } |
| 3199 | None |
| 3200 | } |
| 3201 | |
| 3202 | /// Migrate a session recovered from the legacy single-slot checkpoint to |
| 3203 | /// a per-session checkpoint file. Never overwrites an existing |
| 3204 | /// per-session file and leaves the legacy file in place (older binaries |
| 3205 | /// still read it; the legacy writer is already gone). Returns whether a |
| 3206 | /// file was written. |
| 3207 | pub fn write_session_checkpoint_if_absent( |
| 3208 | &self, |
| 3209 | session: &SavedSession, |
| 3210 | ) -> std::io::Result<bool> { |
| 3211 | let path = self.validated_checkpoint_path(&session.metadata.id)?; |
| 3212 | if path.exists() { |
| 3213 | return Ok(false); |
| 3214 | } |
| 3215 | self.save_checkpoint(session)?; |
| 3216 | Ok(true) |
| 3217 | } |
| 3218 | |
| 3219 | /// Acquire before loading or editing a queue, including on in-process |
| 3220 | /// resume. A per-write lock is insufficient: the second editor's stale |
| 3221 | /// snapshot would overwrite the first as soon as its write completed. |
| 3222 | pub fn acquire_offline_queue_lease( |
| 3223 | &self, |
| 3224 | session_id: &str, |
| 3225 | ) -> io::Result<std::sync::Arc<OfflineQueueLease>> { |
| 3226 | let session_id = self.validated_session_id(session_id)?.to_string(); |
| 3227 | let directory = self.checkpoints_dir(); |
| 3228 | fs::create_dir_all(&directory)?; |
| 3229 | let path = directory.join(format!("{session_id}.offline_queue.lock")); |
| 3230 | let file = fs::OpenOptions::new() |
| 3231 | .create(true) |
| 3232 | .truncate(false) |
| 3233 | .read(true) |
| 3234 | .write(true) |
| 3235 | .open(path)?; |
| 3236 | let mut lock = fd_lock::RwLock::new(file); |
| 3237 | let guard = lock.try_write().map_err(|error| { |
| 3238 | io::Error::new( |
| 3239 | error.kind(), |
| 3240 | format!("Cannot open session {session_id}: its queued input is already open in another window, or its previous writes are still finishing ({error})"), |
| 3241 | ) |
| 3242 | })?; |
| 3243 | // fd-lock's guard borrows its owner. Retain the underlying descriptor |
| 3244 | // instead so this lease can travel with asynchronous writes. Forgetting |
| 3245 | // this non-owning guard keeps the OS lock held; the final Arc explicitly |
| 3246 | // unlocks in Drop. The OS also releases it when the process crashes. |
| 3247 | std::mem::forget(guard); |
| 3248 | Ok(std::sync::Arc::new(OfflineQueueLease { |
| 3249 | session_id, |
| 3250 | _file: lock.into_inner(), |
| 3251 | })) |
| 3252 | } |
| 3253 | |
| 3254 | /// Park this session's offline queue (queued + draft messages). |
| 3255 | /// |
| 3256 | /// Queues are keyed per session (`checkpoints/<session_id>.offline_queue.json`) |
| 3257 | /// for exactly the reason checkpoints are: concurrent Codewhale instances |
| 3258 | /// must never overwrite — or delete — each other's unsent user text. |
| 3259 | /// |
| 3260 | /// A queue with no session id has no owner to restore it to, so parking is |
| 3261 | /// refused rather than written to a shared file where the next boot would |
| 3262 | /// destroy it. |
| 3263 | pub fn save_offline_queue_state( |
| 3264 | &self, |
| 3265 | state: &OfflineQueueState, |
| 3266 | session_id: Option<&str>, |
| 3267 | ) -> std::io::Result<PathBuf> { |
| 3268 | let session_id = session_id.ok_or_else(|| { |
| 3269 | std::io::Error::new( |
| 3270 | std::io::ErrorKind::InvalidInput, |
| 3271 | "Offline queue cannot be parked without a session id", |
| 3272 | ) |
| 3273 | })?; |
| 3274 | let path = self.validated_offline_queue_path(session_id)?; |
| 3275 | fs::create_dir_all(self.checkpoints_dir())?; |
| 3276 | let mut owned = state.clone(); |
| 3277 | // The stamp is redundant with the file name; it stays because the UI's |
| 3278 | // restore path still compares it against the live session id. |
| 3279 | owned.session_id = Some(self.validated_session_id(session_id)?.to_string()); |
| 3280 | let content = serde_json::to_string_pretty(&owned) |
| 3281 | .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?; |
| 3282 | write_atomic(&path, content.as_bytes())?; |
| 3283 | Ok(path) |
| 3284 | } |
| 3285 | |
| 3286 | /// Load one session's parked offline queue if present. |
| 3287 | pub fn load_offline_queue_state( |
| 3288 | &self, |
| 3289 | session_id: &str, |
| 3290 | ) -> std::io::Result<Option<OfflineQueueState>> { |
| 3291 | let path = self.validated_offline_queue_path(session_id)?; |
| 3292 | Ok(match Self::read_offline_queue_file(&path)? { |
| 3293 | Some(state) => Some(state), |
| 3294 | None => self.adopt_legacy_offline_queue(session_id, &path)?, |
| 3295 | }) |
| 3296 | } |
| 3297 | |
| 3298 | /// Remove one named session's parked offline queue. |
| 3299 | pub fn clear_offline_queue_state_for(&self, session_id: &str) -> std::io::Result<()> { |
| 3300 | let path = self.validated_offline_queue_path(session_id)?; |
| 3301 | match fs::remove_file(&path) { |
| 3302 | Ok(()) => {} |
| 3303 | Err(error) if error.kind() == io::ErrorKind::NotFound => {} |
| 3304 | Err(error) => return Err(error), |
| 3305 | } |
| 3306 | Ok(()) |
| 3307 | } |
| 3308 | |
| 3309 | fn validated_offline_queue_path(&self, session_id: &str) -> std::io::Result<PathBuf> { |
| 3310 | let trimmed = self.validated_session_id(session_id)?; |
| 3311 | Ok(self |
| 3312 | .checkpoints_dir() |
| 3313 | .join(format!("{trimmed}{OFFLINE_QUEUE_SUFFIX}"))) |
| 3314 | } |
| 3315 | |
| 3316 | fn read_offline_queue_file(path: &Path) -> std::io::Result<Option<OfflineQueueState>> { |
| 3317 | let content = match fs::read_to_string(path) { |
| 3318 | Ok(content) => content, |
| 3319 | Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None), |
| 3320 | Err(error) => return Err(error), |
| 3321 | }; |
| 3322 | let state: OfflineQueueState = serde_json::from_str(&content) |
| 3323 | .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?; |
| 3324 | if state.schema_version > CURRENT_QUEUE_SCHEMA_VERSION { |
| 3325 | return Err(std::io::Error::new( |
| 3326 | std::io::ErrorKind::InvalidData, |
| 3327 | format!( |
| 3328 | "Offline queue schema v{} is newer than supported v{}", |
| 3329 | state.schema_version, CURRENT_QUEUE_SCHEMA_VERSION |
| 3330 | ), |
| 3331 | )); |
| 3332 | } |
| 3333 | Ok(Some(state)) |
| 3334 | } |
| 3335 | |
| 3336 | /// Migrate the pre-per-session global queue (`checkpoints/offline_queue.json`). |
| 3337 | /// |
| 3338 | /// It holds user-authored text, so it is adopted only by the session it was |
| 3339 | /// stamped for, and it is removed only once this session's copy is durably |
| 3340 | /// written. A queue stamped for someone else — or for nobody — is left |
| 3341 | /// exactly where it is, still readable, for its owner to claim. |
| 3342 | fn adopt_legacy_offline_queue( |
| 3343 | &self, |
| 3344 | session_id: &str, |
| 3345 | path: &Path, |
| 3346 | ) -> std::io::Result<Option<OfflineQueueState>> { |
| 3347 | let legacy = self.checkpoints_dir().join(OFFLINE_QUEUE_FILE); |
| 3348 | // A corrupt or future-schema legacy file must not fail this session's |
| 3349 | // boot: leave it on disk untouched and start with an empty queue. |
| 3350 | let Ok(Some(state)) = Self::read_offline_queue_file(&legacy) else { |
| 3351 | return Ok(None); |
| 3352 | }; |
| 3353 | if state.session_id.as_deref() != Some(self.validated_session_id(session_id)?) { |
| 3354 | return Ok(None); |
| 3355 | } |
| 3356 | let content = serde_json::to_string_pretty(&state) |
| 3357 | .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?; |
| 3358 | fs::create_dir_all(self.checkpoints_dir())?; |
| 3359 | write_atomic(path, content.as_bytes())?; |
| 3360 | match fs::remove_file(&legacy) { |
| 3361 | Ok(()) => {} |
| 3362 | // A second instance of the same session can win the adoption |
| 3363 | // race: both read the legacy file, both write this session's |
| 3364 | // per-session copy, and the twin's remove already retired the |
| 3365 | // legacy one. The queue is durably adopted either way, so a |
| 3366 | // vanished legacy file is success here, not a boot error. |
| 3367 | Err(error) if error.kind() == std::io::ErrorKind::NotFound => {} |
| 3368 | Err(error) => return Err(error), |
| 3369 | } |
| 3370 | Ok(Some(state)) |
| 3371 | } |
| 3372 | |
| 3373 | /// Read a session snapshot without repairing tool call/result pairs. |
| 3374 | /// |
| 3375 | /// This is the correct API for embedding hosts that inspect or update a |
| 3376 | /// durable session while an engine may still be executing a tool call. |
| 3377 | /// A dangling `tool_use` is not proof of a crashed process in that state. |
| 3378 | pub fn load_session_snapshot(&self, id: &str) -> std::io::Result<SavedSession> { |
| 3379 | self.load_session_snapshot_with_limit(id, None) |
| 3380 | } |
| 3381 | |
| 3382 | /// The same protected session parser with an explicit transport byte ceiling. |
| 3383 | /// The file read itself is capped before allocating or parsing its document. |
| 3384 | pub(crate) fn load_session_snapshot_bounded( |
| 3385 | &self, |
| 3386 | id: &str, |
| 3387 | limit: usize, |
| 3388 | ) -> io::Result<SavedSession> { |
| 3389 | self.load_session_snapshot_with_limit(id, Some(limit)) |
| 3390 | } |
| 3391 | |
| 3392 | fn load_session_snapshot_with_limit( |
| 3393 | &self, |
| 3394 | id: &str, |
| 3395 | limit: Option<usize>, |
| 3396 | ) -> io::Result<SavedSession> { |
| 3397 | let path = self.validated_session_path(id)?; |
| 3398 | let content = match limit { |
| 3399 | None => fs::read_to_string(&path)?, |
| 3400 | Some(limit) => { |
| 3401 | use std::io::Read; |
| 3402 | let file = fs::File::open(&path)?; |
| 3403 | if file.metadata()?.len() > limit as u64 { |
| 3404 | return Err(io::Error::new( |
| 3405 | io::ErrorKind::InvalidData, |
| 3406 | "session exceeds full-history transport bound; source retained", |
| 3407 | )); |
| 3408 | } |
| 3409 | let mut bytes = Vec::new(); |
| 3410 | file.take((limit as u64).saturating_add(1)) |
| 3411 | .read_to_end(&mut bytes)?; |
| 3412 | if bytes.len() > limit { |
| 3413 | return Err(io::Error::new( |
| 3414 | io::ErrorKind::InvalidData, |
| 3415 | "session grew beyond full-history transport bound; source retained", |
| 3416 | )); |
| 3417 | } |
| 3418 | String::from_utf8(bytes) |
| 3419 | .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))? |
| 3420 | } |
| 3421 | }; |
| 3422 | let mut session: SavedSession = serde_json::from_str(&content) |
| 3423 | .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?; |
| 3424 | if session.schema_version > CURRENT_SESSION_SCHEMA_VERSION { |
| 3425 | return Err(std::io::Error::new( |
| 3426 | std::io::ErrorKind::InvalidData, |
| 3427 | format!( |
| 3428 | "Session schema v{} is newer than supported v{}", |
| 3429 | session.schema_version, CURRENT_SESSION_SCHEMA_VERSION |
| 3430 | ), |
| 3431 | )); |
| 3432 | } |
| 3433 | |
| 3434 | // The file name is the identity callers asked for. A document that |
| 3435 | // names another session would attach that session's receipts, lease |
| 3436 | // and later saves to the wrong record. |
| 3437 | if session.metadata.id != id.trim() { |
| 3438 | return Err(std::io::Error::new( |
| 3439 | std::io::ErrorKind::InvalidData, |
| 3440 | format!( |
| 3441 | "Session file {} records a different session id", |
| 3442 | path.display() |
| 3443 | ), |
| 3444 | )); |
| 3445 | } |
| 3446 | |
| 3447 | session.system_prompt = strip_legacy_truncation_note(session.system_prompt); |
| 3448 | if let Some(journal) = session.journal.as_ref() { |
| 3449 | validate_saved_journal(journal)?; |
| 3450 | } |
| 3451 | session.ensure_journal(); |
| 3452 | self.hydrate_approval_receipts(&mut session)?; |
| 3453 | self.apply_late_usage_to_metadata(&mut session.metadata); |
| 3454 | |
| 3455 | Ok(session) |
| 3456 | } |
| 3457 | |
| 3458 | /// Load and repair a session after a known process or engine restart. |
| 3459 | /// |
| 3460 | /// The returned repair remains in memory until the caller persists |
| 3461 | /// `recovery.session`. Keeping persistence explicit lets embedding hosts |
| 3462 | /// serialize recovery with their own transcript mutation lock. |
| 3463 | pub fn recover_session_for_resume(&self, id: &str) -> std::io::Result<SessionRecovery> { |
| 3464 | let mut session = self.load_session_snapshot(id)?; |
| 3465 | let repair = repair_recovered_session(&mut session); |
| 3466 | |
| 3467 | Ok(SessionRecovery { |
| 3468 | session, |
| 3469 | changed: !repair.is_empty(), |
| 3470 | repaired_call_count: repair.repaired_call_ids.len(), |
| 3471 | duplicate_result_count: repair.duplicate_result_ids.len(), |
| 3472 | orphan_result_count: repair.orphan_result_ids.len(), |
| 3473 | }) |
| 3474 | } |
| 3475 | |
| 3476 | /// Load, repair, and durably persist a session being resumed. |
| 3477 | /// |
| 3478 | /// Resume is where a crash-repaired history becomes durable: the repaired |
| 3479 | /// record replaces the interrupted one so the same repair does not re-run |
| 3480 | /// on every later load. A persist failure is logged and the repaired |
| 3481 | /// in-memory session is still returned — a failed write-back must not |
| 3482 | /// strand the resume. |
| 3483 | pub fn resume_session(&self, id: &str) -> std::io::Result<SessionRecovery> { |
| 3484 | let recovery = self.recover_session_for_resume(id)?; |
| 3485 | if recovery.changed |
| 3486 | && let Err(error) = self.save_session(&recovery.session) |
| 3487 | { |
| 3488 | tracing::warn!( |
| 3489 | session_id = %recovery.session.metadata.id, |
| 3490 | %error, |
| 3491 | "repaired session history could not be persisted; the repair will re-run on the next load" |
| 3492 | ); |
| 3493 | } |
| 3494 | Ok(recovery) |
| 3495 | } |
| 3496 | |
| 3497 | /// [`Self::resume_session`] for a surface that will own and autosave the |
| 3498 | /// session: it first reserves the session's live lease, and refuses with |
| 3499 | /// `ResourceBusy` when another process has it open |
| 3500 | /// ([`Self::reserve_session_for_attach`]). The caller commits the |
| 3501 | /// returned lease once the session is applied. |
| 3502 | /// |
| 3503 | /// A session that was interrupted mid-turn is attached in its interrupted |
| 3504 | /// state: its crash checkpoint, when newer than the saved document, is |
| 3505 | /// promoted first. Attaching the older document instead dropped the |
| 3506 | /// in-flight turn, and the next autosave then cleared the only record of |
| 3507 | /// it while the turn's file edits stayed on disk. |
| 3508 | pub fn attach_session(&self, id: &str) -> std::io::Result<(SessionRecovery, SessionLease)> { |
| 3509 | let lease = self.reserve_session_for_attach(id)?; |
| 3510 | self.promote_interrupted_checkpoint(id); |
| 3511 | let recovery = self.resume_session(id)?; |
| 3512 | Ok((recovery, lease)) |
| 3513 | } |
| 3514 | |
| 3515 | /// Persist `id`'s crash checkpoint as its saved document when the |
| 3516 | /// checkpoint is the newer of the two (or there is no document yet), then |
| 3517 | /// consume it. Callers hold the session's attach lease. A stale checkpoint |
| 3518 | /// never replaces a newer document; an unreadable document and a failed |
| 3519 | /// save leave both files as they were. |
| 3520 | fn promote_interrupted_checkpoint(&self, id: &str) { |
| 3521 | let Ok(Some(checkpoint)) = self.load_session_checkpoint(id) else { |
| 3522 | return; |
| 3523 | }; |
| 3524 | match self.load_session(id) { |
| 3525 | Ok(saved) if saved.metadata.updated_at >= checkpoint.metadata.updated_at => {} |
| 3526 | Ok(_) => { |
| 3527 | if self.save_session(&checkpoint).is_err() { |
| 3528 | return; |
| 3529 | } |
| 3530 | } |
| 3531 | Err(error) if error.kind() == std::io::ErrorKind::NotFound => { |
| 3532 | if self.save_session(&checkpoint).is_err() { |
| 3533 | return; |
| 3534 | } |
| 3535 | } |
| 3536 | // A document that exists but cannot be read is not ours to |
| 3537 | // replace here; the attach reports it and the checkpoint stays. |
| 3538 | Err(_) => return, |
| 3539 | } |
| 3540 | let _ = self.clear_session_checkpoint(id); |
| 3541 | } |
| 3542 | |
| 3543 | /// [`Self::attach_session`] with a partial-ID prefix. |
| 3544 | pub fn attach_session_by_prefix( |
| 3545 | &self, |
| 3546 | prefix: &str, |
| 3547 | ) -> std::io::Result<(SessionRecovery, SessionLease)> { |
| 3548 | self.attach_session(&self.resolve_session_id_prefix(prefix)?) |
| 3549 | } |
| 3550 | |
| 3551 | /// [`Self::attach_session`] for a session file `/load` has already read. |
| 3552 | /// Either way the surface becomes the writer of the id the file names, so |
| 3553 | /// a managed record and a foreign file take the same contract: the id's |
| 3554 | /// live lease is reserved first and refused while another window has it |
| 3555 | /// open. A managed record then resumes through the manager so its repair |
| 3556 | /// is persisted in place; a foreign file is not ours to rewrite, so its |
| 3557 | /// journal projection and repair stay in memory. |
| 3558 | pub fn attach_session_file( |
| 3559 | &self, |
| 3560 | parsed: SavedSession, |
| 3561 | path: &Path, |
| 3562 | ) -> std::io::Result<(SavedSession, SessionLease)> { |
| 3563 | if self.owns_session_path(&parsed.metadata.id, path) { |
| 3564 | return self |
| 3565 | .attach_session(&parsed.metadata.id) |
| 3566 | .map(|(recovery, lease)| (recovery.session, lease)); |
| 3567 | } |
| 3568 | let lease = self.reserve_session_for_attach(&parsed.metadata.id)?; |
| 3569 | let mut session = parsed; |
| 3570 | if let Some(journal) = session.journal.as_ref() { |
| 3571 | validate_saved_journal(journal)?; |
| 3572 | } |
| 3573 | session.ensure_journal(); |
| 3574 | repair_recovered_session(&mut session); |
| 3575 | Ok((session, lease)) |
| 3576 | } |
| 3577 | |
| 3578 | /// [`Self::resume_session`] with a partial-ID prefix. |
| 3579 | pub fn resume_session_by_prefix(&self, prefix: &str) -> std::io::Result<SessionRecovery> { |
| 3580 | self.resume_session(&self.resolve_session_id_prefix(prefix)?) |
| 3581 | } |
| 3582 | |
| 3583 | /// True when `path` is this store's durable record for `id`. File-based |
| 3584 | /// session loads use it to decide whether a repair may be written back in |
| 3585 | /// place or must stay in memory (a foreign file is not ours to rewrite). |
| 3586 | pub(crate) fn owns_session_path(&self, id: &str, path: &Path) -> bool { |
| 3587 | let Ok(managed) = self.validated_session_path(id) else { |
| 3588 | return false; |
| 3589 | }; |
| 3590 | managed == path |
| 3591 | || managed |
| 3592 | .canonicalize() |
| 3593 | .is_ok_and(|managed| path.canonicalize().is_ok_and(|path| managed == path)) |
| 3594 | } |
| 3595 | |
| 3596 | /// Load a session by ID for the standalone CodeWhale resume flow. |
| 3597 | /// |
| 3598 | /// This preserves the historical recovery behavior for existing callers. |
| 3599 | /// Embedding hosts performing ordinary runtime reads should use |
| 3600 | /// [`Self::load_session_snapshot`] instead. |
| 3601 | pub fn load_session(&self, id: &str) -> std::io::Result<SavedSession> { |
| 3602 | self.recover_session_for_resume(id) |
| 3603 | .map(|recovery| recovery.session) |
| 3604 | } |
| 3605 | |
| 3606 | /// Load a session by partial ID prefix |
| 3607 | pub fn load_session_by_prefix(&self, prefix: &str) -> std::io::Result<SavedSession> { |
| 3608 | self.load_session(&self.resolve_session_id_prefix(prefix)?) |
| 3609 | } |
| 3610 | |
| 3611 | /// Resolve a unique ID without applying resume-time repair to its record. |
| 3612 | pub(crate) fn resolve_session_id_prefix(&self, prefix: &str) -> std::io::Result<String> { |
| 3613 | let sessions = self.list_sessions()?; |
| 3614 | |
| 3615 | // One session listed more than once (a stray copy of its file under |
| 3616 | // another name) is still one session, not an ambiguous prefix. |
| 3617 | let mut matches: Vec<_> = sessions |
| 3618 | .into_iter() |
| 3619 | .filter(|s| s.id.starts_with(prefix)) |
| 3620 | .collect(); |
| 3621 | matches.sort_by(|a, b| a.id.cmp(&b.id)); |
| 3622 | matches.dedup_by(|a, b| a.id == b.id); |
| 3623 | |
| 3624 | match matches.len() { |
| 3625 | 0 => Err(std::io::Error::new( |
| 3626 | std::io::ErrorKind::NotFound, |
| 3627 | format!("No session found with prefix: {prefix}"), |
| 3628 | )), |
| 3629 | 1 => Ok(matches[0].id.clone()), |
| 3630 | _ => Err(std::io::Error::new( |
| 3631 | std::io::ErrorKind::InvalidInput, |
| 3632 | format!( |
| 3633 | "Ambiguous prefix '{}' matches {} sessions", |
| 3634 | prefix, |
| 3635 | matches.len() |
| 3636 | ), |
| 3637 | )), |
| 3638 | } |
| 3639 | } |
| 3640 | |
| 3641 | /// List all saved sessions, sorted by most recently updated |
| 3642 | pub fn list_sessions(&self) -> std::io::Result<Vec<SessionMetadata>> { |
| 3643 | Ok(self |
| 3644 | .list_session_records()? |
| 3645 | .into_iter() |
| 3646 | .map(|(_, metadata)| metadata) |
| 3647 | .collect()) |
| 3648 | } |
| 3649 | |
| 3650 | /// [`Self::list_sessions`], keeping the file each record was read from. |
| 3651 | /// A record's file is not always `<id>.json`: early builds wrote |
| 3652 | /// `session_<timestamp>.json`, and retention must address those by path. |
| 3653 | fn list_session_records(&self) -> std::io::Result<Vec<(PathBuf, SessionMetadata)>> { |
| 3654 | let mut sessions = Vec::new(); |
| 3655 | |
| 3656 | for entry in fs::read_dir(&self.sessions_dir)? { |
| 3657 | let entry = entry?; |
| 3658 | let path = entry.path(); |
| 3659 | |
| 3660 | if path.extension().is_some_and(|ext| ext == "json") |
| 3661 | && let Ok(mut session) = Self::load_session_metadata(&path) |
| 3662 | { |
| 3663 | self.apply_late_usage_to_metadata(&mut session); |
| 3664 | sessions.push((path, session)); |
| 3665 | } |
| 3666 | } |
| 3667 | |
| 3668 | // Sort by updated_at descending (most recent first) |
| 3669 | sessions.sort_by_key(|(_, s)| std::cmp::Reverse(s.updated_at)); |
| 3670 | |
| 3671 | Ok(sessions) |
| 3672 | } |
| 3673 | |
| 3674 | /// Set the durable archive flag on a saved session and return the |
| 3675 | /// resulting metadata. |
| 3676 | /// |
| 3677 | /// This is the single writer for the flag: the picker, the `/sessions` |
| 3678 | /// command, and `PATCH /v1/sessions/{id}` all route through it so the TUI |
| 3679 | /// and the web dashboard cannot drift into two archive notions. A no-op |
| 3680 | /// call (already in the requested state) still returns the metadata and |
| 3681 | /// does not rewrite the file. |
| 3682 | pub fn set_session_archived( |
| 3683 | &self, |
| 3684 | id: &str, |
| 3685 | archived: bool, |
| 3686 | mutator: SessionMutator, |
| 3687 | ) -> std::io::Result<SessionMetadata> { |
| 3688 | let _lease = (mutator == SessionMutator::External) |
| 3689 | .then(|| self.reserve_session_for_external_write(id)) |
| 3690 | .transpose()?; |
| 3691 | let mut session = self.load_session(id)?; |
| 3692 | if session.metadata.archived == archived { |
| 3693 | return Ok(session.metadata); |
| 3694 | } |
| 3695 | session.metadata.archived = archived; |
| 3696 | self.save_session(&session)?; |
| 3697 | Ok(session.metadata) |
| 3698 | } |
| 3699 | |
| 3700 | /// Re-read the durable lifecycle fields for `metadata` from disk. |
| 3701 | /// |
| 3702 | /// This is the autosave-survival guard. A TUI autosave rebuilds the whole |
| 3703 | /// session document from in-memory `App` state; any lifecycle field it |
| 3704 | /// carries from a stale cache would silently revert a rename or archive |
| 3705 | /// that landed in between — including one applied by the picker earlier in |
| 3706 | /// the same event loop, or by `/rename` while a snapshot was already |
| 3707 | /// queued. |
| 3708 | /// |
| 3709 | /// So rather than trusting any cache, the writer re-reads the persisted |
| 3710 | /// values immediately before writing. `title`, `archived`, `created_at`, |
| 3711 | /// and fork lineage are *lifecycle* state owned by the file, not |
| 3712 | /// conversation state owned by the running turn. Reading them back costs |
| 3713 | /// one bounded metadata-prefix read. |
| 3714 | /// |
| 3715 | /// Returns `true` when an existing record was found and merged. A missing |
| 3716 | /// record is not an error: the first save of a new session has nothing to |
| 3717 | /// merge from. |
| 3718 | pub fn merge_persisted_lifecycle(&self, metadata: &mut SessionMetadata) -> bool { |
| 3719 | let Ok(path) = self.validated_session_path(&metadata.id) else { |
| 3720 | return false; |
| 3721 | }; |
| 3722 | let Ok(persisted) = Self::load_session_metadata(&path) else { |
| 3723 | return false; |
| 3724 | }; |
| 3725 | metadata.title = persisted.title; |
| 3726 | metadata.archived = persisted.archived; |
| 3727 | metadata.created_at = persisted.created_at; |
| 3728 | metadata.parent_session_id = persisted.parent_session_id; |
| 3729 | metadata.forked_from_message_count = persisted.forked_from_message_count; |
| 3730 | metadata.runtime_store = persisted.runtime_store; |
| 3731 | true |
| 3732 | } |
| 3733 | |
| 3734 | /// Rename a saved session and return the resulting metadata. |
| 3735 | /// |
| 3736 | /// Titles are trimmed and bounded to [`MAX_SESSION_TITLE_CHARS`] |
| 3737 | /// characters (counted in `char`s, not bytes, so a CJK or emoji title is |
| 3738 | /// not truncated mid-scalar). Created-at and fork lineage are untouched. |
| 3739 | pub fn rename_session( |
| 3740 | &self, |
| 3741 | id: &str, |
| 3742 | title: &str, |
| 3743 | mutator: SessionMutator, |
| 3744 | ) -> std::io::Result<SessionMetadata> { |
| 3745 | let title = normalize_session_title(title)?; |
| 3746 | let _lease = (mutator == SessionMutator::External) |
| 3747 | .then(|| self.reserve_session_for_external_write(id)) |
| 3748 | .transpose()?; |
| 3749 | let mut session = self.load_session(id)?; |
| 3750 | if session.metadata.title == title { |
| 3751 | return Ok(session.metadata); |
| 3752 | } |
| 3753 | session.metadata.title = title; |
| 3754 | self.save_session(&session)?; |
| 3755 | Ok(session.metadata) |
| 3756 | } |
| 3757 | |
| 3758 | /// Load only the metadata from a session file. |
| 3759 | /// |
| 3760 | /// Optimization for #337: previously this called |
| 3761 | /// `serde_json::from_reader` which forces serde to scan every token in |
| 3762 | /// the file just to validate JSON structure — including the |
| 3763 | /// (potentially many MB of) `messages` and `tool_log` arrays we're |
| 3764 | /// going to discard. For a user with hundreds of long sessions, a |
| 3765 | /// single `list_sessions()` call could chew through tens of MB of |
| 3766 | /// JSON per startup. |
| 3767 | /// |
| 3768 | /// We now read at most 64 KB up front and string-extract the |
| 3769 | /// top-level `metadata` object, which is invariably tiny (~500 B) |
| 3770 | /// and appears before any large `messages`/`tool_log` payload. We |
| 3771 | /// fall back to a full-file read only if the prefix doesn't yield a |
| 3772 | /// parseable metadata block (e.g. an oddly-formatted legacy file). |
| 3773 | pub(crate) fn load_session_metadata(path: &Path) -> std::io::Result<SessionMetadata> { |
| 3774 | use std::io::Read; |
| 3775 | |
| 3776 | const PREFIX_BYTES: usize = 64 * 1024; |
| 3777 | let mut file = fs::File::open(path)?; |
| 3778 | let mut buf = Vec::with_capacity(PREFIX_BYTES); |
| 3779 | file.by_ref() |
| 3780 | .take(PREFIX_BYTES as u64) |
| 3781 | .read_to_end(&mut buf)?; |
| 3782 | |
| 3783 | if let Some(mut metadata) = extract_top_level_metadata(&buf) { |
| 3784 | apply_legacy_title_recovery(&mut metadata, &buf); |
| 3785 | return Ok(metadata); |
| 3786 | } |
| 3787 | |
| 3788 | // Metadata wasn't extractable from the prefix: usually a `metadata` |
| 3789 | // block longer than 64 KB (a long-lived session's usage fingerprints, |
| 3790 | // #6842), rarely unusual key ordering. Grow the read geometrically |
| 3791 | // until the block closes, so the cost tracks the metadata's size |
| 3792 | // rather than the transcript's; only past `GROWN_PREFIX_LIMIT` (or |
| 3793 | // for a document that really lacks the block) read everything. |
| 3794 | const GROWN_PREFIX_LIMIT: usize = 16 * 1024 * 1024; |
| 3795 | let mut limit = PREFIX_BYTES; |
| 3796 | while buf.len() == limit && limit < GROWN_PREFIX_LIMIT { |
| 3797 | let next = (limit * 2).min(GROWN_PREFIX_LIMIT); |
| 3798 | file.by_ref() |
| 3799 | .take((next - limit) as u64) |
| 3800 | .read_to_end(&mut buf)?; |
| 3801 | limit = next; |
| 3802 | if let Some(mut metadata) = extract_top_level_metadata(&buf) { |
| 3803 | apply_legacy_title_recovery(&mut metadata, &buf); |
| 3804 | return Ok(metadata); |
| 3805 | } |
| 3806 | } |
| 3807 | let mut rest = Vec::new(); |
| 3808 | file.read_to_end(&mut rest)?; |
| 3809 | buf.extend_from_slice(&rest); |
| 3810 | let mut metadata = extract_top_level_metadata(&buf).ok_or_else(|| { |
| 3811 | std::io::Error::new( |
| 3812 | std::io::ErrorKind::InvalidData, |
| 3813 | "session file missing parseable `metadata` block", |
| 3814 | ) |
| 3815 | })?; |
| 3816 | apply_legacy_title_recovery(&mut metadata, &buf); |
| 3817 | Ok(metadata) |
| 3818 | } |
| 3819 | |
| 3820 | /// Whether `remove_session` has anything of `id`'s to act on: its |
| 3821 | /// document, a recovery checkpoint (its own, or the legacy slot it |
| 3822 | /// originated), or a deletion marker whose cleanup a retry finishes. The |
| 3823 | /// same test `remove_session` repeats under the session's lock, made |
| 3824 | /// before any lease or lock file is created for the id. |
| 3825 | fn session_may_have_records(&self, id: &str, path: &Path) -> io::Result<bool> { |
| 3826 | match fs::symlink_metadata(path) { |
| 3827 | Ok(_) => return Ok(true), |
| 3828 | Err(error) if error.kind() == io::ErrorKind::NotFound => {} |
| 3829 | Err(error) => return Err(error), |
| 3830 | } |
| 3831 | if let Ok(checkpoint) = self.validated_checkpoint_path(id) |
| 3832 | && checkpoint.try_exists()? |
| 3833 | { |
| 3834 | return Ok(true); |
| 3835 | } |
| 3836 | if matches!(self.legacy_checkpoint_origin(), Ok(Some(origin)) if origin == id.trim()) { |
| 3837 | return Ok(true); |
| 3838 | } |
| 3839 | let (late_path, _) = self.late_usage_paths(id)?; |
| 3840 | Self::late_usage_is_deleted(&late_path) |
| 3841 | } |
| 3842 | |
| 3843 | /// Delete a session and its recovery checkpoints, retiring its origin. |
| 3844 | pub fn delete_session(&self, id: &str) -> std::io::Result<()> { |
| 3845 | self.remove_session(id, SessionRemoval::Explicit) |
| 3846 | } |
| 3847 | |
| 3848 | /// Remove the pre-import copy of `session_path` that the first |
| 3849 | /// graph-bearing write kept (see `archive_before_first_graph_write`): it is |
| 3850 | /// the same transcript, so a deleted session must not survive in it. |
| 3851 | /// |
| 3852 | /// A linked archive directory is not followed, as the late-usage store |
| 3853 | /// refuses one: deletion must not reach a same-named file elsewhere. |
| 3854 | fn remove_import_archive_copy(&self, session_path: &Path) -> io::Result<()> { |
| 3855 | let Some(name) = session_path.file_name() else { |
| 3856 | return Ok(()); |
| 3857 | }; |
| 3858 | let dir = self.sessions_dir.join(WORK_GRAPH_IMPORT_ARCHIVE_DIR); |
| 3859 | match fs::symlink_metadata(&dir) { |
| 3860 | Ok(metadata) if crate::plugins::metadata_is_link_or_reparse(&metadata) => { |
| 3861 | tracing::warn!( |
| 3862 | path = %dir.display(), |
| 3863 | "work-graph import archive is a link; its copies were not removed" |
| 3864 | ); |
| 3865 | return Ok(()); |
| 3866 | } |
| 3867 | Ok(_) => {} |
| 3868 | Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()), |
| 3869 | Err(error) => return Err(error), |
| 3870 | } |
| 3871 | match fs::remove_file(dir.join(name)) { |
| 3872 | Ok(()) => Ok(()), |
| 3873 | Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()), |
| 3874 | Err(error) => Err(error), |
| 3875 | } |
| 3876 | } |
| 3877 | |
| 3878 | fn remove_session(&self, id: &str, removal: SessionRemoval) -> std::io::Result<()> { |
| 3879 | let path = self.validated_session_path(id)?; |
| 3880 | // Reserving the lease creates `.late-usage/<id>.live`; an id with |
| 3881 | // nothing to remove must not leave one behind. The authoritative |
| 3882 | // check below repeats this under the session's lock. |
| 3883 | if !self.session_may_have_records(id, &path)? { |
| 3884 | return Err(io::Error::new( |
| 3885 | io::ErrorKind::NotFound, |
| 3886 | format!("Session '{}' not found", id.trim()), |
| 3887 | )); |
| 3888 | } |
| 3889 | let _lease = self.reserve_session_for_external_write(id)?; |
| 3890 | // Older ordinary snapshots may use a name reserved by the checkpoint |
| 3891 | // directory. Such a name must never address its shared legacy files. |
| 3892 | let checkpoint = self.validated_checkpoint_path(id).ok(); |
| 3893 | let legacy_checkpoint = self.checkpoints_dir().join(LEGACY_CHECKPOINT_FILE); |
| 3894 | let (late_path, lock_path) = self.ensure_late_usage_paths(id)?; |
| 3895 | let lock_file = open_private_lock_file(&lock_path)?; |
| 3896 | let mut lock = fd_lock::RwLock::new(lock_file); |
| 3897 | let _guard = lock.write()?; |
| 3898 | let already_deleted = Self::late_usage_is_deleted(&late_path)?; |
| 3899 | let legacy_origin = self.legacy_checkpoint_origin(); |
| 3900 | let owns_legacy_checkpoint = |
| 3901 | matches!(&legacy_origin, Ok(Some(origin)) if origin == id.trim()); |
| 3902 | let has_recovery = match checkpoint.as_ref() { |
| 3903 | Some(path) => path.try_exists()?, |
| 3904 | None => false, |
| 3905 | } || owns_legacy_checkpoint; |
| 3906 | if !already_deleted { |
| 3907 | // An unknown id must not acquire a deletion marker. A prior |
| 3908 | // tombstone, however, lets a retry finish interrupted cleanup. |
| 3909 | match fs::symlink_metadata(&path) { |
| 3910 | Ok(_) => {} |
| 3911 | Err(error) if error.kind() == io::ErrorKind::NotFound && has_recovery => {} |
| 3912 | Err(error) => return Err(error), |
| 3913 | } |
| 3914 | } |
| 3915 | if matches!(removal, SessionRemoval::Retention) && !already_deleted { |
| 3916 | // Retention never acts on an uncertain origin. An unreadable |
| 3917 | // legacy checkpoint may be this id's crash recovery, and a |
| 3918 | // damaged recovery leaves the ordinary snapshot as the only |
| 3919 | // readable copy of the conversation. Fail closed: keep every byte |
| 3920 | // and let the caller skip this record. |
| 3921 | if let Err(error) = &legacy_origin { |
| 3922 | return Err(io::Error::new( |
| 3923 | error.kind(), |
| 3924 | format!( |
| 3925 | "retention left session '{}' untouched: the legacy checkpoint's origin is unreadable ({error})", |
| 3926 | id.trim() |
| 3927 | ), |
| 3928 | )); |
| 3929 | } |
| 3930 | } |
| 3931 | if matches!(removal, SessionRemoval::Retention) && has_recovery && !already_deleted { |
| 3932 | // Retention owns the ordinary snapshot, not crash recovery. Keep |
| 3933 | // the origin and its accounting/evidence writable for resume. |
| 3934 | match fs::remove_file(&path) { |
| 3935 | Ok(()) => {} |
| 3936 | Err(error) if error.kind() == io::ErrorKind::NotFound => {} |
| 3937 | Err(error) => return Err(error), |
| 3938 | } |
| 3939 | // The pre-import copy is an older ordinary snapshot, not recovery. |
| 3940 | return self.remove_import_archive_copy(&path); |
| 3941 | } |
| 3942 | // The store this document is bound to is named by its binding, not by |
| 3943 | // its id: a host store usually sits under another conversation's |
| 3944 | // directory. Read it before the document is gone (#6144 P2). |
| 3945 | let bound_store = Self::load_session_metadata(&path) |
| 3946 | .ok() |
| 3947 | .and_then(|metadata| metadata.runtime_store); |
| 3948 | self.save_session_goal(id, None)?; |
| 3949 | // Publish the tombstone before removing data. A crash or a delayed |
| 3950 | // callback can no longer re-create this session's accounting. The |
| 3951 | // stable lock inode must never be removed or atomically replaced. |
| 3952 | if !already_deleted { |
| 3953 | write_atomic(&late_path.with_extension("deleted"), LATE_USAGE_DELETED)?; |
| 3954 | } |
| 3955 | match fs::remove_file(&path) { |
| 3956 | Ok(()) => {} |
| 3957 | Err(error) if error.kind() == io::ErrorKind::NotFound => {} |
| 3958 | Err(error) => return Err(error), |
| 3959 | } |
| 3960 | match fs::remove_file(&late_path) { |
| 3961 | Ok(()) => {} |
| 3962 | Err(error) if error.kind() == io::ErrorKind::NotFound => {} |
| 3963 | Err(error) => return Err(error), |
| 3964 | } |
| 3965 | if let Some(checkpoint) = checkpoint { |
| 3966 | match fs::remove_file(checkpoint) { |
| 3967 | Ok(()) => {} |
| 3968 | Err(error) if error.kind() == io::ErrorKind::NotFound => {} |
| 3969 | Err(error) => return Err(error), |
| 3970 | } |
| 3971 | } |
| 3972 | if owns_legacy_checkpoint { |
| 3973 | match fs::remove_file(&legacy_checkpoint) { |
| 3974 | Ok(()) => {} |
| 3975 | Err(error) if error.kind() == io::ErrorKind::NotFound => {} |
| 3976 | Err(error) => return Err(error), |
| 3977 | } |
| 3978 | } |
| 3979 | self.remove_import_archive_copy(&path)?; |
| 3980 | self.remove_journal_archive(id)?; |
| 3981 | self.clear_session_boot_owner(id); |
| 3982 | let session_dir = self.sessions_dir.join(id.trim()); |
| 3983 | if session_dir.exists() { |
| 3984 | if crate::plugins::metadata_is_link_or_reparse(&fs::symlink_metadata(&session_dir)?) { |
| 3985 | // Preserve remove_dir_all's existing no-follow behavior. |
| 3986 | fs::remove_dir_all(session_dir)?; |
| 3987 | return Ok(()); |
| 3988 | } |
| 3989 | // Other conversations and automations can share this host's Runtime |
| 3990 | // authority. Deleting a transcript must never delete that store — |
| 3991 | // including a `runtime-recovered-*` sibling, which another |
| 3992 | // document may be bound to. |
| 3993 | for entry in fs::read_dir(&session_dir)? { |
| 3994 | let entry = entry?; |
| 3995 | if is_runtime_store_dir_name(&entry.file_name()) { |
| 3996 | continue; |
| 3997 | } |
| 3998 | if entry.file_type()?.is_dir() { |
| 3999 | fs::remove_dir_all(entry.path())?; |
| 4000 | } else { |
| 4001 | fs::remove_file(entry.path())?; |
| 4002 | } |
| 4003 | } |
| 4004 | } |
| 4005 | self.retire_released_stores(id, bound_store); |
| 4006 | if session_dir.is_dir() && fs::read_dir(&session_dir)?.next().is_none() { |
| 4007 | fs::remove_dir(session_dir)?; |
| 4008 | } |
| 4009 | Ok(()) |
| 4010 | } |
| 4011 | |
| 4012 | /// Set aside the stores a deleted document released — the one its |
| 4013 | /// binding names, and any left under its own directory — when nothing |
| 4014 | /// else binds them and they hold no work. Never unlinks; a store in use, |
| 4015 | /// holding work, or bound elsewhere stays exactly where it is (#6144 P2). |
| 4016 | fn retire_released_stores( |
| 4017 | &self, |
| 4018 | id: &str, |
| 4019 | bound_store: Option<crate::runtime_threads::RuntimeStoreBinding>, |
| 4020 | ) { |
| 4021 | let mut candidates: Vec<PathBuf> = bound_store |
| 4022 | .map(|binding| binding.data_dir) |
| 4023 | .into_iter() |
| 4024 | .collect(); |
| 4025 | if let Ok(entries) = fs::read_dir(self.sessions_dir.join(id.trim())) { |
| 4026 | candidates.extend( |
| 4027 | entries |
| 4028 | .flatten() |
| 4029 | .filter(|entry| is_runtime_store_dir_name(&entry.file_name())) |
| 4030 | .map(|entry| entry.path()), |
| 4031 | ); |
| 4032 | } |
| 4033 | for store in candidates { |
| 4034 | crate::session_reconcile::retire_unbound_store(self, &store, "session deleted"); |
| 4035 | } |
| 4036 | } |
| 4037 | |
| 4038 | /// Clean up old sessions to stay within the active cap. |
| 4039 | pub fn cleanup_old_sessions(&self) -> std::io::Result<()> { |
| 4040 | self.cleanup_old_sessions_keeping(None) |
| 4041 | } |
| 4042 | |
| 4043 | /// As [`Self::cleanup_old_sessions`], but never touches `keep` — the |
| 4044 | /// session being resumed at boot. Without this, a background cleanup that |
| 4045 | /// races session restore can retire the just-resumed session when 50+ |
| 4046 | /// newer records exist (its `updated_at` is not bumped until first save). |
| 4047 | /// |
| 4048 | /// The cap counts *active* transcripts: archived records sit outside it |
| 4049 | /// until the user prunes them, and empty auto-created stubs get their own |
| 4050 | /// small cap so they can never push a real transcript out (#6136, #6137). |
| 4051 | /// A transcript past the cap is archived, never unlinked; the store's |
| 4052 | /// destructive paths stay the explicit user actions. |
| 4053 | pub fn cleanup_old_sessions_keeping(&self, keep: Option<&str>) -> std::io::Result<()> { |
| 4054 | // Archiving saves the record, and every save runs retention again. |
| 4055 | // Drain a backlog in one pass here instead of nesting one cleanup per |
| 4056 | // archived transcript. |
| 4057 | if self |
| 4058 | .retention_in_progress |
| 4059 | .swap(true, std::sync::atomic::Ordering::SeqCst) |
| 4060 | { |
| 4061 | return Ok(()); |
| 4062 | } |
| 4063 | let result = self.cleanup_old_sessions_inner(keep); |
| 4064 | self.retention_in_progress |
| 4065 | .store(false, std::sync::atomic::Ordering::SeqCst); |
| 4066 | result |
| 4067 | } |
| 4068 | |
| 4069 | fn cleanup_old_sessions_inner(&self, keep: Option<&str>) -> std::io::Result<()> { |
| 4070 | let records = self.list_session_records()?; |
| 4071 | |
| 4072 | // What retention owes each class (#6136/#6137): archived records are |
| 4073 | // already outside the cap; empty auto-created stubs are junk the |
| 4074 | // product writes on every boot and are capped apart so they can never |
| 4075 | // occupy a transcript's slot; everything else carries the |
| 4076 | // MAX_SESSIONS window. |
| 4077 | let mut active: Vec<&SessionMetadata> = Vec::new(); |
| 4078 | let mut stubs: Vec<(&Path, &SessionMetadata)> = Vec::new(); |
| 4079 | for (path, session) in &records { |
| 4080 | if session.archived { |
| 4081 | continue; |
| 4082 | } |
| 4083 | if is_empty_auto_created_session(session) { |
| 4084 | stubs.push((path, session)); |
| 4085 | } else { |
| 4086 | active.push(session); |
| 4087 | } |
| 4088 | } |
| 4089 | |
| 4090 | for session in active.iter().skip(MAX_SESSIONS) { |
| 4091 | if keep.is_some_and(|id| id == session.id) { |
| 4092 | continue; |
| 4093 | } |
| 4094 | // `External` keeps the live-session guard honest: a record |
| 4095 | // another process is driving is not ours to retire. |
| 4096 | if let Err(err) = self.set_session_archived(&session.id, true, SessionMutator::External) |
| 4097 | { |
| 4098 | tracing::warn!( |
| 4099 | target: "session", |
| 4100 | session = session.id, |
| 4101 | ?err, |
| 4102 | "retention could not archive a transcript past the cap; it stays active" |
| 4103 | ); |
| 4104 | } |
| 4105 | } |
| 4106 | |
| 4107 | for (listed_path, session) in stubs.iter().skip(MAX_EMPTY_SESSION_STUBS) { |
| 4108 | if keep.is_some_and(|id| id == session.id) { |
| 4109 | continue; |
| 4110 | } |
| 4111 | if let Err(err) = self.retire_empty_stub(listed_path, &session.id) { |
| 4112 | tracing::warn!( |
| 4113 | target: "session", |
| 4114 | session = session.id, |
| 4115 | ?err, |
| 4116 | "retention could not remove an empty session stub" |
| 4117 | ); |
| 4118 | } |
| 4119 | } |
| 4120 | |
| 4121 | // A directory without a top-level session snapshot is not proof of an |
| 4122 | // orphan: runtime threads and automations own independent durable stores, |
| 4123 | // including in other processes and before their first snapshot. Retention |
| 4124 | // only retires records it listed above; never infer authority to delete |
| 4125 | // other directories from an absent transcript or process-local claim. |
| 4126 | |
| 4127 | Ok(()) |
| 4128 | } |
| 4129 | |
| 4130 | /// Retire one empty auto-created stub that retention listed at |
| 4131 | /// `listed_path`. |
| 4132 | /// |
| 4133 | /// Early builds wrote records as `session_<timestamp>.json`, so the file |
| 4134 | /// retention listed is not always `<id>.json`. Removing such a stub by id |
| 4135 | /// could never find it: it was listed again, failed with `NotFound`, and |
| 4136 | /// warned on every launch forever. That record has no id-addressed |
| 4137 | /// accounting, checkpoint, or directory (nothing could reach it by id), so |
| 4138 | /// retiring it is removing exactly the file that was listed. A stub that |
| 4139 | /// is already gone when removal runs (another process's retention got |
| 4140 | /// there first) is already removed, not an error. |
| 4141 | fn retire_empty_stub(&self, listed_path: &Path, id: &str) -> std::io::Result<()> { |
| 4142 | let canonical = self.validated_session_path(id)?; |
| 4143 | let result = if listed_path == canonical { |
| 4144 | self.remove_session(id, SessionRemoval::Retention) |
| 4145 | } else { |
| 4146 | let _lease = self.reserve_session_for_external_write(id)?; |
| 4147 | fs::remove_file(listed_path) |
| 4148 | }; |
| 4149 | match result { |
| 4150 | Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()), |
| 4151 | other => other, |
| 4152 | } |
| 4153 | } |
| 4154 | |
| 4155 | /// Remove session files whose `updated_at` is older than `max_age` |
| 4156 | /// from the persisted-sessions directory. Returns the number of |
| 4157 | /// records pruned. Building block for #406's phase-2 auto-archive |
| 4158 | /// on boot; today the user-facing entry point is the |
| 4159 | /// `/sessions prune <days>` slash command. |
| 4160 | /// |
| 4161 | /// Crash-recovery safety: skips the per-session checkpoint files |
| 4162 | /// (`checkpoints/<session_id>.json`), the legacy single-slot |
| 4163 | /// checkpoint (`checkpoints/latest.json`), and any file under `checkpoints/` |
| 4164 | /// — those are owned by the checkpoint subsystem and live with |
| 4165 | /// stricter durability rules. Only top-level `<session_id>.json` |
| 4166 | /// files are candidates. |
| 4167 | /// |
| 4168 | /// `max_age` is checked against the metadata's `updated_at` |
| 4169 | /// timestamp embedded in the JSON, not the filesystem mtime — the |
| 4170 | /// user may have rsynced their `~/.deepseek` between machines and |
| 4171 | /// fs mtimes can lie. |
| 4172 | #[cfg_attr(not(test), expect(dead_code))] |
| 4173 | pub fn prune_sessions_older_than( |
| 4174 | &self, |
| 4175 | max_age: std::time::Duration, |
| 4176 | ) -> std::io::Result<usize> { |
| 4177 | self.prune_sessions_older_than_keeping(max_age, None) |
| 4178 | } |
| 4179 | |
| 4180 | /// As [`Self::prune_sessions_older_than`], but never deletes `keep` — the |
| 4181 | /// active session. A just-resumed session's `updated_at` is stale until |
| 4182 | /// its first post-resume save, so an age prune could otherwise delete the |
| 4183 | /// live session out from under the TUI. |
| 4184 | pub fn prune_sessions_older_than_keeping( |
| 4185 | &self, |
| 4186 | max_age: std::time::Duration, |
| 4187 | keep: Option<&str>, |
| 4188 | ) -> std::io::Result<usize> { |
| 4189 | // A max_age too large to represent (or to subtract from now) means |
| 4190 | // nothing can be old enough to prune — keep everything instead of |
| 4191 | // panicking on DateTime underflow or inventing a fallback horizon. |
| 4192 | let Some(cutoff) = chrono::Duration::from_std(max_age) |
| 4193 | .ok() |
| 4194 | .and_then(|age| Utc::now().checked_sub_signed(age)) |
| 4195 | else { |
| 4196 | return Ok(0); |
| 4197 | }; |
| 4198 | let sessions = self.list_sessions()?; |
| 4199 | let mut pruned = 0usize; |
| 4200 | for session in sessions { |
| 4201 | if keep.is_some_and(|id| id == session.id) { |
| 4202 | continue; |
| 4203 | } |
| 4204 | if session.updated_at < cutoff { |
| 4205 | if let Err(err) = self.remove_session(&session.id, SessionRemoval::Retention) { |
| 4206 | tracing::warn!( |
| 4207 | target: "session", |
| 4208 | session = session.id, |
| 4209 | ?err, |
| 4210 | "session prune skipped a record", |
| 4211 | ); |
| 4212 | continue; |
| 4213 | } |
| 4214 | pruned += 1; |
| 4215 | } |
| 4216 | } |
| 4217 | Ok(pruned) |
| 4218 | } |
| 4219 | |
| 4220 | /// Get the most recent session scoped to the current workspace. |
| 4221 | /// |
| 4222 | /// Archived sessions are skipped: archiving is the user saying "not this |
| 4223 | /// one", and `--continue` / auto-resume must honour that rather than |
| 4224 | /// dragging a put-away session back. |
| 4225 | pub fn get_latest_session_for_workspace( |
| 4226 | &self, |
| 4227 | workspace: &Path, |
| 4228 | ) -> std::io::Result<Option<SessionMetadata>> { |
| 4229 | let sessions = self.list_sessions()?; |
| 4230 | Ok(sessions.into_iter().find(|session| { |
| 4231 | !session.archived |
| 4232 | && workspace_scope_matches(&session.workspace, workspace) |
| 4233 | && !is_empty_auto_created_session(session) |
| 4234 | })) |
| 4235 | } |
| 4236 | |
| 4237 | /// Search sessions by title |
| 4238 | pub fn search_sessions(&self, query: &str) -> std::io::Result<Vec<SessionMetadata>> { |
| 4239 | let query_lower = query.to_lowercase(); |
| 4240 | let sessions = self.list_sessions()?; |
| 4241 | |
| 4242 | Ok(sessions |
| 4243 | .into_iter() |
| 4244 | .filter(|s| s.title.to_lowercase().contains(&query_lower)) |
| 4245 | .collect()) |
| 4246 | } |
| 4247 | } |
| 4248 | |
| 4249 | /// A Runtime store directory inside a session directory: `runtime`, or a |
| 4250 | /// `runtime-recovered-*` sibling opened for a conversation whose store was |
| 4251 | /// missing. |
| 4252 | pub(crate) fn is_runtime_store_dir_name(name: &std::ffi::OsStr) -> bool { |
| 4253 | name.to_str() |
| 4254 | .is_some_and(|name| name == "runtime" || name.starts_with("runtime-recovered-")) |
| 4255 | } |
| 4256 | |
| 4257 | /// Unicode format characters that never belong in a session title: bidi |
| 4258 | /// embeddings/overrides/isolates and marks, zero-width joiners/spaces, the |
| 4259 | /// soft hyphen, BOM, and line/paragraph separators. Together with |
| 4260 | /// `char::is_control` (C0, DEL, C1 — so ESC, BEL, ST, and OSC introducers) |
| 4261 | /// this is the one character policy for the persisted title, the terminal |
| 4262 | /// tab title, and every plain-text listing that echoes a title. |
| 4263 | pub(crate) fn is_title_format_char(ch: char) -> bool { |
| 4264 | matches!( |
| 4265 | ch, |
| 4266 | '\u{00ad}' |
| 4267 | | '\u{061c}' |
| 4268 | | '\u{200b}'..='\u{200f}' |
| 4269 | | '\u{2028}'..='\u{202e}' |
| 4270 | | '\u{2060}'..='\u{2064}' |
| 4271 | | '\u{2066}'..='\u{2069}' |
| 4272 | | '\u{feff}' |
| 4273 | ) |
| 4274 | } |
| 4275 | |
| 4276 | /// One-line notice that a prior session in `workspace` ended mid-turn |
| 4277 | /// (#5715), for the session-pinned prompt prefix. `current_session_id` is |
| 4278 | /// excluded: an in-flight checkpoint of the live session is current work, |
| 4279 | /// not prior work. Returns `None` when no interrupted session exists. |
| 4280 | pub(crate) fn session_recovery_hint( |
| 4281 | workspace: &Path, |
| 4282 | current_session_id: Option<&str>, |
| 4283 | ) -> Option<String> { |
| 4284 | let manager = SessionManager::default_location().ok()?; |
| 4285 | let meta = manager.interrupted_workspace_session(workspace, current_session_id)?; |
| 4286 | Some(format!( |
| 4287 | "A previous Codewhale session in this workspace (\"{}\", id {}, last active {}) has a recovery checkpoint — it likely ended mid-task. Use session_search/session_get to inspect it and offer to summarize or continue the work; resuming is the user's decision (e.g. /resume).", |
| 4288 | meta.title, |
| 4289 | truncate_id(&meta.id), |
| 4290 | meta.updated_at.format("%Y-%m-%d %H:%M UTC"), |
| 4291 | )) |
| 4292 | } |
| 4293 | |
| 4294 | /// Drop control and bidi/zero-width format characters from a title. |
| 4295 | /// |
| 4296 | /// A session title is user- or content-derived text that later reaches an |
| 4297 | /// OSC 0 terminal title, `codewhale sessions` stdout, and the picker, so the |
| 4298 | /// persisted value must not be able to carry a raw escape sequence. Ordinary |
| 4299 | /// text, punctuation, CJK, and emoji pass through untouched. |
| 4300 | pub fn sanitize_session_title(raw: &str) -> String { |
| 4301 | raw.chars() |
| 4302 | .filter(|ch| !ch.is_control() && !is_title_format_char(*ch)) |
| 4303 | .collect() |
| 4304 | } |
| 4305 | |
| 4306 | /// Sanitize, trim, and bound a user-supplied session title. |
| 4307 | /// |
| 4308 | /// Returns `InvalidInput` for an empty title or one longer than |
| 4309 | /// [`MAX_SESSION_TITLE_CHARS`] so every rename surface (picker, `/rename`, |
| 4310 | /// `PATCH /v1/sessions/{id}`) rejects the same inputs with the same reason. |
| 4311 | pub fn normalize_session_title(title: &str) -> std::io::Result<String> { |
| 4312 | let sanitized = sanitize_session_title(title); |
| 4313 | let trimmed = sanitized.trim(); |
| 4314 | if trimmed.is_empty() { |
| 4315 | return Err(std::io::Error::new( |
| 4316 | std::io::ErrorKind::InvalidInput, |
| 4317 | "Session title cannot be empty", |
| 4318 | )); |
| 4319 | } |
| 4320 | if trimmed.chars().count() > MAX_SESSION_TITLE_CHARS { |
| 4321 | return Err(std::io::Error::new( |
| 4322 | std::io::ErrorKind::InvalidInput, |
| 4323 | format!("Session title cannot exceed {MAX_SESSION_TITLE_CHARS} characters"), |
| 4324 | )); |
| 4325 | } |
| 4326 | Ok(trimmed.to_string()) |
| 4327 | } |
| 4328 | |
| 4329 | pub(crate) fn workspace_scope_matches(saved_workspace: &Path, current_workspace: &Path) -> bool { |
| 4330 | if paths_equivalent(saved_workspace, current_workspace) { |
| 4331 | return true; |
| 4332 | } |
| 4333 | |
| 4334 | // Repository identity comes from the containing checkout itself (Git |
| 4335 | // dir/worktree traversal shared with project-context scope resolution), |
| 4336 | // never from branch names or paths mentioned in conversation. |
| 4337 | let canonical = |path: &Path| fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf()); |
| 4338 | match ( |
| 4339 | find_git_root(&canonical(saved_workspace)), |
| 4340 | find_git_root(&canonical(current_workspace)), |
| 4341 | ) { |
| 4342 | (Some(saved_root), Some(current_root)) => paths_equivalent(&saved_root, ¤t_root), |
| 4343 | _ => false, |
| 4344 | } |
| 4345 | } |
| 4346 | |
| 4347 | pub(crate) fn is_empty_auto_created_session(session: &SessionMetadata) -> bool { |
| 4348 | session.message_count == 0 |
| 4349 | && session |
| 4350 | .title |
| 4351 | .trim() |
| 4352 | .eq_ignore_ascii_case(DEFAULT_SESSION_TITLE) |
| 4353 | } |
| 4354 | |
| 4355 | pub(crate) fn paths_equivalent(lhs: &Path, rhs: &Path) -> bool { |
| 4356 | let lhs_canonical = fs::canonicalize(lhs).ok(); |
| 4357 | let rhs_canonical = fs::canonicalize(rhs).ok(); |
| 4358 | match (lhs_canonical, rhs_canonical) { |
| 4359 | (Some(lhs), Some(rhs)) => lhs == rhs, |
| 4360 | _ => lhs == rhs, |
| 4361 | } |
| 4362 | } |
| 4363 | |
| 4364 | /// Resolve the default session directory path. |
| 4365 | /// |
| 4366 | /// v0.8.44: prefers `~/.codewhale/sessions`, falls back to |
| 4367 | /// `~/.deepseek/sessions` for existing installs. Uses the write-path resolver |
| 4368 | /// so the first access relocates any legacy `~/.deepseek/sessions` into |
| 4369 | /// `~/.codewhale/sessions` when the primary directory is missing (#3240). |
| 4370 | /// If an older build already created an empty primary sessions directory, copy |
| 4371 | /// missing legacy entries into it without overwriting newer CodeWhale data. |
| 4372 | pub fn default_sessions_dir() -> std::io::Result<PathBuf> { |
| 4373 | if let Some(dir) = unsealed_sessions_dir() { |
| 4374 | return Ok(dir); |
| 4375 | } |
| 4376 | let dir = codewhale_config::ensure_state_dir("sessions") |
| 4377 | .map_err(|e| std::io::Error::new(std::io::ErrorKind::NotFound, e.to_string()))?; |
| 4378 | match merge_missing_legacy_session_entries(&dir) { |
| 4379 | Ok(0) => {} |
| 4380 | Ok(count) => { |
| 4381 | tracing::info!( |
| 4382 | target: "session::migration", |
| 4383 | "Copied {count} missing legacy session entries into {}", |
| 4384 | dir.display() |
| 4385 | ); |
| 4386 | } |
| 4387 | Err(err) => { |
| 4388 | tracing::warn!( |
| 4389 | target: "session::migration", |
| 4390 | "Could not copy legacy sessions into {}: {err}", |
| 4391 | dir.display() |
| 4392 | ); |
| 4393 | } |
| 4394 | } |
| 4395 | Ok(dir) |
| 4396 | } |
| 4397 | |
| 4398 | /// `App::new` lists recent sessions, so merely building a fixture reaches the |
| 4399 | /// sessions directory — including the legacy relocation `ensure_state_dir` |
| 4400 | /// performs, which would move a developer's real `~/.deepseek/sessions`. |
| 4401 | /// An unsealed test gets a private directory (and no migration); a sealed one |
| 4402 | /// follows its own environment. |
| 4403 | #[cfg(test)] |
| 4404 | fn unsealed_sessions_dir() -> Option<PathBuf> { |
| 4405 | crate::test_support::unsealed_state_dir("sessions") |
| 4406 | } |
| 4407 | |
| 4408 | #[cfg(not(test))] |
| 4409 | fn unsealed_sessions_dir() -> Option<PathBuf> { |
| 4410 | None |
| 4411 | } |
| 4412 | |
| 4413 | fn merge_missing_legacy_session_entries(primary: &Path) -> io::Result<usize> { |
| 4414 | if codewhale_paths::codewhale_home_is_explicit() { |
| 4415 | return Ok(0); |
| 4416 | } |
| 4417 | |
| 4418 | let legacy = codewhale_config::legacy_deepseek_home() |
| 4419 | .map_err(|e| io::Error::new(io::ErrorKind::NotFound, e.to_string()))? |
| 4420 | .join("sessions"); |
| 4421 | if !legacy.is_dir() || paths_equivalent(primary, &legacy) { |
| 4422 | return Ok(0); |
| 4423 | } |
| 4424 | |
| 4425 | copy_missing_dir_entries(&legacy, primary) |
| 4426 | } |
| 4427 | |
| 4428 | fn copy_missing_dir_entries(src: &Path, dst: &Path) -> io::Result<usize> { |
| 4429 | fs::create_dir_all(dst)?; |
| 4430 | let mut copied = 0; |
| 4431 | for entry in fs::read_dir(src)? { |
| 4432 | let entry = entry?; |
| 4433 | let source = entry.path(); |
| 4434 | let target = dst.join(entry.file_name()); |
| 4435 | |
| 4436 | let file_type = entry.file_type()?; |
| 4437 | if file_type.is_dir() { |
| 4438 | if entry.file_name() == std::ffi::OsStr::new("checkpoints") || target.exists() { |
| 4439 | continue; |
| 4440 | } |
| 4441 | copied += copy_missing_dir_entries(&source, &target)?; |
| 4442 | } else if file_type.is_file() { |
| 4443 | copied += usize::from(copy_file_create_new(&source, &target)?); |
| 4444 | } |
| 4445 | } |
| 4446 | Ok(copied) |
| 4447 | } |
| 4448 | |
| 4449 | fn copy_file_create_new(src: &Path, dst: &Path) -> io::Result<bool> { |
| 4450 | let mut source = fs::File::open(src)?; |
| 4451 | let mut target = match fs::OpenOptions::new() |
| 4452 | .write(true) |
| 4453 | .create_new(true) |
| 4454 | .open(dst) |
| 4455 | { |
| 4456 | Ok(file) => file, |
| 4457 | Err(err) if err.kind() == io::ErrorKind::AlreadyExists => return Ok(false), |
| 4458 | Err(err) => return Err(err), |
| 4459 | }; |
| 4460 | if let Err(err) = io::copy(&mut source, &mut target) { |
| 4461 | let _ = fs::remove_file(dst); |
| 4462 | return Err(err); |
| 4463 | } |
| 4464 | Ok(true) |
| 4465 | } |
| 4466 | |
| 4467 | /// Prune snapshots older than `max_age` for `workspace`. |
| 4468 | /// |
| 4469 | /// Always non-fatal. Returns silently — callers don't need the count |
| 4470 | /// (the underlying repo logs at WARN if anything blew up). |
| 4471 | pub fn prune_workspace_snapshots(workspace: &Path, max_age: std::time::Duration) { |
| 4472 | match crate::snapshot::prune_older_than(workspace, max_age) { |
| 4473 | Ok(0) => {} |
| 4474 | Ok(n) => { |
| 4475 | tracing::debug!(target: "snapshot", "boot prune removed {n} snapshot(s)"); |
| 4476 | } |
| 4477 | Err(e) => { |
| 4478 | tracing::warn!(target: "snapshot", "boot prune failed: {e}"); |
| 4479 | } |
| 4480 | } |
| 4481 | } |
| 4482 | |
| 4483 | /// Create a new `SavedSession` from conversation state |
| 4484 | #[cfg(test)] |
| 4485 | pub fn create_saved_session( |
| 4486 | messages: &[Message], |
| 4487 | model: &str, |
| 4488 | workspace: &Path, |
| 4489 | total_tokens: u64, |
| 4490 | system_prompt: Option<&SystemPrompt>, |
| 4491 | ) -> SavedSession { |
| 4492 | create_saved_session_with_mode( |
| 4493 | messages, |
| 4494 | model, |
| 4495 | workspace, |
| 4496 | total_tokens, |
| 4497 | system_prompt, |
| 4498 | None, |
| 4499 | ) |
| 4500 | } |
| 4501 | |
| 4502 | /// Placeholder title used for a session that has no first user message yet. |
| 4503 | /// `build_session_snapshot` (tui/ui/frame.rs) treats a title equal to this |
| 4504 | /// constant as an auto-generated placeholder and lets the conversation-derived |
| 4505 | /// title win once a user message exists. Keep this string stable on purpose. |
| 4506 | pub(crate) const DEFAULT_SESSION_TITLE: &str = "New Session"; |
| 4507 | |
| 4508 | /// Create a new `SavedSession` from conversation state with optional mode label |
| 4509 | pub fn create_saved_session_with_mode( |
| 4510 | messages: &[Message], |
| 4511 | model: &str, |
| 4512 | workspace: &Path, |
| 4513 | total_tokens: u64, |
| 4514 | system_prompt: Option<&SystemPrompt>, |
| 4515 | mode: Option<&str>, |
| 4516 | ) -> SavedSession { |
| 4517 | create_saved_session_with_id_and_mode( |
| 4518 | Uuid::new_v4().to_string(), |
| 4519 | messages, |
| 4520 | model, |
| 4521 | workspace, |
| 4522 | total_tokens, |
| 4523 | system_prompt, |
| 4524 | mode, |
| 4525 | ) |
| 4526 | } |
| 4527 | |
| 4528 | /// Create a new `SavedSession` using a caller-owned session id. |
| 4529 | pub fn create_saved_session_with_id_and_mode( |
| 4530 | id: String, |
| 4531 | messages: &[Message], |
| 4532 | model: &str, |
| 4533 | workspace: &Path, |
| 4534 | total_tokens: u64, |
| 4535 | system_prompt: Option<&SystemPrompt>, |
| 4536 | mode: Option<&str>, |
| 4537 | ) -> SavedSession { |
| 4538 | create_saved_session_with_id_mode_and_stamps( |
| 4539 | id, |
| 4540 | messages, |
| 4541 | &[], |
| 4542 | model, |
| 4543 | workspace, |
| 4544 | total_tokens, |
| 4545 | system_prompt, |
| 4546 | mode, |
| 4547 | ) |
| 4548 | } |
| 4549 | |
| 4550 | /// Create a new `SavedSession` whose journal entries keep the time each |
| 4551 | /// message actually landed. `message_stamps[i]` is the append time of |
| 4552 | /// `messages[i]`; a missing stamp falls back to now. Callers without a |
| 4553 | /// live stamp log pass `&[]` and get the historical save-time behavior. |
| 4554 | pub fn create_saved_session_with_id_mode_and_stamps( |
| 4555 | id: String, |
| 4556 | messages: &[Message], |
| 4557 | message_stamps: &[DateTime<Utc>], |
| 4558 | model: &str, |
| 4559 | workspace: &Path, |
| 4560 | total_tokens: u64, |
| 4561 | system_prompt: Option<&SystemPrompt>, |
| 4562 | mode: Option<&str>, |
| 4563 | ) -> SavedSession { |
| 4564 | create_saved_session_inner( |
| 4565 | id, |
| 4566 | messages, |
| 4567 | SessionJournal::from_messages_stamped(messages.to_vec(), message_stamps, 0), |
| 4568 | model, |
| 4569 | workspace, |
| 4570 | total_tokens, |
| 4571 | system_prompt, |
| 4572 | mode, |
| 4573 | true, |
| 4574 | ) |
| 4575 | } |
| 4576 | |
| 4577 | /// Create a snapshot whose `messages` projection is left empty (#6214 T3). |
| 4578 | /// |
| 4579 | /// The journal still carries every message, and every serialization path |
| 4580 | /// rehydrates the projection (`serialize_saved_session` runs |
| 4581 | /// `make_storage_compatible`), so the on-disk bytes are identical to the |
| 4582 | /// filled form. The persistence queue drops the projection anyway |
| 4583 | /// (`compact_for_persistence_queue`), so building it first is one full |
| 4584 | /// history copy per debounced flush for nothing. Callers that serialize a |
| 4585 | /// journal-only snapshot directly must run `make_storage_compatible` first. |
| 4586 | pub fn create_saved_session_journal_only( |
| 4587 | id: String, |
| 4588 | messages: &[Message], |
| 4589 | journal: SessionJournal, |
| 4590 | model: &str, |
| 4591 | workspace: &Path, |
| 4592 | total_tokens: u64, |
| 4593 | system_prompt: Option<&SystemPrompt>, |
| 4594 | mode: Option<&str>, |
| 4595 | ) -> SavedSession { |
| 4596 | create_saved_session_inner( |
| 4597 | id, |
| 4598 | messages, |
| 4599 | journal, |
| 4600 | model, |
| 4601 | workspace, |
| 4602 | total_tokens, |
| 4603 | system_prompt, |
| 4604 | mode, |
| 4605 | false, |
| 4606 | ) |
| 4607 | } |
| 4608 | |
| 4609 | fn create_saved_session_inner( |
| 4610 | id: String, |
| 4611 | messages: &[Message], |
| 4612 | journal: SessionJournal, |
| 4613 | model: &str, |
| 4614 | workspace: &Path, |
| 4615 | total_tokens: u64, |
| 4616 | system_prompt: Option<&SystemPrompt>, |
| 4617 | mode: Option<&str>, |
| 4618 | fill_messages: bool, |
| 4619 | ) -> SavedSession { |
| 4620 | let now = Utc::now(); |
| 4621 | |
| 4622 | // Generate title from the first real user message (runtime-owned control |
| 4623 | // traffic is skipped by `conversation_derived_title`). Fall back to the |
| 4624 | // placeholder when no user-authored prompt exists yet. |
| 4625 | let title = |
| 4626 | conversation_derived_title(messages).unwrap_or_else(|| DEFAULT_SESSION_TITLE.to_string()); |
| 4627 | |
| 4628 | let leaf_id = journal.leaf_id.clone(); |
| 4629 | SavedSession { |
| 4630 | schema_version: CURRENT_SESSION_SCHEMA_VERSION, |
| 4631 | metadata: SessionMetadata { |
| 4632 | id, |
| 4633 | title, |
| 4634 | created_at: now, |
| 4635 | updated_at: now, |
| 4636 | message_count: messages.len(), |
| 4637 | total_tokens, |
| 4638 | model: model.to_string(), |
| 4639 | model_provider: default_model_provider(), |
| 4640 | model_provider_id: None, |
| 4641 | workspace: workspace.to_path_buf(), |
| 4642 | mode: mode.map(str::to_string), |
| 4643 | cost: SessionCostSnapshot::default(), |
| 4644 | parent_session_id: None, |
| 4645 | forked_from_message_count: None, |
| 4646 | runtime_store: None, |
| 4647 | cumulative_turn_secs: 0, |
| 4648 | archived: false, |
| 4649 | spawn_depth: journal.spawn_depth, |
| 4650 | }, |
| 4651 | messages: if fill_messages { |
| 4652 | messages.to_vec() |
| 4653 | } else { |
| 4654 | Vec::new() |
| 4655 | }, |
| 4656 | journal: Some(journal), |
| 4657 | leaf_id, |
| 4658 | system_prompt: system_prompt_to_string(system_prompt), |
| 4659 | context_references: Vec::new(), |
| 4660 | artifacts: Vec::new(), |
| 4661 | approval_receipts: Vec::new(), |
| 4662 | work_state: None, |
| 4663 | window_title: None, |
| 4664 | last_auto_route: None, |
| 4665 | turn_outcomes: Vec::new(), |
| 4666 | } |
| 4667 | } |
| 4668 | |
| 4669 | /// Update an existing session with new messages |
| 4670 | pub fn update_session( |
| 4671 | mut session: SavedSession, |
| 4672 | messages: &[Message], |
| 4673 | total_tokens: u64, |
| 4674 | system_prompt: Option<&SystemPrompt>, |
| 4675 | ) -> SavedSession { |
| 4676 | session.schema_version = CURRENT_SESSION_SCHEMA_VERSION; |
| 4677 | session.ensure_journal(); |
| 4678 | let old_len = session.messages.len(); |
| 4679 | let new_len = messages.len(); |
| 4680 | if new_len >= old_len && messages[..old_len] == session.messages[..] { |
| 4681 | if let Some(journal) = session.journal.as_mut() { |
| 4682 | for msg in &messages[old_len..] { |
| 4683 | journal.append_message(msg.clone()); |
| 4684 | } |
| 4685 | session.leaf_id = journal.leaf_id.clone(); |
| 4686 | } |
| 4687 | } else if (new_len != old_len || messages != session.messages.as_slice()) |
| 4688 | && let Some(journal) = session.journal.as_mut() |
| 4689 | { |
| 4690 | let common = messages |
| 4691 | .iter() |
| 4692 | .zip(session.messages.iter()) |
| 4693 | .take_while(|(a, b)| a == b) |
| 4694 | .count(); |
| 4695 | if common > 0 && common <= journal.entries.len() { |
| 4696 | let target_id = journal |
| 4697 | .root_to_leaf() |
| 4698 | .get(common - 1) |
| 4699 | .map(|entry| entry.id.clone()); |
| 4700 | if let Some(target_id) = target_id { |
| 4701 | let _ = journal.branch_to(&target_id); |
| 4702 | } else { |
| 4703 | journal.leaf_id = None; |
| 4704 | } |
| 4705 | } else if common == 0 { |
| 4706 | journal.leaf_id = journal.entries.first().and_then(|e| e.parent_id.clone()); |
| 4707 | if journal.leaf_id.is_none() && !journal.entries.is_empty() { |
| 4708 | journal.leaf_id = None; |
| 4709 | } |
| 4710 | } |
| 4711 | for msg in messages.iter().skip(common) { |
| 4712 | journal.append_message(msg.clone()); |
| 4713 | } |
| 4714 | session.leaf_id = journal.leaf_id.clone(); |
| 4715 | } |
| 4716 | session.messages.clear(); |
| 4717 | session.messages.extend_from_slice(messages); |
| 4718 | session.metadata.updated_at = Utc::now(); |
| 4719 | session.metadata.message_count = messages.len(); |
| 4720 | session.metadata.total_tokens = total_tokens; |
| 4721 | session.system_prompt = system_prompt_to_string(system_prompt); |
| 4722 | session |
| 4723 | } |
| 4724 | |
| 4725 | /// Strip a stale `[Session note]` block that was written by the old |
| 4726 | /// 500-message cap. Only removes notes that contain the specific |
| 4727 | /// "older messages were dropped" phrase — ordinary user-added |
| 4728 | /// `[Session note]` prompts are left untouched. |
| 4729 | fn strip_legacy_truncation_note(system_prompt: Option<String>) -> Option<String> { |
| 4730 | let sp = system_prompt?; |
| 4731 | let Some(trimmed) = sp.strip_prefix("[Session note]\n") else { |
| 4732 | return Some(sp); |
| 4733 | }; |
| 4734 | // Only strip if this is the known cap_messages note. |
| 4735 | if !trimmed.contains("older messages were dropped") { |
| 4736 | return Some(sp); |
| 4737 | } |
| 4738 | // The note block ends with "\n\n---\n\n" (7 chars) followed by the real prompt. |
| 4739 | trimmed |
| 4740 | .find("\n\n---\n\n") |
| 4741 | .map(|pos| trimmed[pos + 7..].to_string()) |
| 4742 | } |
| 4743 | |
| 4744 | /// Byte offset of `key` (a quoted JSON key such as `"metadata"`) outside any |
| 4745 | /// string literal. Brace/string-aware so a key name quoted inside an earlier |
| 4746 | /// message body is never matched. |
| 4747 | fn find_json_key(bytes: &[u8], key: &[u8]) -> Option<usize> { |
| 4748 | let mut idx = 0usize; |
| 4749 | let mut in_string = false; |
| 4750 | let mut escape = false; |
| 4751 | while idx < bytes.len() { |
| 4752 | let c = bytes[idx]; |
| 4753 | if escape { |
| 4754 | escape = false; |
| 4755 | } else if c == b'\\' { |
| 4756 | escape = true; |
| 4757 | } else if c == b'"' { |
| 4758 | if !in_string && bytes[idx..].starts_with(key) { |
| 4759 | return Some(idx); |
| 4760 | } |
| 4761 | in_string = !in_string; |
| 4762 | } |
| 4763 | idx += 1; |
| 4764 | } |
| 4765 | None |
| 4766 | } |
| 4767 | |
| 4768 | /// Offset of the value opening with `open` that follows the key at |
| 4769 | /// `key_offset`. |
| 4770 | fn json_value_start(bytes: &[u8], key_offset: usize, key_len: usize, open: u8) -> Option<usize> { |
| 4771 | let mut idx = key_offset + key_len; |
| 4772 | while idx < bytes.len() && (bytes[idx] as char).is_whitespace() { |
| 4773 | idx += 1; |
| 4774 | } |
| 4775 | if idx >= bytes.len() || bytes[idx] != b':' { |
| 4776 | return None; |
| 4777 | } |
| 4778 | idx += 1; |
| 4779 | while idx < bytes.len() && (bytes[idx] as char).is_whitespace() { |
| 4780 | idx += 1; |
| 4781 | } |
| 4782 | (idx < bytes.len() && bytes[idx] == open).then_some(idx) |
| 4783 | } |
| 4784 | |
| 4785 | /// Exclusive end of the balanced `{...}` starting at `start`, or `None` when |
| 4786 | /// the buffer is truncated before it closes. |
| 4787 | fn json_object_end(bytes: &[u8], start: usize) -> Option<usize> { |
| 4788 | let mut depth = 0i32; |
| 4789 | let mut in_string = false; |
| 4790 | let mut escape = false; |
| 4791 | for (offset, &c) in bytes[start..].iter().enumerate() { |
| 4792 | if escape { |
| 4793 | escape = false; |
| 4794 | continue; |
| 4795 | } |
| 4796 | match c { |
| 4797 | b'\\' => escape = true, |
| 4798 | b'"' => in_string = !in_string, |
| 4799 | b'{' if !in_string => depth += 1, |
| 4800 | b'}' if !in_string => { |
| 4801 | depth -= 1; |
| 4802 | if depth == 0 { |
| 4803 | return Some(start + offset + 1); |
| 4804 | } |
| 4805 | } |
| 4806 | _ => {} |
| 4807 | } |
| 4808 | } |
| 4809 | None |
| 4810 | } |
| 4811 | |
| 4812 | /// The longest valid UTF-8 prefix of `buf`. A fixed-size read prefix can end |
| 4813 | /// inside a multi-byte character (routine for CJK or emoji transcripts, which |
| 4814 | /// serde writes unescaped); rejecting the whole buffer then forced a full-file |
| 4815 | /// read of every such session on each listing. |
| 4816 | /// |
| 4817 | /// Only a character cut at the very end is trimmed. An invalid byte anywhere |
| 4818 | /// else means a damaged file, and an empty prefix sends it down the full read, |
| 4819 | /// which reports it. |
| 4820 | fn utf8_prefix(buf: &[u8]) -> &str { |
| 4821 | match std::str::from_utf8(buf) { |
| 4822 | Ok(s) => s, |
| 4823 | // `valid_up_to` marks a char boundary, so this slice is valid. |
| 4824 | Err(error) if error.error_len().is_none() => { |
| 4825 | std::str::from_utf8(&buf[..error.valid_up_to()]).unwrap_or_default() |
| 4826 | } |
| 4827 | Err(_) => "", |
| 4828 | } |
| 4829 | } |
| 4830 | |
| 4831 | /// String-scan a JSON byte buffer for the top-level `"metadata":{...}` |
| 4832 | /// block and return it parsed. Returns `None` if no balanced metadata |
| 4833 | /// object is present in the buffer. |
| 4834 | /// |
| 4835 | /// Supports the optimisation in `SessionManager::load_session_metadata` |
| 4836 | /// (#337). The scanner is brace-balanced and string-aware so a `{` or |
| 4837 | /// `}` appearing inside a string literal doesn't perturb the depth |
| 4838 | /// count. |
| 4839 | fn extract_top_level_metadata(buf: &[u8]) -> Option<SessionMetadata> { |
| 4840 | let s = utf8_prefix(buf); |
| 4841 | let bytes = s.as_bytes(); |
| 4842 | const KEY: &[u8] = b"\"metadata\""; |
| 4843 | let start = json_value_start(bytes, find_json_key(bytes, KEY)?, KEY.len(), b'{')?; |
| 4844 | let end = json_object_end(bytes, start)?; |
| 4845 | serde_json::from_str::<SessionMetadata>(&s[start..end]).ok() |
| 4846 | } |
| 4847 | |
| 4848 | /// Complete message objects from the front of the `messages` array, plus |
| 4849 | /// whether the array was seen to end. A message the prefix cut in half is |
| 4850 | /// simply absent; nothing is reconstructed. |
| 4851 | fn extract_leading_messages(buf: &[u8], max: usize) -> (Vec<Message>, bool) { |
| 4852 | let s = utf8_prefix(buf); |
| 4853 | let bytes = s.as_bytes(); |
| 4854 | const KEY: &[u8] = b"\"messages\""; |
| 4855 | let Some(key_offset) = find_json_key(bytes, KEY) else { |
| 4856 | return (Vec::new(), false); |
| 4857 | }; |
| 4858 | let Some(array_start) = json_value_start(bytes, key_offset, KEY.len(), b'[') else { |
| 4859 | return (Vec::new(), false); |
| 4860 | }; |
| 4861 | let mut cursor = array_start + 1; |
| 4862 | let mut out = Vec::new(); |
| 4863 | loop { |
| 4864 | while cursor < bytes.len() && matches!(bytes[cursor], b' ' | b'\t' | b'\r' | b'\n' | b',') { |
| 4865 | cursor += 1; |
| 4866 | } |
| 4867 | if cursor < bytes.len() && bytes[cursor] == b']' { |
| 4868 | return (out, true); |
| 4869 | } |
| 4870 | if out.len() >= max || cursor >= bytes.len() || bytes[cursor] != b'{' { |
| 4871 | return (out, false); |
| 4872 | } |
| 4873 | let Some(end) = json_object_end(bytes, cursor) else { |
| 4874 | return (out, false); |
| 4875 | }; |
| 4876 | let Ok(message) = serde_json::from_str::<Message>(&s[cursor..end]) else { |
| 4877 | return (out, false); |
| 4878 | }; |
| 4879 | out.push(message); |
| 4880 | cursor = end; |
| 4881 | } |
| 4882 | } |
| 4883 | |
| 4884 | /// How many leading messages the legacy-title recovery will parse. The |
| 4885 | /// enclosing read is already bounded to a 64 KB prefix (#337); this bounds |
| 4886 | /// the parse inside it. |
| 4887 | const LEGACY_TITLE_SCAN_MESSAGES: usize = 24; |
| 4888 | |
| 4889 | /// Recover a title that a superseded derivation took from runtime control |
| 4890 | /// traffic, using the session's own first real user prompt. |
| 4891 | /// |
| 4892 | /// Provenance is proven, never guessed: the stored title has to be exactly |
| 4893 | /// what the old rule produced — [`truncate_title`] of a message the current |
| 4894 | /// classifier rejects as not a user turn. A renamed session, and a person who |
| 4895 | /// literally typed an envelope as their first message, never match, so their |
| 4896 | /// text is kept. Returns `None` when there is nothing proven to recover. |
| 4897 | fn recovered_legacy_title( |
| 4898 | stored: &str, |
| 4899 | messages: &[Message], |
| 4900 | array_complete: bool, |
| 4901 | ) -> Option<String> { |
| 4902 | let stale = messages.iter().any(|message| { |
| 4903 | crate::runtime_handoff::classify_user_turn_prompt(message) |
| 4904 | == crate::runtime_handoff::UserTurnPromptKind::NotPrompt |
| 4905 | && message.content.iter().any(|block| match block { |
| 4906 | ContentBlock::Text { text, .. } => { |
| 4907 | truncate_title(text, 50) == stored |
| 4908 | || truncate_title(extract_user_prompt(text), 50) == stored |
| 4909 | } |
| 4910 | _ => false, |
| 4911 | }) |
| 4912 | }); |
| 4913 | if !stale { |
| 4914 | return None; |
| 4915 | } |
| 4916 | match conversation_derived_title(messages) { |
| 4917 | Some(title) => Some(title), |
| 4918 | // No user turn in what we read. Only claim the conversation has none |
| 4919 | // when the array actually ended inside the prefix; a truncated read |
| 4920 | // keeps the stored title rather than inventing a neutral one. |
| 4921 | None if array_complete => Some(DEFAULT_SESSION_TITLE.to_string()), |
| 4922 | None => None, |
| 4923 | } |
| 4924 | } |
| 4925 | |
| 4926 | /// Apply [`recovered_legacy_title`] to freshly loaded metadata. In memory |
| 4927 | /// only — the session file is never rewritten, so the stored title (and any |
| 4928 | /// rename) survives on disk. |
| 4929 | fn apply_legacy_title_recovery(metadata: &mut SessionMetadata, buf: &[u8]) { |
| 4930 | // Cost gate, never the rename decision: an envelope title always opens |
| 4931 | // with `<`, so this keeps #337's bounded-parse win for ordinary titles. |
| 4932 | // Whether to rewrite is `recovered_legacy_title`'s proven provenance. |
| 4933 | if !metadata.title.starts_with('<') { |
| 4934 | return; |
| 4935 | } |
| 4936 | let (messages, complete) = extract_leading_messages(buf, LEGACY_TITLE_SCAN_MESSAGES); |
| 4937 | if messages.is_empty() { |
| 4938 | return; |
| 4939 | } |
| 4940 | if let Some(title) = recovered_legacy_title(&metadata.title, &messages, complete) { |
| 4941 | metadata.title = title; |
| 4942 | } |
| 4943 | } |
| 4944 | |
| 4945 | fn system_prompt_to_string(system_prompt: Option<&SystemPrompt>) -> Option<String> { |
| 4946 | match system_prompt { |
| 4947 | Some(SystemPrompt::Text(text)) => Some(text.clone()), |
| 4948 | Some(SystemPrompt::Blocks(blocks)) => Some( |
| 4949 | blocks |
| 4950 | .iter() |
| 4951 | .map(|b| b.text.clone()) |
| 4952 | .collect::<Vec<_>>() |
| 4953 | .join("\n\n---\n\n"), |
| 4954 | ), |
| 4955 | None => None, |
| 4956 | } |
| 4957 | } |
| 4958 | |
| 4959 | /// Truncate a session ID to 8 characters for compact display. |
| 4960 | /// Returns a `&str` borrowing from the input — no allocation. |
| 4961 | pub fn truncate_id(id: &str) -> &str { |
| 4962 | id.get(..8).unwrap_or(id) |
| 4963 | } |
| 4964 | |
| 4965 | /// Strip a leading `<turn_meta>...</turn_meta>` block from saved user text. |
| 4966 | /// |
| 4967 | /// Older sessions can have turn metadata prefixed to the first user message. |
| 4968 | /// The session picker and generated session titles should show the user's |
| 4969 | /// prompt, not the cache/debug envelope. |
| 4970 | pub(crate) fn extract_user_prompt(raw: &str) -> &str { |
| 4971 | let trimmed = raw.trim_start(); |
| 4972 | let Some(after_open) = trimmed.strip_prefix("<turn_meta>") else { |
| 4973 | return trimmed; |
| 4974 | }; |
| 4975 | if let Some(close_pos) = after_open.find("</turn_meta>") { |
| 4976 | return after_open[close_pos + "</turn_meta>".len()..].trim_start(); |
| 4977 | } |
| 4978 | after_open.trim_start() |
| 4979 | } |
| 4980 | |
| 4981 | /// Clean a stored title for display, falling back to a neutral label. |
| 4982 | pub(crate) fn extract_title(raw: &str) -> &str { |
| 4983 | let title = extract_user_prompt(raw); |
| 4984 | if title.is_empty() { "Session" } else { title } |
| 4985 | } |
| 4986 | |
| 4987 | /// Strip common inline thinking/reasoning XML sections from saved assistant |
| 4988 | /// text before it is shown in session previews. |
| 4989 | pub(crate) fn strip_thinking_tags(text: &str) -> String { |
| 4990 | if !text.contains("<think") && !text.contains("<thinking") && !text.contains("<reasoning") { |
| 4991 | return text.to_string(); |
| 4992 | } |
| 4993 | |
| 4994 | let tags = ["think", "thinking", "reasoning"]; |
| 4995 | let mut result = text.to_string(); |
| 4996 | for tag in tags { |
| 4997 | let open = format!("<{tag}>"); |
| 4998 | let close = format!("</{tag}>"); |
| 4999 | while let Some(start) = result.find(&open) { |
| 5000 | let Some(end) = result[start..].find(&close) else { |
| 5001 | break; |
| 5002 | }; |
| 5003 | let end_abs = start + end + close.len(); |
| 5004 | result.replace_range(start..end_abs, ""); |
| 5005 | } |
| 5006 | } |
| 5007 | result |
| 5008 | } |
| 5009 | |
| 5010 | /// Truncate a string to create a title (character-safe for UTF-8) |
| 5011 | fn truncate_title(s: &str, max_len: usize) -> String { |
| 5012 | let s = s.trim(); |
| 5013 | // Older sessions may carry a title saved before sanitization existed; |
| 5014 | // never echo raw controls into stdout or the picker. Take the first |
| 5015 | // line before sanitizing so a legacy multi-line title still shows only |
| 5016 | // its first line. |
| 5017 | let first_line = sanitize_session_title(s.lines().next().unwrap_or(s)); |
| 5018 | let first_line = first_line.trim(); |
| 5019 | |
| 5020 | let char_count = first_line.chars().count(); |
| 5021 | if char_count <= max_len { |
| 5022 | first_line.to_string() |
| 5023 | } else { |
| 5024 | let truncated: String = first_line.chars().take(max_len - 3).collect(); |
| 5025 | format!("{truncated}...") |
| 5026 | } |
| 5027 | } |
| 5028 | |
| 5029 | /// Derive the auto-title from the first real user message of a conversation. |
| 5030 | /// |
| 5031 | /// Returns `None` when no user-authored message exists to name the session |
| 5032 | /// after (an empty transcript, or one holding only runtime-owned control |
| 5033 | /// traffic); callers fall back to [`DEFAULT_SESSION_TITLE`]. |
| 5034 | /// |
| 5035 | /// Chat-template compatibility forces runtime-owned control traffic |
| 5036 | /// (sub-agent handoffs, the Operate contract, restore checkpoints) through |
| 5037 | /// `role = "user"`, but an internal envelope is not what the person typed. |
| 5038 | /// Prompt eligibility comes from the existing user-turn classifier; the live |
| 5039 | /// title fallback shares the same selection through `conversation_title_prompt`. |
| 5040 | fn conversation_derived_title(messages: &[Message]) -> Option<String> { |
| 5041 | conversation_title_prompt(messages).map(|prompt| truncate_title(prompt, 50)) |
| 5042 | } |
| 5043 | |
| 5044 | /// Select the first real user turn's text for persisted and live titles. |
| 5045 | /// Keep an image-only turn as the first user boundary, and strip historical |
| 5046 | /// leading turn metadata without introducing another provenance classifier. |
| 5047 | pub(crate) fn conversation_title_prompt(messages: &[Message]) -> Option<&str> { |
| 5048 | messages |
| 5049 | .iter() |
| 5050 | .find(|message| { |
| 5051 | crate::runtime_handoff::classify_user_turn_prompt(message) |
| 5052 | != crate::runtime_handoff::UserTurnPromptKind::NotPrompt |
| 5053 | }) |
| 5054 | .and_then(|m| { |
| 5055 | m.content.iter().find_map(|block| match block { |
| 5056 | ContentBlock::Text { text, .. } => { |
| 5057 | let prompt = extract_user_prompt(text); |
| 5058 | if prompt.is_empty() { |
| 5059 | None |
| 5060 | } else { |
| 5061 | Some(prompt) |
| 5062 | } |
| 5063 | } |
| 5064 | _ => None, |
| 5065 | }) |
| 5066 | }) |
| 5067 | } |
| 5068 | |
| 5069 | /// Format a session for display in a picker |
| 5070 | pub fn format_session_line(meta: &SessionMetadata) -> String { |
| 5071 | let age = format_age(&meta.updated_at); |
| 5072 | let updated = format_session_updated_at(&meta.updated_at, &age); |
| 5073 | let truncated_title = truncate_title(extract_title(&meta.title), 40); |
| 5074 | let fork_label = if meta.parent_session_id.is_some() { |
| 5075 | " | fork" |
| 5076 | } else { |
| 5077 | "" |
| 5078 | }; |
| 5079 | |
| 5080 | format!( |
| 5081 | "{} | {} | {} msgs{} | {}", |
| 5082 | truncate_id(&meta.id), |
| 5083 | truncated_title, |
| 5084 | meta.message_count, |
| 5085 | fork_label, |
| 5086 | updated |
| 5087 | ) |
| 5088 | } |
| 5089 | |
| 5090 | pub(crate) fn format_session_updated_at(dt: &DateTime<Utc>, age: &str) -> String { |
| 5091 | format!("{} ({age})", dt.format("%Y-%m-%d %H:%M UTC")) |
| 5092 | } |
| 5093 | |
| 5094 | /// Format a datetime as relative age |
| 5095 | fn format_age(dt: &DateTime<Utc>) -> String { |
| 5096 | let now = Utc::now(); |
| 5097 | let duration = now.signed_duration_since(*dt); |
| 5098 | |
| 5099 | if duration.num_minutes() < 1 { |
| 5100 | "just now".to_string() |
| 5101 | } else if duration.num_hours() < 1 { |
| 5102 | format!("{}m ago", duration.num_minutes()) |
| 5103 | } else if duration.num_days() < 1 { |
| 5104 | format!("{}h ago", duration.num_hours()) |
| 5105 | } else if duration.num_weeks() < 1 { |
| 5106 | format!("{}d ago", duration.num_days()) |
| 5107 | } else { |
| 5108 | format!("{}w ago", duration.num_weeks()) |
| 5109 | } |
| 5110 | } |
| 5111 | |
| 5112 | // === Unit Tests === |
| 5113 | |
| 5114 | #[cfg(test)] |
| 5115 | mod tests { |
| 5116 | include!("session_manager/test_cases_01.rs"); |
| 5117 | |
| 5118 | include!("session_manager/test_cases_02.rs"); |
| 5119 | |
| 5120 | include!("session_manager/test_cases_03.rs"); |
| 5121 | |
| 5122 | include!("session_manager/test_cases_04.rs"); |
| 5123 | } |
| 5124 | |
| 5125 | #[cfg(test)] |
| 5126 | mod storage_compatible_tests { |
| 5127 | use super::*; |
| 5128 | |
| 5129 | fn user(text: &str) -> Message { |
| 5130 | Message { |
| 5131 | role: codewhale_models::Role::from("user"), |
| 5132 | content: vec![codewhale_models::ContentBlock::Text { |
| 5133 | text: text.to_string(), |
| 5134 | cache_control: None, |
| 5135 | }], |
| 5136 | } |
| 5137 | } |
| 5138 | |
| 5139 | /// Journal-only snapshots (#6214 T3) skip building the `messages` |
| 5140 | /// projection, but a save still lands full history: serialization |
| 5141 | /// rehydrates the projection from the journal, so a reload is whole. |
| 5142 | #[test] |
| 5143 | fn journal_only_snapshot_saves_and_reloads_full_history() { |
| 5144 | let tmp = tempfile::tempdir().expect("tempdir"); |
| 5145 | let messages = vec![user("first"), user("answer")]; |
| 5146 | let sparse = create_saved_session_journal_only( |
| 5147 | "roundtrip".to_string(), |
| 5148 | &messages, |
| 5149 | SessionJournal::from_messages(messages.clone(), 0), |
| 5150 | "test-model", |
| 5151 | tmp.path(), |
| 5152 | 7, |
| 5153 | None, |
| 5154 | None, |
| 5155 | ); |
| 5156 | assert!( |
| 5157 | sparse.messages.is_empty(), |
| 5158 | "journal-only snapshots carry no messages projection" |
| 5159 | ); |
| 5160 | assert_eq!( |
| 5161 | sparse.journal.as_ref().expect("journal").to_messages(), |
| 5162 | messages, |
| 5163 | "the journal still carries every message" |
| 5164 | ); |
| 5165 | let manager = SessionManager::new(tmp.path().join("sessions")).expect("manager"); |
| 5166 | manager.save_session_owned(sparse).expect("save sparse"); |
| 5167 | let reloaded = manager.load_session("roundtrip").expect("reload"); |
| 5168 | assert_eq!(reloaded.messages, messages); |
| 5169 | } |
| 5170 | |
| 5171 | /// The no-op cases must stay no-ops, byte for byte. |
| 5172 | /// |
| 5173 | /// `make_storage_compatible` replaced a clone-and-return-`Option` helper |
| 5174 | /// (#6214 T3). Two of that helper's paths returned `None`, and the caller |
| 5175 | /// then serialized the *original* — so a `metadata.message_count` that |
| 5176 | /// disagrees with `messages.len()` survived untouched. Rewriting it in |
| 5177 | /// place would silently edit live data on every save, and nothing else in |
| 5178 | /// the suite catches that. |
| 5179 | #[test] |
| 5180 | fn make_storage_compatible_leaves_the_no_op_cases_byte_identical() { |
| 5181 | let workspace = std::env::temp_dir(); |
| 5182 | let messages = vec![user("one"), user("two")]; |
| 5183 | |
| 5184 | // 1. No journal at all (legacy, pre-journal files): untouched. |
| 5185 | let mut legacy = create_saved_session(&messages, "test-model", &workspace, 0, None); |
| 5186 | legacy.journal = None; |
| 5187 | legacy.metadata.message_count = 99; // deliberately disagrees |
| 5188 | let before = serde_json::to_string_pretty(&legacy).expect("legacy json"); |
| 5189 | let mut after_session = legacy.clone(); |
| 5190 | after_session.make_storage_compatible(); |
| 5191 | assert_eq!( |
| 5192 | serde_json::to_string_pretty(&after_session).expect("legacy json"), |
| 5193 | before, |
| 5194 | "a session with no journal must serialize exactly as it arrived" |
| 5195 | ); |
| 5196 | assert_eq!(after_session.metadata.message_count, 99); |
| 5197 | |
| 5198 | // 2. Journal present and `messages` already equals its active branch: |
| 5199 | // still untouched, including the disagreeing count. |
| 5200 | let mut settled = create_saved_session(&messages, "test-model", &workspace, 0, None); |
| 5201 | assert!(settled.journal.is_some(), "fixture must carry a journal"); |
| 5202 | settled.messages = settled |
| 5203 | .journal |
| 5204 | .as_ref() |
| 5205 | .expect("journal present") |
| 5206 | .to_messages(); |
| 5207 | settled.metadata.message_count = 99; |
| 5208 | let before = serde_json::to_string_pretty(&settled).expect("settled json"); |
| 5209 | let mut after_session = settled.clone(); |
| 5210 | after_session.make_storage_compatible(); |
| 5211 | assert_eq!( |
| 5212 | serde_json::to_string_pretty(&after_session).expect("settled json"), |
| 5213 | before, |
| 5214 | "an already-consistent session must not be rewritten" |
| 5215 | ); |
| 5216 | assert_eq!( |
| 5217 | after_session.metadata.message_count, 99, |
| 5218 | "the early return must happen before message_count is recomputed" |
| 5219 | ); |
| 5220 | } |
| 5221 | |
| 5222 | /// The queued (journal-only) path is what the debounced flush actually |
| 5223 | /// writes: `compact_for_persistence_queue` empties `messages` first. |
| 5224 | #[test] |
| 5225 | fn make_storage_compatible_rehydrates_a_compacted_queue_snapshot() { |
| 5226 | let workspace = std::env::temp_dir(); |
| 5227 | let messages = vec![user("one"), user("two"), user("three")]; |
| 5228 | let mut session = create_saved_session(&messages, "test-model", &workspace, 0, None); |
| 5229 | let expected = session |
| 5230 | .journal |
| 5231 | .as_ref() |
| 5232 | .expect("journal present") |
| 5233 | .to_messages(); |
| 5234 | |
| 5235 | session.compact_for_persistence_queue(); |
| 5236 | assert!( |
| 5237 | session.messages.is_empty(), |
| 5238 | "the queued snapshot is journal-only" |
| 5239 | ); |
| 5240 | |
| 5241 | session.make_storage_compatible(); |
| 5242 | assert_eq!( |
| 5243 | session.messages, expected, |
| 5244 | "the compat projection is rebuilt from the journal" |
| 5245 | ); |
| 5246 | assert_eq!(session.metadata.message_count, expected.len()); |
| 5247 | } |
| 5248 | } |
| 5249 | |
| 5250 | #[cfg(test)] |
| 5251 | mod canonical_journal_admission_tests { |
| 5252 | use super::*; |
| 5253 | #[test] |
| 5254 | fn corrupted_saved_branch_refuses_read_and_write_without_truncating_source() { |
| 5255 | let tmp = tempfile::tempdir().unwrap(); |
| 5256 | let manager = SessionManager::new(tmp.path().join("sessions")).unwrap(); |
| 5257 | let message = Message { |
| 5258 | role: codewhale_models::Role::User, |
| 5259 | content: vec![ContentBlock::Text { |
| 5260 | text: "retained".into(), |
| 5261 | cache_control: None, |
| 5262 | }], |
| 5263 | }; |
| 5264 | let mut session = create_saved_session_with_id_and_mode( |
| 5265 | "corrupt-graph".into(), |
| 5266 | &[message], |
| 5267 | "test-model", |
| 5268 | tmp.path(), |
| 5269 | 0, |
| 5270 | None, |
| 5271 | None, |
| 5272 | ); |
| 5273 | let mut empty_root = session.clone(); |
| 5274 | empty_root.metadata.id = "empty-root-graph".into(); |
| 5275 | let mut alternative = session.messages.clone(); |
| 5276 | alternative[0].role = codewhale_models::Role::Assistant; |
| 5277 | let graph = empty_root.journal.as_mut().unwrap(); |
| 5278 | graph.rebranch_active_messages(&alternative); |
| 5279 | let retained_entries = graph.entries.clone(); |
| 5280 | assert_eq!(retained_entries.len(), 2); |
| 5281 | graph.rebranch_active_messages(&[]); |
| 5282 | empty_root.leaf_id = None; |
| 5283 | empty_root.messages.clear(); |
| 5284 | empty_root.metadata.message_count = 0; |
| 5285 | manager.save_session(&empty_root).unwrap(); |
| 5286 | let observed = manager |
| 5287 | .load_session_snapshot_bounded( |
| 5288 | "empty-root-graph", |
| 5289 | codewhale_protocol::MAX_CANONICAL_HISTORY_BYTES, |
| 5290 | ) |
| 5291 | .unwrap(); |
| 5292 | assert!(observed.messages.is_empty()); |
| 5293 | assert_eq!(observed.leaf_id, None); |
| 5294 | assert_eq!(observed.journal.as_ref().unwrap().leaf_id, None); |
| 5295 | assert_eq!(observed.journal.as_ref().unwrap().entries, retained_entries); |
| 5296 | let graph = session.journal.as_mut().unwrap(); |
| 5297 | graph.entries[0].parent_id = Some(graph.entries[0].id.clone()); |
| 5298 | let bytes = serde_json::to_vec(&session).unwrap(); |
| 5299 | let path = manager.validated_session_path("corrupt-graph").unwrap(); |
| 5300 | fs::write(&path, &bytes).unwrap(); |
| 5301 | assert_eq!( |
| 5302 | manager |
| 5303 | .load_session_snapshot("corrupt-graph") |
| 5304 | .unwrap_err() |
| 5305 | .kind(), |
| 5306 | io::ErrorKind::InvalidData |
| 5307 | ); |
| 5308 | assert!(manager.save_session(&session).is_err()); |
| 5309 | assert_eq!(fs::read(&path).unwrap(), bytes); |
| 5310 | assert!( |
| 5311 | SavedSession::import_foreign( |
| 5312 | session.export_container("fixture"), |
| 5313 | tmp.path().into(), |
| 5314 | "test-model".into() |
| 5315 | ) |
| 5316 | .is_err() |
| 5317 | ); |
| 5318 | } |
| 5319 | } |
| 5320 | |
| 5321 | #[cfg(test)] |
| 5322 | mod journal_bound_tests { |
| 5323 | //! #6842: bounded autosaves archive superseded journal entries before |
| 5324 | //! dropping them from the session document. |
| 5325 | use super::*; |
| 5326 | use std::collections::HashSet; |
| 5327 | |
| 5328 | fn text(text: &str) -> Message { |
| 5329 | Message { |
| 5330 | role: codewhale_models::Role::User, |
| 5331 | content: vec![ContentBlock::Text { |
| 5332 | text: text.into(), |
| 5333 | cache_control: None, |
| 5334 | }], |
| 5335 | } |
| 5336 | } |
| 5337 | |
| 5338 | /// A session whose journal holds `rounds` superseded copies of a |
| 5339 | /// 10-message tail, as repeated compaction leaves it. |
| 5340 | fn compacted_session(id: &str, workspace: &Path, rounds: usize) -> SavedSession { |
| 5341 | let mut session = create_saved_session_with_id_and_mode( |
| 5342 | id.into(), |
| 5343 | &[text("start")], |
| 5344 | "test-model", |
| 5345 | workspace, |
| 5346 | 0, |
| 5347 | None, |
| 5348 | None, |
| 5349 | ); |
| 5350 | let journal = session.journal.as_mut().unwrap(); |
| 5351 | for round in 0..rounds { |
| 5352 | let tail: Vec<Message> = (0..10) |
| 5353 | .map(|i| text(&format!("round {round} message {i}"))) |
| 5354 | .collect(); |
| 5355 | journal.rebranch_active_messages(&tail); |
| 5356 | } |
| 5357 | session.leaf_id = journal.leaf_id.clone(); |
| 5358 | session.messages = journal.to_messages(); |
| 5359 | session.metadata.message_count = session.messages.len(); |
| 5360 | session |
| 5361 | } |
| 5362 | |
| 5363 | fn all_ids(entries: &[SessionEntry]) -> HashSet<String> { |
| 5364 | entries.iter().map(|e| e.id.clone()).collect() |
| 5365 | } |
| 5366 | |
| 5367 | #[test] |
| 5368 | fn bounded_save_archives_then_prunes_and_live_journal_follows() { |
| 5369 | let tmp = tempfile::tempdir().unwrap(); |
| 5370 | let manager = SessionManager::new(tmp.path().join("sessions")).unwrap(); |
| 5371 | let session = compacted_session("bounded-a", tmp.path(), 460); |
| 5372 | let original = session.journal.clone().unwrap(); |
| 5373 | let active = original.to_messages(); |
| 5374 | |
| 5375 | manager.save_session_bounded(session).unwrap(); |
| 5376 | |
| 5377 | let saved = manager.load_session("bounded-a").unwrap(); |
| 5378 | let journal = saved.journal.as_ref().unwrap(); |
| 5379 | assert!( |
| 5380 | journal.len() <= 10 + JOURNAL_RETAINED_DEAD_ENTRIES, |
| 5381 | "{}", |
| 5382 | journal.len() |
| 5383 | ); |
| 5384 | assert_eq!(journal.to_messages(), active); |
| 5385 | assert_eq!(saved.leaf_id, original.leaf_id); |
| 5386 | let archive = manager.load_journal_archive("bounded-a").unwrap(); |
| 5387 | // Nothing lost: document + archive is exactly the original journal. |
| 5388 | let mut union = all_ids(&journal.entries); |
| 5389 | union.extend(all_ids(&archive)); |
| 5390 | assert_eq!(union, all_ids(&original.entries)); |
| 5391 | assert_eq!(archive.len() + journal.len(), original.len()); |
| 5392 | |
| 5393 | // The live journal drops exactly what the save archived. |
| 5394 | let mut live = original.clone(); |
| 5395 | let archived = take_archived_journal_ids("bounded-a"); |
| 5396 | assert_eq!(archived, all_ids(&archive)); |
| 5397 | live.remove_entries(&archived).unwrap(); |
| 5398 | assert_eq!(all_ids(&live.entries), all_ids(&journal.entries)); |
| 5399 | assert!(take_archived_journal_ids("bounded-a").is_empty()); |
| 5400 | |
| 5401 | // A plain save never prunes. |
| 5402 | let small = compacted_session("bounded-plain", tmp.path(), 460); |
| 5403 | let len = small.journal.as_ref().unwrap().len(); |
| 5404 | manager.save_session(&small).unwrap(); |
| 5405 | let reloaded = manager.load_session("bounded-plain").unwrap(); |
| 5406 | assert_eq!(reloaded.journal.unwrap().len(), len); |
| 5407 | } |
| 5408 | |
| 5409 | #[test] |
| 5410 | fn archive_failure_prunes_nothing() { |
| 5411 | let tmp = tempfile::tempdir().unwrap(); |
| 5412 | let sessions = tmp.path().join("sessions"); |
| 5413 | let manager = SessionManager::new(sessions.clone()).unwrap(); |
| 5414 | // A regular file where the archive directory belongs. |
| 5415 | fs::write(sessions.join(JOURNAL_ARCHIVE_DIR), b"not a dir").unwrap(); |
| 5416 | let session = compacted_session("bounded-b", tmp.path(), 460); |
| 5417 | let len = session.journal.as_ref().unwrap().len(); |
| 5418 | manager.save_session_bounded(session).unwrap(); |
| 5419 | let saved = manager.load_session("bounded-b").unwrap(); |
| 5420 | assert_eq!(saved.journal.unwrap().len(), len); |
| 5421 | assert!(take_archived_journal_ids("bounded-b").is_empty()); |
| 5422 | } |
| 5423 | |
| 5424 | #[cfg(unix)] |
| 5425 | #[test] |
| 5426 | fn archive_written_then_document_write_fails_loses_nothing() { |
| 5427 | use std::os::unix::fs::PermissionsExt as _; |
| 5428 | let tmp = tempfile::tempdir().unwrap(); |
| 5429 | let sessions = tmp.path().join("sessions"); |
| 5430 | let manager = SessionManager::new(sessions.clone()).unwrap(); |
| 5431 | let session = compacted_session("bounded-c", tmp.path(), 460); |
| 5432 | let original = session.journal.clone().unwrap(); |
| 5433 | manager.save_session(&session).unwrap(); |
| 5434 | fs::create_dir_all(sessions.join(JOURNAL_ARCHIVE_DIR)).unwrap(); |
| 5435 | // The archive can be appended, but the document's temp file cannot |
| 5436 | // be created: the crash window between the two writes. |
| 5437 | fs::set_permissions(&sessions, fs::Permissions::from_mode(0o500)).unwrap(); |
| 5438 | let failed = manager.save_session_bounded(session.clone()); |
| 5439 | fs::set_permissions(&sessions, fs::Permissions::from_mode(0o700)).unwrap(); |
| 5440 | assert!(failed.is_err(), "document write must fail in this setup"); |
| 5441 | let on_disk = manager.load_session("bounded-c").unwrap(); |
| 5442 | assert_eq!(on_disk.journal.as_ref().unwrap(), &original); |
| 5443 | assert!( |
| 5444 | !manager |
| 5445 | .load_journal_archive("bounded-c") |
| 5446 | .unwrap() |
| 5447 | .is_empty() |
| 5448 | ); |
| 5449 | assert!(take_archived_journal_ids("bounded-c").is_empty()); |
| 5450 | |
| 5451 | // Retrying archives the same entries again; the reader dedupes and |
| 5452 | // the union is still the whole journal. |
| 5453 | manager.save_session_bounded(session).unwrap(); |
| 5454 | let saved = manager.load_session("bounded-c").unwrap(); |
| 5455 | let archive = manager.load_journal_archive("bounded-c").unwrap(); |
| 5456 | let mut union = all_ids(&saved.journal.unwrap().entries); |
| 5457 | union.extend(all_ids(&archive)); |
| 5458 | assert_eq!(union, all_ids(&original.entries)); |
| 5459 | let _ = take_archived_journal_ids("bounded-c"); |
| 5460 | } |
| 5461 | |
| 5462 | #[test] |
| 5463 | fn branch_to_an_archived_entry_restores_its_chain() { |
| 5464 | let tmp = tempfile::tempdir().unwrap(); |
| 5465 | let manager = SessionManager::new(tmp.path().join("sessions")).unwrap(); |
| 5466 | let session = compacted_session("bounded-d", tmp.path(), 460); |
| 5467 | let original = session.journal.clone().unwrap(); |
| 5468 | manager.save_session_bounded(session).unwrap(); |
| 5469 | let archived = take_archived_journal_ids("bounded-d"); |
| 5470 | // The deepest archived entry of the oldest round. |
| 5471 | let target = original |
| 5472 | .entries |
| 5473 | .iter() |
| 5474 | .filter(|e| archived.contains(&e.id)) |
| 5475 | .find(|e| { |
| 5476 | e.kind |
| 5477 | .as_message() |
| 5478 | .is_some_and(|m| m.content == text("round 0 message 9").content) |
| 5479 | }) |
| 5480 | .unwrap() |
| 5481 | .id |
| 5482 | .clone(); |
| 5483 | |
| 5484 | let mut session = manager.load_session("bounded-d").unwrap(); |
| 5485 | assert!(session.journal_branch_to(&target).is_err()); |
| 5486 | assert!( |
| 5487 | manager |
| 5488 | .restore_archived_journal_chain(&mut session, &target) |
| 5489 | .unwrap() |
| 5490 | ); |
| 5491 | session.journal_branch_to(&target).unwrap(); |
| 5492 | manager.save_session(&session).unwrap(); |
| 5493 | let reloaded = manager.load_session("bounded-d").unwrap(); |
| 5494 | assert_eq!(reloaded.leaf_id.as_deref(), Some(target.as_str())); |
| 5495 | // Each round replaced the whole tail, so round 0 is its own root chain. |
| 5496 | let expected: Vec<Message> = (0..10) |
| 5497 | .map(|i| text(&format!("round 0 message {i}"))) |
| 5498 | .collect(); |
| 5499 | assert_eq!(reloaded.messages, expected); |
| 5500 | |
| 5501 | let mut missing = manager.load_session("bounded-d").unwrap(); |
| 5502 | assert!( |
| 5503 | !manager |
| 5504 | .restore_archived_journal_chain(&mut missing, "no-such-entry") |
| 5505 | .unwrap() |
| 5506 | ); |
| 5507 | } |
| 5508 | |
| 5509 | #[test] |
| 5510 | fn deleting_a_session_removes_its_journal_archive() { |
| 5511 | let tmp = tempfile::tempdir().unwrap(); |
| 5512 | let sessions = tmp.path().join("sessions"); |
| 5513 | let manager = SessionManager::new(sessions.clone()).unwrap(); |
| 5514 | manager |
| 5515 | .save_session_bounded(compacted_session("bounded-e", tmp.path(), 460)) |
| 5516 | .unwrap(); |
| 5517 | let _ = take_archived_journal_ids("bounded-e"); |
| 5518 | let path = journal_archive_path(&sessions, "bounded-e").unwrap(); |
| 5519 | assert!(path.exists()); |
| 5520 | manager.delete_session("bounded-e").unwrap(); |
| 5521 | assert!(!path.exists()); |
| 5522 | } |
| 5523 | |
| 5524 | #[test] |
| 5525 | fn metadata_larger_than_the_prefix_is_read_without_the_transcript() { |
| 5526 | let tmp = tempfile::tempdir().unwrap(); |
| 5527 | let manager = SessionManager::new(tmp.path().join("sessions")).unwrap(); |
| 5528 | let mut session = compacted_session("bounded-f", tmp.path(), 10); |
| 5529 | for i in 0..9_000 { |
| 5530 | session.metadata.cost.usage_source_fingerprints.insert( |
| 5531 | crate::cost_status::usage_source_fingerprint(&format!("call-{i}")), |
| 5532 | ); |
| 5533 | } |
| 5534 | let path = manager.save_session(&session).unwrap(); |
| 5535 | // Replace the transcript with 1 MB of non-JSON: listing must still |
| 5536 | // find the >64 KB metadata block (the grown read stops once it |
| 5537 | // closes; this checks correctness, not how many bytes were read). |
| 5538 | let bytes = fs::read(&path).unwrap(); |
| 5539 | let start = find_json_key(&bytes, b"\"metadata\"").unwrap(); |
| 5540 | let open = json_value_start(&bytes, start, 10, b'{').unwrap(); |
| 5541 | let end = json_object_end(&bytes, open).unwrap(); |
| 5542 | assert!( |
| 5543 | end > 64 * 1024, |
| 5544 | "metadata must exceed the old prefix: {end}" |
| 5545 | ); |
| 5546 | let mut damaged = bytes[..end].to_vec(); |
| 5547 | damaged.extend(std::iter::repeat_n(b'x', 1 << 20)); |
| 5548 | fs::write(&path, &damaged).unwrap(); |
| 5549 | let metadata = SessionManager::load_session_metadata(&path).unwrap(); |
| 5550 | assert_eq!(metadata.id, "bounded-f"); |
| 5551 | assert_eq!(metadata.cost.usage_source_fingerprints.len(), 9_000); |
| 5552 | } |
| 5553 | } |
| 5554 |