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_session(record)?;
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 = self
42                .controller
43                .lock()
44                .unwrap_or_else(PoisonError::into_inner);
45            let source = controller
46                .state
47                .sessions
48                .get(&selection.session_id)
49                .context("Move session is missing")?;
50            (
51                source.harness_kind,
52                matches!(
53                    source.state,
54                    SessionState::Running | SessionState::Disconnected
55                ),
56            )
57        };
58        let snapshot = if source_active {
59            crate::controller::move_session::refresh_move_source(
60                &self.session_manager,
61                &selection.session_id,
62            )
63            .await?
64        } else {
65            None
66        };
67        let runtime = tokio::runtime::Handle::current();
68        let mut preparation = blocking(move || {
69            let controller = Controller::load()?;
70            runtime
71                .block_on(controller.prepare_move_session_controlled(selection, &ProcessExecutor))
72        })
73        .await?;
74        preparation.source_unavailable = source_active
75            && snapshot
76                .as_ref()
77                .is_none_or(|snapshot| !snapshot.operational.native_session_is_ready());
78        preparation.active |= preparation.source_unavailable;
79        if let Some(snapshot) = snapshot {
80            let mut operational = snapshot.operational;
81            operational.queued_prompts.clear();
82            operational.checkpoint_barrier = None;
83            preparation.active |= !operational.safe_to_replace(source_harness);
84        }
85        Ok(preparation)
86    }
87
88    pub async fn move_session(
89        self: &Arc<Self>,
90        request: MoveSessionRequest,
91    ) -> Result<MoveOutcome> {
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.start_or_join_lifecycle_with_key(
100            session_id.clone(),
101            LifecycleKind::Move,
102            None,
103            Some(key),
104            move |state, session_id, cancelled| async move {
105                let runtime = tokio::runtime::Handle::current();
106                blocking(move || {
107                    let _reservation = reserve_recovery_or_cancel(
108                        &state.recovery_observer,
109                        &session_id,
110                        &cancelled,
111                    )?;
112                    let mut controller = Controller::load()?;
113                    let executor = DaemonStageReportingExecutor::new(
114                        CancellableProcessExecutor::new(cancelled),
115                        state.clone(),
116                        session_id,
117                    );
118                    let outcome = runtime.block_on(controller.move_session_managed_controlled(
119                        request,
120                        &executor,
121                        &state.session_manager,
122                    ))?;
123                    Ok(DaemonLifecycleResult::Move(outcome))
124                })
125                .await
126            },
127        )?;
128        self.set_lifecycle_resume_destination(
129            &session_id,
130            selection.profile_id.clone().unwrap_or_default(),
131            selection.target_template_id.clone().unwrap_or_default(),
132        );
133        let channel = result.clone();
134        let result = Self::wait_lifecycle_result(result).await;
135        self.remove_completed_lifecycle(&channel);
136        match result {
137            Ok(DaemonLifecycleResult::Move(outcome)) => Ok(outcome),
138            Ok(_) => bail!("move returned an unrelated lifecycle result"),
139            Err(error) => Ok(MoveOutcome {
140                operation_id, session_id, profile_id: selection.profile_id.unwrap_or_default(),
141                target_template_id: selection.target_template_id.unwrap_or_default(), outcome: "failed".into(),
142                error: Some(format!("{error:#}")), recovery: Some("Inspect session status and prepare Move again; any verified checkpoint is retained.".into()),
143            }),
144        }
145    }
146
147    pub(super) fn recover_moves(
148        self: &Arc<Self>,
149        operations: Vec<MoveOperation>,
150    ) -> Result<BTreeSet<String>> {
151        let mut owned = BTreeSet::new();
152        for operation in &operations {
153            crate::controller::move_session::restore_move_queue_hold(operation);
154        }
155        for operation in operations.into_iter().filter(|op| {
156            op.is_active()
157                || self
158                    .controller
159                    .lock()
160                    .unwrap_or_else(PoisonError::into_inner)
161                    .state
162                    .sessions
163                    .get(&op.selection.session_id)
164                    .is_some_and(|session| {
165                        matches!(
166                            session.state,
167                            SessionState::Closing | SessionState::Destroying
168                        )
169                    })
170        }) {
171            let id = operation.selection.session_id.clone();
172            let key = format!("recovery:{}", operation.operation_id);
173            let result = self.start_or_join_lifecycle_with_key(
174                id.clone(),
175                LifecycleKind::Move,
176                None,
177                Some(key),
178                move |state, session_id, cancelled| async move {
179                    let runtime = tokio::runtime::Handle::current();
180                    blocking(move || {
181                        let _reservation = reserve_recovery_or_cancel(
182                            &state.recovery_observer,
183                            &session_id,
184                            &cancelled,
185                        )?;
186                        let mut controller = Controller::load()?;
187                        let executor = DaemonStageReportingExecutor::new(
188                            CancellableProcessExecutor::new(cancelled),
189                            state.clone(),
190                            session_id,
191                        );
192                        runtime
193                            .block_on(controller.recover_move_managed_controlled(
194                                operation,
195                                &executor,
196                                &state.session_manager,
197                            ))
198                            .map(DaemonLifecycleResult::Move)
199                    })
200                    .await
201                },
202            )?;
203            owned.insert(id.clone());
204            let state = self.clone();
205            tokio::spawn(async move {
206                let channel = result.clone();
207                match Self::wait_lifecycle_result(result).await {
208                    Ok(DaemonLifecycleResult::Move(outcome)) if outcome.outcome != "completed" => {
209                        state.push_notice(
210                            &id,
211                            outcome
212                                .error
213                                .unwrap_or_else(|| "Move recovery needs attention".into()),
214                        );
215                    }
216                    Err(error) => {
217                        state.push_notice(&id, format!("Move recovery failed: {error:#}"))
218                    }
219                    _ => {}
220                }
221                state.remove_completed_lifecycle(&channel);
222            });
223        }
224        Ok(owned)
225    }
226}