返回 CodeWhale
runtime_store_binding.rs
根目录 / crates / tui / src / tui / ui / tests / runtime_store_binding.rs
1 use super::*;
2 use crate::automation_manager::{
3 AutomationManager, AutomationStatus, CreateAutomationRequest, run_now_shared,
4 };
5 use crate::runtime_threads::{RuntimeThreadManager, RuntimeThreadManagerConfig};
6 use crate::task_manager::{TaskManager, TaskManagerConfig};
7
8 fn fixture_config() -> Config {
9 let mut config = Config {
10 ..Config::default()
11 }
12 .with_legacy_root(
13 Some("local-runtime-binding-fixture".into()),
14 Some("http://127.0.0.1:1/v1".into()),
15 );
16 config.set_feature("mcp", false).unwrap();
17 config.set_feature("subagents", false).unwrap();
18 config
19 }
20
21 #[tokio::test]
22 async fn runtime_store_binding_persists_on_exit_without_a_model_turn() -> anyhow::Result<()> {
23 let _environment = crate::test_support::lock_test_env();
24 let root = tempfile::tempdir()?;
25 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
26 let _runtime = crate::test_support::EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
27 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
28 let explicit_store = crate::test_support::EnvVarGuard::set(
29 "CODEWHALE_RUNTIME_DIR",
30 root.path().join("original-runtime"),
31 );
32 let mut config = fixture_config();
33 let sessions = SessionManager::default_location()?;
34 let original = crate::session_manager::create_saved_session_with_id_and_mode(
35 "legacy-conversation".into(),
36 &[text_message("user", "retain this earlier conversation")],
37 "deepseek-v4-pro",
38 root.path(),
39 0,
40 None,
41 None,
42 );
43 assert!(original.metadata.runtime_store.is_none());
44 sessions.save_session(&original)?;
45 sessions.save_checkpoint(&original)?;
46 let mut app = Box::new(create_test_app());
47 apply_loaded_session_with_goal(&mut app, &mut config, original.clone(), None)
48 .map_err(anyhow::Error::msg)?;
49 let task_config = TaskManagerConfig::from_runtime(&config, root.path().into(), None, Some(1));
50 let tasks = TaskManager::start(
51 task_config.clone(),
52 config.clone(),
53 app.plugin_registry.clone(),
54 &original.metadata.id,
55 None,
56 )
57 .await?;
58 let binding = tasks
59 .session_store_binding()
60 .expect("attached Runtime store");
61 app.runtime_services.task_manager = Some(tasks.clone());
62 let (handle, actor) =
63 persistence_actor::spawn_persistence_actor(SessionManager::default_location()?);
64 // Match clean exit ordering: no Engine turn, checkpoint or snapshot has
65 // been queued by this host before its TaskManager stops.
66 tasks.shutdown_and_wait().await?;
67 assert!(
68 super::super::event_loop::persist_settled_session_on_shutdown(&mut app, &handle)
69 .map_err(anyhow::Error::msg)?
70 );
71 assert!(handle.try_send(PersistRequest::Shutdown));
72 actor.await?;
73 let saved = sessions.load_session(&original.metadata.id)?;
74 assert_eq!(saved.metadata.runtime_store.as_ref(), Some(&binding));
75 assert_eq!(saved.metadata.title, original.metadata.title);
76 assert_eq!(saved.messages, original.messages);
77 assert!(
78 sessions
79 .load_session_checkpoint(&original.metadata.id)?
80 .is_none()
81 );
82 drop(app);
83 drop(tasks);
84 drop(explicit_store);
85 // Ordinary resume now reopens the same authority without an env override.
86 let resumed = TaskManager::start(
87 task_config,
88 config,
89 Arc::new(crate::plugins::PluginRegistry::empty(root.path())),
90 &saved.metadata.id,
91 saved.metadata.runtime_store.as_ref(),
92 )
93 .await?;
94 assert_eq!(resumed.execution_scope(), binding.execution_scope);
95 resumed.shutdown_and_wait().await?;
96 Ok(())
97 }
98
99 #[tokio::test]
100 async fn runtime_store_binding_exit_preserves_inflight_recovery() -> anyhow::Result<()> {
101 let _environment = crate::test_support::lock_test_env();
102 let root = tempfile::tempdir()?;
103 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
104 let sessions = SessionManager::default_location()?;
105 // U02-09: a locally cancelled turn still owes its terminal event.
106 for (loading, dispatch, cancelled) in [
107 (true, false, false),
108 (false, true, false),
109 (false, false, true),
110 ] {
111 let original = crate::session_manager::create_saved_session_with_mode(
112 &[],
113 "deepseek-v4-pro",
114 root.path(),
115 0,
116 None,
117 None,
118 );
119 let path = sessions.save_session(&original)?;
120 let checkpoint = sessions.save_checkpoint(&original)?;
121 let saved_before = std::fs::read(&path)?;
122 let checkpoint_before = std::fs::read(&checkpoint)?;
123 let mut app = Box::new(create_test_app());
124 app.current_session_id = Some(original.metadata.id.clone());
125 app.is_loading = loading;
126 app.dispatch_in_flight = dispatch;
127 app.suppress_stream_events_until_turn_complete = cancelled;
128 let (handle, actor) =
129 persistence_actor::spawn_persistence_actor(SessionManager::default_location()?);
130 assert!(
131 !super::super::event_loop::persist_settled_session_on_shutdown(&mut app, &handle)
132 .map_err(anyhow::Error::msg)?
133 );
134 assert!(handle.try_send(PersistRequest::Shutdown));
135 actor.await?;
136 assert_eq!(std::fs::read(path)?, saved_before);
137 assert_eq!(std::fs::read(checkpoint)?, checkpoint_before);
138 }
139 Ok(())
140 }
141
142 /// #6362: this test used to be one async body. Every debug-build temporary
143 /// of that body — the boxed `App` returns, the cloned `Config`s, two session
144 /// snapshots, the task-manager futures — got its own slot in a single poll
145 /// frame, which alone measured 1.1 MiB on the 2 MiB stack libtest gives a
146 /// test thread (gdb frame attribution, 2026-09-20). The phases below are
147 /// built and boxed through `boxed_phase`, so each phase's temporaries die
148 /// with its own poll frame and the outer body holds pointers; inline
149 /// `async {}` phases measured 806 KiB of never-reused slots on the outer
150 /// frame and still overflowed. The test pins the default budget explicitly
151 /// instead of inheriting CI's 16 MiB `RUST_MIN_STACK`, which is what masked
152 /// the overflow.
153 #[test]
154 fn runtime_store_binding_survives_launch_snapshot_and_resume() -> anyhow::Result<()> {
155 use crate::test_support::boxed_phase;
156
157 crate::test_support::block_on_default_test_stack(|| async {
158 let _environment = crate::test_support::lock_test_env();
159 let root = tempfile::tempdir()?;
160 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
161 let _runtime = crate::test_support::EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
162 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
163 let config = fixture_config();
164 let sessions = SessionManager::default_location()?;
165 let task_config =
166 TaskManagerConfig::from_runtime(&config, root.path().into(), None, Some(1));
167 let (root, config, task_config, sessions) = (&root, &config, &task_config, &sessions);
168
169 // Phase 1: launch over a saved conversation and bind its Runtime store.
170 let (initial_id, saved_id, binding, automation) = boxed_phase(move || async move {
171 let mut app = Box::new(create_test_app());
172 app.workspace = root.path().into();
173 let initial_id = super::super::event_loop::ensure_runtime_session_id(&mut app);
174 let tasks = TaskManager::start(
175 task_config.clone(),
176 config.clone(),
177 app.plugin_registry.clone(),
178 &initial_id,
179 None,
180 )
181 .await?;
182 app.runtime_services.task_manager = Some(tasks.clone());
183 // A saved initial conversation may later be deleted while another
184 // launch still refers to its Runtime store.
185 let initial = build_session_snapshot(&mut app, sessions).map_err(anyhow::Error::msg)?;
186 sessions.save_session(&initial)?;
187 let launch = begin_launch_session(&mut app, None);
188 assert!(!launch.is_error, "{:?}", launch.message);
189 assert_ne!(app.current_session_id.as_deref(), Some(initial_id.as_str()));
190 let saved = build_session_snapshot(&mut app, sessions).map_err(anyhow::Error::msg)?;
191 let binding = saved
192 .metadata
193 .runtime_store
194 .clone()
195 .expect("attached host binding");
196 assert_eq!(binding.execution_scope, tasks.execution_scope());
197 sessions.save_session(&saved)?;
198 let mut automations = AutomationManager::open(root.path().join("automations"))?;
199 automations.bind_task_manager(&tasks)?;
200 let automation = automations.create_automation(CreateAutomationRequest {
201 name: "resumed ownership fixture".into(),
202 prompt: "local fixture only".into(),
203 rrule: "FREQ=HOURLY;INTERVAL=1".into(),
204 cwds: vec![root.path().into()],
205 model: None,
206 model_provider: None,
207 model_provider_id: None,
208 mode: None,
209 allow_shell: Some(false),
210 trust_mode: Some(false),
211 auto_approve: Some(false),
212 delivery_mode: None,
213 status: Some(AutomationStatus::Paused),
214 })?;
215 tasks.shutdown_and_wait().await?;
216 drop(app);
217 drop(tasks);
218 drop(automations);
219 Ok::<_, anyhow::Error>((initial_id, saved.metadata.id.clone(), binding, automation))
220 })
221 .await?;
222 sessions.delete_session(&initial_id)?;
223 assert!(
224 binding.data_dir.is_dir(),
225 "transcript deletion cannot erase Runtime authority"
226 );
227 let loaded = sessions.load_session(&saved_id)?;
228 assert_eq!(loaded.metadata.runtime_store.as_ref(), Some(&binding));
229 let mut resumed_config = config.clone();
230 let (loaded, automation) = (&loaded, &automation);
231
232 // Phase 2: resume with the binding and run the automation for real.
233 {
234 let resumed_config = &mut resumed_config;
235 boxed_phase(move || async move {
236 let mut resumed = Box::new(create_test_app());
237 apply_loaded_session_with_goal(&mut resumed, resumed_config, loaded.clone(), None)
238 .map_err(anyhow::Error::msg)?;
239 let tasks = TaskManager::start(
240 task_config.clone(),
241 config.clone(),
242 resumed.plugin_registry.clone(),
243 &loaded.metadata.id,
244 loaded.metadata.runtime_store.as_ref(),
245 )
246 .await?;
247 assert_eq!(
248 tasks.execution_scope(),
249 automation.execution_scope.as_deref().unwrap()
250 );
251 resumed.runtime_services.task_manager = Some(tasks.clone());
252 let automations = Arc::new(tokio::sync::Mutex::new(AutomationManager::open(
253 root.path().join("automations"),
254 )?));
255 // The real Run-now admission must now create its durable
256 // receipt. The configured endpoint is closed loopback and no
257 // shell/tool is authorized.
258 let run = run_now_shared(&automations, &automation.id, &tasks).await?;
259 assert!(run.task_id.is_some(), "{run:?}");
260 assert_eq!(
261 automations
262 .lock()
263 .await
264 .list_runs(&automation.id, None)?
265 .len(),
266 1
267 );
268 assert_eq!(
269 automations
270 .lock()
271 .await
272 .get_automation(&automation.id)?
273 .execution_scope,
274 automation.execution_scope
275 );
276 tasks.shutdown_and_wait().await?;
277 drop(resumed);
278 drop(tasks);
279 Ok::<_, anyhow::Error>(())
280 })
281 .await?;
282 }
283
284 // Phase 3: reproduce the old resume path — deriving a store from the
285 // saved conversation id without its binding opens a foreign scope and
286 // cannot run.
287 let foreign = TaskManager::start(
288 task_config.clone(),
289 config.clone(),
290 Arc::new(crate::plugins::PluginRegistry::empty(root.path())),
291 &loaded.metadata.id,
292 None,
293 )
294 .await?;
295 let foreign = &foreign;
296 boxed_phase(move || async move {
297 let foreign_automations = Arc::new(tokio::sync::Mutex::new(AutomationManager::open(
298 root.path().join("automations"),
299 )?));
300 let definition_path = root
301 .path()
302 .join("automations/automations")
303 .join(format!("{}.json", automation.id));
304 let before_foreign_run = std::fs::read(&definition_path)?;
305 let error = run_now_shared(&foreign_automations, &automation.id, foreign)
306 .await
307 .unwrap_err();
308 assert!(
309 error
310 .to_string()
311 .contains("another Runtime execution scope"),
312 "{error:#}"
313 );
314 assert_eq!(
315 foreign_automations
316 .lock()
317 .await
318 .list_runs(&automation.id, None)?
319 .len(),
320 1
321 );
322 assert_eq!(std::fs::read(definition_path)?, before_foreign_run);
323 Ok::<_, anyhow::Error>(())
324 })
325 .await?;
326
327 // Phase 4: a host that already owns a foreign scope refuses the switch
328 // and keeps its pending input.
329 {
330 let resumed_config = &mut resumed_config;
331 boxed_phase(move || async move {
332 let mut other_app = Box::new(create_test_app());
333 other_app.runtime_services.task_manager = Some(foreign.clone());
334 other_app.input = "preserve pending input".into();
335 let old_id = other_app.current_session_id.clone();
336 let error = apply_loaded_session_with_goal(
337 &mut other_app,
338 resumed_config,
339 loaded.clone(),
340 None,
341 )
342 .unwrap_err();
343 // The refusal must name the route that actually works. "Resume
344 // it in a new Codewhale process" was true but unactionable:
345 // starting a new process and then picking the session from
346 // `/resume` returns here, because that is this same switch
347 // path (#6207, #6225).
348 assert!(
349 error.contains("codewhale resume"),
350 "the refusal must point at the direct-open path: {error}"
351 );
352 assert!(
353 error.contains(&loaded.metadata.id),
354 "the refusal must name the session to open: {error}"
355 );
356 assert_eq!(other_app.current_session_id, old_id);
357 assert_eq!(other_app.input, "preserve pending input");
358 Ok::<_, anyhow::Error>(())
359 })
360 .await?;
361 }
362 foreign.shutdown_and_wait().await?;
363 Ok(())
364 })
365 }
366
367 #[cfg(unix)]
368 #[test]
369 fn runtime_store_binding_retention_does_not_follow_session_directory_symlinks() -> anyhow::Result<()>
370 {
371 let _environment = crate::test_support::lock_test_env();
372 let root = tempfile::tempdir()?;
373 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
374 let sessions = SessionManager::default_location()?;
375 let saved = crate::session_manager::create_saved_session_with_id_and_mode(
376 "linked-session".into(),
377 &[],
378 "fixture",
379 root.path(),
380 0,
381 None,
382 None,
383 );
384 sessions.save_session(&saved)?;
385 let target = root.path().join("unrelated-directory");
386 std::fs::create_dir_all(&target)?;
387 std::fs::write(target.join("keep.txt"), "preserve user data")?;
388 let link = root.path().join("sessions/linked-session");
389 std::os::unix::fs::symlink(&target, &link)?;
390 sessions.delete_session("linked-session")?;
391 assert!(!link.exists());
392 assert_eq!(
393 std::fs::read_to_string(target.join("keep.txt"))?,
394 "preserve user data"
395 );
396 Ok(())
397 }
398
399 #[test]
400 fn runtime_store_binding_rejects_foreign_missing_or_overridden_store() -> anyhow::Result<()> {
401 let _environment = crate::test_support::lock_test_env();
402 let root = tempfile::tempdir()?;
403 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
404 let _runtime = crate::test_support::EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
405 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
406 let config = fixture_config();
407 let cfg = RuntimeThreadManagerConfig::for_session(root.path().join("tasks"), "original");
408 let runtime = RuntimeThreadManager::open(config.clone(), root.path().into(), cfg.clone())?;
409 let binding = runtime.session_store_binding();
410 drop(runtime);
411 let state_path = binding.data_dir.join("state.json");
412 let before = std::fs::read(&state_path)?;
413 let mut wrong = binding.clone();
414 wrong.execution_scope = "0".repeat(64);
415 let open = |binding: &crate::runtime_threads::RuntimeStoreBinding| {
416 RuntimeThreadManager::open_for_session(
417 config.clone(),
418 root.path().into(),
419 cfg.clone(),
420 Arc::new(crate::plugins::PluginRegistry::empty(root.path())),
421 Some(binding),
422 )
423 };
424 assert!(
425 open(&wrong)
426 .err()
427 .unwrap()
428 .to_string()
429 .contains("ownership does not match")
430 );
431 assert_eq!(std::fs::read(&state_path)?, before);
432 wrong.data_dir = root.path().join("missing-store");
433 assert!(open(&wrong).is_err());
434 assert!(
435 !wrong.data_dir.exists(),
436 "saved binding cannot create a replacement authority"
437 );
438 let _override = crate::test_support::EnvVarGuard::set(
439 "CODEWHALE_RUNTIME_DIR",
440 root.path().join("foreign-override"),
441 );
442 assert!(
443 open(&binding)
444 .err()
445 .unwrap()
446 .to_string()
447 .contains("override conflicts")
448 );
449 Ok(())
450 }
451
452 #[test]
453 fn missing_runtime_store_recovers_without_reusing_authority_or_resurrecting_stale_binding()
454 -> anyhow::Result<()> {
455 let _environment = crate::test_support::lock_test_env();
456 let root = tempfile::tempdir()?;
457 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
458 let _runtime = crate::test_support::EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
459 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
460 let sessions = SessionManager::default_location()?;
461 let mut saved = crate::session_manager::create_saved_session_with_id_and_mode(
462 "interrupted".into(),
463 &[text_message("user", "retain my work")],
464 "deepseek-v4-pro",
465 root.path(),
466 0,
467 None,
468 None,
469 );
470 let missing = crate::runtime_threads::RuntimeStoreBinding {
471 data_dir: root.path().join("sessions/previous/runtime"),
472 execution_scope: "0".repeat(64),
473 };
474 saved.metadata.runtime_store = Some(missing.clone());
475 sessions.save_session(&saved)?;
476 let stale = saved.clone();
477 let manager = RuntimeThreadManager::open_for_session(
478 fixture_config(),
479 root.path().into(),
480 RuntimeThreadManagerConfig::for_session(root.path().join("tasks"), "interrupted"),
481 Arc::new(crate::plugins::PluginRegistry::empty(root.path())),
482 Some(&missing),
483 )?;
484 let recovered = manager.session_store_binding();
485 assert_ne!(recovered.execution_scope, missing.execution_scope);
486 assert_ne!(recovered.data_dir, missing.data_dir);
487 assert!(!missing.data_dir.exists());
488 assert!(
489 RuntimeThreadManager::open_for_session(
490 fixture_config(),
491 root.path().into(),
492 RuntimeThreadManagerConfig::for_session(root.path().join("tasks"), "interrupted"),
493 Arc::new(crate::plugins::PluginRegistry::empty(root.path())),
494 Some(&missing),
495 )
496 .is_err(),
497 "concurrent recovery cannot mint a second owner"
498 );
499 // Losing the process before the repaired binding is saved must leave a
500 // retryable, session-scoped store, not an orphan or a second authority.
501 drop(manager);
502 let manager = RuntimeThreadManager::open_for_session(
503 fixture_config(),
504 root.path().into(),
505 RuntimeThreadManagerConfig::for_session(root.path().join("tasks"), "interrupted"),
506 Arc::new(crate::plugins::PluginRegistry::empty(root.path())),
507 Some(&missing),
508 )?;
509 assert_eq!(manager.session_store_binding(), recovered);
510 let other_recovery = RuntimeThreadManager::open_for_session(
511 fixture_config(),
512 root.path().into(),
513 RuntimeThreadManagerConfig::for_session(root.path().join("tasks"), "other-interrupted"),
514 Arc::new(crate::plugins::PluginRegistry::empty(root.path())),
515 Some(&missing),
516 )?;
517 assert_ne!(
518 other_recovery.session_store_binding().data_dir,
519 recovered.data_dir
520 );
521 assert_ne!(
522 other_recovery.session_store_binding().execution_scope,
523 recovered.execution_scope
524 );
525 saved.metadata.runtime_store = Some(recovered.clone());
526 sessions.save_session(&saved)?;
527 sessions.save_session(&stale)?;
528 sessions.save_checkpoint(&stale)?;
529 let durable = sessions.load_session("interrupted")?;
530 assert_eq!(durable.metadata.runtime_store, Some(recovered.clone()));
531 assert_eq!(durable.messages, stale.messages);
532 let competing = RuntimeThreadManager::open_for_session(
533 fixture_config(),
534 root.path().into(),
535 RuntimeThreadManagerConfig::for_session(root.path().join("tasks"), "competing"),
536 Arc::new(crate::plugins::PluginRegistry::empty(root.path())),
537 None,
538 )?;
539 let mut competing_snapshot = stale.clone();
540 competing_snapshot.metadata.runtime_store = Some(competing.session_store_binding());
541 assert!(sessions.save_session(&competing_snapshot).is_err());
542 assert!(sessions.save_checkpoint(&competing_snapshot).is_err());
543 assert_eq!(
544 sessions.load_session("interrupted")?.metadata.runtime_store,
545 Some(recovered),
546 "another valid owner cannot overwrite the completed recovery"
547 );
548 Ok(())
549 }
550
551 #[tokio::test]
552 async fn picker_recovers_missing_store_into_the_idle_host_and_persists_before_returning()
553 -> anyhow::Result<()> {
554 let _environment = crate::test_support::lock_test_env();
555 let root = tempfile::tempdir()?;
556 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
557 let _runtime = crate::test_support::EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
558 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
559 let sessions = SessionManager::default_location()?;
560 let mut config = fixture_config();
561 let mut saved = crate::session_manager::create_saved_session_with_id_and_mode(
562 "picker-interrupted".into(),
563 &[text_message("user", "retain my work")],
564 "deepseek-v4-pro",
565 root.path(),
566 0,
567 None,
568 None,
569 );
570 saved.metadata.runtime_store = Some(crate::runtime_threads::RuntimeStoreBinding {
571 data_dir: root.path().join("sessions/previous/runtime"),
572 execution_scope: "0".repeat(64),
573 });
574 sessions.save_session(&saved)?;
575 let mut app = Box::new(create_test_app());
576 let tasks = TaskManager::start(
577 TaskManagerConfig::from_runtime(&config, root.path().into(), None, Some(1)),
578 config.clone(),
579 app.plugin_registry.clone(),
580 "picker-current",
581 None,
582 )
583 .await?;
584 let binding = tasks.session_store_binding().expect("current host");
585 app.runtime_services.task_manager = Some(tasks.clone());
586 app.current_session_id = Some("picker-current".into());
587 app.api_messages_mut()
588 .push(text_message("user", "current conversation"));
589 let current_messages = app.api_messages.clone();
590 let plan_state = app.plan_state.clone();
591 let held = plan_state
592 .try_lock()
593 .expect("hold Work state during recovery");
594 assert!(apply_loaded_session_with_goal(&mut app, &mut config, saved.clone(), None).is_err());
595 assert_eq!(app.current_session_id.as_deref(), Some("picker-current"));
596 assert_eq!(app.api_messages, current_messages);
597 assert_eq!(
598 sessions
599 .load_session("picker-interrupted")?
600 .metadata
601 .runtime_store
602 .as_ref(),
603 Some(&binding),
604 "binding repair survives a contended UI restore"
605 );
606 drop(held);
607 apply_loaded_session_with_goal(&mut app, &mut config, saved.clone(), None)
608 .map_err(anyhow::Error::msg)?;
609 assert_eq!(
610 app.current_session_id.as_deref(),
611 Some("picker-interrupted")
612 );
613 let durable = sessions.load_session("picker-interrupted")?;
614 assert_eq!(durable.metadata.runtime_store.as_ref(), Some(&binding));
615 assert_eq!(durable.messages, saved.messages);
616 assert_eq!(
617 app.current_session_metadata.as_ref().unwrap().runtime_store,
618 Some(binding)
619 );
620 sessions.save_session(&saved)?;
621 assert_eq!(
622 sessions
623 .load_session("picker-interrupted")?
624 .metadata
625 .runtime_store,
626 durable.metadata.runtime_store
627 );
628 tasks.shutdown_and_wait().await?;
629 Ok(())
630 }
631
632 /// #6207: a store that exists but is empty, unheld, and scope-free holds
633 /// nothing a session switch could abandon. A force-quit leaves exactly that
634 /// shape — the directory is on disk, ownerless, with zero events — and
635 /// refusing it left the session unopenable while protecting nothing.
636 #[test]
637 fn adoptable_empty_store_reports_nothing_to_abandon() -> anyhow::Result<()> {
638 let _environment = crate::test_support::lock_test_env();
639 let root = tempfile::tempdir()?;
640 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
641 let _runtime = crate::test_support::EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
642 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
643
644 let store_dir = root.path().join("sessions/interrupted/runtime");
645 // Open a real store, so the layout under test is the product's rather than
646 // the test's idea of it.
647 drop(crate::runtime_threads::RuntimeThreadStore::open(
648 store_dir.clone(),
649 )?);
650
651 let binding = crate::runtime_threads::RuntimeStoreBinding {
652 data_dir: store_dir.clone(),
653 execution_scope: "0".repeat(64),
654 };
655 assert!(
656 !binding.is_missing_session_store()?,
657 "the store exists, so the old predicate cannot recover it"
658 );
659 assert_eq!(
660 binding.adoption_refusal()?,
661 None,
662 "a freshly opened store holds nothing to abandon"
663 );
664 assert!(
665 !binding.has_live_holder()?,
666 "nobody holds the freshly opened store"
667 );
668 assert!(
669 !binding.has_scope_pinned_automation()?,
670 "no automations exist under the fixture home"
671 );
672 assert!(
673 binding.is_adoptable_empty_store()?,
674 "empty, unheld, scope-free: adoptable"
675 );
676
677 // Each work directory must be load-bearing on its own. If `open` gains a
678 // directory that RUNTIME_STORE_WORK_DIRS misses, this is the assertion that
679 // notices, instead of the miss silently widening what a switch will adopt.
680 for dir in [
681 "threads",
682 "turns",
683 "items",
684 "events",
685 "goals",
686 "agent-mail",
687 "turn-operations",
688 ] {
689 let marker = store_dir.join(dir).join("work.json");
690 std::fs::write(&marker, "{}")?;
691 assert_eq!(
692 binding.adoption_refusal()?,
693 Some(crate::runtime_threads::StoreAdoptionRefusal::HasDurableWork { dir }),
694 "{dir} holds work; the store must not be adopted"
695 );
696 assert!(
697 !binding.is_adoptable_empty_store()?,
698 "{dir} blocks the adopt"
699 );
700 std::fs::remove_file(&marker)?;
701 }
702 assert!(
703 binding.is_adoptable_empty_store()?,
704 "markers removed: adoptable again"
705 );
706 Ok(())
707 }
708
709 /// #6418: the binding a live host records is its store's canonical root,
710 /// while the configured sessions root is spelled lexically. When those differ
711 /// (a Windows verbatim prefix or short name, a symlinked home on Unix), the
712 /// hand-built bindings above still pass but every *recorded* binding read as
713 /// unconfined, so an empty, unheld store could never be adopted in-session.
714 #[test]
715 fn recorded_binding_is_confined_under_a_non_canonical_home() -> anyhow::Result<()> {
716 let _environment = crate::test_support::lock_test_env();
717 let root = tempfile::tempdir()?;
718 let real_home = root.path().join("real-home");
719 std::fs::create_dir_all(&real_home)?;
720 #[cfg(unix)]
721 let home = {
722 let linked = root.path().join("linked-home");
723 std::os::unix::fs::symlink(&real_home, &linked)?;
724 linked
725 };
726 // Windows needs no fixture: canonical paths there carry `\\?\`.
727 #[cfg(not(unix))]
728 let home = real_home.clone();
729 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", &home);
730 let _runtime = crate::test_support::EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
731 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
732
733 let runtime = RuntimeThreadManager::open(
734 fixture_config(),
735 home.clone(),
736 RuntimeThreadManagerConfig::for_session(home.join("tasks"), "previous"),
737 )?;
738 let binding = runtime.session_store_binding();
739 drop(runtime);
740
741 #[cfg(unix)]
742 assert!(
743 !binding.data_dir.starts_with(&home),
744 "fixture must record a binding spelled differently from the home"
745 );
746 assert!(
747 !binding.is_missing_session_store()?,
748 "the recorded store exists"
749 );
750 assert!(
751 binding.is_adoptable_empty_store()?,
752 "a recorded, empty, unheld store is confined and adoptable"
753 );
754
755 // Confinement still refuses a store outside the sessions root.
756 let outside = crate::runtime_threads::RuntimeStoreBinding {
757 data_dir: root.path().join("elsewhere/previous/runtime"),
758 execution_scope: binding.execution_scope.clone(),
759 };
760 std::fs::create_dir_all(&outside.data_dir)?;
761 assert!(!outside.is_adoptable_empty_store()?);
762 Ok(())
763 }
764
765 /// #6207: the race that reverted the first fix — a live foreign TaskManager
766 /// holds the store while its disk state is still empty, so emptiness alone
767 /// cannot tell abandonment from a holder that has not flushed yet. The
768 /// process-owner lock is held for the manager's lifetime, which is what
769 /// distinguishes the two.
770 #[test]
771 fn adoptable_empty_store_refuses_a_live_holder() -> anyhow::Result<()> {
772 let _environment = crate::test_support::lock_test_env();
773 let root = tempfile::tempdir()?;
774 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
775 let _runtime = crate::test_support::EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
776 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
777
778 let store_dir = root.path().join("sessions/held/runtime");
779 drop(crate::runtime_threads::RuntimeThreadStore::open(
780 store_dir.clone(),
781 )?);
782 let binding = crate::runtime_threads::RuntimeStoreBinding {
783 data_dir: store_dir.clone(),
784 execution_scope: "0".repeat(64),
785 };
786 assert!(binding.is_adoptable_empty_store()?);
787
788 let _held = crate::runtime_threads::RuntimeProcessOwnerLock::acquire(&store_dir)?;
789 assert!(
790 binding.has_live_holder()?,
791 "the held owner lock reads as held"
792 );
793 assert!(
794 !binding.is_adoptable_empty_store()?,
795 "a held store must refuse even while its disk state is empty"
796 );
797 drop(_held);
798 assert!(
799 !binding.has_live_holder()?,
800 "releasing the lock releases the hold"
801 );
802 assert!(
803 binding.is_adoptable_empty_store()?,
804 "unheld again: adoptable"
805 );
806 Ok(())
807 }
808
809 /// #6207: scope-pinned automations are recorded outside the store
810 /// directories, so an otherwise empty store with one bound to its scope
811 /// still refuses — adopting it would orphan their scheduled work.
812 #[test]
813 fn adoptable_empty_store_refuses_a_scope_pinned_automation() -> anyhow::Result<()> {
814 let _environment = crate::test_support::lock_test_env();
815 let root = tempfile::tempdir()?;
816 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
817 let _runtime = crate::test_support::EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
818 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
819
820 let store_dir = root.path().join("sessions/pinned/runtime");
821 drop(crate::runtime_threads::RuntimeThreadStore::open(
822 store_dir.clone(),
823 )?);
824 let scope = "ab".repeat(32);
825 let binding = crate::runtime_threads::RuntimeStoreBinding {
826 data_dir: store_dir.clone(),
827 execution_scope: scope.clone(),
828 };
829 assert!(binding.is_adoptable_empty_store()?);
830
831 let automations = AutomationManager::open(root.path().join("automations"))?;
832 let created = automations.create_automation(CreateAutomationRequest {
833 name: "scope fixture".into(),
834 prompt: "local fixture only".into(),
835 rrule: "FREQ=HOURLY;INTERVAL=1".into(),
836 cwds: vec![root.path().into()],
837 model: None,
838 model_provider: None,
839 model_provider_id: None,
840 mode: None,
841 allow_shell: Some(false),
842 trust_mode: Some(false),
843 auto_approve: Some(false),
844 delivery_mode: None,
845 status: Some(AutomationStatus::Paused),
846 })?;
847 automations.edit_automation(&created.id, |record| {
848 let mut record = record.ok_or_else(|| anyhow::anyhow!("fresh automation must exist"))?;
849 record.execution_scope = Some(scope.clone());
850 Ok(Some(record))
851 })?;
852 assert!(
853 binding.has_scope_pinned_automation()?,
854 "the bound automation is visible from the binding's scope"
855 );
856 assert!(
857 !binding.is_adoptable_empty_store()?,
858 "a scope-pinned automation blocks the adopt"
859 );
860 Ok(())
861 }
862
863 /// #6207 end to end: the picker adopts an existing-but-empty unheld store
864 /// into the idle host and persists the repaired binding, the way the
865 /// missing-store path already does.
866 #[tokio::test]
867 async fn picker_adopts_existing_empty_unheld_store() -> anyhow::Result<()> {
868 let _environment = crate::test_support::lock_test_env();
869 let root = tempfile::tempdir()?;
870 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
871 let _runtime = crate::test_support::EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
872 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
873 let sessions = SessionManager::default_location()?;
874 let mut config = fixture_config();
875
876 let store_dir = root.path().join("sessions/previous/runtime");
877 drop(crate::runtime_threads::RuntimeThreadStore::open(
878 store_dir.clone(),
879 )?);
880 let mut saved = crate::session_manager::create_saved_session_with_id_and_mode(
881 "picker-adoptable".into(),
882 &[text_message("user", "retain my work")],
883 "deepseek-v4-pro",
884 root.path(),
885 0,
886 None,
887 None,
888 );
889 saved.metadata.runtime_store = Some(crate::runtime_threads::RuntimeStoreBinding {
890 data_dir: store_dir,
891 execution_scope: "0".repeat(64),
892 });
893 sessions.save_session(&saved)?;
894
895 let mut app = Box::new(create_test_app());
896 let tasks = TaskManager::start(
897 TaskManagerConfig::from_runtime(&config, root.path().into(), None, Some(1)),
898 config.clone(),
899 app.plugin_registry.clone(),
900 "picker-current",
901 None,
902 )
903 .await?;
904 let binding = tasks.session_store_binding().expect("current host");
905 app.runtime_services.task_manager = Some(tasks.clone());
906 app.current_session_id = Some("picker-current".into());
907
908 apply_loaded_session_with_goal(&mut app, &mut config, saved.clone(), None)
909 .map_err(anyhow::Error::msg)?;
910 assert_eq!(app.current_session_id.as_deref(), Some("picker-adoptable"));
911 let durable = sessions.load_session("picker-adoptable")?;
912 assert_eq!(durable.metadata.runtime_store.as_ref(), Some(&binding));
913 assert_eq!(durable.messages, saved.messages);
914 // #6144 P1a: the store the conversation left is set aside where it was
915 // abandoned, not left on disk with nothing pointing at it. The switch
916 // hands that off the UI runtime. Moving the store precedes writing its
917 // manifest, so wait for both rather than treating the move as completion.
918 let abandoned = root.path().join("sessions/previous/runtime");
919 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
920 let set_aside = loop {
921 let manifest = std::fs::read_dir(root.path().join("sessions/.set-aside"))
922 .into_iter()
923 .flatten()
924 .flatten()
925 .map(|run| {
926 std::fs::read_to_string(run.path().join("MANIFEST.jsonl")).unwrap_or_default()
927 })
928 .collect::<String>();
929 if (!abandoned.exists() && manifest.contains("previous"))
930 || std::time::Instant::now() >= deadline
931 {
932 break manifest;
933 }
934 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
935 };
936 assert!(!abandoned.exists());
937 assert!(set_aside.contains("previous"), "{set_aside}");
938 tasks.shutdown_and_wait().await?;
939 Ok(())
940 }
941
942 /// #6144 P1a: an abandoned store another document still binds is kept.
943 #[tokio::test]
944 async fn picker_adoption_keeps_a_store_another_document_binds() -> anyhow::Result<()> {
945 let _environment = crate::test_support::lock_test_env();
946 let root = tempfile::tempdir()?;
947 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path());
948 let _runtime = crate::test_support::EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
949 let _legacy = crate::test_support::EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
950 let sessions = SessionManager::default_location()?;
951 let mut config = fixture_config();
952
953 let store_dir = root.path().join("sessions/previous/runtime");
954 drop(crate::runtime_threads::RuntimeThreadStore::open(
955 store_dir.clone(),
956 )?);
957 let binding = crate::runtime_threads::RuntimeStoreBinding::for_store_dir(&store_dir)?;
958 let mut saved = crate::session_manager::create_saved_session_with_id_and_mode(
959 "picker-adoptable".into(),
960 &[text_message("user", "retain my work")],
961 "deepseek-v4-pro",
962 root.path(),
963 0,
964 None,
965 None,
966 );
967 saved.metadata.runtime_store = Some(binding.clone());
968 sessions.save_session(&saved)?;
969 let mut sibling = crate::session_manager::create_saved_session_with_id_and_mode(
970 "same-host-sibling".into(),
971 &[text_message("user", "saved in the same host")],
972 "deepseek-v4-pro",
973 root.path(),
974 0,
975 None,
976 None,
977 );
978 sibling.metadata.runtime_store = Some(binding);
979 sessions.save_session(&sibling)?;
980
981 let mut app = Box::new(create_test_app());
982 let tasks = TaskManager::start(
983 TaskManagerConfig::from_runtime(&config, root.path().into(), None, Some(1)),
984 config.clone(),
985 app.plugin_registry.clone(),
986 "picker-current",
987 None,
988 )
989 .await?;
990 app.runtime_services.task_manager = Some(tasks.clone());
991 app.current_session_id = Some("picker-current".into());
992 apply_loaded_session_with_goal(&mut app, &mut config, saved, None)
993 .map_err(anyhow::Error::msg)?;
994 // Give the background retirement time to (wrongly) act.
995 tokio::time::sleep(std::time::Duration::from_millis(300)).await;
996 assert!(store_dir.is_dir(), "the sibling still binds it");
997 tasks.shutdown_and_wait().await?;
998 Ok(())
999 }
1000
1000 lines RUST