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 | 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; 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}