返回 CodeWhale
session_reconcile.rs
根目录 / crates / tui / src / session_reconcile.rs
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
1050 lines RUST