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 fn stop_subagents_before_close(self: &Arc<Self>, parent_session_id: &str) -> BeforeClose {
201 let state = Arc::clone(self);
202 let parent_session_id = parent_session_id.to_owned();
203 Box::pin(async move { state.stop_subagents_for_suspend(&parent_session_id).await })
204 }
205
206 pub(super) async fn tell_live_parent_about_stopped_subagents(
217 self: &Arc<Self>,
218 parent_session_id: &str,
219 ) {
220 let loaded = blocking({
221 let parent_session_id = parent_session_id.to_owned();
222 move || {
223 let live = Controller::load()?
224 .state
225 .sessions
226 .get(&parent_session_id)
227 .is_some_and(|session| {
228 matches!(
229 session.state,
230 SessionState::Running | SessionState::Disconnected
231 )
232 });
233 if !live {
234 return Ok(Vec::new());
235 }
236 crate::database::load_stopped_subagents(&parent_session_id)
237 }
238 })
239 .await;
240 let stopped = match loaded {
241 Ok(stopped) => stopped,
242 Err(error) => {
243 tracing::warn!(
244 session_id = %parent_session_id,
245 error = format!("{error:#}"),
246 "could not read the sub-agents a failed suspend stopped"
247 );
248 return;
249 }
250 };
251 let Some(context) = mj_core::subagent::stopped_subagents_prompt_context(&stopped) else {
252 return;
253 };
254 let delivered = async {
255 let handle = self
256 .session_manager
257 .session(parent_session_id.to_owned())
258 .await?;
259 handle.install_prompt_context(context).await?;
260 if let Some(text) = mj_core::subagent::stopped_subagents_notice(&stopped) {
261 handle
262 .submit(
263 new_command_id("stopped-subagents")?,
264 RelayCommand::RecordNotice { text },
265 )
266 .await?;
267 }
268 anyhow::Ok(())
269 }
270 .await;
271 if let Err(error) = delivered {
272 tracing::warn!(
273 session_id = %parent_session_id,
274 error = format!("{error:#}"),
275 "could not tell a live parent which sub-agents a failed suspend stopped; its next resume will"
276 );
277 return;
278 }
279 let delivered = stopped
280 .into_iter()
281 .map(|child| child.child_session_id)
282 .collect::<Vec<_>>();
283 if let Err(error) = blocking({
286 let parent_session_id = parent_session_id.to_owned();
287 move || crate::database::clear_stopped_subagents(&parent_session_id, &delivered)
288 })
289 .await
290 {
291 tracing::warn!(
292 session_id = %parent_session_id,
293 error = format!("{error:#}"),
294 "could not clear the stopped sub-agents after telling the live parent"
295 );
296 }
297 }
298
299 async fn remove_subagent_records(self: &Arc<Self>, child_id: &str) -> Result<()> {
303 self.preempt_active_lifecycle(child_id).await?;
304 blocking({
305 let child_id = child_id.to_owned();
306 move || {
307 if let Err(error) = mj_core::attachment::AttachmentStore::controller(&child_id)
308 .and_then(|store| store.remove_session_data())
309 {
310 tracing::warn!(
311 %child_id,
312 error = format!("{error:#}"),
313 "could not remove a stopped sub-agent's attachments"
314 );
315 }
316 crate::database::delete_session(&child_id)
317 }
318 })
319 .await?;
320 self.reload_controller().await?;
321 self.publish_revision();
322 Ok(())
323 }
324
325 pub(super) async fn wait_before_close(self: &Arc<Self>, session_id: &str) -> Result<()> {
326 let pending = {
329 let mut owner = self.owner();
330 let operations = &mut owner.lifecycle;
331 operations
332 .get_mut(session_id)
333 .filter(|operation| {
334 !matches!(
335 operation.kind,
336 LifecycleKind::Suspend | LifecycleKind::Cleanup
337 )
338 })
339 .map(|operation| {
340 operation.request_cancel();
341 operation.result.clone()
342 })
343 };
344 if let Some(pending) = pending {
345 if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
346 tracing::debug!(%session_id, %error, "previous lifecycle ended before close");
347 }
348 self.remove_completed_lifecycle(&pending);
349 }
350
351 Ok(())
352 }
353
354 pub async fn close_subagent_request(
358 self: &Arc<Self>,
359 session_id: String,
360 parent_session_id: String,
361 request_id: String,
362 ) -> Result<()> {
363 let key = serde_json::to_string(&(parent_session_id, request_id))?;
364 let result = self.start_or_join_lifecycle_with_key(
365 session_id,
366 LifecycleKind::Suspend,
367 None,
368 Some(key),
369 move |state, session_id, cancelled| async move {
370 state.suspend_admitted(session_id, cancelled, true).await
371 },
372 )?;
373 let completed = result.clone();
374 let outcome = Self::wait_lifecycle_result(result).await;
375 self.remove_completed_lifecycle(&completed);
376 outcome?;
377 Ok(())
378 }
379
380 async fn close_requested_session_with_ack(
381 self: &Arc<Self>,
382 session_id: String,
383 acknowledge_unpublished_work: bool,
384 ) -> Result<()> {
385 self.wait_before_close(&session_id).await?;
386 self.run_lifecycle(
387 session_id,
388 LifecycleKind::Suspend,
389 move |state, session_id, cancelled| async move {
390 state
391 .suspend_admitted(session_id, cancelled, acknowledge_unpublished_work)
392 .await
393 },
394 )
395 .await?;
396 Ok(())
397 }
398
399 async fn suspend_admitted(
400 self: &Arc<Self>,
401 session_id: String,
402 cancelled: Arc<AtomicBool>,
403 acknowledge_unpublished_work: bool,
404 ) -> Result<DaemonLifecycleResult> {
405 let _recovery_reservation = tokio::task::spawn_blocking({
406 let observer = self.recovery_observer.clone();
407 let session_id = session_id.clone();
408 let cancelled = cancelled.clone();
409 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
410 })
411 .await
412 .context("reserve recovery for daemon close task")??;
413 let mut controller = tokio::task::spawn_blocking(Controller::load)
414 .await
415 .context("load controller for daemon close task")??;
416 let route = close_route(
417 controller.state.sessions.get(&session_id),
418 controller.state.subagents.contains_key(&session_id),
419 );
420 if matches!(route, CloseRoute::Done | CloseRoute::DeferredCleanup) {
421 self.stop_subagents_for_suspend(&session_id).await?;
422 return Ok(if route == CloseRoute::DeferredCleanup {
423 DaemonLifecycleResult::DeferredCleanup
424 } else {
425 DaemonLifecycleResult::Done
426 });
427 }
428 let executor = DaemonStageReportingExecutor::new(
429 CancellableProcessExecutor::new(cancelled),
430 self.clone(),
431 session_id.clone(),
432 );
433 let deferred = match route {
434 CloseRoute::RecoverInterrupted => {
437 controller
438 .recover_interrupted_close_managed(
439 &session_id,
440 &executor,
441 &self.session_manager,
442 acknowledge_unpublished_work,
443 Some(self.stop_subagents_before_close(&session_id)),
444 )
445 .await?
446 }
447 CloseRoute::SettleWithoutCheckpoint => {
453 self.stop_subagents_for_suspend(&session_id).await?;
455 controller.suspend_session_without_checkpoint(&session_id, &executor)?
456 }
457 _ => {
458 controller
459 .suspend_session_managed_controlled(
460 &session_id,
461 &executor,
462 &self.session_manager,
463 acknowledge_unpublished_work,
464 Some(self.stop_subagents_before_close(&session_id)),
465 )
466 .await?
467 }
468 };
469 Ok(if deferred {
470 DaemonLifecycleResult::DeferredCleanup
471 } else {
472 DaemonLifecycleResult::Done
473 })
474 }
475
476 pub(super) fn start_deferred_cleanup(
477 self: &Arc<Self>,
478 session_id: String,
479 ) -> Result<LifecycleWatch> {
480 let result = self.start_or_join_lifecycle(
481 session_id.clone(),
482 LifecycleKind::Cleanup,
483 |state, session_id, cancelled| async move {
484 blocking(move || {
485 let mut controller = Controller::load()?;
486 let executor = DaemonStageReportingExecutor::new(
487 CancellableProcessExecutor::new(cancelled),
488 state,
489 session_id.clone(),
490 );
491 controller.cleanup_stopped_target(&session_id, &executor)?;
492 Ok(DaemonLifecycleResult::Done)
493 })
494 .await
495 },
496 )?;
497 let caller_result = result.clone();
498 let channel = result.clone();
499 let state = Arc::clone(self);
500 tokio::spawn(async move {
501 if let Err(error) = Self::wait_lifecycle_result(result).await {
502 tracing::warn!(%session_id, error = format!("{error:#}"), "deferred Podman cleanup failed");
503 state.push_notice(
504 &session_id,
505 "Container storage cleanup failed; the stopped session retains its target for retry.",
506 );
507 }
508 state.remove_completed_lifecycle(&channel);
509 });
510 Ok(caller_result)
511 }
512
513 pub(super) fn resume_startup_cleanups(self: &Arc<Self>, immediately: bool) {
515 self.resume_move_destination_cleanups(immediately);
516 let ids = {
517 let owner = self.owner();
518 owner
519 .controller()
520 .state
521 .sessions
522 .values()
523 .filter(|session| {
524 session.state == SessionState::StartupCleanup
525 && !owner
526 .lifecycle
527 .get(&session.id)
528 .is_some_and(|operation| operation.is_running())
529 && (immediately
530 || chrono::DateTime::parse_from_rfc3339(&session.updated_at)
531 .map(|at| {
532 (chrono::Utc::now() - at.with_timezone(&chrono::Utc))
533 .num_seconds()
534 >= 30
535 })
536 .unwrap_or(true))
537 })
538 .map(|session| session.id.clone())
539 .collect::<Vec<_>>()
540 };
541 for session_id in ids {
542 let result = self.start_or_join_lifecycle(
543 session_id.clone(),
544 LifecycleKind::StartupCleanup,
545 |_state, session_id, _cancelled| async move {
546 blocking(move || {
547 let mut controller = Controller::load()?;
548 if controller
551 .state
552 .sessions
553 .get(&session_id)
554 .is_none_or(|record| record.state != SessionState::StartupCleanup)
555 {
556 return Ok(DaemonLifecycleResult::Done);
557 }
558 let executor = crate::controller::failed_launch_cleanup_executor();
559 controller.cleanup_failed_startup_controlled(&session_id, &executor)?;
560 Ok(DaemonLifecycleResult::Done)
561 })
562 .await
563 },
564 );
565 match result {
566 Ok(result) => {
567 let state = self.clone();
568 let completed = result.clone();
569 tokio::spawn(async move {
570 if let Err(error) = Self::wait_lifecycle_result(result).await {
571 tracing::warn!(%session_id, error=format!("{error:#}"), "failed startup cleanup remains pending; retry in 30s");
572 }
573 state.remove_completed_lifecycle(&completed);
574 });
575 }
576 Err(error) => {
577 tracing::debug!(%session_id, error=format!("{error:#}"), "startup cleanup admission deferred")
578 }
579 }
580 }
581 }
582
583 pub(super) fn resume_retained_cleanups(self: &Arc<Self>) {
584 let session_ids = self
585 .owner()
586 .controller()
587 .state
588 .sessions
589 .iter()
590 .filter(|(_, session)| {
591 session.state == SessionState::Stopped && session.target.is_some()
592 })
593 .map(|(session_id, _)| session_id.clone())
594 .collect::<Vec<_>>();
595 for session_id in session_ids {
596 if crate::controller::move_session::move_owns_session(&session_id) {
597 continue;
598 }
599 if let Err(error) = self.start_deferred_cleanup(session_id.clone()) {
600 tracing::warn!(%session_id, error = format!("{error:#}"), "could not resume deferred Podman cleanup");
601 self.push_notice(
602 &session_id,
603 format!("Could not resume container storage cleanup: {error:#}"),
604 );
605 }
606 }
607 }
608
609 pub(super) async fn wait_for_deferred_cleanup(
610 self: &Arc<Self>,
611 session_id: &str,
612 ) -> Result<()> {
613 let existing = {
614 let lifecycle_owner = self.owner();
615 let lifecycle = &lifecycle_owner.lifecycle;
616 lifecycle.get(session_id).and_then(|active| {
617 (active.kind == LifecycleKind::Cleanup).then(|| active.result.clone())
618 })
619 };
620 let result = match existing {
621 Some(result) => result,
622 None => {
623 let needs_cleanup = blocking({
624 let session_id = session_id.to_owned();
625 move || {
626 let controller = Controller::load()?;
627 Ok(controller
628 .state
629 .sessions
630 .get(&session_id)
631 .is_some_and(|session| {
632 session.state == SessionState::Stopped && session.target.is_some()
633 }))
634 }
635 })
636 .await?;
637 if !needs_cleanup {
638 return Ok(());
639 }
640 self.start_deferred_cleanup(session_id.to_owned())?
641 }
642 };
643 let channel = result.clone();
644 let outcome = Self::wait_lifecycle_result(result).await;
645 self.remove_completed_lifecycle(&channel);
646 match outcome? {
647 DaemonLifecycleResult::Done => Ok(()),
648 DaemonLifecycleResult::Move(_) | DaemonLifecycleResult::Park(_) => {
649 unreachable!("cleanup cannot return a move outcome")
650 }
651 DaemonLifecycleResult::DeferredCleanup => {
652 unreachable!("cleanup cannot schedule another cleanup")
653 }
654 }
655 }
656}