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 | CloseRoute::RecoverInterrupted => {}
99        }
100        let operation_session_id = session_id.clone();
101        let result = self
102            .run_lifecycle(
103                operation_session_id,
104                LifecycleKind::Close,
105                move |state, session_id, cancelled| async move {
106                    let _recovery_reservation = tokio::task::spawn_blocking({
107                        let observer = state.recovery_observer.clone();
108                        let session_id = session_id.clone();
109                        let cancelled = cancelled.clone();
110                        move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
111                    })
112                    .await
113                    .context("reserve recovery for daemon close task")??;
114                    let mut controller = tokio::task::spawn_blocking(Controller::load)
115                        .await
116                        .context("load controller for daemon close task")??;
117                    let executor = DaemonStageReportingExecutor::new(
118                        CancellableProcessExecutor::new(cancelled),
119                        state.clone(),
120                        session_id.clone(),
121                    );
122                    let deferred = if route == CloseRoute::RecoverInterrupted {
123                        controller
124                            .recover_interrupted_close_managed(
125                                &session_id,
126                                &executor,
127                                &state.session_manager,
128                            )
129                            .await?
130                    } else {
131                        controller
132                            .close_session_managed_controlled(
133                                &session_id,
134                                &executor,
135                                &state.session_manager,
136                            )
137                            .await?
138                    };
139                    Ok(if deferred {
140                        DaemonLifecycleResult::DeferredCleanup
141                    } else {
142                        DaemonLifecycleResult::Done
143                    })
144                },
145            )
146            .await?;
147        let _ = result; // Deferred cleanup is handed off by the daemon-owned supervisor.
148        Ok(())
149    }
150
151    pub(super) fn start_deferred_cleanup(
152        self: &Arc<Self>,
153        session_id: String,
154    ) -> Result<
155        tokio::sync::watch::Receiver<Option<std::result::Result<DaemonLifecycleResult, String>>>,
156    > {
157        let result = self.start_or_join_lifecycle(
158            session_id.clone(),
159            LifecycleKind::Cleanup,
160            |state, session_id, cancelled| async move {
161                blocking(move || {
162                    let mut controller = Controller::load()?;
163                    let executor = DaemonStageReportingExecutor::new(
164                        CancellableProcessExecutor::new(cancelled),
165                        state,
166                        session_id.clone(),
167                    );
168                    controller.cleanup_stopped_target(&session_id, &executor)?;
169                    Ok(DaemonLifecycleResult::Done)
170                })
171                .await
172            },
173        )?;
174        let caller_result = result.clone();
175        let channel = result.clone();
176        let state = Arc::clone(self);
177        tokio::spawn(async move {
178            if let Err(error) = Self::wait_lifecycle_result(result).await {
179                tracing::warn!(%session_id, error = format!("{error:#}"), "deferred Podman cleanup failed");
180                state.push_notice(
181                    &session_id,
182                    "Container storage cleanup failed; the stopped session retains its target for retry.",
183                );
184            }
185            state.remove_completed_lifecycle(&channel);
186        });
187        Ok(caller_result)
188    }
189
190    pub(super) fn resume_retained_cleanups(self: &Arc<Self>) {
191        let session_ids = self
192            .controller
193            .lock()
194            .unwrap_or_else(PoisonError::into_inner)
195            .state
196            .sessions
197            .iter()
198            .filter(|(_, session)| {
199                session.state == SessionState::Stopped && session.target.is_some()
200            })
201            .map(|(session_id, _)| session_id.clone())
202            .collect::<Vec<_>>();
203        for session_id in session_ids {
204            if crate::controller::move_session::move_owns_session(&session_id) {
205                continue;
206            }
207            if let Err(error) = self.start_deferred_cleanup(session_id.clone()) {
208                tracing::warn!(%session_id, error = format!("{error:#}"), "could not resume deferred Podman cleanup");
209                self.push_notice(
210                    &session_id,
211                    format!("Could not resume container storage cleanup: {error:#}"),
212                );
213            }
214        }
215    }
216
217    pub(super) async fn wait_for_deferred_cleanup(
218        self: &Arc<Self>,
219        session_id: &str,
220    ) -> Result<()> {
221        let existing = {
222            let lifecycle = self
223                .lifecycle
224                .lock()
225                .unwrap_or_else(PoisonError::into_inner);
226            lifecycle.get(session_id).and_then(|active| {
227                (active.kind == LifecycleKind::Cleanup).then(|| active.result.clone())
228            })
229        };
230        let result = match existing {
231            Some(result) => result,
232            None => {
233                let needs_cleanup = blocking({
234                    let session_id = session_id.to_owned();
235                    move || {
236                        let controller = Controller::load()?;
237                        Ok(controller
238                            .state
239                            .sessions
240                            .get(&session_id)
241                            .is_some_and(|session| {
242                                session.state == SessionState::Stopped && session.target.is_some()
243                            }))
244                    }
245                })
246                .await?;
247                if !needs_cleanup {
248                    return Ok(());
249                }
250                self.start_deferred_cleanup(session_id.to_owned())?
251            }
252        };
253        let channel = result.clone();
254        let outcome = Self::wait_lifecycle_result(result).await;
255        self.remove_completed_lifecycle(&channel);
256        match outcome? {
257            DaemonLifecycleResult::Done => Ok(()),
258            DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
259            DaemonLifecycleResult::DeferredCleanup => {
260                unreachable!("cleanup cannot schedule another cleanup")
261            }
262        }
263    }
264}