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        self.request_close(&session_id);
83        let result = self.suspend_with_children(&session_id).await;
84        if let Err(error) = &result {
85            let reference = new_command_id("suspension").unwrap_or_else(|_| "suspension".into());
86            tracing::warn!(%session_id, %reference, error = format!("{error:#}"), "session suspension failed");
87            self.record_failed_close(&session_id, &reference, &LifecycleFailure::of(error))
88                .await;
89        }
90        self.clear_close_request(&session_id);
91        result
92    }
93
94    async fn suspend_with_children(self: &Arc<Self>, session_id: &str) -> Result<()> {
95        self.prepare_suspension(session_id).await?;
96        let children = blocking({
97            let session_id = session_id.to_owned();
98            move || {
99                let controller = Controller::load()?;
100                Ok(active_child_session_ids(&controller.state, &session_id))
101            }
102        })
103        .await?;
104        for child_id in children {
105            if let Err(error) = Box::pin(self.suspend_session(child_id.clone())).await {
106                tracing::warn!(%child_id, error = format!("{error:#}"), "child suspension failed");
107                return Err(mj_core::refusal::Refusal::precondition(format!(
108                    "sub-agent {child_id} could not be suspended; inspect that session and retry"
109                ))
110                .into());
111            }
112        }
113        self.close_requested_session(session_id.to_owned()).await
114    }
115
116    pub(super) async fn wait_before_close(self: &Arc<Self>, session_id: &str) -> Result<()> {
117        // Cancellation is a request: the old owner must actually finish before
118        // close acquires the target, including an irreversible create commit.
119        let pending = {
120            let operations = self
121                .lifecycle
122                .lock()
123                .unwrap_or_else(PoisonError::into_inner);
124            operations
125                .get(session_id)
126                .filter(|operation| {
127                    !matches!(
128                        operation.kind,
129                        LifecycleKind::Suspend | LifecycleKind::Cleanup
130                    )
131                })
132                .map(|operation| {
133                    operation.request_cancel();
134                    operation.result.clone()
135                })
136        };
137        if let Some(pending) = pending {
138            if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
139                tracing::debug!(%session_id, %error, "previous lifecycle ended before close");
140            }
141            self.remove_completed_lifecycle(&pending);
142        }
143
144        Ok(())
145    }
146
147    pub(super) async fn close_requested_session(
148        self: &Arc<Self>,
149        session_id: String,
150    ) -> Result<()> {
151        self.wait_before_close(&session_id).await?;
152        let route = blocking({
153            let session_id = session_id.clone();
154            move || {
155                let controller = Controller::load()?;
156                Ok(close_route(controller.state.sessions.get(&session_id)))
157            }
158        })
159        .await?;
160        match route {
161            CloseRoute::Done => return Ok(()),
162            CloseRoute::DeferredCleanup => {
163                self.start_deferred_cleanup(session_id)?;
164                return Ok(());
165            }
166            CloseRoute::Graceful
167            | CloseRoute::RecoverInterrupted
168            | CloseRoute::SettleWithoutCheckpoint => {}
169        }
170        let operation_session_id = session_id.clone();
171        let result = self
172            .run_lifecycle(
173                operation_session_id,
174                LifecycleKind::Suspend,
175                move |state, session_id, cancelled| async move {
176                    let _recovery_reservation = tokio::task::spawn_blocking({
177                        let observer = state.recovery_observer.clone();
178                        let session_id = session_id.clone();
179                        let cancelled = cancelled.clone();
180                        move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
181                    })
182                    .await
183                    .context("reserve recovery for daemon close task")??;
184                    let mut controller = tokio::task::spawn_blocking(Controller::load)
185                        .await
186                        .context("load controller for daemon close task")??;
187                    let executor = DaemonStageReportingExecutor::new(
188                        CancellableProcessExecutor::new(cancelled),
189                        state.clone(),
190                        session_id.clone(),
191                    );
192                    let deferred = match route {
193                        CloseRoute::RecoverInterrupted => {
194                            controller
195                                .recover_interrupted_close_managed(
196                                    &session_id,
197                                    &executor,
198                                    &state.session_manager,
199                                )
200                                .await?
201                        }
202                        // Nothing to archive and no relay to latch, so this
203                        // close only tears down and settles. The route was
204                        // decided after `wait_before_close` let any live
205                        // create or resume finish, so a session that is still
206                        // genuinely provisioning is not caught here.
207                        CloseRoute::SettleWithoutCheckpoint => {
208                            controller.suspend_session_without_checkpoint(&session_id, &executor)?
209                        }
210                        _ => {
211                            controller
212                                .suspend_session_managed_controlled(
213                                    &session_id,
214                                    &executor,
215                                    &state.session_manager,
216                                )
217                                .await?
218                        }
219                    };
220                    Ok(if deferred {
221                        DaemonLifecycleResult::DeferredCleanup
222                    } else {
223                        DaemonLifecycleResult::Done
224                    })
225                },
226            )
227            .await?;
228        let _ = result; // Deferred cleanup is handed off by the daemon-owned supervisor.
229        Ok(())
230    }
231
232    pub(super) fn start_deferred_cleanup(
233        self: &Arc<Self>,
234        session_id: String,
235    ) -> Result<LifecycleWatch> {
236        let result = self.start_or_join_lifecycle(
237            session_id.clone(),
238            LifecycleKind::Cleanup,
239            |state, session_id, cancelled| async move {
240                blocking(move || {
241                    let mut controller = Controller::load()?;
242                    let executor = DaemonStageReportingExecutor::new(
243                        CancellableProcessExecutor::new(cancelled),
244                        state,
245                        session_id.clone(),
246                    );
247                    controller.cleanup_stopped_target(&session_id, &executor)?;
248                    Ok(DaemonLifecycleResult::Done)
249                })
250                .await
251            },
252        )?;
253        let caller_result = result.clone();
254        let channel = result.clone();
255        let state = Arc::clone(self);
256        tokio::spawn(async move {
257            if let Err(error) = Self::wait_lifecycle_result(result).await {
258                tracing::warn!(%session_id, error = format!("{error:#}"), "deferred Podman cleanup failed");
259                state.push_notice(
260                    &session_id,
261                    "Container storage cleanup failed; the stopped session retains its target for retry.",
262                );
263            }
264            state.remove_completed_lifecycle(&channel);
265        });
266        Ok(caller_result)
267    }
268
269    pub(super) fn resume_retained_cleanups(self: &Arc<Self>) {
270        let session_ids = self
271            .controller
272            .lock()
273            .unwrap_or_else(PoisonError::into_inner)
274            .state
275            .sessions
276            .iter()
277            .filter(|(_, session)| {
278                session.state == SessionState::Stopped && session.target.is_some()
279            })
280            .map(|(session_id, _)| session_id.clone())
281            .collect::<Vec<_>>();
282        for session_id in session_ids {
283            if crate::controller::move_session::move_owns_session(&session_id) {
284                continue;
285            }
286            if let Err(error) = self.start_deferred_cleanup(session_id.clone()) {
287                tracing::warn!(%session_id, error = format!("{error:#}"), "could not resume deferred Podman cleanup");
288                self.push_notice(
289                    &session_id,
290                    format!("Could not resume container storage cleanup: {error:#}"),
291                );
292            }
293        }
294    }
295
296    pub(super) async fn wait_for_deferred_cleanup(
297        self: &Arc<Self>,
298        session_id: &str,
299    ) -> Result<()> {
300        let existing = {
301            let lifecycle = self
302                .lifecycle
303                .lock()
304                .unwrap_or_else(PoisonError::into_inner);
305            lifecycle.get(session_id).and_then(|active| {
306                (active.kind == LifecycleKind::Cleanup).then(|| active.result.clone())
307            })
308        };
309        let result = match existing {
310            Some(result) => result,
311            None => {
312                let needs_cleanup = blocking({
313                    let session_id = session_id.to_owned();
314                    move || {
315                        let controller = Controller::load()?;
316                        Ok(controller
317                            .state
318                            .sessions
319                            .get(&session_id)
320                            .is_some_and(|session| {
321                                session.state == SessionState::Stopped && session.target.is_some()
322                            }))
323                    }
324                })
325                .await?;
326                if !needs_cleanup {
327                    return Ok(());
328                }
329                self.start_deferred_cleanup(session_id.to_owned())?
330            }
331        };
332        let channel = result.clone();
333        let outcome = Self::wait_lifecycle_result(result).await;
334        self.remove_completed_lifecycle(&channel);
335        match outcome? {
336            DaemonLifecycleResult::Done => Ok(()),
337            DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
338            DaemonLifecycleResult::DeferredCleanup => {
339                unreachable!("cleanup cannot schedule another cleanup")
340            }
341        }
342    }
343}