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                // A sub-agent borrows its parent's environment, which Move
111                // replaces, so its children stop exactly as they do when the
112                // parent is suspended. The destination's resume tells the
113                // model which ones stopped.
114                let stopping = std::time::Instant::now();
115                state.stop_subagents_for_suspend(&session_id).await?;
116                tracing::info!(
117                    %session_id,
118                    phase = "stop sub-agents",
119                    elapsed_ms = stopping.elapsed().as_millis() as u64,
120                    "move phase finished"
121                );
122                let result = blocking({
123                    let state = state.clone();
124                    let session_id = session_id.clone();
125                    move || {
126                        let reserving = std::time::Instant::now();
127                        let _reservation = reserve_recovery_or_cancel(
128                            &state.recovery_observer,
129                            &session_id,
130                            &cancelled,
131                        )?;
132                        tracing::info!(
133                            %session_id,
134                            phase = "recovery reservation",
135                            elapsed_ms = reserving.elapsed().as_millis() as u64,
136                            "move phase finished"
137                        );
138                        let loading = std::time::Instant::now();
139                        let mut controller = Controller::load()?;
140                        tracing::info!(
141                            %session_id,
142                            phase = "load controller state",
143                            elapsed_ms = loading.elapsed().as_millis() as u64,
144                            "move phase finished"
145                        );
146                        let executor = DaemonStageReportingExecutor::new(
147                            CancellableProcessExecutor::new(cancelled),
148                            state.clone(),
149                            session_id,
150                        );
151                        let outcome = mj_core::runtime::block_on(
152                            controller.move_session_managed_controlled(
153                                request,
154                                &executor,
155                                &state.session_manager,
156                            ),
157                        )??;
158                        Ok(DaemonLifecycleResult::Move(outcome))
159                    }
160                })
161                .await;
162                if result.is_err() {
163                    state
164                        .tell_live_parent_about_stopped_subagents(&session_id)
165                        .await;
166                }
167                result
168            },
169        )?;
170        self.set_lifecycle_resume_destination(
171            &session_id,
172            selection.profile_id.clone().unwrap_or_default(),
173            selection.target_template_id.clone().unwrap_or_default(),
174        );
175        Ok(result)
176    }
177
178    pub async fn move_session(
179        self: &Arc<Self>,
180        request: MoveSessionRequest,
181    ) -> Result<MoveOutcome> {
182        let selection = request.preparation.selection.clone();
183        let operation_id = request.preparation.operation_id.clone();
184        let session_id = selection.session_id.clone();
185        let result = self.admit_move_session(request)?;
186        let channel = result.clone();
187        let result = Self::wait_lifecycle_result(result).await;
188        self.remove_completed_lifecycle(&channel);
189        match result {
190            Ok(DaemonLifecycleResult::Move(outcome)) => Ok(outcome),
191            Ok(_) => bail!("move returned an unrelated lifecycle result"),
192            Err(error) => Ok(MoveOutcome {
193                operation_id, session_id, profile_id: selection.profile_id.unwrap_or_default(),
194                target_template_id: selection.target_template_id.unwrap_or_default(), outcome: "failed".into(),
195                error: Some(format!("{error:#}")), recovery: Some("Inspect session status and prepare Move again; any verified checkpoint is retained.".into()),
196            }),
197        }
198    }
199
200    pub(super) fn recover_moves(
201        self: &Arc<Self>,
202        operations: Vec<MoveOperation>,
203    ) -> Result<BTreeSet<String>> {
204        let mut owned = BTreeSet::new();
205        for operation in &operations {
206            crate::controller::move_session::restore_move_queue_hold(operation);
207        }
208        for operation in operations.into_iter().filter(|op| {
209            op.is_active()
210                || self
211                    .owner()
212                    .controller()
213                    .state
214                    .sessions
215                    .get(&op.selection.session_id)
216                    .is_some_and(|session| {
217                        matches!(
218                            session.state,
219                            SessionState::Closing | SessionState::Destroying
220                        )
221                    })
222        }) {
223            let id = operation.selection.session_id.clone();
224            let key = format!("recovery:{}", operation.operation_id);
225            let phase = if operation.recovery_session.is_some() {
226                LifecyclePhase::MovingDestination
227            } else {
228                LifecyclePhase::Executing
229            };
230            let result = self.admit_lifecycle(
231                id.clone(),
232                LifecycleKind::Move,
233                super::lifecycle::LifecycleStart {
234                    resume_workspace_id: None,
235                    request_key: Some(key),
236                    create_control: None,
237                    phase,
238                    move_operation_id: None,
239                },
240                move |state, session_id, cancelled| async move {
241                    blocking(move || {
242                        let _reservation = reserve_recovery_or_cancel(
243                            &state.recovery_observer,
244                            &session_id,
245                            &cancelled,
246                        )?;
247                        let mut controller = Controller::load()?;
248                        let executor = DaemonStageReportingExecutor::new(
249                            CancellableProcessExecutor::new(cancelled),
250                            state.clone(),
251                            session_id,
252                        );
253                        mj_core::runtime::block_on(controller.recover_move_managed_controlled(
254                            operation,
255                            &executor,
256                            &state.session_manager,
257                        ))?
258                        .map(DaemonLifecycleResult::Move)
259                    })
260                    .await
261                },
262            )?;
263            owned.insert(id.clone());
264            let state = self.clone();
265            tokio::spawn(async move {
266                let channel = result.clone();
267                match Self::wait_lifecycle_result(result).await {
268                    Ok(DaemonLifecycleResult::Move(outcome))
269                        if !matches!(outcome.outcome.as_str(), "completed" | "interrupted") =>
270                    {
271                        state.push_notice(
272                            &id,
273                            outcome
274                                .error
275                                .unwrap_or_else(|| "Move recovery needs attention".into()),
276                        );
277                    }
278                    Err(error) => {
279                        state.push_notice(&id, format!("Move recovery failed: {error:#}"))
280                    }
281                    _ => {}
282                }
283                state.remove_completed_lifecycle(&channel);
284            });
285        }
286        Ok(owned)
287    }
288}