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