| 1 | //! Repair of the saved-session store (#6144). |
| 2 | //! |
| 3 | //! The session document is the one authority for a conversation. Two links |
| 4 | //! hang off it, each with exactly one home: document → Runtime store in |
| 5 | //! `metadata.runtime_store`, and Runtime thread → document in |
| 6 | //! `ThreadRecord.session_id` plus its checkpoint. A directory's *name* is not |
| 7 | //! ownership: a host store sits under whichever id the process booted with, |
| 8 | //! and every conversation that process saved binds it. |
| 9 | //! |
| 10 | //! The write paths now keep those links consistent (stores are retired where |
| 11 | //! they are abandoned, export is idempotent, deleting a document unbinds its |
| 12 | //! threads, checkpoints cover a prefix). What they cannot prevent — a crash |
| 13 | //! or kill between two writes, and everything earlier builds left behind — is |
| 14 | //! repaired here, on load: |
| 15 | //! |
| 16 | //! - **R1** an unreadable document is set aside (a newer-schema one is left in |
| 17 | //! place and reported); |
| 18 | //! - **R2/R5** a store no document binds, holding no work, is set aside while |
| 19 | //! its process-owner lock is held, so no process can open it mid-move; a |
| 20 | //! store another process holds is skipped and counted as in use; |
| 21 | //! - **R3** a store no document binds that holds threads is re-indexed: one |
| 22 | //! "Recovered:" document per thread, bound to the store, with the thread |
| 23 | //! bound back to it; |
| 24 | //! - **R4** a thread bound to a document that no longer exists is unbound |
| 25 | //! (the old binding goes in the receipts) and loads from its own turns; |
| 26 | //! - **R6** a directory with no document, no store and only artifacts or |
| 27 | //! approval receipts is set aside when no document, thread or derived |
| 28 | //! thread id names it and it has been untouched for a day. |
| 29 | //! |
| 30 | //! Repair only ever *moves* things, into `sessions/.set-aside/<run>/`, with a |
| 31 | //! `MANIFEST.jsonl` naming each original path; it never unlinks. Every action |
| 32 | //! appends a line to `sessions/.reconcile/receipts.jsonl` (a dropped thread |
| 33 | //! binding is recorded in its store's `session-unbind-receipts.jsonl`), and |
| 34 | //! the run's summary is kept in `sessions/.reconcile/last.json` for |
| 35 | //! `codewhale doctor`. |
| 36 | |
| 37 | use std::collections::HashSet; |
| 38 | use std::fs; |
| 39 | use std::io::{self, Write as _}; |
| 40 | use std::path::{Path, PathBuf}; |
| 41 | use std::time::{Duration, SystemTime}; |
| 42 | |
| 43 | use chrono::{DateTime, Utc}; |
| 44 | use serde::{Deserialize, Serialize}; |
| 45 | use serde_json::{Value, json}; |
| 46 | use sha2::{Digest, Sha256}; |
| 47 | |
| 48 | use crate::runtime_threads::{RuntimeStoreBinding, ThreadRecord}; |
| 49 | use crate::session_manager::{SavedSession, SessionManager, is_runtime_store_dir_name}; |
| 50 | |
| 51 | pub(crate) const SET_ASIDE_DIR: &str = ".set-aside"; |
| 52 | const RECONCILE_DIR: &str = ".reconcile"; |
| 53 | const RECONCILE_LOCK_FILE: &str = ".reconcile.lock"; |
| 54 | const RECEIPTS_FILE: &str = "receipts.jsonl"; |
| 55 | const LAST_RUN_FILE: &str = "last.json"; |
| 56 | const MANIFEST_FILE: &str = "MANIFEST.jsonl"; |
| 57 | const LATE_USAGE_DIR: &str = ".late-usage"; |
| 58 | const CHECKPOINTS_DIR: &str = "checkpoints"; |
| 59 | const SESSION_BOOT_OWNERS_FILE: &str = "session_boot_owners.json"; |
| 60 | /// Upper bound on actions per run; the rest waits for the next launch. |
| 61 | pub(crate) const DEFAULT_ACTION_LIMIT: usize = 500; |
| 62 | /// A document-less artifact directory has no lock to prove nobody is still |
| 63 | /// writing it (an older build's engine, say), so it must also be idle. |
| 64 | const ARTIFACT_DIR_IDLE: Duration = Duration::from_secs(24 * 60 * 60); |
| 65 | |
| 66 | /// What one reconcile run did (or, dry, would do). |
| 67 | #[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] |
| 68 | pub struct ReconcileSummary { |
| 69 | pub ran_at: Option<DateTime<Utc>>, |
| 70 | #[serde(default)] |
| 71 | pub dry_run: bool, |
| 72 | /// Another process was already reconciling; nothing was examined. |
| 73 | #[serde(default)] |
| 74 | pub skipped_concurrent: bool, |
| 75 | /// R1: unreadable documents set aside. |
| 76 | #[serde(default)] |
| 77 | pub documents_set_aside: usize, |
| 78 | /// Documents written by a newer build, left in place. |
| 79 | #[serde(default)] |
| 80 | pub documents_newer_schema: usize, |
| 81 | /// R2/R5: unbound, empty Runtime stores set aside. |
| 82 | #[serde(default)] |
| 83 | pub stores_set_aside: usize, |
| 84 | /// Unbound stores that hold work with no thread to re-index; kept. |
| 85 | #[serde(default)] |
| 86 | pub stores_kept_with_work: usize, |
| 87 | /// Unbound stores a live process holds; left for a later run. |
| 88 | #[serde(default)] |
| 89 | pub stores_in_use: usize, |
| 90 | /// R3: "Recovered:" documents written for threads in unbound stores. |
| 91 | #[serde(default)] |
| 92 | pub sessions_recovered: usize, |
| 93 | /// R4: threads whose document no longer exists, unbound. |
| 94 | #[serde(default)] |
| 95 | pub threads_unbound: usize, |
| 96 | /// R6: document-less artifact directories set aside. |
| 97 | #[serde(default)] |
| 98 | pub artifact_dirs_set_aside: usize, |
| 99 | /// Document-less directories kept because something names them, they |
| 100 | /// are recent, or they hold something this repair does not recognise. |
| 101 | #[serde(default)] |
| 102 | pub artifact_dirs_kept: usize, |
| 103 | /// The per-run action limit stopped this run; the next one continues. |
| 104 | #[serde(default)] |
| 105 | pub limit_reached: bool, |
| 106 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 107 | pub set_aside_path: Option<PathBuf>, |
| 108 | #[serde(default, skip_serializing_if = "Vec::is_empty")] |
| 109 | pub errors: Vec<String>, |
| 110 | } |
| 111 | |
| 112 | impl ReconcileSummary { |
| 113 | fn actions(&self) -> usize { |
| 114 | self.documents_set_aside |
| 115 | + self.stores_set_aside |
| 116 | + self.sessions_recovered |
| 117 | + self.threads_unbound |
| 118 | + self.artifact_dirs_set_aside |
| 119 | } |
| 120 | |
| 121 | /// Whether this run changed anything on disk. |
| 122 | #[must_use] |
| 123 | pub fn changed(&self) -> bool { |
| 124 | !self.dry_run && self.actions() > 0 |
| 125 | } |
| 126 | |
| 127 | /// One line for the first frame, only when something changed. |
| 128 | #[must_use] |
| 129 | pub fn notice(&self) -> Option<String> { |
| 130 | if !self.changed() { |
| 131 | return None; |
| 132 | } |
| 133 | let mut parts = Vec::new(); |
| 134 | if self.sessions_recovered > 0 { |
| 135 | parts.push(format!("{} recovered", self.sessions_recovered)); |
| 136 | } |
| 137 | if self.threads_unbound > 0 { |
| 138 | parts.push(format!("{} threads re-linked", self.threads_unbound)); |
| 139 | } |
| 140 | let set_aside = |
| 141 | self.stores_set_aside + self.artifact_dirs_set_aside + self.documents_set_aside; |
| 142 | if set_aside > 0 { |
| 143 | parts.push(format!("{set_aside} unused items set aside")); |
| 144 | } |
| 145 | let mut line = format!("Sessions repaired: {}", parts.join(", ")); |
| 146 | if let Some(path) = &self.set_aside_path { |
| 147 | line.push_str(&format!(" → {}", path.display())); |
| 148 | } |
| 149 | Some(line) |
| 150 | } |
| 151 | |
| 152 | /// The `codewhale doctor` detail for this run. |
| 153 | #[must_use] |
| 154 | pub fn doctor_detail(&self) -> String { |
| 155 | let when = self |
| 156 | .ran_at |
| 157 | .map_or_else(|| "never".to_string(), |at| at.to_rfc3339()); |
| 158 | let mut line = format!( |
| 159 | "last repair {when}: {} recovered, {} threads unbound, {} stores + {} artifact dirs + {} unreadable documents set aside, {} newer-schema documents left in place, {} unbound stores in use by another process, {} kept with work", |
| 160 | self.sessions_recovered, |
| 161 | self.threads_unbound, |
| 162 | self.stores_set_aside, |
| 163 | self.artifact_dirs_set_aside, |
| 164 | self.documents_set_aside, |
| 165 | self.documents_newer_schema, |
| 166 | self.stores_in_use, |
| 167 | self.stores_kept_with_work, |
| 168 | ); |
| 169 | if self.limit_reached { |
| 170 | line.push_str("; more remains for the next run"); |
| 171 | } |
| 172 | if let Some(path) = &self.set_aside_path { |
| 173 | line.push_str(&format!("; set aside under {}", path.display())); |
| 174 | } |
| 175 | if !self.errors.is_empty() { |
| 176 | line.push_str(&format!("; {} errors", self.errors.len())); |
| 177 | } |
| 178 | line |
| 179 | } |
| 180 | } |
| 181 | |
| 182 | #[derive(Debug, Clone)] |
| 183 | pub struct ReconcileOptions { |
| 184 | /// Examine and count, but move and write nothing. |
| 185 | pub dry_run: bool, |
| 186 | pub limit: usize, |
| 187 | /// A document being resumed right now; never touched. |
| 188 | pub skip_session: Option<String>, |
| 189 | /// How long a document-less artifact directory must sit untouched. |
| 190 | pub artifact_idle: Duration, |
| 191 | } |
| 192 | |
| 193 | impl Default for ReconcileOptions { |
| 194 | fn default() -> Self { |
| 195 | Self { |
| 196 | dry_run: false, |
| 197 | limit: DEFAULT_ACTION_LIMIT, |
| 198 | skip_session: None, |
| 199 | artifact_idle: ARTIFACT_DIR_IDLE, |
| 200 | } |
| 201 | } |
| 202 | } |
| 203 | |
| 204 | /// The summary of the last completed (non-dry) run, if any. |
| 205 | #[must_use] |
| 206 | pub fn last_run(sessions_dir: &Path) -> Option<ReconcileSummary> { |
| 207 | let raw = fs::read(sessions_dir.join(RECONCILE_DIR).join(LAST_RUN_FILE)).ok()?; |
| 208 | serde_json::from_slice(&raw).ok() |
| 209 | } |
| 210 | |
| 211 | static PENDING_NOTICE: std::sync::Mutex<Option<String>> = std::sync::Mutex::new(None); |
| 212 | |
| 213 | /// The first-frame notice a background run left, taken once. |
| 214 | pub(crate) fn take_pending_notice() -> Option<String> { |
| 215 | PENDING_NOTICE.lock().ok()?.take() |
| 216 | } |
| 217 | |
| 218 | /// Reconcile the default sessions store on a blocking worker. Called once a |
| 219 | /// host holds its own Runtime store, so that store and any store the launch |
| 220 | /// just opened read as in use rather than as candidates. |
| 221 | pub(crate) fn spawn_background_reconcile(skip_session: Option<String>) { |
| 222 | tokio::task::spawn_blocking(move || { |
| 223 | let Ok(manager) = SessionManager::default_location() else { |
| 224 | return; |
| 225 | }; |
| 226 | let options = ReconcileOptions { |
| 227 | skip_session, |
| 228 | ..ReconcileOptions::default() |
| 229 | }; |
| 230 | match reconcile(&manager, &options) { |
| 231 | Ok(summary) => { |
| 232 | if let Some(notice) = summary.notice() |
| 233 | && let Ok(mut pending) = PENDING_NOTICE.lock() |
| 234 | { |
| 235 | *pending = Some(notice); |
| 236 | } |
| 237 | } |
| 238 | Err(error) => tracing::warn!(%error, "session store repair did not run"), |
| 239 | } |
| 240 | }); |
| 241 | } |
| 242 | |
| 243 | /// Run one bounded reconcile pass over `sessions`' store. |
| 244 | pub fn reconcile( |
| 245 | sessions: &SessionManager, |
| 246 | options: &ReconcileOptions, |
| 247 | ) -> io::Result<ReconcileSummary> { |
| 248 | let sessions_dir = sessions.sessions_dir().to_path_buf(); |
| 249 | let mut summary = ReconcileSummary { |
| 250 | ran_at: Some(Utc::now()), |
| 251 | dry_run: options.dry_run, |
| 252 | ..ReconcileSummary::default() |
| 253 | }; |
| 254 | let lock_file = |
| 255 | crate::session_manager::open_private_lock_file(&sessions_dir.join(RECONCILE_LOCK_FILE))?; |
| 256 | if !crate::runtime_threads::try_lock_file_exclusive(&lock_file)? { |
| 257 | summary.skipped_concurrent = true; |
| 258 | return Ok(summary); |
| 259 | } |
| 260 | let mut run = Run { |
| 261 | sessions, |
| 262 | canonical_root: sessions_dir |
| 263 | .canonicalize() |
| 264 | .unwrap_or_else(|_| sessions_dir.clone()), |
| 265 | sessions_dir, |
| 266 | options, |
| 267 | set_aside: SetAside::default(), |
| 268 | summary, |
| 269 | }; |
| 270 | run.execute(); |
| 271 | let summary = run.summary; |
| 272 | if !options.dry_run { |
| 273 | let dir = sessions.sessions_dir().join(RECONCILE_DIR); |
| 274 | fs::create_dir_all(&dir)?; |
| 275 | let bytes = serde_json::to_vec_pretty(&summary).map_err(io::Error::other)?; |
| 276 | crate::utils::write_atomic(&dir.join(LAST_RUN_FILE), &bytes)?; |
| 277 | } |
| 278 | drop(lock_file); |
| 279 | Ok(summary) |
| 280 | } |
| 281 | |
| 282 | /// Receipts of thread bindings dropped in a store, kept in that store's root |
| 283 | /// (outside its work directories, so they never make it look busy). |
| 284 | pub(crate) const THREAD_UNBIND_RECEIPTS_FILE: &str = "session-unbind-receipts.jsonl"; |
| 285 | |
| 286 | /// Record the binding a thread is losing, next to the thread itself: the |
| 287 | /// store is the one place every writer that unbinds (the thread's own host, |
| 288 | /// the Runtime API, reconcile) already knows. |
| 289 | pub(crate) fn record_thread_unbound(store_dir: &Path, thread: &ThreadRecord, reason: &str) { |
| 290 | let receipt = json!({ |
| 291 | "action": "thread_unbound", |
| 292 | "thread_id": thread.id, |
| 293 | "session_id": thread.session_id, |
| 294 | "saved_session_checkpoint": thread.saved_session_checkpoint, |
| 295 | "reason": reason, |
| 296 | "at": Utc::now(), |
| 297 | }); |
| 298 | let result = fs::OpenOptions::new() |
| 299 | .create(true) |
| 300 | .append(true) |
| 301 | .open(store_dir.join(THREAD_UNBIND_RECEIPTS_FILE)) |
| 302 | .and_then(|mut file| { |
| 303 | let mut line = serde_json::to_vec(&receipt).map_err(io::Error::other)?; |
| 304 | line.push(b'\n'); |
| 305 | file.write_all(&line) |
| 306 | }); |
| 307 | if let Err(error) = result { |
| 308 | tracing::warn!(%error, thread_id = %thread.id, "thread unbind receipt was not written"); |
| 309 | } |
| 310 | } |
| 311 | |
| 312 | /// Set aside `store` when no document or checkpoint binds it, no process |
| 313 | /// holds it, and it holds no work. Returns where it went. Used where a store |
| 314 | /// is abandoned — a switch that rebinds its conversation elsewhere, a delete |
| 315 | /// — so the abandonment does not wait for the next reconcile. |
| 316 | pub(crate) fn retire_unbound_store( |
| 317 | sessions: &SessionManager, |
| 318 | store: &Path, |
| 319 | reason: &str, |
| 320 | ) -> Option<PathBuf> { |
| 321 | if !store.is_dir() { |
| 322 | return None; |
| 323 | } |
| 324 | let references = collect_document_references(sessions.sessions_dir(), false); |
| 325 | if references.binds(store) { |
| 326 | return None; |
| 327 | } |
| 328 | let binding = RuntimeStoreBinding::for_store_dir(store).ok()?; |
| 329 | let held = binding.try_hold().ok().flatten()?; |
| 330 | if !matches!(held.keep_reason(), Ok(None)) { |
| 331 | return None; |
| 332 | } |
| 333 | let canonical_root = sessions |
| 334 | .sessions_dir() |
| 335 | .canonicalize() |
| 336 | .unwrap_or_else(|_| sessions.sessions_dir().to_path_buf()); |
| 337 | let mut set_aside = SetAside::default(); |
| 338 | let original = held.binding().data_dir.clone(); |
| 339 | let target = set_aside |
| 340 | .target_for(sessions.sessions_dir(), &canonical_root, &original) |
| 341 | .ok()?; |
| 342 | if let Err(error) = held.move_to(&target) { |
| 343 | tracing::debug!(store = %original.display(), %error, "released store was not set aside"); |
| 344 | return None; |
| 345 | } |
| 346 | set_aside.manifest(&original, &target, reason, json!({"kind": "runtime_store"})); |
| 347 | append_receipt( |
| 348 | sessions.sessions_dir(), |
| 349 | json!({"action": "store_set_aside", "path": original, "moved_to": target, "reason": reason}), |
| 350 | ); |
| 351 | remove_if_empty(original.parent()); |
| 352 | Some(target) |
| 353 | } |
| 354 | |
| 355 | /// [`retire_unbound_store`] off the calling thread when a Tokio runtime is |
| 356 | /// running (a session switch runs on the UI runtime and must not wait on a |
| 357 | /// scan of every document), inline otherwise. |
| 358 | pub(crate) fn retire_unbound_store_in_background( |
| 359 | sessions: SessionManager, |
| 360 | store: PathBuf, |
| 361 | reason: &'static str, |
| 362 | ) { |
| 363 | match tokio::runtime::Handle::try_current() { |
| 364 | Ok(runtime) => { |
| 365 | runtime.spawn_blocking(move || retire_unbound_store(&sessions, &store, reason)); |
| 366 | } |
| 367 | Err(_) => { |
| 368 | retire_unbound_store(&sessions, &store, reason); |
| 369 | } |
| 370 | } |
| 371 | } |
| 372 | |
| 373 | fn append_receipt(sessions_dir: &Path, mut receipt: Value) { |
| 374 | if let Value::Object(map) = &mut receipt { |
| 375 | map.insert("at".into(), json!(Utc::now())); |
| 376 | } |
| 377 | let dir = sessions_dir.join(RECONCILE_DIR); |
| 378 | let result = fs::create_dir_all(&dir).and_then(|()| { |
| 379 | let mut file = fs::OpenOptions::new() |
| 380 | .create(true) |
| 381 | .append(true) |
| 382 | .open(dir.join(RECEIPTS_FILE))?; |
| 383 | let mut line = serde_json::to_vec(&receipt).map_err(io::Error::other)?; |
| 384 | line.push(b'\n'); |
| 385 | file.write_all(&line) |
| 386 | }); |
| 387 | if let Err(error) = result { |
| 388 | tracing::warn!(%error, "session repair receipt was not written"); |
| 389 | } |
| 390 | } |
| 391 | |
| 392 | fn remove_if_empty(dir: Option<&Path>) { |
| 393 | if let Some(dir) = dir |
| 394 | && fs::read_dir(dir).is_ok_and(|mut entries| entries.next().is_none()) |
| 395 | { |
| 396 | let _ = fs::remove_dir(dir); |
| 397 | } |
| 398 | } |
| 399 | |
| 400 | /// One `.set-aside/<run>/` directory, created on first use. |
| 401 | #[derive(Default)] |
| 402 | struct SetAside { |
| 403 | run_dir: Option<PathBuf>, |
| 404 | } |
| 405 | |
| 406 | impl SetAside { |
| 407 | fn target_for( |
| 408 | &mut self, |
| 409 | sessions_dir: &Path, |
| 410 | canonical_root: &Path, |
| 411 | original: &Path, |
| 412 | ) -> io::Result<PathBuf> { |
| 413 | let run_dir = match &self.run_dir { |
| 414 | Some(dir) => dir.clone(), |
| 415 | None => { |
| 416 | let stamp = Utc::now().format("%Y%m%dT%H%M%SZ"); |
| 417 | let suffix = &uuid::Uuid::new_v4().simple().to_string()[..6]; |
| 418 | let dir = sessions_dir |
| 419 | .join(SET_ASIDE_DIR) |
| 420 | .join(format!("{stamp}-{suffix}")); |
| 421 | fs::create_dir_all(&dir)?; |
| 422 | self.run_dir = Some(dir.clone()); |
| 423 | dir |
| 424 | } |
| 425 | }; |
| 426 | let relative = original |
| 427 | .strip_prefix(canonical_root) |
| 428 | .or_else(|_| original.strip_prefix(sessions_dir)) |
| 429 | .map(Path::to_path_buf) |
| 430 | .unwrap_or_else(|_| { |
| 431 | original |
| 432 | .file_name() |
| 433 | .map(PathBuf::from) |
| 434 | .unwrap_or_else(|| PathBuf::from("unnamed")) |
| 435 | }); |
| 436 | Ok(run_dir.join(relative)) |
| 437 | } |
| 438 | |
| 439 | fn manifest(&self, original: &Path, moved_to: &Path, reason: &str, detail: Value) { |
| 440 | let Some(run_dir) = &self.run_dir else { |
| 441 | return; |
| 442 | }; |
| 443 | let line = json!({ |
| 444 | "original": original, |
| 445 | "moved_to": moved_to, |
| 446 | "reason": reason, |
| 447 | "detail": detail, |
| 448 | "at": Utc::now(), |
| 449 | }); |
| 450 | let result = fs::OpenOptions::new() |
| 451 | .create(true) |
| 452 | .append(true) |
| 453 | .open(run_dir.join(MANIFEST_FILE)) |
| 454 | .and_then(|mut file| { |
| 455 | let mut bytes = serde_json::to_vec(&line).map_err(io::Error::other)?; |
| 456 | bytes.push(b'\n'); |
| 457 | file.write_all(&bytes) |
| 458 | }); |
| 459 | if let Err(error) = result { |
| 460 | tracing::warn!(%error, "set-aside manifest entry was not written"); |
| 461 | } |
| 462 | } |
| 463 | } |
| 464 | |
| 465 | /// Everything that names a session id or a store. |
| 466 | #[derive(Default)] |
| 467 | struct References { |
| 468 | /// Ids with a readable document or a crash-recovery checkpoint. |
| 469 | documents: HashSet<String>, |
| 470 | /// Stores some document or checkpoint binds (canonical where possible). |
| 471 | bound_stores: HashSet<PathBuf>, |
| 472 | /// Ids named by a Runtime thread, or derived from one. |
| 473 | thread_named: HashSet<String>, |
| 474 | /// Every uuid-shaped token found in any document's text. |
| 475 | document_tokens: HashSet<String>, |
| 476 | /// Document texts, kept only to answer non-uuid ids (rare). |
| 477 | document_texts: Vec<String>, |
| 478 | } |
| 479 | |
| 480 | impl References { |
| 481 | fn binds(&self, store: &Path) -> bool { |
| 482 | self.bound_stores.contains(&canonical(store)) || self.bound_stores.contains(store) |
| 483 | } |
| 484 | |
| 485 | fn names(&self, id: &str) -> bool { |
| 486 | self.documents.contains(id) |
| 487 | || self.thread_named.contains(id) |
| 488 | || self.document_tokens.contains(id) |
| 489 | || (!looks_like_uuid(id) && self.document_texts.iter().any(|text| text.contains(id))) |
| 490 | } |
| 491 | } |
| 492 | |
| 493 | fn canonical(path: &Path) -> PathBuf { |
| 494 | path.canonicalize().unwrap_or_else(|_| path.to_path_buf()) |
| 495 | } |
| 496 | |
| 497 | fn looks_like_uuid(id: &str) -> bool { |
| 498 | id.len() == 36 && uuid::Uuid::parse_str(id).is_ok() |
| 499 | } |
| 500 | |
| 501 | fn uuid_tokens(bytes: &[u8]) -> impl Iterator<Item = String> + '_ { |
| 502 | static UUID: std::sync::LazyLock<regex::bytes::Regex> = std::sync::LazyLock::new(|| { |
| 503 | regex::bytes::Regex::new( |
| 504 | "[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}", |
| 505 | ) |
| 506 | .expect("uuid pattern") |
| 507 | }); |
| 508 | UUID.find_iter(bytes) |
| 509 | .map(|found| String::from_utf8_lossy(found.as_bytes()).to_ascii_lowercase()) |
| 510 | } |
| 511 | |
| 512 | fn document_id(path: &Path) -> Option<String> { |
| 513 | if path.extension().is_none_or(|ext| ext != "json") |
| 514 | || path |
| 515 | .file_name() |
| 516 | .is_some_and(|name| name == SESSION_BOOT_OWNERS_FILE) |
| 517 | { |
| 518 | return None; |
| 519 | } |
| 520 | let id = path.file_stem()?.to_str()?; |
| 521 | crate::artifacts::is_valid_session_id(id).then(|| id.to_string()) |
| 522 | } |
| 523 | |
| 524 | /// Store bindings and ids of every readable document and checkpoint. With |
| 525 | /// `keep_text`, also collect what R6 needs to know which ids documents mention. |
| 526 | fn collect_document_references(sessions_dir: &Path, keep_text: bool) -> References { |
| 527 | let mut references = References::default(); |
| 528 | let scan = |dir: &Path, references: &mut References, keep_text: bool| { |
| 529 | let Ok(entries) = fs::read_dir(dir) else { |
| 530 | return; |
| 531 | }; |
| 532 | for entry in entries.flatten() { |
| 533 | let path = entry.path(); |
| 534 | let Some(id) = document_id(&path) else { |
| 535 | continue; |
| 536 | }; |
| 537 | let Ok(metadata) = SessionManager::load_session_metadata(&path) else { |
| 538 | continue; |
| 539 | }; |
| 540 | references.documents.insert(id); |
| 541 | if let Some(binding) = metadata.runtime_store { |
| 542 | references.bound_stores.insert(canonical(&binding.data_dir)); |
| 543 | references.bound_stores.insert(binding.data_dir); |
| 544 | } |
| 545 | if keep_text && let Ok(bytes) = fs::read(&path) { |
| 546 | references.document_tokens.extend(uuid_tokens(&bytes)); |
| 547 | references |
| 548 | .document_texts |
| 549 | .push(String::from_utf8_lossy(&bytes).into_owned()); |
| 550 | } |
| 551 | } |
| 552 | }; |
| 553 | scan(sessions_dir, &mut references, keep_text); |
| 554 | scan( |
| 555 | &sessions_dir.join(CHECKPOINTS_DIR), |
| 556 | &mut references, |
| 557 | keep_text, |
| 558 | ); |
| 559 | references |
| 560 | } |
| 561 | |
| 562 | /// One store directory under a session directory. |
| 563 | struct StoreDir { |
| 564 | owner_id: String, |
| 565 | path: PathBuf, |
| 566 | } |
| 567 | |
| 568 | struct Run<'a> { |
| 569 | sessions: &'a SessionManager, |
| 570 | sessions_dir: PathBuf, |
| 571 | canonical_root: PathBuf, |
| 572 | options: &'a ReconcileOptions, |
| 573 | set_aside: SetAside, |
| 574 | summary: ReconcileSummary, |
| 575 | } |
| 576 | |
| 577 | impl Run<'_> { |
| 578 | fn budget_left(&mut self) -> bool { |
| 579 | if self.summary.actions() >= self.options.limit { |
| 580 | self.summary.limit_reached = true; |
| 581 | return false; |
| 582 | } |
| 583 | true |
| 584 | } |
| 585 | |
| 586 | fn error(&mut self, context: &str, error: impl std::fmt::Display) { |
| 587 | self.summary.errors.push(format!("{context}: {error}")); |
| 588 | } |
| 589 | |
| 590 | fn receipt(&self, receipt: Value) { |
| 591 | if !self.options.dry_run { |
| 592 | append_receipt(&self.sessions_dir, receipt); |
| 593 | } |
| 594 | } |
| 595 | |
| 596 | fn execute(&mut self) { |
| 597 | self.repair_documents(); |
| 598 | let mut references = collect_document_references(&self.sessions_dir, true); |
| 599 | let stores = self.store_dirs(); |
| 600 | self.collect_thread_references(&stores, &mut references); |
| 601 | for store in &stores { |
| 602 | if !self.budget_left() { |
| 603 | return; |
| 604 | } |
| 605 | self.repair_store(store, &references); |
| 606 | } |
| 607 | self.repair_artifact_dirs(&references); |
| 608 | if let Some(dir) = &self.set_aside.run_dir { |
| 609 | self.summary.set_aside_path = Some(dir.clone()); |
| 610 | } |
| 611 | } |
| 612 | |
| 613 | /// R1: set aside documents that cannot be read at all. |
| 614 | fn repair_documents(&mut self) { |
| 615 | let Ok(entries) = fs::read_dir(&self.sessions_dir) else { |
| 616 | return; |
| 617 | }; |
| 618 | let mut paths: Vec<PathBuf> = entries.flatten().map(|entry| entry.path()).collect(); |
| 619 | paths.sort(); |
| 620 | for path in paths { |
| 621 | let Some(id) = document_id(&path) else { |
| 622 | continue; |
| 623 | }; |
| 624 | if self.options.skip_session.as_deref() == Some(id.as_str()) |
| 625 | || SessionManager::load_session_metadata(&path).is_ok() |
| 626 | { |
| 627 | continue; |
| 628 | } |
| 629 | match classify_unlisted_document(&path) { |
| 630 | DocumentState::Readable => continue, |
| 631 | DocumentState::NewerSchema => { |
| 632 | self.summary.documents_newer_schema += 1; |
| 633 | continue; |
| 634 | } |
| 635 | DocumentState::Unreadable(reason) => { |
| 636 | if !self.budget_left() { |
| 637 | return; |
| 638 | } |
| 639 | self.set_aside_path(&path, &format!("unreadable document: {reason}"), true); |
| 640 | self.summary.documents_set_aside += 1; |
| 641 | } |
| 642 | } |
| 643 | } |
| 644 | } |
| 645 | |
| 646 | fn store_dirs(&self) -> Vec<StoreDir> { |
| 647 | let mut out = Vec::new(); |
| 648 | let Ok(entries) = fs::read_dir(&self.sessions_dir) else { |
| 649 | return out; |
| 650 | }; |
| 651 | for entry in entries.flatten() { |
| 652 | let owner = entry.file_name(); |
| 653 | let Some(owner_id) = owner.to_str() else { |
| 654 | continue; |
| 655 | }; |
| 656 | if !crate::artifacts::is_valid_session_id(owner_id) || !entry.path().is_dir() { |
| 657 | continue; |
| 658 | } |
| 659 | let Ok(children) = fs::read_dir(entry.path()) else { |
| 660 | continue; |
| 661 | }; |
| 662 | for child in children.flatten() { |
| 663 | if is_runtime_store_dir_name(&child.file_name()) && child.path().is_dir() { |
| 664 | out.push(StoreDir { |
| 665 | owner_id: owner_id.to_string(), |
| 666 | path: child.path(), |
| 667 | }); |
| 668 | } |
| 669 | } |
| 670 | } |
| 671 | out.sort_by(|left, right| left.path.cmp(&right.path)); |
| 672 | out |
| 673 | } |
| 674 | |
| 675 | /// Every thread in every store (read-only; no locks) names its bound |
| 676 | /// document and the id its engine writes under. Plus the host stores |
| 677 | /// outside the sessions directory (`codewhale serve`, task runtimes). |
| 678 | fn collect_thread_references(&self, stores: &[StoreDir], references: &mut References) { |
| 679 | let mut roots: Vec<PathBuf> = stores.iter().map(|store| store.path.clone()).collect(); |
| 680 | let tasks = crate::task_manager::default_tasks_dir().join("runtime"); |
| 681 | roots.push(tasks.clone()); |
| 682 | if let Ok(entries) = fs::read_dir(&tasks) { |
| 683 | roots.extend(entries.flatten().map(|entry| entry.path())); |
| 684 | } |
| 685 | for root in roots { |
| 686 | let Ok(entries) = fs::read_dir(root.join("threads")) else { |
| 687 | continue; |
| 688 | }; |
| 689 | for entry in entries.flatten() { |
| 690 | let path = entry.path(); |
| 691 | if path.extension().is_none_or(|ext| ext != "json") { |
| 692 | continue; |
| 693 | } |
| 694 | if let Some(thread_id) = path.file_stem().and_then(|stem| stem.to_str()) { |
| 695 | references |
| 696 | .thread_named |
| 697 | .insert(crate::runtime_threads::thread_session_id(thread_id)); |
| 698 | } |
| 699 | if let Ok(raw) = fs::read(&path) |
| 700 | && let Ok(value) = serde_json::from_slice::<Value>(&raw) |
| 701 | && let Some(session_id) = value.get("session_id").and_then(Value::as_str) |
| 702 | { |
| 703 | references.thread_named.insert(session_id.to_string()); |
| 704 | } |
| 705 | } |
| 706 | } |
| 707 | } |
| 708 | |
| 709 | fn repair_store(&mut self, store: &StoreDir, references: &References) { |
| 710 | let binding = match RuntimeStoreBinding::for_store_dir(&store.path) { |
| 711 | Ok(binding) => binding, |
| 712 | Err(error) => return self.error(&format!("store {}", store.path.display()), error), |
| 713 | }; |
| 714 | let held = match binding.try_hold() { |
| 715 | Ok(Some(held)) => held, |
| 716 | Ok(None) => { |
| 717 | if !references.binds(&store.path) { |
| 718 | self.summary.stores_in_use += 1; |
| 719 | } |
| 720 | return; |
| 721 | } |
| 722 | Err(error) => return self.error(&format!("store {}", store.path.display()), error), |
| 723 | }; |
| 724 | if self |
| 725 | .options |
| 726 | .skip_session |
| 727 | .as_deref() |
| 728 | .is_some_and(|skip| skip == store.owner_id) |
| 729 | { |
| 730 | return; |
| 731 | } |
| 732 | let exists = |id: &str| references.documents.contains(id); |
| 733 | if references.binds(&store.path) { |
| 734 | // R4 for a bound store nobody holds. |
| 735 | if self.options.dry_run { |
| 736 | return; |
| 737 | } |
| 738 | match held.unbind_threads_without_documents(exists) { |
| 739 | Ok(count) => self.summary.threads_unbound += count, |
| 740 | Err(error) => self.error(&format!("threads in {}", store.path.display()), error), |
| 741 | } |
| 742 | return; |
| 743 | } |
| 744 | match held.keep_reason() { |
| 745 | Ok(None) => { |
| 746 | let tombstoned = self |
| 747 | .sessions_dir |
| 748 | .join(LATE_USAGE_DIR) |
| 749 | .join(format!("{}.deleted", store.owner_id)) |
| 750 | .exists(); |
| 751 | let reason = if tombstoned { |
| 752 | "empty Runtime store of a deleted session" |
| 753 | } else { |
| 754 | "empty Runtime store no document binds" |
| 755 | }; |
| 756 | self.summary.stores_set_aside += 1; |
| 757 | if self.options.dry_run { |
| 758 | return; |
| 759 | } |
| 760 | let target = match self.set_aside.target_for( |
| 761 | &self.sessions_dir, |
| 762 | &self.canonical_root, |
| 763 | &held.binding().data_dir.clone(), |
| 764 | ) { |
| 765 | Ok(target) => target, |
| 766 | Err(error) => { |
| 767 | self.summary.stores_set_aside -= 1; |
| 768 | return self.error("set-aside directory", error); |
| 769 | } |
| 770 | }; |
| 771 | let original = held.binding().data_dir.clone(); |
| 772 | if let Err(error) = held.move_to(&target) { |
| 773 | self.summary.stores_set_aside -= 1; |
| 774 | return self.error(&format!("store {}", original.display()), error); |
| 775 | } |
| 776 | self.set_aside.manifest( |
| 777 | &original, |
| 778 | &target, |
| 779 | reason, |
| 780 | json!({"kind": "runtime_store"}), |
| 781 | ); |
| 782 | self.receipt(json!({"action": "store_set_aside", "path": original, "moved_to": target, "reason": reason})); |
| 783 | remove_if_empty(store.path.parent()); |
| 784 | } |
| 785 | Ok(Some(_)) => self.recover_store(held, &store.path, exists), |
| 786 | Err(error) => self.error(&format!("store {}", store.path.display()), error), |
| 787 | } |
| 788 | } |
| 789 | |
| 790 | /// R3: give each conversation in an unbound store a document of its own. |
| 791 | fn recover_store( |
| 792 | &mut self, |
| 793 | held: crate::runtime_threads::HeldRuntimeStore, |
| 794 | path: &Path, |
| 795 | exists: impl Fn(&str) -> bool, |
| 796 | ) { |
| 797 | // R4 applies here too. Deleting a document from the session picker |
| 798 | // does not unbind the threads naming it, so an unbound store can hold |
| 799 | // threads bound to a document that is gone. Left bound, they are |
| 800 | // neither loadable through that document nor recoverable below, and |
| 801 | // the store is kept forever. Unbinding them (with a receipt) makes |
| 802 | // them recoverable in this same pass. |
| 803 | if !self.options.dry_run { |
| 804 | match held.unbind_threads_without_documents(&exists) { |
| 805 | Ok(count) => self.summary.threads_unbound += count, |
| 806 | Err(error) => return self.error(&format!("threads in {}", path.display()), error), |
| 807 | } |
| 808 | } |
| 809 | let threads = match held.recoverable_threads() { |
| 810 | Ok(threads) => threads, |
| 811 | Err(error) => return self.error(&format!("threads in {}", path.display()), error), |
| 812 | }; |
| 813 | if threads.is_empty() { |
| 814 | self.summary.stores_kept_with_work += 1; |
| 815 | return; |
| 816 | } |
| 817 | for recovered in threads { |
| 818 | if !self.budget_left() { |
| 819 | return; |
| 820 | } |
| 821 | self.summary.sessions_recovered += 1; |
| 822 | if self.options.dry_run { |
| 823 | continue; |
| 824 | } |
| 825 | let session = match self.recovered_document(&held, &recovered, &exists) { |
| 826 | Ok(session) => session, |
| 827 | Err(error) => { |
| 828 | self.summary.sessions_recovered -= 1; |
| 829 | self.error(&format!("recover thread {}", recovered.thread.id), error); |
| 830 | continue; |
| 831 | } |
| 832 | }; |
| 833 | if let Err(error) = held.bind_recovered_thread(&recovered.thread.id, &session) { |
| 834 | self.error(&format!("bind thread {}", recovered.thread.id), error); |
| 835 | continue; |
| 836 | } |
| 837 | self.receipt(json!({ |
| 838 | "action": "session_recovered", |
| 839 | "store": path, |
| 840 | "thread_id": recovered.thread.id, |
| 841 | "session_id": session.metadata.id, |
| 842 | "title": session.metadata.title, |
| 843 | })); |
| 844 | } |
| 845 | } |
| 846 | |
| 847 | fn recovered_document( |
| 848 | &self, |
| 849 | held: &crate::runtime_threads::HeldRuntimeStore, |
| 850 | recovered: &crate::runtime_threads::RecoverableThread, |
| 851 | exists: &impl Fn(&str) -> bool, |
| 852 | ) -> anyhow::Result<SavedSession> { |
| 853 | let thread = &recovered.thread; |
| 854 | let id = crate::runtime_threads::thread_session_id(&thread.id); |
| 855 | // A run that stopped between saving the document and binding the |
| 856 | // thread left the document behind at this same id; reuse it. |
| 857 | if exists(&id) || self.sessions.session_document_exists(&id) { |
| 858 | return Ok(self.sessions.load_session(&id)?); |
| 859 | } |
| 860 | let mut session = crate::session_manager::create_saved_session_with_id_and_mode( |
| 861 | id, |
| 862 | &recovered.messages, |
| 863 | &thread.model, |
| 864 | &thread.workspace, |
| 865 | 0, |
| 866 | thread |
| 867 | .system_prompt |
| 868 | .clone() |
| 869 | .map(codewhale_models::SystemPrompt::Text) |
| 870 | .as_ref(), |
| 871 | Some(thread.mode.as_str()), |
| 872 | ); |
| 873 | if let Some(provider) = thread.model_provider.as_deref() { |
| 874 | session |
| 875 | .metadata |
| 876 | .set_model_provider_route(provider, thread.model_provider_id.as_deref()); |
| 877 | } |
| 878 | let label = thread |
| 879 | .title |
| 880 | .clone() |
| 881 | .filter(|title| !title.trim().is_empty()) |
| 882 | .unwrap_or_else(|| session.metadata.title.clone()); |
| 883 | session.metadata.title = recovered_title(&label); |
| 884 | session.metadata.runtime_store = Some(held.binding().clone()); |
| 885 | self.sessions.save_session(&session)?; |
| 886 | Ok(session) |
| 887 | } |
| 888 | |
| 889 | /// R6: directories with no document and no store that hold only |
| 890 | /// artifacts or approval receipts nothing names. |
| 891 | fn repair_artifact_dirs(&mut self, references: &References) { |
| 892 | let Ok(entries) = fs::read_dir(&self.sessions_dir) else { |
| 893 | return; |
| 894 | }; |
| 895 | let mut dirs: Vec<(String, PathBuf)> = entries |
| 896 | .flatten() |
| 897 | .filter_map(|entry| { |
| 898 | let name = entry.file_name().to_str()?.to_string(); |
| 899 | (crate::artifacts::is_valid_session_id(&name) && entry.path().is_dir()) |
| 900 | .then(|| (name, entry.path())) |
| 901 | }) |
| 902 | .collect(); |
| 903 | dirs.sort(); |
| 904 | let idle_before = SystemTime::now() |
| 905 | .checked_sub(self.options.artifact_idle) |
| 906 | .unwrap_or(SystemTime::UNIX_EPOCH); |
| 907 | for (id, path) in dirs { |
| 908 | if references.documents.contains(&id) |
| 909 | || self.options.skip_session.as_deref() == Some(id.as_str()) |
| 910 | || !holds_only_artifacts(&path) |
| 911 | { |
| 912 | continue; |
| 913 | } |
| 914 | if references.names(&id.to_ascii_lowercase()) |
| 915 | || references.names(&id) |
| 916 | || newest_mtime(&path).is_none_or(|mtime| mtime > idle_before) |
| 917 | { |
| 918 | self.summary.artifact_dirs_kept += 1; |
| 919 | continue; |
| 920 | } |
| 921 | if !self.budget_left() { |
| 922 | return; |
| 923 | } |
| 924 | self.set_aside_path( |
| 925 | &path, |
| 926 | "artifacts of a conversation no document or thread names", |
| 927 | false, |
| 928 | ); |
| 929 | self.summary.artifact_dirs_set_aside += 1; |
| 930 | } |
| 931 | } |
| 932 | |
| 933 | fn set_aside_path(&mut self, path: &Path, reason: &str, hash: bool) { |
| 934 | if self.options.dry_run { |
| 935 | return; |
| 936 | } |
| 937 | let detail = if hash { |
| 938 | fs::read(path).map_or( |
| 939 | Value::Null, |
| 940 | |bytes| json!({"sha256": file_sha256(&bytes), "bytes": bytes.len()}), |
| 941 | ) |
| 942 | } else { |
| 943 | json!({"kind": "directory"}) |
| 944 | }; |
| 945 | let target = match self |
| 946 | .set_aside |
| 947 | .target_for(&self.sessions_dir, &self.canonical_root, path) |
| 948 | { |
| 949 | Ok(target) => target, |
| 950 | Err(error) => return self.error("set-aside directory", error), |
| 951 | }; |
| 952 | if let Some(parent) = target.parent() |
| 953 | && let Err(error) = fs::create_dir_all(parent) |
| 954 | { |
| 955 | return self.error("set-aside directory", error); |
| 956 | } |
| 957 | if let Err(error) = fs::rename(path, &target) { |
| 958 | return self.error(&format!("set aside {}", path.display()), error); |
| 959 | } |
| 960 | self.set_aside.manifest(path, &target, reason, detail); |
| 961 | self.receipt( |
| 962 | json!({"action": "set_aside", "path": path, "moved_to": target, "reason": reason}), |
| 963 | ); |
| 964 | } |
| 965 | } |
| 966 | |
| 967 | enum DocumentState { |
| 968 | Readable, |
| 969 | NewerSchema, |
| 970 | Unreadable(String), |
| 971 | } |
| 972 | |
| 973 | /// A document `list_sessions` could not list: newer, readable after all, or |
| 974 | /// genuinely unreadable. |
| 975 | fn classify_unlisted_document(path: &Path) -> DocumentState { |
| 976 | let bytes = match fs::read(path) { |
| 977 | Ok(bytes) => bytes, |
| 978 | Err(error) => return DocumentState::Unreadable(error.to_string()), |
| 979 | }; |
| 980 | if let Ok(value) = serde_json::from_slice::<Value>(&bytes) |
| 981 | && value |
| 982 | .get("schema_version") |
| 983 | .and_then(Value::as_u64) |
| 984 | .is_some_and(|version| { |
| 985 | version > u64::from(crate::session_manager::CURRENT_SESSION_SCHEMA_VERSION) |
| 986 | }) |
| 987 | { |
| 988 | return DocumentState::NewerSchema; |
| 989 | } |
| 990 | match serde_json::from_slice::<SavedSession>(&bytes) { |
| 991 | Ok(_) => DocumentState::Readable, |
| 992 | Err(error) => DocumentState::Unreadable(error.to_string()), |
| 993 | } |
| 994 | } |
| 995 | |
| 996 | fn holds_only_artifacts(dir: &Path) -> bool { |
| 997 | let Ok(entries) = fs::read_dir(dir) else { |
| 998 | return false; |
| 999 | }; |
| 1000 | let mut any = false; |
| 1001 | for entry in entries.flatten() { |
| 1002 | any = true; |
| 1003 | let name = entry.file_name(); |
| 1004 | let Some(name) = name.to_str() else { |
| 1005 | return false; |
| 1006 | }; |
| 1007 | if !matches!( |
| 1008 | name, |
| 1009 | "artifacts" | "approval_receipts.jsonl" | "approval_receipts.lock" |
| 1010 | ) { |
| 1011 | return false; |
| 1012 | } |
| 1013 | } |
| 1014 | any |
| 1015 | } |
| 1016 | |
| 1017 | fn newest_mtime(path: &Path) -> Option<SystemTime> { |
| 1018 | let metadata = fs::symlink_metadata(path).ok()?; |
| 1019 | let mut newest = metadata.modified().ok()?; |
| 1020 | if metadata.is_dir() { |
| 1021 | for entry in fs::read_dir(path).ok()?.flatten() { |
| 1022 | if let Some(child) = newest_mtime(&entry.path()) { |
| 1023 | newest = newest.max(child); |
| 1024 | } |
| 1025 | } |
| 1026 | } |
| 1027 | Some(newest) |
| 1028 | } |
| 1029 | |
| 1030 | fn recovered_title(label: &str) -> String { |
| 1031 | let title = |
| 1032 | crate::session_manager::sanitize_session_title(&format!("Recovered: {}", label.trim())); |
| 1033 | title |
| 1034 | .chars() |
| 1035 | .take(crate::session_manager::MAX_SESSION_TITLE_CHARS) |
| 1036 | .collect() |
| 1037 | } |
| 1038 | |
| 1039 | /// SHA-256 of a file's bytes, for the manifest. |
| 1040 | fn file_sha256(bytes: &[u8]) -> String { |
| 1041 | Sha256::digest(bytes) |
| 1042 | .iter() |
| 1043 | .map(|byte| format!("{byte:02x}")) |
| 1044 | .collect() |
| 1045 | } |
| 1046 | |
| 1047 | #[cfg(test)] |
| 1048 | #[path = "session_reconcile/tests.rs"] |
| 1049 | mod tests; |
| 1050 |