1use super::*;
4use mj_core::state::MoveOperation;
5
6pub(super) fn load_controller_for_resume(request: &ResumeSessionRequest) -> Result<Controller> {
7 let mut controller = Controller::load()?;
8 if let Some(operation) = crate::database::load_move_operation(&request.session_id)?
9 && matches!(
10 operation.phase,
11 mj_core::state::MovePhase::Failed | mj_core::state::MovePhase::Cancelled
12 )
13 && !operation.queue_admission_started
14 && operation.source_profile_id == request.profile_id
15 && operation.source_target_template_id == request.target_template_id
16 {
17 let record = controller
20 .state
21 .sessions
22 .get_mut(&request.session_id)
23 .context("Move recovery session is missing")?;
24 record.resource_allocation = operation.source_resource_allocation;
25 record.additional_mounts = operation.source_additional_mounts;
26 if let Some(previous) = operation.recovery_session {
27 record.container_cpus = previous.container_cpus;
28 record.container_memory = previous.container_memory;
29 }
30 crate::database::save_resumed_session(record, None)?;
31 }
32 Ok(controller)
33}
34
35impl RuntimeState {
36 pub async fn prepare_move_session(
37 self: &Arc<Self>,
38 selection: MoveSelection,
39 ) -> Result<MovePreparation> {
40 let (source_harness, source_active) = {
41 let controller_owner = self.owner();
42 let controller = controller_owner.controller();
43 let source = controller
44 .state
45 .sessions
46 .get(&selection.session_id)
47 .context("Move session is missing")?;
48 (
49 source.harness_kind,
50 matches!(
51 source.state,
52 SessionState::Running | SessionState::Disconnected
53 ),
54 )
55 };
56 let snapshot = if source_active {
57 crate::controller::move_session::refresh_move_source(
58 &self.session_manager,
59 &selection.session_id,
60 )
61 .await?
62 } else {
63 None
64 };
65 let mut preparation = blocking(move || {
66 let controller = Controller::load()?;
67 mj_core::runtime::block_on(
68 controller.prepare_move_session_controlled(selection, &ProcessExecutor),
69 )?
70 })
71 .await?;
72 preparation.source_unavailable = source_active
73 && snapshot
74 .as_ref()
75 .is_none_or(|snapshot| !snapshot.operational.native_session_is_ready());
76 preparation.active |= preparation.source_unavailable;
77 if let Some(snapshot) = snapshot {
78 let mut operational = snapshot.operational;
79 operational.queued_prompts.clear();
80 operational.checkpoint_barrier = None;
81 preparation.active |= !operational.safe_to_replace(source_harness);
82 }
83 Ok(preparation)
84 }
85
86 pub(crate) fn start_move_session(self: &Arc<Self>, request: MoveSessionRequest) -> Result<()> {
88 self.admit_move_session(request).map(|_| ())
89 }
90
91 fn admit_move_session(self: &Arc<Self>, request: MoveSessionRequest) -> Result<LifecycleWatch> {
92 let selection = request.preparation.selection.clone();
93 let operation_id = request.preparation.operation_id.clone();
94 let key =
97 serde_json::to_string(&(&selection, request.queue, request.acknowledge_interruption))?;
98 let session_id = selection.session_id.clone();
99 let result = self.admit_lifecycle(
100 session_id.clone(),
101 LifecycleKind::Move,
102 super::lifecycle::LifecycleStart {
103 resume_workspace_id: None,
104 request_key: Some(key),
105 create_control: None,
106 phase: LifecyclePhase::Executing,
107 move_operation_id: Some(operation_id.clone()),
108 },
109 move |state, session_id, cancelled| async move {
110 let result = state
111 .clone()
112 .run_move_controller_work(
113 session_id.clone(),
114 cancelled,
115 |mut controller, executor, manager| async move {
116 controller
117 .move_session_managed_controlled(request, &executor, &manager)
118 .await
119 },
120 )
121 .await
122 .map(DaemonLifecycleResult::Move);
123 if result.is_err() {
124 state
125 .tell_live_parent_about_stopped_subagents(&session_id)
126 .await;
127 }
128 result
129 },
130 )?;
131 self.set_lifecycle_resume_destination(
132 &session_id,
133 selection.profile_id.clone().unwrap_or_default(),
134 selection.target_template_id.clone().unwrap_or_default(),
135 );
136 Ok(result)
137 }
138
139 pub async fn move_session(
140 self: &Arc<Self>,
141 request: MoveSessionRequest,
142 ) -> Result<MoveOutcome> {
143 let selection = request.preparation.selection.clone();
144 let operation_id = request.preparation.operation_id.clone();
145 let session_id = selection.session_id.clone();
146 let result = self.admit_move_session(request)?;
147 let channel = result.clone();
148 let result = Self::wait_lifecycle_result(result).await;
149 self.remove_completed_lifecycle(&channel);
150 match result {
151 Ok(DaemonLifecycleResult::Move(outcome)) => Ok(outcome),
152 Ok(_) => bail!("move returned an unrelated lifecycle result"),
153 Err(error) => Ok(MoveOutcome {
154 operation_id, session_id, profile_id: selection.profile_id.unwrap_or_default(),
155 target_template_id: selection.target_template_id.unwrap_or_default(), outcome: "failed".into(),
156 error: Some(format!("{error:#}")), recovery: Some("Inspect session status and prepare Move again; any verified checkpoint is retained.".into()),
157 }),
158 }
159 }
160
161 pub(super) fn recover_moves(
162 self: &Arc<Self>,
163 operations: Vec<MoveOperation>,
164 ) -> Result<BTreeSet<String>> {
165 let mut owned = BTreeSet::new();
166 for operation in &operations {
167 crate::controller::move_session::restore_move_queue_hold(operation);
168 }
169 for operation in operations.into_iter().filter(|op| {
170 op.is_active()
171 || self
172 .owner()
173 .controller()
174 .state
175 .sessions
176 .get(&op.selection.session_id)
177 .is_some_and(|session| {
178 matches!(
179 session.state,
180 SessionState::Closing | SessionState::Destroying
181 )
182 })
183 }) {
184 let id = operation.selection.session_id.clone();
185 let key = format!("recovery:{}", operation.operation_id);
186 let phase = if operation.recovery_session.is_some() {
187 LifecyclePhase::MovingDestination
188 } else {
189 LifecyclePhase::Executing
190 };
191 let result = self.admit_lifecycle(
192 id.clone(),
193 LifecycleKind::Move,
194 super::lifecycle::LifecycleStart {
195 resume_workspace_id: None,
196 request_key: Some(key),
197 create_control: None,
198 phase,
199 move_operation_id: None,
200 },
201 move |state, session_id, cancelled| async move {
202 state
203 .run_move_controller_work(
204 session_id,
205 cancelled,
206 |mut controller, executor, manager| async move {
207 controller
208 .recover_move_managed_controlled(operation, &executor, &manager)
209 .await
210 },
211 )
212 .await
213 .map(DaemonLifecycleResult::Move)
214 },
215 )?;
216 owned.insert(id.clone());
217 let state = self.clone();
218 tokio::spawn(async move {
219 let channel = result.clone();
220 match Self::wait_lifecycle_result(result).await {
221 Ok(DaemonLifecycleResult::Move(outcome))
222 if !matches!(outcome.outcome.as_str(), "completed" | "interrupted") =>
223 {
224 state.push_notice(
225 &id,
226 outcome
227 .error
228 .unwrap_or_else(|| "Move recovery needs attention".into()),
229 );
230 }
231 Err(error) => {
232 state.push_notice(&id, format!("Move recovery failed: {error:#}"))
233 }
234 _ => {}
235 }
236 state.remove_completed_lifecycle(&channel);
237 });
238 }
239 Ok(owned)
240 }
241
242 pub(super) fn resume_move_destination_cleanups(self: &Arc<Self>, immediately: bool) {
243 let ids = {
244 let owner = self.owner();
245 owner
246 .committed()
247 .into_iter()
248 .flat_map(|committed| committed.moves.values())
249 .filter(|op| {
250 op.prepared_destination.as_ref().is_some_and(|d| {
251 matches!(
252 d.state,
253 mj_core::state::PreparedDestinationState::CleanupPending { .. }
254 )
255 })
256 })
257 .filter(|op| {
258 !owner
259 .lifecycle
260 .get(&op.selection.session_id)
261 .is_some_and(|active| active.is_running())
262 })
263 .filter(|op| {
264 immediately
265 || chrono::DateTime::parse_from_rfc3339(&op.updated_at)
266 .map(|at| {
267 (chrono::Utc::now() - at.with_timezone(&chrono::Utc)).num_seconds()
268 >= 30
269 })
270 .unwrap_or(true)
271 })
272 .map(|op| op.selection.session_id.clone())
273 .collect::<Vec<_>>()
274 };
275 for id in ids {
276 let result = self.start_or_join_lifecycle(
277 id.clone(),
278 LifecycleKind::Cleanup,
279 |state, id, cancelled| async move {
280 state
281 .run_move_controller_work(
282 id.clone(),
283 cancelled,
284 move |controller, executor, _manager| async move {
285 let mut operation = crate::database::load_move_operation(&id)?
286 .context("Move cleanup intent missing")?;
287 executor.begin_resumable_move_work()?;
288 let result = controller
289 .cleanup_prepared_move_destination(&mut operation, &executor);
290 operation.updated_at = chrono::Utc::now().to_rfc3339();
291 if result.is_ok() {
292 crate::controller::move_session::record_finished_move_recovery(
293 &controller.state,
294 &mut operation,
295 )?;
296 } else {
297 crate::database::save_move_operation(&operation)?;
298 }
299 executor.end_resumable_move_work()?;
300 result?;
301 Ok(MoveOutcome {
302 operation_id: operation.operation_id,
303 session_id: id,
304 profile_id: operation.selection.profile_id.unwrap_or_default(),
305 target_template_id: operation
306 .selection
307 .target_template_id
308 .unwrap_or_default(),
309 outcome: "cleaned_up".into(),
310 error: None,
311 recovery: None,
312 })
313 },
314 )
315 .await
316 .map(DaemonLifecycleResult::Move)
317 },
318 );
319 match result {
320 Ok(result) => {
321 let state = self.clone();
322 tokio::spawn(async move {
323 let channel = result.clone();
324 if let Err(error) = Self::wait_lifecycle_result(result).await {
325 tracing::warn!(session_id=%id, %error, "EC2 Move destination cleanup remains pending; retry in 30s");
326 state.push_notice(
327 &id,
328 format!(
329 "EC2 destination cleanup will retry automatically: {error:#}"
330 ),
331 );
332 }
333 state.remove_completed_lifecycle(&channel);
334 });
335 }
336 Err(error) => {
337 tracing::warn!(session_id=%id, %error, "EC2 Move destination cleanup could not start")
338 }
339 }
340 }
341 }
342
343 pub(super) async fn run_move_controller_work<W, Fut>(
357 self: Arc<Self>,
358 session_id: String,
359 cancelled: Arc<AtomicBool>,
360 work: W,
361 ) -> Result<MoveOutcome>
362 where
363 W: FnOnce(
364 Controller,
365 DaemonStageReportingExecutor<CancellableProcessExecutor>,
366 SessionManagerControl,
367 ) -> Fut
368 + Send
369 + 'static,
370 Fut: std::future::Future<Output = Result<MoveOutcome>>,
371 {
372 tokio::task::spawn_blocking(move || {
373 let reserving = std::time::Instant::now();
374 let _reservation =
375 reserve_recovery_or_cancel(&self.recovery_observer, &session_id, &cancelled)?;
376 tracing::info!(
377 %session_id,
378 phase = "recovery reservation",
379 elapsed_ms = reserving.elapsed().as_millis() as u64,
380 "move phase finished"
381 );
382 let loading = std::time::Instant::now();
383 let controller = (self.controller_loader)()?;
384 tracing::info!(
385 %session_id,
386 phase = "load controller state",
387 elapsed_ms = loading.elapsed().as_millis() as u64,
388 "move phase finished"
389 );
390 let manager = self.session_manager.clone();
391 let executor = DaemonStageReportingExecutor::new(
392 CancellableProcessExecutor::new(cancelled),
393 self,
394 session_id,
395 );
396 mj_core::runtime::block_on(Box::pin(work(controller, executor, manager)))?
397 })
398 .await
399 .context("Move controller task failed")?
400 }
401}
402
403#[cfg(all(test, unix))]
404mod admission_tests {
405 use super::*;
406 use crate::controller::test_support::{IsolatedTest, test_name};
407
408 const BOUND: Duration = Duration::from_secs(5);
409
410 fn metadata() -> DaemonMetadata {
411 DaemonMetadata {
412 protocol_version: PROTOCOL_VERSION,
413 pid: 1,
414 address: SocketAddr::from((Ipv4Addr::LOCALHOST, 0)),
415 token: "right-token".into(),
416 started_at: "now".into(),
417 build_version: "test".into(),
418 }
419 }
420
421 fn outcome(status: &str) -> MoveOutcome {
422 MoveOutcome {
423 operation_id: "move-operation".into(),
424 session_id: "moving-session".into(),
425 profile_id: String::new(),
426 target_template_id: String::new(),
427 outcome: status.into(),
428 error: None,
429 recovery: None,
430 }
431 }
432
433 fn held_command(directory: &Path, purpose: &str) -> CommandSpec {
435 CommandSpec::new(
436 "sh",
437 [
438 "-c".to_owned(),
439 format!(
440 "echo $$ > '{0}/pid'; while [ ! -e '{0}/release' ]; do sleep 0.05; done",
441 directory.display()
442 ),
443 ],
444 )
445 .purpose(purpose)
446 }
447
448 struct ReleaseOnDrop(PathBuf);
451
452 impl Drop for ReleaseOnDrop {
453 fn drop(&mut self) {
454 let _ = std::fs::write(&self.0, b"");
455 }
456 }
457
458 async fn wait_for_file(path: &Path) {
459 let deadline = std::time::Instant::now() + BOUND;
460 while !path.exists() {
461 assert!(
462 std::time::Instant::now() < deadline,
463 "{} never appeared",
464 path.display()
465 );
466 tokio::time::sleep(Duration::from_millis(10)).await;
467 }
468 }
469
470 fn process_is_running(pid: i32) -> bool {
471 unsafe { libc::kill(pid, 0) == 0 }
473 }
474
475 async fn handoff(state: &Arc<RuntimeState>) -> DaemonReply {
476 super::super::actions::handle_action(
477 DaemonAction::PrepareUpgrade,
478 &metadata(),
479 state,
480 &CancellationToken::new(),
481 )
482 .await
483 .unwrap()
484 }
485
486 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
493 async fn a_handoff_does_not_wait_for_a_slow_path_move_copy() {
494 const NAME: &str = "a_handoff_does_not_wait_for_a_slow_path_move_copy";
495 const CHILD: &str = "MJ_TEST_SLOW_MOVE_HANDOFF";
496 if std::env::var_os(CHILD).is_none() {
497 let root = tempfile::tempdir().unwrap();
498 IsolatedTest::new(test_name(module_path!(), NAME))
499 .env(CHILD, "1")
500 .env("MJ_INSTANCE", "slow-move-handoff")
501 .isolated_store(root.path())
502 .run();
503 return;
504 }
505 let directory = tempfile::tempdir().unwrap();
506 let _release = ReleaseOnDrop(directory.path().join("release"));
507 let copy = held_command(directory.path(), "copy Move workspace");
508 let state = super::super::tests::test_runtime_state();
509 let result = state
510 .start_or_join_lifecycle(
511 "moving-session".into(),
512 LifecycleKind::Move,
513 move |state, session_id, cancelled| async move {
514 state
515 .run_move_controller_work(
516 session_id,
517 cancelled,
518 |_, executor, _| async move {
519 executor.begin_resumable_move_work()?;
520 let copied = executor.execute(©);
521 let resumed = executor.end_resumable_move_work();
522 if copied.is_ok() && resumed.is_ok() {
523 return Ok(outcome("completed"));
524 }
525 ensure!(
526 !crate::upgrade::gate().is_open(),
527 "the copy stopped without a handoff"
528 );
529 Ok(outcome("interrupted"))
530 },
531 )
532 .await
533 .map(DaemonLifecycleResult::Move)
534 },
535 )
536 .unwrap();
537 wait_for_file(&directory.path().join("pid")).await;
538 let pid: i32 = std::fs::read_to_string(directory.path().join("pid"))
539 .unwrap()
540 .trim()
541 .parse()
542 .unwrap();
543
544 let labels = crate::upgrade::active_labels();
545 assert!(
546 !labels
547 .iter()
548 .any(|label| label == "session lifecycle" || label == "database operation"),
549 "the copy holds handoff admission: {labels:?}"
550 );
551 let deadline = std::time::Instant::now() + BOUND;
552 while !matches!(handoff(&state).await, DaemonReply::Done) {
553 assert!(
554 std::time::Instant::now() < deadline,
555 "the handoff waited for the copy: {:?}",
556 crate::upgrade::active_labels()
557 );
558 tokio::time::sleep(Duration::from_millis(10)).await;
559 }
560 assert!(
561 process_is_running(pid),
562 "the handoff itself cancels nothing"
563 );
564
565 let stopping = std::time::Instant::now();
567 state.cancel_and_wait_lifecycles().await.unwrap();
568 assert!(stopping.elapsed() < BOUND);
569 match RuntimeState::wait_lifecycle_result(result).await.unwrap() {
570 DaemonLifecycleResult::Move(outcome) => assert_eq!(outcome.outcome, "interrupted"),
571 _ => panic!("a Move lifecycle returns a Move outcome"),
572 }
573 let deadline = std::time::Instant::now() + BOUND;
574 while process_is_running(pid) {
575 assert!(
576 std::time::Instant::now() < deadline,
577 "the cancelled copy kept running"
578 );
579 tokio::time::sleep(Duration::from_millis(10)).await;
580 }
581 assert!(!directory.path().join("release").exists());
582 }
583
584 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
587 async fn a_handoff_waits_for_a_fast_path_move() {
588 const NAME: &str = "a_handoff_waits_for_a_fast_path_move";
589 const CHILD: &str = "MJ_TEST_FAST_MOVE_HANDOFF";
590 if std::env::var_os(CHILD).is_none() {
591 let root = tempfile::tempdir().unwrap();
592 IsolatedTest::new(test_name(module_path!(), NAME))
593 .env(CHILD, "1")
594 .env("MJ_INSTANCE", "fast-move-handoff")
595 .isolated_store(root.path())
596 .run();
597 return;
598 }
599 let directory = tempfile::tempdir().unwrap();
600 let _release = ReleaseOnDrop(directory.path().join("release"));
601 let swap = held_command(directory.path(), "swap the harness in place");
602 let state = super::super::tests::test_runtime_state();
603 let result = state
604 .start_or_join_lifecycle(
605 "moving-session".into(),
606 LifecycleKind::Move,
607 move |state, session_id, cancelled| async move {
608 state
609 .run_move_controller_work(
610 session_id,
611 cancelled,
612 |_, executor, _| async move {
613 executor.execute(&swap)?;
614 Ok(outcome("completed"))
615 },
616 )
617 .await
618 .map(DaemonLifecycleResult::Move)
619 },
620 )
621 .unwrap();
622 wait_for_file(&directory.path().join("pid")).await;
623 assert!(
624 crate::upgrade::active_labels()
625 .iter()
626 .any(|label| label == "session lifecycle")
627 );
628 assert!(matches!(handoff(&state).await, DaemonReply::UpgradePending));
629 std::fs::write(directory.path().join("release"), b"").unwrap();
630 match RuntimeState::wait_lifecycle_result(result).await.unwrap() {
631 DaemonLifecycleResult::Move(outcome) => assert_eq!(outcome.outcome, "completed"),
632 _ => panic!("a Move lifecycle returns a Move outcome"),
633 }
634 let deadline = std::time::Instant::now() + BOUND;
635 loop {
636 if matches!(handoff(&state).await, DaemonReply::Done) {
637 break;
638 }
639 assert!(
640 std::time::Instant::now() < deadline,
641 "the finished Move still holds the handoff"
642 );
643 tokio::time::sleep(Duration::from_millis(10)).await;
644 }
645 }
646}