mj_controller/daemon/
close.rs1use 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 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 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; 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}