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    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    /// Tell a parent that is still live which of its sub-agents a suspend or
207    /// a discard stopped before it failed.
208    ///
209    /// The list otherwise waits for the parent's next resume, and a parent
210    /// that never stopped may not resume for a long time; meanwhile its model
211    /// would wait on children that are gone, and the person would not know
212    /// why. So the relay takes the note for the parent's next prompt now, and
213    /// the conversation gets the line a resume records. A parent that is not
214    /// live keeps the list for its resume, and so does one whose relay
215    /// refuses the note.
216    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        // The relay owns the note now. Failing to forget the list only means
284        // a later resume tells the model again.
285        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    /// Forget a sub-agent whose stop failed: its record, its relation to the
300    /// parent, its conversation and its attachments. Its worker lives on the
301    /// parent's target, which the parent's suspend releases next.
302    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        // Cancellation is a request: the old owner must actually finish before
327        // close acquires the target, including an irreversible create commit.
328        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    /// A parent's close request, admitted once per (parent, request) so a
355    /// replay after a daemon restart joins the close already in flight
356    /// instead of starting another.
357    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            // `prepare_suspension` marks a live session `Closing`,
435            // so this is also the route of every live suspend.
436            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            // Nothing to archive and no relay to latch, so this
448            // close only tears down and settles. The route was
449            // decided after `wait_before_close` let any live
450            // create or resume finish, so a session that is still
451            // genuinely provisioning is not caught here.
452            CloseRoute::SettleWithoutCheckpoint => {
453                // No checkpoint can fail after the sub-agents stop.
454                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    /// Recovery is teardown only: never reconnect or reprovision a failed worker.
514    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                        // The decision follows lifecycle admission, so a competing
549                        // close or destroy cannot settle this snapshot under us.
550                        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}