1use super::*;
2
3impl RuntimeState {
4 pub fn request_close(&self, session_id: &str) {
5 self.owner().close_requested.insert(session_id.to_owned());
6 self.publish_revision();
7 }
8
9 pub fn clear_close_request(&self, session_id: &str) {
10 self.owner().close_requested.remove(session_id);
11 self.publish_revision();
12 }
13
14 pub fn close_is_requested(&self, session_id: &str) -> bool {
15 self.owner().close_requested.contains(session_id)
16 }
17
18 pub async fn prepare_suspension(self: &Arc<Self>, session_id: &str) -> Result<()> {
21 self.wait_before_close(session_id).await?;
22 if self
23 .owner()
24 .lifecycle
25 .get(session_id)
26 .is_some_and(|operation| {
27 operation.kind == LifecycleKind::Suspend && operation.is_running()
28 })
29 {
30 return Ok(());
31 }
32 blocking({
33 let session_id = session_id.to_owned();
34 let observer = self.recovery_observer.clone();
35 move || {
36 let cancelled = AtomicBool::new(false);
37 let _reservation = reserve_recovery_or_cancel(&observer, &session_id, &cancelled)?;
38 ensure!(
39 !crate::controller::move_session::move_has_pending_queue(&session_id),
40 "Move queue admission is incomplete; retry Move on the same destination before suspending"
41 );
42 let mut controller = Controller::load()?;
43 let record = controller
44 .state
45 .sessions
46 .get_mut(&session_id)
47 .with_context(|| format!("unknown session {session_id}"))?;
48 if matches!(
49 record.state,
50 SessionState::Running
51 | SessionState::Disconnected
52 | SessionState::Checkpointing
53 ) {
54 record.state = SessionState::Closing;
55 record.last_error = None;
56 record.updated_at = chrono::Utc::now().to_rfc3339();
57 crate::database::save_lifecycle_session(record)?;
58 } else if record.public_error().is_some() {
59 record.last_error = None;
60 crate::database::save_lifecycle_session(record)?;
61 }
62 Ok(())
63 }
64 })
65 .await?;
66 self.reload_controller().await?;
67 self.publish_revision();
68 Ok(())
69 }
70
71 pub async fn suspend_session(self: &Arc<Self>, session_id: String) -> Result<()> {
72 self.suspend_session_with_ack(session_id, true).await
76 }
77
78 pub async fn suspend_session_with_ack(
79 self: &Arc<Self>,
80 session_id: String,
81 acknowledge_unpublished_work: bool,
82 ) -> Result<()> {
83 self.request_close(&session_id);
84 let result = self
85 .suspend_with_children(&session_id, acknowledge_unpublished_work)
86 .await;
87 if let Err(error) = &result {
88 let reference = new_command_id("suspension").unwrap_or_else(|_| "suspension".into());
89 tracing::warn!(%session_id, %reference, error = format!("{error:#}"), "session suspension failed");
90 self.record_failed_close(&session_id, &reference, &LifecycleFailure::of(error))
91 .await;
92 }
93 self.clear_close_request(&session_id);
94 if result.is_err() {
95 self.tell_live_parent_about_stopped_subagents(&session_id)
96 .await;
97 }
98 result
99 }
100
101 async fn suspend_with_children(
104 self: &Arc<Self>,
105 session_id: &str,
106 acknowledge_unpublished_work: bool,
107 ) -> Result<()> {
108 self.prepare_suspension(session_id).await?;
109 self.close_requested_session_with_ack(session_id.to_owned(), acknowledge_unpublished_work)
110 .await
111 }
112
113 pub(super) async fn stop_subagents_for_suspend(
134 self: &Arc<Self>,
135 parent_session_id: &str,
136 ) -> Result<()> {
137 let stopped = blocking({
138 let parent_session_id = parent_session_id.to_owned();
139 move || {
140 let controller = Controller::load()?;
141 let stopped = active_child_session_ids(&controller.state, &parent_session_id)
142 .iter()
143 .map(|child_id| {
144 crate::controller::stopped_subagent(&controller.state, child_id)
145 })
146 .collect::<Result<Vec<_>>>()?;
147 if !stopped.is_empty() {
148 crate::database::record_stopped_subagents(&parent_session_id, &stopped)?;
149 }
150 Ok(stopped)
151 }
152 })
153 .await
154 .context("list the sub-agents this suspend stops")?;
155 if stopped.is_empty() {
156 return Ok(());
157 }
158 let not_handed_back = stopped.iter().filter(|child| !child.handed_back).count();
159 tracing::info!(
160 session_id = %parent_session_id,
161 stopped = stopped.len(),
162 not_handed_back,
163 "stopping sub-agents before suspending their parent"
164 );
165 self.index_before_destroy(parent_session_id).await;
168 for child in stopped {
169 let child_id = child.child_session_id;
170 let Err(error) = Box::pin(self.force_destroy_indexed_session(
172 child_id.clone(),
173 BranchDisposition::Keep,
174 LifecycleKind::StopSubagent,
175 ))
176 .await
177 else {
178 continue;
179 };
180 tracing::warn!(
181 session_id = %parent_session_id,
182 %child_id,
183 error = format!("{error:#}"),
184 "could not stop a sub-agent for its parent's suspend; removing its records"
185 );
186 if let Err(error) = self.remove_subagent_records(&child_id).await {
187 tracing::warn!(
188 session_id = %parent_session_id,
189 %child_id,
190 error = format!("{error:#}"),
191 "could not remove the records of a sub-agent its parent's suspend stopped"
192 );
193 }
194 }
195 Ok(())
196 }
197
198 pub(super) fn stop_subagents_before_close(
201 self: &Arc<Self>,
202 parent_session_id: &str,
203 ) -> BeforeClose {
204 let state = Arc::clone(self);
205 let parent_session_id = parent_session_id.to_owned();
206 Box::pin(async move { state.stop_subagents_for_suspend(&parent_session_id).await })
207 }
208
209 pub(super) async fn tell_live_parent_about_stopped_subagents(
220 self: &Arc<Self>,
221 parent_session_id: &str,
222 ) {
223 let loaded = blocking({
224 let parent_session_id = parent_session_id.to_owned();
225 move || {
226 let live = Controller::load()?
227 .state
228 .sessions
229 .get(&parent_session_id)
230 .is_some_and(|session| {
231 matches!(
232 session.state,
233 SessionState::Running | SessionState::Disconnected
234 )
235 });
236 if !live {
237 return Ok(Vec::new());
238 }
239 crate::database::load_stopped_subagents(&parent_session_id)
240 }
241 })
242 .await;
243 let stopped = match loaded {
244 Ok(stopped) => stopped,
245 Err(error) => {
246 tracing::warn!(
247 session_id = %parent_session_id,
248 error = format!("{error:#}"),
249 "could not read the sub-agents a failed suspend stopped"
250 );
251 return;
252 }
253 };
254 let Some(context) = mj_core::subagent::stopped_subagents_prompt_context(&stopped) else {
255 return;
256 };
257 let delivered = async {
258 let handle = self
259 .session_manager
260 .session(parent_session_id.to_owned())
261 .await?;
262 handle.install_prompt_context(context).await?;
263 if let Some(text) = mj_core::subagent::stopped_subagents_notice(&stopped) {
264 handle
265 .submit(
266 new_command_id("stopped-subagents")?,
267 RelayCommand::RecordNotice { text },
268 )
269 .await?;
270 }
271 anyhow::Ok(())
272 }
273 .await;
274 if let Err(error) = delivered {
275 tracing::warn!(
276 session_id = %parent_session_id,
277 error = format!("{error:#}"),
278 "could not tell a live parent which sub-agents a failed suspend stopped; its next resume will"
279 );
280 return;
281 }
282 let delivered = stopped
283 .into_iter()
284 .map(|child| child.child_session_id)
285 .collect::<Vec<_>>();
286 if let Err(error) = blocking({
289 let parent_session_id = parent_session_id.to_owned();
290 move || crate::database::clear_stopped_subagents(&parent_session_id, &delivered)
291 })
292 .await
293 {
294 tracing::warn!(
295 session_id = %parent_session_id,
296 error = format!("{error:#}"),
297 "could not clear the stopped sub-agents after telling the live parent"
298 );
299 }
300 }
301
302 async fn remove_subagent_records(self: &Arc<Self>, child_id: &str) -> Result<()> {
306 self.preempt_active_lifecycle(child_id).await?;
307 blocking({
308 let child_id = child_id.to_owned();
309 move || {
310 if let Err(error) = mj_core::attachment::AttachmentStore::controller(&child_id)
311 .and_then(|store| store.remove_session_data())
312 {
313 tracing::warn!(
314 %child_id,
315 error = format!("{error:#}"),
316 "could not remove a stopped sub-agent's attachments"
317 );
318 }
319 crate::database::delete_session(&child_id)
320 }
321 })
322 .await?;
323 self.reload_controller().await?;
324 self.publish_revision();
325 Ok(())
326 }
327
328 pub(super) async fn wait_before_close(self: &Arc<Self>, session_id: &str) -> Result<()> {
329 let pending = {
332 let mut owner = self.owner();
333 let operations = &mut owner.lifecycle;
334 operations
335 .get_mut(session_id)
336 .filter(|operation| {
337 !matches!(
338 operation.kind,
339 LifecycleKind::Suspend | LifecycleKind::Cleanup
340 )
341 })
342 .map(|operation| {
343 operation.request_cancel();
344 (
345 operation.kind == LifecycleKind::Restart,
346 operation.result.clone(),
347 )
348 })
349 };
350 if let Some((cancelled_restart, pending)) = pending {
351 if cancelled_restart {
352 blocking({
353 let session_id = session_id.to_owned();
354 move || crate::database::cancel_session_restart(&session_id)
355 })
356 .await?;
357 }
358 if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
359 tracing::debug!(%session_id, %error, "previous lifecycle ended before close");
360 }
361 self.remove_completed_lifecycle(&pending);
362 }
363
364 Ok(())
365 }
366
367 pub async fn close_subagent_request(
371 self: &Arc<Self>,
372 session_id: String,
373 parent_session_id: String,
374 request_id: String,
375 ) -> Result<()> {
376 let key = serde_json::to_string(&(parent_session_id, request_id))?;
377 let result = self.start_or_join_lifecycle_with_key(
378 session_id,
379 LifecycleKind::Suspend,
380 None,
381 Some(key),
382 move |state, session_id, cancelled| async move {
383 state.suspend_admitted(session_id, cancelled, true).await
384 },
385 )?;
386 let completed = result.clone();
387 let outcome = Self::wait_lifecycle_result(result).await;
388 self.remove_completed_lifecycle(&completed);
389 outcome?;
390 Ok(())
391 }
392
393 async fn close_requested_session_with_ack(
394 self: &Arc<Self>,
395 session_id: String,
396 acknowledge_unpublished_work: bool,
397 ) -> Result<()> {
398 self.wait_before_close(&session_id).await?;
399 self.run_lifecycle(
400 session_id,
401 LifecycleKind::Suspend,
402 move |state, session_id, cancelled| async move {
403 state
404 .suspend_admitted(session_id, cancelled, acknowledge_unpublished_work)
405 .await
406 },
407 )
408 .await?;
409 Ok(())
410 }
411
412 async fn suspend_admitted(
413 self: &Arc<Self>,
414 session_id: String,
415 cancelled: Arc<AtomicBool>,
416 acknowledge_unpublished_work: bool,
417 ) -> Result<DaemonLifecycleResult> {
418 let _recovery_reservation = tokio::task::spawn_blocking({
419 let observer = self.recovery_observer.clone();
420 let session_id = session_id.clone();
421 let cancelled = cancelled.clone();
422 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
423 })
424 .await
425 .context("reserve recovery for daemon close task")??;
426 let mut controller = tokio::task::spawn_blocking(Controller::load)
427 .await
428 .context("load controller for daemon close task")??;
429 let executor = DaemonStageReportingExecutor::new(
430 CancellableProcessExecutor::new(cancelled),
431 self.clone(),
432 session_id.clone(),
433 );
434 self.suspend_with_loaded_controller(
435 &session_id,
436 &mut controller,
437 &executor,
438 acknowledge_unpublished_work,
439 )
440 .await
441 }
442
443 pub(super) async fn suspend_with_loaded_controller(
444 self: &Arc<Self>,
445 session_id: &str,
446 controller: &mut Controller,
447 executor: &(impl CommandExecutor + Sync),
448 acknowledge_unpublished_work: bool,
449 ) -> Result<DaemonLifecycleResult> {
450 let route = close_route(
451 controller.state.sessions.get(session_id),
452 controller.state.subagents.contains_key(session_id),
453 );
454 if matches!(route, CloseRoute::Done | CloseRoute::DeferredCleanup) {
455 self.stop_subagents_for_suspend(session_id).await?;
456 return Ok(if route == CloseRoute::DeferredCleanup {
457 DaemonLifecycleResult::DeferredCleanup
458 } else {
459 DaemonLifecycleResult::Done
460 });
461 }
462 let deferred = match route {
463 CloseRoute::RecoverInterrupted => {
466 controller
467 .recover_interrupted_close_managed(
468 session_id,
469 executor,
470 &self.session_manager,
471 acknowledge_unpublished_work,
472 Some(self.stop_subagents_before_close(session_id)),
473 )
474 .await?
475 }
476 CloseRoute::SettleWithoutCheckpoint => {
482 self.stop_subagents_for_suspend(session_id).await?;
484 controller.suspend_session_without_checkpoint(session_id, executor)?
485 }
486 _ => {
487 controller
488 .suspend_session_managed_controlled(
489 session_id,
490 executor,
491 &self.session_manager,
492 acknowledge_unpublished_work,
493 Some(self.stop_subagents_before_close(session_id)),
494 )
495 .await?
496 }
497 };
498 Ok(if deferred {
499 DaemonLifecycleResult::DeferredCleanup
500 } else {
501 DaemonLifecycleResult::Done
502 })
503 }
504
505 pub(super) fn start_deferred_cleanup(
506 self: &Arc<Self>,
507 session_id: String,
508 ) -> Result<LifecycleWatch> {
509 let result = self.start_or_join_lifecycle(
510 session_id.clone(),
511 LifecycleKind::Cleanup,
512 |state, session_id, cancelled| async move {
513 blocking(move || {
514 let mut controller = Controller::load()?;
515 let executor = DaemonStageReportingExecutor::new(
516 CancellableProcessExecutor::new(cancelled),
517 state,
518 session_id.clone(),
519 );
520 controller.cleanup_stopped_target(&session_id, &executor)?;
521 Ok(DaemonLifecycleResult::Done)
522 })
523 .await
524 },
525 )?;
526 let caller_result = result.clone();
527 let channel = result.clone();
528 let state = Arc::clone(self);
529 tokio::spawn(async move {
530 if let Err(error) = Self::wait_lifecycle_result(result).await {
531 tracing::warn!(%session_id, error = format!("{error:#}"), "deferred Podman cleanup failed");
532 state.push_notice(
533 &session_id,
534 "Container storage cleanup failed; the stopped session retains its target for retry.",
535 );
536 }
537 state.remove_completed_lifecycle(&channel);
538 });
539 Ok(caller_result)
540 }
541
542 pub(super) fn resume_startup_cleanups(self: &Arc<Self>, immediately: bool) {
544 self.resume_move_destination_cleanups(immediately);
545 let ids = {
546 let owner = self.owner();
547 owner
548 .controller()
549 .state
550 .sessions
551 .values()
552 .filter(|session| {
553 session.state == SessionState::StartupCleanup
554 && !owner
555 .lifecycle
556 .get(&session.id)
557 .is_some_and(|operation| operation.is_running())
558 && (immediately
559 || chrono::DateTime::parse_from_rfc3339(&session.updated_at)
560 .map(|at| {
561 (chrono::Utc::now() - at.with_timezone(&chrono::Utc))
562 .num_seconds()
563 >= 30
564 })
565 .unwrap_or(true))
566 })
567 .map(|session| session.id.clone())
568 .collect::<Vec<_>>()
569 };
570 for session_id in ids {
571 let result = self.start_or_join_lifecycle(
572 session_id.clone(),
573 LifecycleKind::StartupCleanup,
574 |_state, session_id, _cancelled| async move {
575 blocking(move || {
576 let mut controller = Controller::load()?;
577 if controller
580 .state
581 .sessions
582 .get(&session_id)
583 .is_none_or(|record| record.state != SessionState::StartupCleanup)
584 {
585 return Ok(DaemonLifecycleResult::Done);
586 }
587 let executor = crate::controller::failed_launch_cleanup_executor();
588 controller.cleanup_failed_startup_controlled(&session_id, &executor)?;
589 Ok(DaemonLifecycleResult::Done)
590 })
591 .await
592 },
593 );
594 match result {
595 Ok(result) => {
596 let state = self.clone();
597 let completed = result.clone();
598 tokio::spawn(async move {
599 if let Err(error) = Self::wait_lifecycle_result(result).await {
600 tracing::warn!(%session_id, error=format!("{error:#}"), "failed startup cleanup remains pending; retry in 30s");
601 }
602 state.remove_completed_lifecycle(&completed);
603 });
604 }
605 Err(error) => {
606 tracing::debug!(%session_id, error=format!("{error:#}"), "startup cleanup admission deferred")
607 }
608 }
609 }
610 }
611
612 pub(super) fn resume_retained_cleanups(self: &Arc<Self>) {
613 let session_ids = self
614 .owner()
615 .controller()
616 .state
617 .sessions
618 .iter()
619 .filter(|(_, session)| {
620 session.state == SessionState::Stopped && session.target.is_some()
621 })
622 .map(|(session_id, _)| session_id.clone())
623 .collect::<Vec<_>>();
624 for session_id in session_ids {
625 if crate::controller::move_session::move_owns_session(&session_id) {
626 continue;
627 }
628 if self
629 .owner()
630 .lifecycle
631 .get(&session_id)
632 .is_some_and(|active| active.kind == LifecycleKind::Restart && active.is_running())
633 {
634 continue;
635 }
636 if let Err(error) = self.start_deferred_cleanup(session_id.clone()) {
637 tracing::warn!(%session_id, error = format!("{error:#}"), "could not resume deferred Podman cleanup");
638 self.push_notice(
639 &session_id,
640 format!("Could not resume container storage cleanup: {error:#}"),
641 );
642 }
643 }
644 }
645
646 pub(super) async fn wait_for_deferred_cleanup(
647 self: &Arc<Self>,
648 session_id: &str,
649 ) -> Result<()> {
650 let existing = {
651 let lifecycle_owner = self.owner();
652 let lifecycle = &lifecycle_owner.lifecycle;
653 lifecycle.get(session_id).and_then(|active| {
654 (active.kind == LifecycleKind::Cleanup).then(|| active.result.clone())
655 })
656 };
657 let result = match existing {
658 Some(result) => result,
659 None => {
660 let needs_cleanup = blocking({
661 let session_id = session_id.to_owned();
662 move || {
663 let controller = Controller::load()?;
664 Ok(controller
665 .state
666 .sessions
667 .get(&session_id)
668 .is_some_and(|session| {
669 session.state == SessionState::Stopped && session.target.is_some()
670 }))
671 }
672 })
673 .await?;
674 if !needs_cleanup {
675 return Ok(());
676 }
677 self.start_deferred_cleanup(session_id.to_owned())?
678 }
679 };
680 let channel = result.clone();
681 let outcome = Self::wait_lifecycle_result(result).await;
682 self.remove_completed_lifecycle(&channel);
683 match outcome? {
684 DaemonLifecycleResult::Done => Ok(()),
685 DaemonLifecycleResult::Move(_) | DaemonLifecycleResult::Park(_) => {
686 unreachable!("cleanup cannot return a move outcome")
687 }
688 DaemonLifecycleResult::DeferredCleanup => {
689 unreachable!("cleanup cannot schedule another cleanup")
690 }
691 }
692 }
693}