Skip to main content

mj_controller/daemon/
restart.rs

1use super::*;
2
3impl RuntimeState {
4    /// Restart a session while keeping its target when the retained environment
5    /// passes the same checks used by an in-place Move.
6    pub async fn restart_session(self: &Arc<Self>, session_id: String) -> Result<()> {
7        self.restart_session_controlled(session_id, None).await
8    }
9
10    pub async fn restart_session_controlled(
11        self: &Arc<Self>,
12        session_id: String,
13        control: Option<CreateSessionControl>,
14    ) -> Result<()> {
15        let joining_restart = self
16            .owner()
17            .lifecycle
18            .get(&session_id)
19            .is_some_and(|active| active.kind == LifecycleKind::Restart && active.is_running());
20        if !joining_restart {
21            self.wait_before_close(&session_id).await?;
22            self.wait_for_active_suspend(&session_id).await?;
23            let cleanup_running = self
24                .owner()
25                .lifecycle
26                .get(&session_id)
27                .is_some_and(|active| active.kind == LifecycleKind::Cleanup && active.is_running());
28            if cleanup_running {
29                self.wait_for_deferred_cleanup(&session_id).await?;
30            }
31        }
32        let workspace_id = blocking({
33            let session_id = session_id.clone();
34            move || {
35                Controller::load()?
36                    .state
37                    .sessions
38                    .get(&session_id)
39                    .map(|session| session.workspace_id.clone())
40                    .with_context(|| format!("unknown session {session_id}"))
41            }
42        })
43        .await?;
44        let _workspace_admission = self.workspace_resume_gate(&workspace_id).read_owned().await;
45
46        let operation = self.start_or_join_lifecycle_controlled(
47            session_id.clone(),
48            LifecycleKind::Restart,
49            Some(workspace_id),
50            None,
51            control,
52            move |state, session_id, cancelled| async move {
53                state.restart_admitted(session_id, cancelled, None).await
54            },
55        )?;
56        let completed = operation.clone();
57        let outcome = Self::wait_lifecycle_result(operation).await;
58        self.remove_completed_lifecycle(&completed);
59        outcome?;
60        Ok(())
61    }
62
63    async fn wait_for_active_suspend(self: &Arc<Self>, session_id: &str) -> Result<()> {
64        let pending = {
65            let owner = self.owner();
66            owner.lifecycle.get(session_id).and_then(|active| {
67                (active.kind == LifecycleKind::Suspend && active.is_running())
68                    .then(|| active.result.clone())
69            })
70        };
71        if let Some(pending) = pending {
72            let completed = pending.clone();
73            Self::wait_lifecycle_result(pending).await?;
74            self.remove_completed_lifecycle(&completed);
75        }
76        Ok(())
77    }
78
79    /// Recover an accepted Restart after a daemon replacement. Its durable
80    /// phase distinguishes a completed restore from a request that had only
81    /// begun stopping the old worker.
82    pub(super) fn recover_restarts(
83        self: &Arc<Self>,
84        intents: Vec<crate::database::SessionRestartIntent>,
85    ) -> Result<BTreeSet<String>> {
86        let mut owned = BTreeSet::new();
87        for intent in intents {
88            let session_id = intent.session_id.clone();
89            let current_state = self
90                .owner()
91                .controller()
92                .state
93                .sessions
94                .get(&session_id)
95                .map(|record| record.state);
96            if restart_intent_was_superseded(current_state, intent.phase) {
97                crate::database::cancel_session_restart(&session_id)?;
98                tracing::info!(
99                    %session_id,
100                    state = ?current_state,
101                    "discarded stale Restart intent superseded by a later lifecycle"
102                );
103                continue;
104            }
105            let key = format!("recovery:{}", intent.operation_id);
106            let result = self.admit_lifecycle(
107                session_id.clone(),
108                LifecycleKind::Restart,
109                super::lifecycle::LifecycleStart {
110                    resume_workspace_id: None,
111                    request_key: Some(key),
112                    create_control: None,
113                    phase: LifecyclePhase::Executing,
114                    move_operation_id: None,
115                },
116                move |state, session_id, cancelled| async move {
117                    state
118                        .restart_admitted(session_id, cancelled, Some(intent))
119                        .await
120                },
121            )?;
122            owned.insert(session_id.clone());
123            let state = self.clone();
124            tokio::spawn(async move {
125                let channel = result.clone();
126                match Self::wait_lifecycle_result(result).await {
127                    Ok(DaemonLifecycleResult::Done) => {}
128                    Ok(_) => unreachable!("restart recovery returned another lifecycle result"),
129                    Err(error) => {
130                        tracing::warn!(
131                            %session_id,
132                            error = format!("{error:#}"),
133                            "interrupted Restart needs attention"
134                        );
135                        state.push_notice(
136                            &session_id,
137                            format!("Restart could not recover automatically: {error:#}"),
138                        );
139                    }
140                }
141                state.remove_completed_lifecycle(&channel);
142            });
143        }
144        Ok(owned)
145    }
146
147    async fn restart_admitted(
148        self: &Arc<Self>,
149        session_id: String,
150        cancelled: Arc<AtomicBool>,
151        recovered: Option<crate::database::SessionRestartIntent>,
152    ) -> Result<DaemonLifecycleResult> {
153        self.request_close(&session_id);
154        let result = self
155            .restart_admitted_inner(&session_id, cancelled, recovered)
156            .await;
157        self.clear_close_request(&session_id);
158        result
159    }
160
161    async fn restart_admitted_inner(
162        self: &Arc<Self>,
163        session_id: &str,
164        cancelled: Arc<AtomicBool>,
165        recovered: Option<crate::database::SessionRestartIntent>,
166    ) -> Result<DaemonLifecycleResult> {
167        let _recovery_reservation = tokio::task::spawn_blocking({
168            let observer = self.recovery_observer.clone();
169            let session_id = session_id.to_owned();
170            let cancelled = cancelled.clone();
171            move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
172        })
173        .await
174        .context("reserve recovery for daemon restart task")??;
175        let mut controller = tokio::task::spawn_blocking(Controller::load)
176            .await
177            .context("load controller for daemon restart task")??;
178        let move_operation = crate::database::load_move_operation(session_id)?;
179        ensure!(
180            !move_operation
181                .as_ref()
182                .is_some_and(mj_core::state::MoveOperation::holds_source_environment),
183            "a Move still owns this environment; retry Move on its recorded destination"
184        );
185
186        let mut intent = match recovered {
187            Some(intent) => intent,
188            None => {
189                let operation_id = new_command_id("restart")?;
190                blocking({
191                    let session_id = session_id.to_owned();
192                    let cancelled = cancelled.clone();
193                    move || {
194                        crate::database::begin_session_restart(
195                            &session_id,
196                            &operation_id,
197                            cancelled,
198                        )
199                    }
200                })
201                .await?
202            }
203        };
204        ensure!(
205            intent.session_id == session_id,
206            "Restart intent belongs to another session"
207        );
208        tracing::info!(
209            %session_id,
210            operation_id = %intent.operation_id,
211            phase = ?intent.phase,
212            "starting daemon-owned Restart"
213        );
214
215        // A completed in-place restore or fallback resume can only be Running
216        // after the durable phase advanced beyond Stopping. The latter remains
217        // retryable if the daemon stopped before it sealed a checkpoint.
218        if intent.phase != crate::database::SessionRestartPhase::Stopping
219            && controller
220                .state
221                .sessions
222                .get(session_id)
223                .is_some_and(|record| record.state == SessionState::Running)
224        {
225            finish_restart_intent(session_id, &intent.operation_id).await?;
226            tracing::info!(%session_id, "completed recovered Restart intent for an already running session");
227            return Ok(DaemonLifecycleResult::Done);
228        }
229
230        let executor = DaemonStageReportingExecutor::new(
231            CancellableProcessExecutor::new(cancelled.clone()),
232            self.clone(),
233            session_id.to_owned(),
234        );
235
236        if intent.phase == crate::database::SessionRestartPhase::Fallback {
237            ensure_target_unavailable_for_fallback(&controller, session_id, &executor)?;
238            tracing::info!(%session_id, "Restart selected suspend-and-resume fallback");
239            self.restart_fallback(session_id, &mut controller, &executor, &cancelled)
240                .await?;
241            finish_restart_intent(session_id, &intent.operation_id).await?;
242            return Ok(DaemonLifecycleResult::Done);
243        }
244
245        let mut record = controller
246            .state
247            .sessions
248            .get(session_id)
249            .with_context(|| format!("unknown session {session_id}"))?
250            .clone();
251
252        if !restart_checkpoint_matches(&record, &intent) {
253            let target = restart_target_availability(&controller, session_id, &executor)?;
254            if fallback_is_safe(target) {
255                ensure!(
256                    !intent.checkpoint_started,
257                    "Restart target became unavailable after checkpointing began; retaining the target and refusing destructive fallback"
258                );
259                advance_restart_phase(
260                    session_id,
261                    &intent.operation_id,
262                    crate::database::SessionRestartPhase::Fallback,
263                )
264                .await?;
265                tracing::info!(
266                    %session_id,
267                    ?target,
268                    "Restart selected suspend-and-resume because the target was unavailable before checkpointing"
269                );
270                self.restart_fallback(session_id, &mut controller, &executor, &cancelled)
271                    .await?;
272                finish_restart_intent(session_id, &intent.operation_id).await?;
273                return Ok(DaemonLifecycleResult::Done);
274            }
275
276            blocking({
277                let session_id = session_id.to_owned();
278                let operation_id = intent.operation_id.clone();
279                move || {
280                    crate::database::mark_session_restart_checkpoint_started(
281                        &session_id,
282                        &operation_id,
283                    )
284                }
285            })
286            .await?;
287            tracing::info!(%session_id, "Restart selected in-place checkpoint and restore");
288            controller
289                .suspend_session_for_restart(
290                    session_id,
291                    &executor,
292                    &self.session_manager,
293                    Some(self.stop_subagents_before_close(session_id)),
294                )
295                .await
296                .context("checkpoint and seal session for Restart")?;
297            controller = tokio::task::spawn_blocking(Controller::load)
298                .await
299                .context("reload controller after Restart checkpoint")??;
300            record = controller
301                .state
302                .sessions
303                .get(session_id)
304                .with_context(|| format!("unknown session {session_id}"))?
305                .clone();
306            let checkpoint_sha256 = record
307                .checkpoint
308                .as_ref()
309                .context("Restart checkpoint completed without a checkpoint")?
310                .sha256
311                .clone();
312            let checkpoint_sha256_for_record = checkpoint_sha256.clone();
313            blocking({
314                let session_id = session_id.to_owned();
315                let operation_id = intent.operation_id.clone();
316                move || {
317                    crate::database::record_session_restart_checkpoint(
318                        &session_id,
319                        &operation_id,
320                        &checkpoint_sha256_for_record,
321                    )
322                }
323            })
324            .await?;
325            intent.checkpoint_sha256 = Some(checkpoint_sha256);
326        }
327
328        record = controller
329            .state
330            .sessions
331            .get(session_id)
332            .with_context(|| format!("unknown session {session_id}"))?
333            .clone();
334        ensure!(
335            restart_checkpoint_matches(&record, &intent),
336            "Restart checkpoint identity is absent or no longer matches the sealed session checkpoint; retaining the target"
337        );
338        advance_restart_phase(
339            session_id,
340            &intent.operation_id,
341            crate::database::SessionRestartPhase::RestoringInPlace,
342        )
343        .await?;
344        match controller
345            .restore_session_in_place_for_restart(
346                session_id,
347                &record.last_profile,
348                &record.target_template_id,
349                &executor,
350            )
351            .await
352        {
353            Ok(Ok(_)) => {
354                tracing::info!(%session_id, "Restart restored the worker in the retained target");
355                finish_restart_intent(session_id, &intent.operation_id).await?;
356                Ok(DaemonLifecycleResult::Done)
357            }
358            Ok(Err(crate::controller::InPlaceRestartError::Preflight(error))) => {
359                Err(error.context("Restart preflight failed; retained target was not cleaned up"))
360            }
361            Ok(Err(crate::controller::InPlaceRestartError::Restore(error)))
362            | Ok(Err(crate::controller::InPlaceRestartError::Cancelled(error))) => Err(error),
363            Err(error) => Err(error.context("restore Restart in the retained target")),
364        }
365    }
366
367    async fn restart_fallback(
368        self: &Arc<Self>,
369        session_id: &str,
370        controller: &mut Controller,
371        executor: &(impl CommandExecutor + Sync),
372        cancelled: &AtomicBool,
373    ) -> Result<()> {
374        ensure!(
375            !cancelled.load(Ordering::Acquire),
376            "Restart cancelled before fallback cleanup"
377        );
378        ensure_target_unavailable_for_fallback(controller, session_id, executor)?;
379        let outcome = self
380            .suspend_with_loaded_controller(session_id, controller, executor, true)
381            .await?;
382        match outcome {
383            DaemonLifecycleResult::Done => {}
384            DaemonLifecycleResult::DeferredCleanup => {
385                ensure_target_unavailable_for_fallback(controller, session_id, executor)?;
386                controller.cleanup_stopped_target(session_id, executor)?;
387            }
388            DaemonLifecycleResult::Move(_) | DaemonLifecycleResult::Park(_) => {
389                unreachable!("Restart fallback suspension returned another lifecycle result")
390            }
391        }
392        ensure!(
393            !cancelled.load(Ordering::Acquire),
394            "Restart cancelled before fallback resume"
395        );
396        let record = controller
397            .state
398            .sessions
399            .get(session_id)
400            .with_context(|| format!("unknown session {session_id}"))?
401            .clone();
402        let materialized = controller
403            .resume_session_controlled_with_repository_preflight(
404                session_id,
405                &record.last_profile,
406                &record.target_template_id,
407                SessionResumeOptions {
408                    additional_mounts: Some(record.additional_mounts),
409                    resource_allocation: record.resource_allocation,
410                    discard_queue: false,
411                },
412                None,
413                executor,
414            )
415            .await?;
416        let _ = materialized;
417        Ok(())
418    }
419}
420
421#[derive(Debug, Clone, Copy, PartialEq, Eq)]
422pub(super) enum RestartTargetAvailability {
423    Reachable,
424    Missing,
425    Unreachable,
426}
427
428pub(super) fn restart_checkpoint_matches(
429    record: &SessionRecord,
430    intent: &crate::database::SessionRestartIntent,
431) -> bool {
432    record.target.is_some()
433        && record.checkpoint.as_ref().is_some_and(|checkpoint| {
434            intent.checkpoint_sha256.as_deref() == Some(checkpoint.sha256.as_str())
435        })
436        && matches!(
437            record.state,
438            SessionState::Closing
439                | SessionState::Error
440                | SessionState::Stopped
441                | SessionState::Provisioning
442        )
443}
444
445pub(super) fn restart_intent_was_superseded(
446    state: Option<SessionState>,
447    phase: crate::database::SessionRestartPhase,
448) -> bool {
449    state.is_none()
450        || matches!(state, Some(SessionState::Destroying | SessionState::Parked))
451        || (state == Some(SessionState::Stopped)
452            && phase != crate::database::SessionRestartPhase::Fallback)
453}
454
455pub(super) fn restart_target_availability(
456    controller: &Controller,
457    session_id: &str,
458    executor: &impl CommandExecutor,
459) -> Result<RestartTargetAvailability> {
460    let record = controller
461        .state
462        .sessions
463        .get(session_id)
464        .with_context(|| format!("unknown session {session_id}"))?;
465    let Some(locator) = record.target.as_ref() else {
466        return Ok(RestartTargetAvailability::Missing);
467    };
468    let locator = crate::controller::backend_locator(locator, record, &controller.config)?;
469
470    if let Some(plan) = crate::targets::target_recovery_plan(&locator, session_id)? {
471        let output = executor
472            .execute(&plan.exists)
473            .context("check retained target before Restart")?;
474        match output.status {
475            0 => match crate::targets::ensure_recovery_target_running(executor, Some(&plan))? {
476                crate::targets::TargetRecoveryOutcome::Missing => {
477                    Ok(RestartTargetAvailability::Missing)
478                }
479                crate::targets::TargetRecoveryOutcome::AlreadyRunning
480                | crate::targets::TargetRecoveryOutcome::Started
481                | crate::targets::TargetRecoveryOutcome::NotRequired => {
482                    Ok(RestartTargetAvailability::Reachable)
483                }
484            },
485            1 => Ok(RestartTargetAvailability::Missing),
486            125 => Ok(RestartTargetAvailability::Unreachable),
487            255 if locator_is_remote(&locator) => Ok(RestartTargetAvailability::Unreachable),
488            status => bail!(
489                "could not confirm retained target availability before Restart (status {status}): {}",
490                String::from_utf8_lossy(&output.stderr).trim()
491            ),
492        }
493    } else if let mj_core::targets::TargetLocator::AppleContainer { container_id, .. } = &locator {
494        let list = CommandSpec::new("container", ["list", "--all", "--quiet"])
495            .purpose("check retained Apple container before Restart");
496        let output = executor
497            .execute(&list)
498            .context("list retained Apple target before Restart")?;
499        ensure!(
500            output.status == 0,
501            "could not confirm Apple target availability before Restart (status {}): {}",
502            output.status,
503            String::from_utf8_lossy(&output.stderr).trim()
504        );
505        let listed = String::from_utf8(output.stdout)
506            .context("decode Apple container list before Restart")?;
507        if listed.lines().any(|id| id.trim() == container_id) {
508            Ok(RestartTargetAvailability::Reachable)
509        } else {
510            Ok(RestartTargetAvailability::Missing)
511        }
512    } else {
513        let probe = crate::targets::command_on_locator(
514            &locator,
515            session_id,
516            vec!["true".into()],
517            "probe retained target before Restart",
518        )?;
519        let output = executor
520            .execute(&probe)
521            .context("probe retained target before Restart")?;
522        match output.status {
523            0 => Ok(RestartTargetAvailability::Reachable),
524            255 if locator_is_remote(&locator) => Ok(RestartTargetAvailability::Unreachable),
525            status => bail!(
526                "could not confirm retained target availability before Restart (status {status}); refusing destructive fallback: {}",
527                String::from_utf8_lossy(&output.stderr).trim()
528            ),
529        }
530    }
531}
532
533fn locator_is_remote(locator: &mj_core::targets::TargetLocator) -> bool {
534    matches!(
535        locator,
536        mj_core::targets::TargetLocator::AwsEc2 { .. }
537            | mj_core::targets::TargetLocator::SshBare { .. }
538            | mj_core::targets::TargetLocator::SshPodman { .. }
539            | mj_core::targets::TargetLocator::SshDocker { .. }
540    )
541}
542
543fn ensure_target_unavailable_for_fallback(
544    controller: &Controller,
545    session_id: &str,
546    executor: &impl CommandExecutor,
547) -> Result<()> {
548    let availability = restart_target_availability(controller, session_id, executor)?;
549    ensure!(
550        fallback_is_safe(availability),
551        "Restart target is reachable; refusing suspend-and-resume target cleanup"
552    );
553    Ok(())
554}
555
556pub(super) fn fallback_is_safe(availability: RestartTargetAvailability) -> bool {
557    matches!(
558        availability,
559        RestartTargetAvailability::Missing | RestartTargetAvailability::Unreachable
560    )
561}
562
563async fn advance_restart_phase(
564    session_id: &str,
565    operation_id: &str,
566    phase: crate::database::SessionRestartPhase,
567) -> Result<()> {
568    blocking({
569        let session_id = session_id.to_owned();
570        let operation_id = operation_id.to_owned();
571        move || crate::database::advance_session_restart(&session_id, &operation_id, phase)
572    })
573    .await
574}
575
576async fn finish_restart_intent(session_id: &str, operation_id: &str) -> Result<()> {
577    blocking({
578        let session_id = session_id.to_owned();
579        let operation_id = operation_id.to_owned();
580        move || crate::database::finish_session_restart(&session_id, &operation_id)
581    })
582    .await
583}
584
585#[cfg(test)]
586mod tests {
587    use super::*;
588
589    #[test]
590    fn restart_fallback_requires_confirmed_target_absence() {
591        assert!(fallback_is_safe(RestartTargetAvailability::Missing));
592        assert!(fallback_is_safe(RestartTargetAvailability::Unreachable));
593        assert!(!fallback_is_safe(RestartTargetAvailability::Reachable));
594    }
595}