返回 CodeWhale
tests.rs
1 use super::*;
2 use crate::config::Config;
3 use crate::runtime_threads::{
4 CreateThreadRequest, RuntimeProcessOwnerLock, RuntimeThreadManager, RuntimeThreadManagerConfig,
5 RuntimeThreadStore,
6 };
7 use crate::session_manager::create_saved_session_with_id_and_mode;
8 use crate::test_support::{EnvVarGuard, lock_test_env};
9 use codewhale_models::{ContentBlock, Message, Role};
10
11 fn text(role: Role, text: &str) -> Message {
12 Message {
13 role,
14 content: vec![ContentBlock::Text {
15 text: text.to_string(),
16 cache_control: None,
17 }],
18 }
19 }
20
21 fn fixture_config() -> Config {
22 let mut config = Config::default().with_legacy_root(
23 Some("local-reconcile-fixture".into()),
24 Some("http://127.0.0.1:1/v1".into()),
25 );
26 config.set_feature("mcp", false).unwrap();
27 config.set_feature("subagents", false).unwrap();
28 config
29 }
30
31 /// An isolated Codewhale home whose sessions directory is the configured
32 /// one, so store confinement holds for fixtures created under it.
33 struct Fixture {
34 sessions: SessionManager,
35 home: tempfile::TempDir,
36 _runtime: EnvVarGuard,
37 _legacy_runtime: EnvVarGuard,
38 _tasks: EnvVarGuard,
39 _home: EnvVarGuard,
40 }
41
42 impl Fixture {
43 fn new() -> Self {
44 let home = tempfile::tempdir().expect("home");
45 let guard = EnvVarGuard::set("CODEWHALE_HOME", home.path());
46 let runtime = EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
47 let legacy_runtime = EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
48 let tasks = EnvVarGuard::remove("CODEWHALE_TASKS_DIR");
49 let sessions = SessionManager::default_location().expect("sessions");
50 Self {
51 sessions,
52 home,
53 _runtime: runtime,
54 _legacy_runtime: legacy_runtime,
55 _tasks: tasks,
56 _home: guard,
57 }
58 }
59
60 fn dir(&self) -> PathBuf {
61 self.sessions.sessions_dir().to_path_buf()
62 }
63
64 /// An opened, empty store at `sessions/<owner>/runtime`.
65 fn empty_store(&self, owner: &str) -> PathBuf {
66 let path = self.dir().join(owner).join("runtime");
67 RuntimeThreadStore::open(path.clone()).expect("open store");
68 path
69 }
70
71 fn document(&self, id: &str, store: Option<&Path>) -> SavedSession {
72 let mut session = create_saved_session_with_id_and_mode(
73 id.to_string(),
74 &[text(Role::User, "hello"), text(Role::Assistant, "hi")],
75 "deepseek-v4-pro",
76 self.home.path(),
77 0,
78 None,
79 None,
80 );
81 session.metadata.runtime_store =
82 store.map(|store| RuntimeStoreBinding::for_store_dir(store).expect("binding"));
83 self.sessions.save_session(&session).expect("save");
84 session
85 }
86
87 /// A store holding one task-bound thread with a two-message turn, and no
88 /// document bound to it. Returns the store and the thread id.
89 async fn store_with_thread(&self, owner: &str) -> (PathBuf, String) {
90 let path = self.dir().join(owner).join("runtime");
91 let manager = RuntimeThreadManager::open(
92 fixture_config(),
93 self.home.path().to_path_buf(),
94 RuntimeThreadManagerConfig {
95 data_dir: path.clone(),
96 task_data_dir: self.home.path().join("tasks"),
97 sessions_dir: None,
98 max_active_threads: 2,
99 },
100 )
101 .expect("manager");
102 let thread = manager
103 .create_thread(CreateThreadRequest {
104 task_id: Some("task_orphaned".into()),
105 workspace: Some(self.home.path().to_path_buf()),
106 ..CreateThreadRequest::default()
107 })
108 .await
109 .expect("thread");
110 manager
111 .seed_thread_from_messages(
112 &thread.id,
113 &[
114 text(Role::User, "rebuild the parser"),
115 text(Role::Assistant, "parser rebuilt"),
116 ],
117 )
118 .await
119 .expect("seed");
120 drop(manager);
121 (path, thread.id)
122 }
123
124 fn run(&self) -> ReconcileSummary {
125 self.run_with(ReconcileOptions {
126 artifact_idle: Duration::ZERO,
127 ..ReconcileOptions::default()
128 })
129 }
130
131 fn run_with(&self, options: ReconcileOptions) -> ReconcileSummary {
132 reconcile(&self.sessions, &options).expect("reconcile")
133 }
134
135 fn set_aside_root(&self) -> PathBuf {
136 self.dir().join(SET_ASIDE_DIR)
137 }
138
139 fn receipts(&self) -> String {
140 fs::read_to_string(self.dir().join(RECONCILE_DIR).join(RECEIPTS_FILE)).unwrap_or_default()
141 }
142 }
143
144 fn manifests(root: &Path) -> String {
145 let mut out = String::new();
146 if let Ok(runs) = fs::read_dir(root) {
147 for run in runs.flatten() {
148 out.push_str(&fs::read_to_string(run.path().join(MANIFEST_FILE)).unwrap_or_default());
149 }
150 }
151 out
152 }
153
154 /// The measured shape (#6144): documents are bound to a store under
155 /// *another* id's directory, and that directory has no document of its own.
156 /// Only the store nothing binds may move.
157 /// The reconcile lock sidecar is owner-only, even when an older release left
158 /// it world-readable.
159 #[cfg(unix)]
160 #[test]
161 fn reconcile_lock_file_is_owner_only() {
162 use std::os::unix::fs::PermissionsExt as _;
163 let _lock = lock_test_env();
164 let fixture = Fixture::new();
165 let lock = fixture.dir().join(RECONCILE_LOCK_FILE);
166 let mode = || fs::metadata(&lock).unwrap().permissions().mode() & 0o777;
167
168 fixture.run();
169 assert_eq!(mode(), 0o600);
170
171 fs::set_permissions(&lock, fs::Permissions::from_mode(0o644)).unwrap();
172 fixture.run();
173 assert_eq!(mode(), 0o600);
174 }
175
176 #[test]
177 fn reconcile_sets_aside_only_stores_no_document_binds() {
178 let _env = lock_test_env();
179 let fx = Fixture::new();
180 let host_store = fx.empty_store("host-anchor");
181 fx.document("conversation-a", Some(&host_store));
182 fx.document("conversation-b", Some(&host_store));
183 let orphan = fx.empty_store("crashed-anchor");
184
185 let summary = fx.run();
186
187 assert_eq!(summary.stores_set_aside, 1, "{summary:?}");
188 assert!(host_store.is_dir(), "a bound store must never move");
189 assert!(!orphan.exists());
190 assert!(
191 !fx.dir().join("crashed-anchor").exists(),
192 "the emptied directory goes"
193 );
194 let manifest = manifests(&fx.set_aside_root());
195 assert!(manifest.contains("crashed-anchor"), "{manifest}");
196 assert!(fx.receipts().contains("store_set_aside"));
197 assert_eq!(summary.notice().is_some(), summary.changed());
198 assert_eq!(last_run(&fx.dir()), Some(summary));
199
200 // Idempotent: nothing left to do.
201 let again = fx.run();
202 assert!(!again.changed(), "{again:?}");
203 }
204
205 /// A store whose owner lock is held is in use by a live process: skipped and
206 /// reported, never moved. Holding that lock through the move is also what
207 /// makes a concurrent opener fail its lock instead of racing the move.
208 #[test]
209 fn reconcile_skips_a_store_a_live_process_holds() {
210 let _env = lock_test_env();
211 let fx = Fixture::new();
212 let held = fx.empty_store("live-host");
213 let lock = RuntimeProcessOwnerLock::acquire(&held).expect("hold");
214
215 let summary = fx.run();
216 assert_eq!(summary.stores_set_aside, 0);
217 assert_eq!(summary.stores_in_use, 1);
218 assert!(held.is_dir());
219
220 drop(lock);
221 assert_eq!(fx.run().stores_set_aside, 1);
222 }
223
224 /// R3: a store no document binds but that holds a conversation is
225 /// re-indexed, not set aside, and a second run is a no-op.
226 #[tokio::test]
227 async fn reconcile_reindexes_threads_in_an_unbound_store() {
228 let _env = lock_test_env();
229 let fx = Fixture::new();
230 let (store, thread_id) = fx.store_with_thread("worker-anchor").await;
231
232 let summary = fx.run();
233 assert_eq!(summary.sessions_recovered, 1, "{summary:?}");
234 assert_eq!(summary.stores_set_aside, 0);
235 assert!(store.is_dir());
236
237 let id = crate::runtime_threads::thread_session_id(&thread_id);
238 let recovered = fx.sessions.load_session(&id).expect("recovered document");
239 assert!(recovered.metadata.title.starts_with("Recovered: "));
240 assert_eq!(
241 recovered.messages,
242 vec![
243 text(Role::User, "rebuild the parser"),
244 text(Role::Assistant, "parser rebuilt")
245 ]
246 );
247 let binding = recovered
248 .metadata
249 .runtime_store
250 .expect("bound to its store");
251 assert_eq!(binding.data_dir, store.canonicalize().unwrap());
252 let thread = RuntimeThreadStore::open(store.clone())
253 .unwrap()
254 .load_thread(&thread_id)
255 .unwrap();
256 assert_eq!(thread.session_id.as_deref(), Some(id.as_str()));
257 assert_eq!(
258 thread
259 .saved_session_checkpoint
260 .as_ref()
261 .and_then(|checkpoint| checkpoint.messages_len),
262 Some(2)
263 );
264
265 let again = fx.run();
266 assert_eq!(again.sessions_recovered, 0, "{again:?}");
267 assert_eq!(again.stores_set_aside, 0);
268 }
269
270 /// R4: a thread in a bound store whose document is gone is unbound, with
271 /// the old binding in the receipts.
272 #[tokio::test]
273 async fn reconcile_unbinds_threads_whose_document_is_gone() {
274 let _env = lock_test_env();
275 let fx = Fixture::new();
276 let (store, thread_id) = fx.store_with_thread("host").await;
277 fx.document("keeper", Some(&store));
278 {
279 let opened = RuntimeThreadStore::open(store.clone()).unwrap();
280 let mut thread = opened.load_thread(&thread_id).unwrap();
281 thread.session_id = Some("deleted-document".into());
282 opened.save_thread(&thread).unwrap();
283 }
284
285 let summary = fx.run();
286 assert_eq!(summary.threads_unbound, 1, "{summary:?}");
287 let thread = RuntimeThreadStore::open(store.clone())
288 .unwrap()
289 .load_thread(&thread_id)
290 .unwrap();
291 assert_eq!(thread.session_id, None);
292 let receipts = fs::read_to_string(store.join(THREAD_UNBIND_RECEIPTS_FILE)).unwrap();
293 assert!(receipts.contains("thread_unbound") && receipts.contains("deleted-document"));
294 }
295
296 /// R4 then R3 in a store no document binds: deleting a document from the
297 /// session picker leaves the threads naming it bound. They are unbound (with
298 /// a receipt) and then recovered, instead of keeping the store forever.
299 #[tokio::test]
300 async fn reconcile_recovers_threads_in_an_unbound_store_bound_to_a_missing_document() {
301 let _env = lock_test_env();
302 let fx = Fixture::new();
303 let (store, thread_id) = fx.store_with_thread("picker-deleted").await;
304 {
305 let opened = RuntimeThreadStore::open(store.clone()).unwrap();
306 let mut thread = opened.load_thread(&thread_id).unwrap();
307 thread.session_id = Some("deleted-from-picker".into());
308 opened.save_thread(&thread).unwrap();
309 }
310
311 let summary = fx.run();
312 assert_eq!(summary.threads_unbound, 1, "{summary:?}");
313 assert_eq!(summary.sessions_recovered, 1, "{summary:?}");
314 assert_eq!(summary.stores_kept_with_work, 0, "{summary:?}");
315
316 let id = crate::runtime_threads::thread_session_id(&thread_id);
317 let recovered = fx.sessions.load_session(&id).expect("recovered document");
318 assert!(recovered.metadata.title.starts_with("Recovered: "));
319 assert_eq!(
320 recovered
321 .metadata
322 .runtime_store
323 .expect("bound to its store")
324 .data_dir,
325 store.canonicalize().unwrap()
326 );
327 let thread = RuntimeThreadStore::open(store.clone())
328 .unwrap()
329 .load_thread(&thread_id)
330 .unwrap();
331 assert_eq!(thread.session_id.as_deref(), Some(id.as_str()));
332 let receipts = fs::read_to_string(store.join(THREAD_UNBIND_RECEIPTS_FILE)).unwrap();
333 assert!(receipts.contains("deleted-from-picker"), "{receipts}");
334
335 let again = fx.run();
336 assert!(!again.changed(), "{again:?}");
337 }
338
339 /// #6555: reconcile judges seed journals exactly as Runtime startup does,
340 /// but never settles them. An unpublished seed's turns are not history, and
341 /// a thread whose journal cannot be settled is not recovered.
342 #[tokio::test]
343 async fn reconcile_does_not_recover_an_unpublished_or_unsettled_seed() {
344 let _env = lock_test_env();
345 for journal_state in ["uncommitted", "unreadable"] {
346 let fx = Fixture::new();
347 let (store, thread_id) = fx.store_with_thread("seeding-anchor").await;
348 let journal_path = store.join("threads").join(format!("{thread_id}.seed"));
349 {
350 let opened = RuntimeThreadStore::open(store.clone()).unwrap();
351 let turns = opened.list_turns_for_thread(&thread_id).unwrap();
352 let items = opened.list_items_for_turn(&turns[0].id).unwrap();
353 // The process stopped between the seed's records and its commit.
354 let mut thread = opened.load_thread(&thread_id).unwrap();
355 thread.latest_turn_id = None;
356 opened.save_thread(&thread).unwrap();
357 let journal = serde_json::json!({
358 "thread_id": thread_id,
359 "previous_latest_turn_id": null,
360 "turn_ids": turns.iter().map(|turn| turn.id.clone()).collect::<Vec<_>>(),
361 "item_ids": items.iter().map(|item| item.id.clone()).collect::<Vec<_>>(),
362 });
363 let bytes = match journal_state {
364 "uncommitted" => serde_json::to_vec(&journal).unwrap(),
365 _ => b"{truncated seed intent".to_vec(),
366 };
367 fs::write(&journal_path, bytes).unwrap();
368 }
369
370 let summary = fx.run();
371 assert_eq!(
372 summary.sessions_recovered, 0,
373 "{journal_state}: {summary:?}"
374 );
375 assert_eq!(
376 summary.stores_kept_with_work, 1,
377 "{journal_state}: {summary:?}"
378 );
379 assert!(
380 journal_path.is_file(),
381 "{journal_state}: reconcile leaves the journal for Runtime startup"
382 );
383 let thread = RuntimeThreadStore::open(store.clone())
384 .unwrap()
385 .load_thread(&thread_id)
386 .unwrap();
387 assert_eq!(thread.session_id, None, "{journal_state}");
388 }
389 }
390
391 /// R1: an unreadable document is set aside with its hash; a document from a
392 /// newer build is left alone.
393 #[test]
394 fn reconcile_sets_aside_unreadable_documents_and_keeps_newer_ones() {
395 let _env = lock_test_env();
396 let fx = Fixture::new();
397 fs::write(fx.dir().join("torn.json"), b"{\"metadata\": {\"id\": ").unwrap();
398 fs::write(
399 fx.dir().join("future.json"),
400 br#"{"schema_version": 999, "shape": "unknown"}"#,
401 )
402 .unwrap();
403
404 let summary = fx.run();
405 assert_eq!(summary.documents_set_aside, 1, "{summary:?}");
406 assert_eq!(summary.documents_newer_schema, 1);
407 assert!(!fx.dir().join("torn.json").exists());
408 assert!(fx.dir().join("future.json").exists());
409 let manifest = manifests(&fx.set_aside_root());
410 assert!(
411 manifest.contains("torn.json") && manifest.contains("sha256"),
412 "{manifest}"
413 );
414 }
415
416 /// R6: a document-less artifact directory nothing names is set aside; one a
417 /// document mentions, or one still being written, is kept.
418 #[test]
419 fn reconcile_sets_aside_unreferenced_artifact_directories() {
420 let _env = lock_test_env();
421 let fx = Fixture::new();
422 let lost = "0f0f0f0f-1111-4222-8333-444444444444";
423 let named = "0a0a0a0a-1111-4222-8333-555555555555";
424 for id in [lost, named] {
425 let artifacts = fx.dir().join(id).join("artifacts");
426 fs::create_dir_all(&artifacts).unwrap();
427 fs::write(artifacts.join("context-transfer-x.json"), b"[]").unwrap();
428 }
429 let mut mentions = create_saved_session_with_id_and_mode(
430 "mentions".into(),
431 &[text(Role::User, &format!("see sessions/{named}/artifacts"))],
432 "deepseek-v4-pro",
433 fx.home.path(),
434 0,
435 None,
436 None,
437 );
438 mentions.metadata.title = "Mentions".into();
439 fx.sessions.save_session(&mentions).unwrap();
440
441 // Recent directories are kept until they have been idle long enough.
442 let recent = fx.run_with(ReconcileOptions::default());
443 assert_eq!(recent.artifact_dirs_set_aside, 0);
444
445 let summary = fx.run();
446 assert_eq!(summary.artifact_dirs_set_aside, 1, "{summary:?}");
447 assert_eq!(summary.artifact_dirs_kept, 1);
448 assert!(!fx.dir().join(lost).exists());
449 assert!(fx.dir().join(named).exists());
450 }
451
452 /// The per-run limit stops a run and the next continues; a second process
453 /// arriving while one runs does nothing.
454 #[test]
455 fn reconcile_is_bounded_and_single_flight() {
456 let _env = lock_test_env();
457 let fx = Fixture::new();
458 fx.empty_store("orphan-one");
459 fx.empty_store("orphan-two");
460
461 let first = fx.run_with(ReconcileOptions {
462 limit: 1,
463 ..ReconcileOptions::default()
464 });
465 assert_eq!(first.stores_set_aside, 1);
466 assert!(first.limit_reached);
467 let second = fx.run();
468 assert_eq!(second.stores_set_aside, 1);
469 assert!(!second.limit_reached);
470
471 let lock = fs::OpenOptions::new()
472 .create(true)
473 .truncate(false)
474 .write(true)
475 .open(fx.dir().join(RECONCILE_LOCK_FILE))
476 .unwrap();
477 assert!(crate::runtime_threads::try_lock_file_exclusive(&lock).unwrap());
478 let skipped = fx.run();
479 assert!(skipped.skipped_concurrent);
480 }
481
482 /// P2: deleting a document retires the store its *binding* names — not
483 /// `sessions/<id>/runtime` — unless another document binds it or it holds
484 /// work. A `runtime-recovered-*` store under the deleted id is never removed.
485 #[tokio::test]
486 async fn deleting_a_document_retires_the_store_it_released() {
487 let _env = lock_test_env();
488 let fx = Fixture::new();
489 let shared = fx.empty_store("host-one");
490 let sole = fx.empty_store("host-two");
491 fx.document("keeps-shared", Some(&shared));
492 fx.document("uses-shared", Some(&shared));
493 fx.document("uses-sole", Some(&sole));
494 let (busy, _) = fx.store_with_thread("host-three").await;
495 fx.document("uses-busy", Some(&busy));
496 let recovered = fx.dir().join("uses-sole").join("runtime-recovered-session");
497 RuntimeThreadStore::open(recovered.clone()).unwrap();
498 fx.document("bound-to-recovered", Some(&recovered));
499
500 fx.sessions.delete_session("uses-shared").unwrap();
501 fx.sessions.delete_session("uses-sole").unwrap();
502 fx.sessions.delete_session("uses-busy").unwrap();
503
504 assert!(shared.is_dir(), "another document still binds it");
505 assert!(!sole.exists(), "released and empty: set aside");
506 assert!(busy.is_dir(), "holds work: kept");
507 assert!(recovered.is_dir(), "bound elsewhere: never deleted");
508 assert!(manifests(&fx.set_aside_root()).contains("host-two"));
509 }
510
511 /// P8: concurrent boot-owner stamps from separate handles lose no entry.
512 #[test]
513 fn concurrent_boot_owner_stamps_keep_every_entry() {
514 let _env = lock_test_env();
515 let fx = Fixture::new();
516 let ids: Vec<String> = (0..8).map(|index| format!("stamped-{index}")).collect();
517 for id in &ids {
518 fx.document(id, None);
519 }
520 let dir = fx.dir();
521 std::thread::scope(|scope| {
522 for id in &ids {
523 let dir = dir.clone();
524 scope.spawn(move || {
525 let manager = SessionManager::new(dir).unwrap();
526 manager.record_session_boot_owner(id, "boot_test").unwrap();
527 });
528 }
529 });
530 for id in &ids {
531 assert_eq!(
532 fx.sessions.session_boot_owner(id).as_deref(),
533 Some("boot_test")
534 );
535 }
536 }
537
538 #[test]
539 fn summary_notice_and_doctor_line_name_what_changed() {
540 let quiet = ReconcileSummary::default();
541 assert_eq!(quiet.notice(), None);
542 assert!(quiet.doctor_detail().contains("never"));
543
544 let changed = ReconcileSummary {
545 ran_at: Some(Utc::now()),
546 sessions_recovered: 2,
547 threads_unbound: 5,
548 stores_set_aside: 140,
549 stores_in_use: 1,
550 set_aside_path: Some(PathBuf::from("/tmp/set-aside")),
551 ..ReconcileSummary::default()
552 };
553 let notice = changed.notice().unwrap();
554 assert!(notice.contains("2 recovered"), "{notice}");
555 assert!(notice.contains("5 threads re-linked"), "{notice}");
556 assert!(notice.contains("140 unused items set aside"), "{notice}");
557 let detail = changed.doctor_detail();
558 assert!(detail.contains("140 stores"), "{detail}");
559 assert!(detail.contains("1 unbound stores in use"), "{detail}");
560
561 let dry = ReconcileSummary {
562 dry_run: true,
563 ..changed
564 };
565 assert_eq!(dry.notice(), None, "a dry run changed nothing");
566 }
567
567 lines RUST