Skip to main content

mj_controller/daemon/
close.rs

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    /// Make the intent durable before an HTTP caller receives acceptance.
19    /// Waiting and store writes happen in the action task, never the UI loop.
20    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        // Internal suspension (workspace close, recovery, child teardown) has
73        // already been admitted by its parent operation. Its checkpoint is
74        // still verified before any owned checkout is released.
75        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    /// The parent's sub-agents are stopped inside its close, once its
102    /// checkpoint is verified (see [`Self::stop_subagents_for_suspend`]).
103    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    /// Stop every active Mjolnir sub-agent of a parent whose suspend or
114    /// discard is about to close it, so only the parent is checkpointed.
115    ///
116    /// A close calls this after the parent's checkpoint is verified and
117    /// recorded, and before the parent's relay is sealed: the children do not
118    /// have to stop for the checkpoint, and a close that fails at its
119    /// checkpoint (a full disk, an archive directory it cannot write, a dirty
120    /// submodule) then leaves them running, with nothing to tell anyone.
121    ///
122    /// A child's durable output is the report it hands back, so a child is
123    /// stopped and removed the way a destroy removes one, with no checkpoint
124    /// of its own; its conversation is put into SessionWiki first, so it stays
125    /// searchable. The children are listed on the parent's record before any
126    /// of them is stopped, so the parent's model can be told about them when
127    /// the parent resumes, or at once when the close fails after this and
128    /// leaves the parent live.
129    ///
130    /// Only reading and recording that list can fail this. A child that
131    /// cannot be stopped is logged and its records are removed anyway: it
132    /// never fails the parent's suspend.
133    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        // One pass indexes the parent's whole tree, so each child is
166        // destroyed below without indexing it again.
167        self.index_before_destroy(parent_session_id).await;
168        for child in stopped {
169            let child_id = child.child_session_id;
170            // A sub-agent borrows its parent's worker and owns no branch.
171            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    /// [`Self::stop_subagents_for_suspend`] as the step a close runs once the
199    /// parent's checkpoint is verified.
200    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    /// Tell a parent that is still live which of its sub-agents a suspend or
210    /// a discard stopped before it failed.
211    ///
212    /// The list otherwise waits for the parent's next resume, and a parent
213    /// that never stopped may not resume for a long time; meanwhile its model
214    /// would wait on children that are gone, and the person would not know
215    /// why. So the relay takes the note for the parent's next prompt now, and
216    /// the conversation gets the line a resume records. A parent that is not
217    /// live keeps the list for its resume, and so does one whose relay
218    /// refuses the note.
219    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        // The relay owns the note now. Failing to forget the list only means
287        // a later resume tells the model again.
288        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    /// Forget a sub-agent whose stop failed: its record, its relation to the
303    /// parent, its conversation and its attachments. Its worker lives on the
304    /// parent's target, which the parent's suspend releases next.
305    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        // Cancellation is a request: the old owner must actually finish before
330        // close acquires the target, including an irreversible create commit.
331        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    /// A parent's close request, admitted once per (parent, request) so a
368    /// replay after a daemon restart joins the close already in flight
369    /// instead of starting another.
370    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            // `prepare_suspension` marks a live session `Closing`,
464            // so this is also the route of every live suspend.
465            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            // Nothing to archive and no relay to latch, so this
477            // close only tears down and settles. The route was
478            // decided after `wait_before_close` let any live
479            // create or resume finish, so a session that is still
480            // genuinely provisioning is not caught here.
481            CloseRoute::SettleWithoutCheckpoint => {
482                // No checkpoint can fail after the sub-agents stop.
483                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    /// Recovery is teardown only: never reconnect or reprovision a failed worker.
543    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                        // The decision follows lifecycle admission, so a competing
578                        // close or destroy cannot settle this snapshot under us.
579                        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}