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.suspend_session_with_ack(session_id, true).await
86 }
87
88 pub async fn suspend_session_with_ack(
89 self: &Arc<Self>,
90 session_id: String,
91 acknowledge_unpublished_work: bool,
92 ) -> Result<()> {
93 self.request_close(&session_id);
94 let result = self
95 .suspend_with_children(&session_id, acknowledge_unpublished_work)
96 .await;
97 if let Err(error) = &result {
98 let reference = new_command_id("suspension").unwrap_or_else(|_| "suspension".into());
99 tracing::warn!(%session_id, %reference, error = format!("{error:#}"), "session suspension failed");
100 self.record_failed_close(&session_id, &reference, &LifecycleFailure::of(error))
101 .await;
102 }
103 self.clear_close_request(&session_id);
104 result
105 }
106
107 async fn suspend_with_children(
108 self: &Arc<Self>,
109 session_id: &str,
110 acknowledge_unpublished_work: bool,
111 ) -> Result<()> {
112 self.prepare_suspension(session_id).await?;
113 let children = blocking({
114 let session_id = session_id.to_owned();
115 move || {
116 let controller = Controller::load()?;
117 Ok(active_child_session_ids(&controller.state, &session_id))
118 }
119 })
120 .await?;
121 for child_id in children {
122 if let Err(error) = Box::pin(self.suspend_session(child_id.clone())).await {
123 tracing::warn!(%child_id, error = format!("{error:#}"), "child suspension failed");
124 return Err(mj_core::refusal::Refusal::precondition(format!(
125 "sub-agent {child_id} could not be suspended; inspect that session and retry"
126 ))
127 .into());
128 }
129 }
130 self.close_requested_session_with_ack(session_id.to_owned(), acknowledge_unpublished_work)
131 .await
132 }
133
134 pub(super) async fn wait_before_close(self: &Arc<Self>, session_id: &str) -> Result<()> {
135 let pending = {
138 let operations = self
139 .lifecycle
140 .lock()
141 .unwrap_or_else(PoisonError::into_inner);
142 operations
143 .get(session_id)
144 .filter(|operation| {
145 !matches!(
146 operation.kind,
147 LifecycleKind::Suspend | LifecycleKind::Cleanup
148 )
149 })
150 .map(|operation| {
151 operation.request_cancel();
152 operation.result.clone()
153 })
154 };
155 if let Some(pending) = pending {
156 if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
157 tracing::debug!(%session_id, %error, "previous lifecycle ended before close");
158 }
159 self.remove_completed_lifecycle(&pending);
160 }
161
162 Ok(())
163 }
164
165 async fn close_requested_session_with_ack(
166 self: &Arc<Self>,
167 session_id: String,
168 acknowledge_unpublished_work: bool,
169 ) -> Result<()> {
170 self.wait_before_close(&session_id).await?;
171 let route = blocking({
172 let session_id = session_id.clone();
173 move || {
174 let controller = Controller::load()?;
175 Ok(close_route(controller.state.sessions.get(&session_id)))
176 }
177 })
178 .await?;
179 match route {
180 CloseRoute::Done => return Ok(()),
181 CloseRoute::DeferredCleanup => {
182 self.start_deferred_cleanup(session_id)?;
183 return Ok(());
184 }
185 CloseRoute::Graceful
186 | CloseRoute::RecoverInterrupted
187 | CloseRoute::SettleWithoutCheckpoint => {}
188 }
189 let operation_session_id = session_id.clone();
190 let result = self
191 .run_lifecycle(
192 operation_session_id,
193 LifecycleKind::Suspend,
194 move |state, session_id, cancelled| async move {
195 let _recovery_reservation = tokio::task::spawn_blocking({
196 let observer = state.recovery_observer.clone();
197 let session_id = session_id.clone();
198 let cancelled = cancelled.clone();
199 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
200 })
201 .await
202 .context("reserve recovery for daemon close task")??;
203 let mut controller = tokio::task::spawn_blocking(Controller::load)
204 .await
205 .context("load controller for daemon close task")??;
206 let executor = DaemonStageReportingExecutor::new(
207 CancellableProcessExecutor::new(cancelled),
208 state.clone(),
209 session_id.clone(),
210 );
211 let deferred = match route {
212 CloseRoute::RecoverInterrupted => {
213 controller
214 .recover_interrupted_close_managed(
215 &session_id,
216 &executor,
217 &state.session_manager,
218 )
219 .await?
220 }
221 CloseRoute::SettleWithoutCheckpoint => {
227 controller.suspend_session_without_checkpoint(&session_id, &executor)?
228 }
229 _ => {
230 controller
231 .suspend_session_managed_controlled(
232 &session_id,
233 &executor,
234 &state.session_manager,
235 acknowledge_unpublished_work,
236 )
237 .await?
238 }
239 };
240 Ok(if deferred {
241 DaemonLifecycleResult::DeferredCleanup
242 } else {
243 DaemonLifecycleResult::Done
244 })
245 },
246 )
247 .await?;
248 let _ = result; Ok(())
250 }
251
252 pub(super) fn start_deferred_cleanup(
253 self: &Arc<Self>,
254 session_id: String,
255 ) -> Result<LifecycleWatch> {
256 let result = self.start_or_join_lifecycle(
257 session_id.clone(),
258 LifecycleKind::Cleanup,
259 |state, session_id, cancelled| async move {
260 blocking(move || {
261 let mut controller = Controller::load()?;
262 let executor = DaemonStageReportingExecutor::new(
263 CancellableProcessExecutor::new(cancelled),
264 state,
265 session_id.clone(),
266 );
267 controller.cleanup_stopped_target(&session_id, &executor)?;
268 Ok(DaemonLifecycleResult::Done)
269 })
270 .await
271 },
272 )?;
273 let caller_result = result.clone();
274 let channel = result.clone();
275 let state = Arc::clone(self);
276 tokio::spawn(async move {
277 if let Err(error) = Self::wait_lifecycle_result(result).await {
278 tracing::warn!(%session_id, error = format!("{error:#}"), "deferred Podman cleanup failed");
279 state.push_notice(
280 &session_id,
281 "Container storage cleanup failed; the stopped session retains its target for retry.",
282 );
283 }
284 state.remove_completed_lifecycle(&channel);
285 });
286 Ok(caller_result)
287 }
288
289 pub(super) fn resume_retained_cleanups(self: &Arc<Self>) {
290 let session_ids = self
291 .controller
292 .lock()
293 .unwrap_or_else(PoisonError::into_inner)
294 .state
295 .sessions
296 .iter()
297 .filter(|(_, session)| {
298 session.state == SessionState::Stopped && session.target.is_some()
299 })
300 .map(|(session_id, _)| session_id.clone())
301 .collect::<Vec<_>>();
302 for session_id in session_ids {
303 if crate::controller::move_session::move_owns_session(&session_id) {
304 continue;
305 }
306 if let Err(error) = self.start_deferred_cleanup(session_id.clone()) {
307 tracing::warn!(%session_id, error = format!("{error:#}"), "could not resume deferred Podman cleanup");
308 self.push_notice(
309 &session_id,
310 format!("Could not resume container storage cleanup: {error:#}"),
311 );
312 }
313 }
314 }
315
316 pub(super) async fn wait_for_deferred_cleanup(
317 self: &Arc<Self>,
318 session_id: &str,
319 ) -> Result<()> {
320 let existing = {
321 let lifecycle = self
322 .lifecycle
323 .lock()
324 .unwrap_or_else(PoisonError::into_inner);
325 lifecycle.get(session_id).and_then(|active| {
326 (active.kind == LifecycleKind::Cleanup).then(|| active.result.clone())
327 })
328 };
329 let result = match existing {
330 Some(result) => result,
331 None => {
332 let needs_cleanup = blocking({
333 let session_id = session_id.to_owned();
334 move || {
335 let controller = Controller::load()?;
336 Ok(controller
337 .state
338 .sessions
339 .get(&session_id)
340 .is_some_and(|session| {
341 session.state == SessionState::Stopped && session.target.is_some()
342 }))
343 }
344 })
345 .await?;
346 if !needs_cleanup {
347 return Ok(());
348 }
349 self.start_deferred_cleanup(session_id.to_owned())?
350 }
351 };
352 let channel = result.clone();
353 let outcome = Self::wait_lifecycle_result(result).await;
354 self.remove_completed_lifecycle(&channel);
355 match outcome? {
356 DaemonLifecycleResult::Done => Ok(()),
357 DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
358 DaemonLifecycleResult::DeferredCleanup => {
359 unreachable!("cleanup cannot schedule another cleanup")
360 }
361 }
362 }
363}