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 mut preparation = blocking(move || {
68            let controller = Controller::load()?;
69            mj_core::runtime::block_on(
70                controller.prepare_move_session_controlled(selection, &ProcessExecutor),
71            )?
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                blocking(move || {
106                    let _reservation = reserve_recovery_or_cancel(
107                        &state.recovery_observer,
108                        &session_id,
109                        &cancelled,
110                    )?;
111                    let mut controller = Controller::load()?;
112                    let executor = DaemonStageReportingExecutor::new(
113                        CancellableProcessExecutor::new(cancelled),
114                        state.clone(),
115                        session_id,
116                    );
117                    let outcome =
118                        mj_core::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                    blocking(move || {
180                        let _reservation = reserve_recovery_or_cancel(
181                            &state.recovery_observer,
182                            &session_id,
183                            &cancelled,
184                        )?;
185                        let mut controller = Controller::load()?;
186                        let executor = DaemonStageReportingExecutor::new(
187                            CancellableProcessExecutor::new(cancelled),
188                            state.clone(),
189                            session_id,
190                        );
191                        mj_core::runtime::block_on(controller.recover_move_managed_controlled(
192                            operation,
193                            &executor,
194                            &state.session_manager,
195                        ))?
196                        .map(DaemonLifecycleResult::Move)
197                    })
198                    .await
199                },
200            )?;
201            owned.insert(id.clone());
202            let state = self.clone();
203            tokio::spawn(async move {
204                let channel = result.clone();
205                match Self::wait_lifecycle_result(result).await {
206                    Ok(DaemonLifecycleResult::Move(outcome)) if outcome.outcome != "completed" => {
207                        state.push_notice(
208                            &id,
209                            outcome
210                                .error
211                                .unwrap_or_else(|| "Move recovery needs attention".into()),
212                        );
213                    }
214                    Err(error) => {
215                        state.push_notice(&id, format!("Move recovery failed: {error:#}"))
216                    }
217                    _ => {}
218                }
219                state.remove_completed_lifecycle(&channel);
220            });
221        }
222        Ok(owned)
223    }
224}