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 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 if result.is_err() {
105 self.tell_live_parent_about_stopped_subagents(&session_id)
106 .await;
107 }
108 result
109 }
110
111 async fn suspend_with_children(
114 self: &Arc<Self>,
115 session_id: &str,
116 acknowledge_unpublished_work: bool,
117 ) -> Result<()> {
118 self.prepare_suspension(session_id).await?;
119 self.close_requested_session_with_ack(session_id.to_owned(), acknowledge_unpublished_work)
120 .await
121 }
122
123 pub(super) async fn stop_subagents_for_suspend(
144 self: &Arc<Self>,
145 parent_session_id: &str,
146 ) -> Result<()> {
147 let stopped = blocking({
148 let parent_session_id = parent_session_id.to_owned();
149 move || {
150 let controller = Controller::load()?;
151 let stopped = active_child_session_ids(&controller.state, &parent_session_id)
152 .iter()
153 .map(|child_id| {
154 crate::controller::stopped_subagent(&controller.state, child_id)
155 })
156 .collect::<Result<Vec<_>>>()?;
157 if !stopped.is_empty() {
158 crate::database::record_stopped_subagents(&parent_session_id, &stopped)?;
159 }
160 Ok(stopped)
161 }
162 })
163 .await
164 .context("list the sub-agents this suspend stops")?;
165 if stopped.is_empty() {
166 return Ok(());
167 }
168 let not_handed_back = stopped.iter().filter(|child| !child.handed_back).count();
169 tracing::info!(
170 session_id = %parent_session_id,
171 stopped = stopped.len(),
172 not_handed_back,
173 "stopping sub-agents before suspending their parent"
174 );
175 self.index_before_destroy(parent_session_id).await;
178 for child in stopped {
179 let child_id = child.child_session_id;
180 let Err(error) = Box::pin(self.force_destroy_indexed_session(
182 child_id.clone(),
183 BranchDisposition::Keep,
184 LifecycleKind::StopSubagent,
185 ))
186 .await
187 else {
188 continue;
189 };
190 tracing::warn!(
191 session_id = %parent_session_id,
192 %child_id,
193 error = format!("{error:#}"),
194 "could not stop a sub-agent for its parent's suspend; removing its records"
195 );
196 if let Err(error) = self.remove_subagent_records(&child_id).await {
197 tracing::warn!(
198 session_id = %parent_session_id,
199 %child_id,
200 error = format!("{error:#}"),
201 "could not remove the records of a sub-agent its parent's suspend stopped"
202 );
203 }
204 }
205 Ok(())
206 }
207
208 fn stop_subagents_before_close(self: &Arc<Self>, parent_session_id: &str) -> BeforeClose {
211 let state = Arc::clone(self);
212 let parent_session_id = parent_session_id.to_owned();
213 Box::pin(async move { state.stop_subagents_for_suspend(&parent_session_id).await })
214 }
215
216 pub(super) async fn tell_live_parent_about_stopped_subagents(
227 self: &Arc<Self>,
228 parent_session_id: &str,
229 ) {
230 let loaded = blocking({
231 let parent_session_id = parent_session_id.to_owned();
232 move || {
233 let live = Controller::load()?
234 .state
235 .sessions
236 .get(&parent_session_id)
237 .is_some_and(|session| {
238 matches!(
239 session.state,
240 SessionState::Running | SessionState::Disconnected
241 )
242 });
243 if !live {
244 return Ok(Vec::new());
245 }
246 crate::database::load_stopped_subagents(&parent_session_id)
247 }
248 })
249 .await;
250 let stopped = match loaded {
251 Ok(stopped) => stopped,
252 Err(error) => {
253 tracing::warn!(
254 session_id = %parent_session_id,
255 error = format!("{error:#}"),
256 "could not read the sub-agents a failed suspend stopped"
257 );
258 return;
259 }
260 };
261 let Some(context) = mj_core::subagent::stopped_subagents_prompt_context(&stopped) else {
262 return;
263 };
264 let delivered = async {
265 let handle = self
266 .session_manager
267 .session(parent_session_id.to_owned())
268 .await?;
269 handle.install_prompt_context(context).await?;
270 if let Some(text) = mj_core::subagent::stopped_subagents_notice(&stopped) {
271 handle
272 .submit(
273 new_command_id("stopped-subagents")?,
274 RelayCommand::RecordNotice { text },
275 )
276 .await?;
277 }
278 anyhow::Ok(())
279 }
280 .await;
281 if let Err(error) = delivered {
282 tracing::warn!(
283 session_id = %parent_session_id,
284 error = format!("{error:#}"),
285 "could not tell a live parent which sub-agents a failed suspend stopped; its next resume will"
286 );
287 return;
288 }
289 let delivered = stopped
290 .into_iter()
291 .map(|child| child.child_session_id)
292 .collect::<Vec<_>>();
293 if let Err(error) = blocking({
296 let parent_session_id = parent_session_id.to_owned();
297 move || crate::database::clear_stopped_subagents(&parent_session_id, &delivered)
298 })
299 .await
300 {
301 tracing::warn!(
302 session_id = %parent_session_id,
303 error = format!("{error:#}"),
304 "could not clear the stopped sub-agents after telling the live parent"
305 );
306 }
307 }
308
309 async fn remove_subagent_records(self: &Arc<Self>, child_id: &str) -> Result<()> {
313 self.preempt_active_lifecycle(child_id).await?;
314 blocking({
315 let child_id = child_id.to_owned();
316 move || {
317 if let Err(error) = mj_core::attachment::AttachmentStore::controller(&child_id)
318 .and_then(|store| store.remove_session_data())
319 {
320 tracing::warn!(
321 %child_id,
322 error = format!("{error:#}"),
323 "could not remove a stopped sub-agent's attachments"
324 );
325 }
326 crate::database::delete_session(&child_id)
327 }
328 })
329 .await?;
330 self.reload_controller().await?;
331 self.publish_revision();
332 Ok(())
333 }
334
335 pub(super) async fn wait_before_close(self: &Arc<Self>, session_id: &str) -> Result<()> {
336 let pending = {
339 let operations = self
340 .lifecycle
341 .lock()
342 .unwrap_or_else(PoisonError::into_inner);
343 operations
344 .get(session_id)
345 .filter(|operation| {
346 !matches!(
347 operation.kind,
348 LifecycleKind::Suspend | LifecycleKind::Cleanup
349 )
350 })
351 .map(|operation| {
352 operation.request_cancel();
353 operation.result.clone()
354 })
355 };
356 if let Some(pending) = pending {
357 if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
358 tracing::debug!(%session_id, %error, "previous lifecycle ended before close");
359 }
360 self.remove_completed_lifecycle(&pending);
361 }
362
363 Ok(())
364 }
365
366 async fn close_requested_session_with_ack(
367 self: &Arc<Self>,
368 session_id: String,
369 acknowledge_unpublished_work: bool,
370 ) -> Result<()> {
371 self.wait_before_close(&session_id).await?;
372 let route = blocking({
373 let session_id = session_id.clone();
374 move || {
375 let controller = Controller::load()?;
376 Ok(close_route(controller.state.sessions.get(&session_id)))
377 }
378 })
379 .await?;
380 match route {
381 CloseRoute::Done | CloseRoute::DeferredCleanup => {
382 self.stop_subagents_for_suspend(&session_id).await?;
385 if route == CloseRoute::DeferredCleanup {
386 self.start_deferred_cleanup(session_id)?;
387 }
388 return Ok(());
389 }
390 CloseRoute::Graceful
391 | CloseRoute::RecoverInterrupted
392 | CloseRoute::SettleWithoutCheckpoint => {}
393 }
394 let operation_session_id = session_id.clone();
395 let result = self
396 .run_lifecycle(
397 operation_session_id,
398 LifecycleKind::Suspend,
399 move |state, session_id, cancelled| async move {
400 let _recovery_reservation = tokio::task::spawn_blocking({
401 let observer = state.recovery_observer.clone();
402 let session_id = session_id.clone();
403 let cancelled = cancelled.clone();
404 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
405 })
406 .await
407 .context("reserve recovery for daemon close task")??;
408 let mut controller = tokio::task::spawn_blocking(Controller::load)
409 .await
410 .context("load controller for daemon close task")??;
411 let executor = DaemonStageReportingExecutor::new(
412 CancellableProcessExecutor::new(cancelled),
413 state.clone(),
414 session_id.clone(),
415 );
416 let deferred = match route {
417 CloseRoute::RecoverInterrupted => {
420 controller
421 .recover_interrupted_close_managed(
422 &session_id,
423 &executor,
424 &state.session_manager,
425 acknowledge_unpublished_work,
426 Some(state.stop_subagents_before_close(&session_id)),
427 )
428 .await?
429 }
430 CloseRoute::SettleWithoutCheckpoint => {
436 state.stop_subagents_for_suspend(&session_id).await?;
438 controller.suspend_session_without_checkpoint(&session_id, &executor)?
439 }
440 _ => {
441 controller
442 .suspend_session_managed_controlled(
443 &session_id,
444 &executor,
445 &state.session_manager,
446 acknowledge_unpublished_work,
447 Some(state.stop_subagents_before_close(&session_id)),
448 )
449 .await?
450 }
451 };
452 Ok(if deferred {
453 DaemonLifecycleResult::DeferredCleanup
454 } else {
455 DaemonLifecycleResult::Done
456 })
457 },
458 )
459 .await?;
460 let _ = result; Ok(())
462 }
463
464 pub(super) fn start_deferred_cleanup(
465 self: &Arc<Self>,
466 session_id: String,
467 ) -> Result<LifecycleWatch> {
468 let result = self.start_or_join_lifecycle(
469 session_id.clone(),
470 LifecycleKind::Cleanup,
471 |state, session_id, cancelled| async move {
472 blocking(move || {
473 let mut controller = Controller::load()?;
474 let executor = DaemonStageReportingExecutor::new(
475 CancellableProcessExecutor::new(cancelled),
476 state,
477 session_id.clone(),
478 );
479 controller.cleanup_stopped_target(&session_id, &executor)?;
480 Ok(DaemonLifecycleResult::Done)
481 })
482 .await
483 },
484 )?;
485 let caller_result = result.clone();
486 let channel = result.clone();
487 let state = Arc::clone(self);
488 tokio::spawn(async move {
489 if let Err(error) = Self::wait_lifecycle_result(result).await {
490 tracing::warn!(%session_id, error = format!("{error:#}"), "deferred Podman cleanup failed");
491 state.push_notice(
492 &session_id,
493 "Container storage cleanup failed; the stopped session retains its target for retry.",
494 );
495 }
496 state.remove_completed_lifecycle(&channel);
497 });
498 Ok(caller_result)
499 }
500
501 pub(super) fn resume_retained_cleanups(self: &Arc<Self>) {
502 let session_ids = self
503 .controller
504 .lock()
505 .unwrap_or_else(PoisonError::into_inner)
506 .state
507 .sessions
508 .iter()
509 .filter(|(_, session)| {
510 session.state == SessionState::Stopped && session.target.is_some()
511 })
512 .map(|(session_id, _)| session_id.clone())
513 .collect::<Vec<_>>();
514 for session_id in session_ids {
515 if crate::controller::move_session::move_owns_session(&session_id) {
516 continue;
517 }
518 if let Err(error) = self.start_deferred_cleanup(session_id.clone()) {
519 tracing::warn!(%session_id, error = format!("{error:#}"), "could not resume deferred Podman cleanup");
520 self.push_notice(
521 &session_id,
522 format!("Could not resume container storage cleanup: {error:#}"),
523 );
524 }
525 }
526 }
527
528 pub(super) async fn wait_for_deferred_cleanup(
529 self: &Arc<Self>,
530 session_id: &str,
531 ) -> Result<()> {
532 let existing = {
533 let lifecycle = self
534 .lifecycle
535 .lock()
536 .unwrap_or_else(PoisonError::into_inner);
537 lifecycle.get(session_id).and_then(|active| {
538 (active.kind == LifecycleKind::Cleanup).then(|| active.result.clone())
539 })
540 };
541 let result = match existing {
542 Some(result) => result,
543 None => {
544 let needs_cleanup = blocking({
545 let session_id = session_id.to_owned();
546 move || {
547 let controller = Controller::load()?;
548 Ok(controller
549 .state
550 .sessions
551 .get(&session_id)
552 .is_some_and(|session| {
553 session.state == SessionState::Stopped && session.target.is_some()
554 }))
555 }
556 })
557 .await?;
558 if !needs_cleanup {
559 return Ok(());
560 }
561 self.start_deferred_cleanup(session_id.to_owned())?
562 }
563 };
564 let channel = result.clone();
565 let outcome = Self::wait_lifecycle_result(result).await;
566 self.remove_completed_lifecycle(&channel);
567 match outcome? {
568 DaemonLifecycleResult::Done => Ok(()),
569 DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
570 DaemonLifecycleResult::DeferredCleanup => {
571 unreachable!("cleanup cannot schedule another cleanup")
572 }
573 }
574 }
575}