Skip to main content

mj_controller/daemon/
close.rs

1use super::*;
2
3impl RuntimeState {
4    pub fn request_close(&self, session_id: &str) {
5        self.close_requested
6            .lock()
7            .unwrap_or_else(PoisonError::into_inner)
8            .insert(session_id.to_owned());
9        self.publish_revision();
10    }
11
12    pub fn clear_close_request(&self, session_id: &str) {
13        self.close_requested
14            .lock()
15            .unwrap_or_else(PoisonError::into_inner)
16            .remove(session_id);
17        self.publish_revision();
18    }
19
20    pub fn close_is_requested(&self, session_id: &str) -> bool {
21        self.close_requested
22            .lock()
23            .unwrap_or_else(PoisonError::into_inner)
24            .contains(session_id)
25    }
26
27    /// Make the intent durable before an HTTP caller receives acceptance.
28    /// Waiting and store writes happen in the action task, never the UI loop.
29    pub async fn prepare_suspension(self: &Arc<Self>, session_id: &str) -> Result<()> {
30        self.wait_before_close(session_id).await?;
31        if self
32            .lifecycle
33            .lock()
34            .unwrap_or_else(PoisonError::into_inner)
35            .get(session_id)
36            .is_some_and(|operation| {
37                operation.kind == LifecycleKind::Suspend && operation.result.borrow().is_none()
38            })
39        {
40            return Ok(());
41        }
42        blocking({
43            let session_id = session_id.to_owned();
44            let observer = self.recovery_observer.clone();
45            move || {
46                let cancelled = AtomicBool::new(false);
47                let _reservation = reserve_recovery_or_cancel(&observer, &session_id, &cancelled)?;
48                ensure!(
49                    !crate::controller::move_session::move_has_pending_queue(&session_id),
50                    "Move queue admission is incomplete; retry Move on the same destination before suspending"
51                );
52                let mut controller = Controller::load()?;
53                let record = controller
54                    .state
55                    .sessions
56                    .get_mut(&session_id)
57                    .with_context(|| format!("unknown session {session_id}"))?;
58                if matches!(
59                    record.state,
60                    SessionState::Running
61                        | SessionState::Disconnected
62                        | SessionState::Checkpointing
63                ) {
64                    record.state = SessionState::Closing;
65                    record.last_error = None;
66                    record.updated_at = chrono::Utc::now().to_rfc3339();
67                    crate::database::save_lifecycle_session(record)?;
68                } else if record.public_error().is_some() {
69                    record.last_error = None;
70                    crate::database::save_lifecycle_session(record)?;
71                }
72                Ok(())
73            }
74        })
75        .await?;
76        self.reload_controller().await?;
77        self.publish_revision();
78        Ok(())
79    }
80
81    pub async fn suspend_session(self: &Arc<Self>, session_id: String) -> Result<()> {
82        // Internal suspension (workspace close, recovery, child teardown) has
83        // already been admitted by its parent operation. Its checkpoint is
84        // still verified before any owned checkout is released.
85        self.suspend_session_with_ack(session_id, true).await
86    }
87
88    pub async fn suspend_session_with_ack(
89        self: &Arc<Self>,
90        session_id: String,
91        acknowledge_unpublished_work: bool,
92    ) -> Result<()> {
93        self.request_close(&session_id);
94        let result = self
95            .suspend_with_children(&session_id, acknowledge_unpublished_work)
96            .await;
97        if let Err(error) = &result {
98            let reference = new_command_id("suspension").unwrap_or_else(|_| "suspension".into());
99            tracing::warn!(%session_id, %reference, error = format!("{error:#}"), "session suspension failed");
100            self.record_failed_close(&session_id, &reference, &LifecycleFailure::of(error))
101                .await;
102        }
103        self.clear_close_request(&session_id);
104        result
105    }
106
107    async fn suspend_with_children(
108        self: &Arc<Self>,
109        session_id: &str,
110        acknowledge_unpublished_work: bool,
111    ) -> Result<()> {
112        self.prepare_suspension(session_id).await?;
113        let children = blocking({
114            let session_id = session_id.to_owned();
115            move || {
116                let controller = Controller::load()?;
117                Ok(active_child_session_ids(&controller.state, &session_id))
118            }
119        })
120        .await?;
121        for child_id in children {
122            if let Err(error) = Box::pin(self.suspend_session(child_id.clone())).await {
123                tracing::warn!(%child_id, error = format!("{error:#}"), "child suspension failed");
124                return Err(mj_core::refusal::Refusal::precondition(format!(
125                    "sub-agent {child_id} could not be suspended; inspect that session and retry"
126                ))
127                .into());
128            }
129        }
130        self.close_requested_session_with_ack(session_id.to_owned(), acknowledge_unpublished_work)
131            .await
132    }
133
134    pub(super) async fn wait_before_close(self: &Arc<Self>, session_id: &str) -> Result<()> {
135        // Cancellation is a request: the old owner must actually finish before
136        // close acquires the target, including an irreversible create commit.
137        let pending = {
138            let operations = self
139                .lifecycle
140                .lock()
141                .unwrap_or_else(PoisonError::into_inner);
142            operations
143                .get(session_id)
144                .filter(|operation| {
145                    !matches!(
146                        operation.kind,
147                        LifecycleKind::Suspend | LifecycleKind::Cleanup
148                    )
149                })
150                .map(|operation| {
151                    operation.request_cancel();
152                    operation.result.clone()
153                })
154        };
155        if let Some(pending) = pending {
156            if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
157                tracing::debug!(%session_id, %error, "previous lifecycle ended before close");
158            }
159            self.remove_completed_lifecycle(&pending);
160        }
161
162        Ok(())
163    }
164
165    async fn close_requested_session_with_ack(
166        self: &Arc<Self>,
167        session_id: String,
168        acknowledge_unpublished_work: bool,
169    ) -> Result<()> {
170        self.wait_before_close(&session_id).await?;
171        let route = blocking({
172            let session_id = session_id.clone();
173            move || {
174                let controller = Controller::load()?;
175                Ok(close_route(controller.state.sessions.get(&session_id)))
176            }
177        })
178        .await?;
179        match route {
180            CloseRoute::Done => return Ok(()),
181            CloseRoute::DeferredCleanup => {
182                self.start_deferred_cleanup(session_id)?;
183                return Ok(());
184            }
185            CloseRoute::Graceful
186            | CloseRoute::RecoverInterrupted
187            | CloseRoute::SettleWithoutCheckpoint => {}
188        }
189        let operation_session_id = session_id.clone();
190        let result = self
191            .run_lifecycle(
192                operation_session_id,
193                LifecycleKind::Suspend,
194                move |state, session_id, cancelled| async move {
195                    let _recovery_reservation = tokio::task::spawn_blocking({
196                        let observer = state.recovery_observer.clone();
197                        let session_id = session_id.clone();
198                        let cancelled = cancelled.clone();
199                        move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
200                    })
201                    .await
202                    .context("reserve recovery for daemon close task")??;
203                    let mut controller = tokio::task::spawn_blocking(Controller::load)
204                        .await
205                        .context("load controller for daemon close task")??;
206                    let executor = DaemonStageReportingExecutor::new(
207                        CancellableProcessExecutor::new(cancelled),
208                        state.clone(),
209                        session_id.clone(),
210                    );
211                    let deferred = match route {
212                        CloseRoute::RecoverInterrupted => {
213                            controller
214                                .recover_interrupted_close_managed(
215                                    &session_id,
216                                    &executor,
217                                    &state.session_manager,
218                                )
219                                .await?
220                        }
221                        // Nothing to archive and no relay to latch, so this
222                        // close only tears down and settles. The route was
223                        // decided after `wait_before_close` let any live
224                        // create or resume finish, so a session that is still
225                        // genuinely provisioning is not caught here.
226                        CloseRoute::SettleWithoutCheckpoint => {
227                            controller.suspend_session_without_checkpoint(&session_id, &executor)?
228                        }
229                        _ => {
230                            controller
231                                .suspend_session_managed_controlled(
232                                    &session_id,
233                                    &executor,
234                                    &state.session_manager,
235                                    acknowledge_unpublished_work,
236                                )
237                                .await?
238                        }
239                    };
240                    Ok(if deferred {
241                        DaemonLifecycleResult::DeferredCleanup
242                    } else {
243                        DaemonLifecycleResult::Done
244                    })
245                },
246            )
247            .await?;
248        let _ = result; // Deferred cleanup is handed off by the daemon-owned supervisor.
249        Ok(())
250    }
251
252    pub(super) fn start_deferred_cleanup(
253        self: &Arc<Self>,
254        session_id: String,
255    ) -> Result<LifecycleWatch> {
256        let result = self.start_or_join_lifecycle(
257            session_id.clone(),
258            LifecycleKind::Cleanup,
259            |state, session_id, cancelled| async move {
260                blocking(move || {
261                    let mut controller = Controller::load()?;
262                    let executor = DaemonStageReportingExecutor::new(
263                        CancellableProcessExecutor::new(cancelled),
264                        state,
265                        session_id.clone(),
266                    );
267                    controller.cleanup_stopped_target(&session_id, &executor)?;
268                    Ok(DaemonLifecycleResult::Done)
269                })
270                .await
271            },
272        )?;
273        let caller_result = result.clone();
274        let channel = result.clone();
275        let state = Arc::clone(self);
276        tokio::spawn(async move {
277            if let Err(error) = Self::wait_lifecycle_result(result).await {
278                tracing::warn!(%session_id, error = format!("{error:#}"), "deferred Podman cleanup failed");
279                state.push_notice(
280                    &session_id,
281                    "Container storage cleanup failed; the stopped session retains its target for retry.",
282                );
283            }
284            state.remove_completed_lifecycle(&channel);
285        });
286        Ok(caller_result)
287    }
288
289    pub(super) fn resume_retained_cleanups(self: &Arc<Self>) {
290        let session_ids = self
291            .controller
292            .lock()
293            .unwrap_or_else(PoisonError::into_inner)
294            .state
295            .sessions
296            .iter()
297            .filter(|(_, session)| {
298                session.state == SessionState::Stopped && session.target.is_some()
299            })
300            .map(|(session_id, _)| session_id.clone())
301            .collect::<Vec<_>>();
302        for session_id in session_ids {
303            if crate::controller::move_session::move_owns_session(&session_id) {
304                continue;
305            }
306            if let Err(error) = self.start_deferred_cleanup(session_id.clone()) {
307                tracing::warn!(%session_id, error = format!("{error:#}"), "could not resume deferred Podman cleanup");
308                self.push_notice(
309                    &session_id,
310                    format!("Could not resume container storage cleanup: {error:#}"),
311                );
312            }
313        }
314    }
315
316    pub(super) async fn wait_for_deferred_cleanup(
317        self: &Arc<Self>,
318        session_id: &str,
319    ) -> Result<()> {
320        let existing = {
321            let lifecycle = self
322                .lifecycle
323                .lock()
324                .unwrap_or_else(PoisonError::into_inner);
325            lifecycle.get(session_id).and_then(|active| {
326                (active.kind == LifecycleKind::Cleanup).then(|| active.result.clone())
327            })
328        };
329        let result = match existing {
330            Some(result) => result,
331            None => {
332                let needs_cleanup = blocking({
333                    let session_id = session_id.to_owned();
334                    move || {
335                        let controller = Controller::load()?;
336                        Ok(controller
337                            .state
338                            .sessions
339                            .get(&session_id)
340                            .is_some_and(|session| {
341                                session.state == SessionState::Stopped && session.target.is_some()
342                            }))
343                    }
344                })
345                .await?;
346                if !needs_cleanup {
347                    return Ok(());
348                }
349                self.start_deferred_cleanup(session_id.to_owned())?
350            }
351        };
352        let channel = result.clone();
353        let outcome = Self::wait_lifecycle_result(result).await;
354        self.remove_completed_lifecycle(&channel);
355        match outcome? {
356            DaemonLifecycleResult::Done => Ok(()),
357            DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
358            DaemonLifecycleResult::DeferredCleanup => {
359                unreachable!("cleanup cannot schedule another cleanup")
360            }
361        }
362    }
363}