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    pub async fn close_session(self: &Arc<Self>, session_id: String) -> Result<()> {
28        let children = blocking({
29            let session_id = session_id.clone();
30            move || {
31                let controller = Controller::load()?;
32                Ok(active_child_session_ids(&controller.state, &session_id))
33            }
34        })
35        .await?;
36        for child_id in children {
37            self.request_close(&child_id);
38            let result = self.close_requested_session(child_id.clone()).await;
39            self.clear_close_request(&child_id);
40            result.with_context(|| format!("stop sub-agent {child_id} before its parent"))?;
41        }
42        self.request_close(&session_id);
43        let result = self.close_requested_session(session_id.clone()).await;
44        self.clear_close_request(&session_id);
45        result
46    }
47
48    pub(super) async fn wait_before_close(self: &Arc<Self>, session_id: &str) -> Result<()> {
49        // Cancellation is a request: the old owner must actually finish before
50        // close acquires the target, including an irreversible create commit.
51        let pending = {
52            let operations = self
53                .lifecycle
54                .lock()
55                .unwrap_or_else(PoisonError::into_inner);
56            operations
57                .get(session_id)
58                .filter(|operation| {
59                    !matches!(
60                        operation.kind,
61                        LifecycleKind::Close | LifecycleKind::Cleanup
62                    )
63                })
64                .map(|operation| {
65                    operation.request_cancel();
66                    operation.result.clone()
67                })
68        };
69        if let Some(pending) = pending {
70            if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
71                tracing::debug!(%session_id, %error, "previous lifecycle ended before close");
72            }
73            self.remove_completed_lifecycle(&pending);
74        }
75
76        Ok(())
77    }
78
79    pub(super) async fn close_requested_session(
80        self: &Arc<Self>,
81        session_id: String,
82    ) -> Result<()> {
83        self.wait_before_close(&session_id).await?;
84        let route = blocking({
85            let session_id = session_id.clone();
86            move || {
87                let controller = Controller::load()?;
88                Ok(close_route(controller.state.sessions.get(&session_id)))
89            }
90        })
91        .await?;
92        match route {
93            CloseRoute::Done => return Ok(()),
94            CloseRoute::DeferredCleanup => {
95                self.start_deferred_cleanup(session_id)?;
96                return Ok(());
97            }
98            CloseRoute::Graceful
99            | CloseRoute::RecoverInterrupted
100            | CloseRoute::SettleWithoutCheckpoint => {}
101        }
102        let operation_session_id = session_id.clone();
103        let result = self
104            .run_lifecycle(
105                operation_session_id,
106                LifecycleKind::Close,
107                move |state, session_id, cancelled| async move {
108                    let _recovery_reservation = tokio::task::spawn_blocking({
109                        let observer = state.recovery_observer.clone();
110                        let session_id = session_id.clone();
111                        let cancelled = cancelled.clone();
112                        move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
113                    })
114                    .await
115                    .context("reserve recovery for daemon close task")??;
116                    let mut controller = tokio::task::spawn_blocking(Controller::load)
117                        .await
118                        .context("load controller for daemon close task")??;
119                    let executor = DaemonStageReportingExecutor::new(
120                        CancellableProcessExecutor::new(cancelled),
121                        state.clone(),
122                        session_id.clone(),
123                    );
124                    let deferred = match route {
125                        CloseRoute::RecoverInterrupted => {
126                            controller
127                                .recover_interrupted_close_managed(
128                                    &session_id,
129                                    &executor,
130                                    &state.session_manager,
131                                )
132                                .await?
133                        }
134                        // Nothing to archive and no relay to latch, so this
135                        // close only tears down and settles. The route was
136                        // decided after `wait_before_close` let any live
137                        // create or resume finish, so a session that is still
138                        // genuinely provisioning is not caught here.
139                        CloseRoute::SettleWithoutCheckpoint => {
140                            controller.close_session_without_checkpoint(&session_id, &executor)?
141                        }
142                        _ => {
143                            controller
144                                .close_session_managed_controlled(
145                                    &session_id,
146                                    &executor,
147                                    &state.session_manager,
148                                )
149                                .await?
150                        }
151                    };
152                    Ok(if deferred {
153                        DaemonLifecycleResult::DeferredCleanup
154                    } else {
155                        DaemonLifecycleResult::Done
156                    })
157                },
158            )
159            .await?;
160        let _ = result; // Deferred cleanup is handed off by the daemon-owned supervisor.
161        Ok(())
162    }
163
164    pub(super) fn start_deferred_cleanup(
165        self: &Arc<Self>,
166        session_id: String,
167    ) -> Result<LifecycleWatch> {
168        let result = self.start_or_join_lifecycle(
169            session_id.clone(),
170            LifecycleKind::Cleanup,
171            |state, session_id, cancelled| async move {
172                blocking(move || {
173                    let mut controller = Controller::load()?;
174                    let executor = DaemonStageReportingExecutor::new(
175                        CancellableProcessExecutor::new(cancelled),
176                        state,
177                        session_id.clone(),
178                    );
179                    controller.cleanup_stopped_target(&session_id, &executor)?;
180                    Ok(DaemonLifecycleResult::Done)
181                })
182                .await
183            },
184        )?;
185        let caller_result = result.clone();
186        let channel = result.clone();
187        let state = Arc::clone(self);
188        tokio::spawn(async move {
189            if let Err(error) = Self::wait_lifecycle_result(result).await {
190                tracing::warn!(%session_id, error = format!("{error:#}"), "deferred Podman cleanup failed");
191                state.push_notice(
192                    &session_id,
193                    "Container storage cleanup failed; the stopped session retains its target for retry.",
194                );
195            }
196            state.remove_completed_lifecycle(&channel);
197        });
198        Ok(caller_result)
199    }
200
201    pub(super) fn resume_retained_cleanups(self: &Arc<Self>) {
202        let session_ids = self
203            .controller
204            .lock()
205            .unwrap_or_else(PoisonError::into_inner)
206            .state
207            .sessions
208            .iter()
209            .filter(|(_, session)| {
210                session.state == SessionState::Stopped && session.target.is_some()
211            })
212            .map(|(session_id, _)| session_id.clone())
213            .collect::<Vec<_>>();
214        for session_id in session_ids {
215            if crate::controller::move_session::move_owns_session(&session_id) {
216                continue;
217            }
218            if let Err(error) = self.start_deferred_cleanup(session_id.clone()) {
219                tracing::warn!(%session_id, error = format!("{error:#}"), "could not resume deferred Podman cleanup");
220                self.push_notice(
221                    &session_id,
222                    format!("Could not resume container storage cleanup: {error:#}"),
223                );
224            }
225        }
226    }
227
228    pub(super) async fn wait_for_deferred_cleanup(
229        self: &Arc<Self>,
230        session_id: &str,
231    ) -> Result<()> {
232        let existing = {
233            let lifecycle = self
234                .lifecycle
235                .lock()
236                .unwrap_or_else(PoisonError::into_inner);
237            lifecycle.get(session_id).and_then(|active| {
238                (active.kind == LifecycleKind::Cleanup).then(|| active.result.clone())
239            })
240        };
241        let result = match existing {
242            Some(result) => result,
243            None => {
244                let needs_cleanup = blocking({
245                    let session_id = session_id.to_owned();
246                    move || {
247                        let controller = Controller::load()?;
248                        Ok(controller
249                            .state
250                            .sessions
251                            .get(&session_id)
252                            .is_some_and(|session| {
253                                session.state == SessionState::Stopped && session.target.is_some()
254                            }))
255                    }
256                })
257                .await?;
258                if !needs_cleanup {
259                    return Ok(());
260                }
261                self.start_deferred_cleanup(session_id.to_owned())?
262            }
263        };
264        let channel = result.clone();
265        let outcome = Self::wait_lifecycle_result(result).await;
266        self.remove_completed_lifecycle(&channel);
267        match outcome? {
268            DaemonLifecycleResult::Done => Ok(()),
269            DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
270            DaemonLifecycleResult::DeferredCleanup => {
271                unreachable!("cleanup cannot schedule another cleanup")
272            }
273        }
274    }
275}