Skip to main content

mj_controller/daemon/
session_move.rs

1//! Daemon-owned move admission, supervision, and restart reconciliation.
2
3use 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        // None in ordinary Resume means inheritance. Restore that inheritance
18        // baseline first so a partial conversion cannot change source sizing.
19        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    /// Reserve lifecycle ownership before a web request is acknowledged.
87    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        // Preparation resolves inherited settings. Compare those settings,
95        // never just the verb or a freshly generated preparation identifier.
96        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    /// The controller half of every Move lifecycle, started or recovered.
344    ///
345    /// It takes no upgrade admission of its own; the lifecycle's hold is the
346    /// only one. That is what lets the slow path work as intended: it
347    /// releases the lifecycle's hold around each resumable workspace copy
348    /// (`begin_resumable_move_work`), so a handoff does not wait for a copy,
349    /// the closing daemon cancels it, and the next daemon resumes the Move
350    /// from its durable phase. An in-place Move copies nothing and never
351    /// releases the hold, so a handoff waits for it to finish.
352    ///
353    /// The Move's controller future holds a standard mutex guard across an
354    /// await, so it is not `Send` and cannot be the lifecycle future itself.
355    /// It is built and awaited on a blocking-pool thread of its own.
356    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    /// A command that records its process id and waits for `release`.
434    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    /// Releases a held command when a test ends, so a failing assertion does
449    /// not leave the runtime waiting for it forever.
450    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        // SAFETY: signal 0 only checks that the process exists.
472        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    /// A slow-path Move holds no handoff admission while it copies a
487    /// workspace: a handoff completes without waiting for the copy, the stop
488    /// that follows kills the copy, and the Move reports that it was
489    /// interrupted for the next daemon to resume, as `finish_move_result`
490    /// does for a Move with a workspace transfer.
491    // Hard-won: 571e8aa0: a handoff waited through a resumable Move copy because duplicate admission stayed held.
492    #[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(&copy);
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        // The closing daemon cancels the lifecycles it leaves behind.
566        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    /// An in-place Move copies nothing and keeps its admission for its whole
585    /// run, so a handoff waits for it and then proceeds.
586    #[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}