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.close_requested
6            .lock()
7            .unwrap_or_else(PoisonError::into_inner)
8            .insert(session_id.to_owned());
9        self.publish_revision();
10    }
11
12    pub fn clear_close_request(&self, session_id: &str) {
13        self.close_requested
14            .lock()
15            .unwrap_or_else(PoisonError::into_inner)
16            .remove(session_id);
17        self.publish_revision();
18    }
19
20    pub fn close_is_requested(&self, session_id: &str) -> bool {
21        self.close_requested
22            .lock()
23            .unwrap_or_else(PoisonError::into_inner)
24            .contains(session_id)
25    }
26
27    /// Make the intent durable before an HTTP caller receives acceptance.
28    /// Waiting and store writes happen in the action task, never the UI loop.
29    pub async fn prepare_suspension(self: &Arc<Self>, session_id: &str) -> Result<()> {
30        self.wait_before_close(session_id).await?;
31        if self
32            .lifecycle
33            .lock()
34            .unwrap_or_else(PoisonError::into_inner)
35            .get(session_id)
36            .is_some_and(|operation| {
37                operation.kind == LifecycleKind::Suspend && operation.result.borrow().is_none()
38            })
39        {
40            return Ok(());
41        }
42        blocking({
43            let session_id = session_id.to_owned();
44            let observer = self.recovery_observer.clone();
45            move || {
46                let cancelled = AtomicBool::new(false);
47                let _reservation = reserve_recovery_or_cancel(&observer, &session_id, &cancelled)?;
48                ensure!(
49                    !crate::controller::move_session::move_has_pending_queue(&session_id),
50                    "Move queue admission is incomplete; retry Move on the same destination before suspending"
51                );
52                let mut controller = Controller::load()?;
53                let record = controller
54                    .state
55                    .sessions
56                    .get_mut(&session_id)
57                    .with_context(|| format!("unknown session {session_id}"))?;
58                if matches!(
59                    record.state,
60                    SessionState::Running
61                        | SessionState::Disconnected
62                        | SessionState::Checkpointing
63                ) {
64                    record.state = SessionState::Closing;
65                    record.last_error = None;
66                    record.updated_at = chrono::Utc::now().to_rfc3339();
67                    crate::database::save_lifecycle_session(record)?;
68                } else if record.public_error().is_some() {
69                    record.last_error = None;
70                    crate::database::save_lifecycle_session(record)?;
71                }
72                Ok(())
73            }
74        })
75        .await?;
76        self.reload_controller().await?;
77        self.publish_revision();
78        Ok(())
79    }
80
81    pub async fn suspend_session(self: &Arc<Self>, session_id: String) -> Result<()> {
82        // Internal suspension (workspace close, recovery, child teardown) has
83        // already been admitted by its parent operation. Its checkpoint is
84        // still verified before any owned checkout is released.
85        self.suspend_session_with_ack(session_id, true).await
86    }
87
88    pub async fn suspend_session_with_ack(
89        self: &Arc<Self>,
90        session_id: String,
91        acknowledge_unpublished_work: bool,
92    ) -> Result<()> {
93        self.request_close(&session_id);
94        let result = self
95            .suspend_with_children(&session_id, acknowledge_unpublished_work)
96            .await;
97        if let Err(error) = &result {
98            let reference = new_command_id("suspension").unwrap_or_else(|_| "suspension".into());
99            tracing::warn!(%session_id, %reference, error = format!("{error:#}"), "session suspension failed");
100            self.record_failed_close(&session_id, &reference, &LifecycleFailure::of(error))
101                .await;
102        }
103        self.clear_close_request(&session_id);
104        if result.is_err() {
105            self.tell_live_parent_about_stopped_subagents(&session_id)
106                .await;
107        }
108        result
109    }
110
111    /// The parent's sub-agents are stopped inside its close, once its
112    /// checkpoint is verified (see [`Self::stop_subagents_for_suspend`]).
113    async fn suspend_with_children(
114        self: &Arc<Self>,
115        session_id: &str,
116        acknowledge_unpublished_work: bool,
117    ) -> Result<()> {
118        self.prepare_suspension(session_id).await?;
119        self.close_requested_session_with_ack(session_id.to_owned(), acknowledge_unpublished_work)
120            .await
121    }
122
123    /// Stop every active Mjolnir sub-agent of a parent whose suspend or
124    /// discard is about to close it, so only the parent is checkpointed.
125    ///
126    /// A close calls this after the parent's checkpoint is verified and
127    /// recorded, and before the parent's relay is sealed: the children do not
128    /// have to stop for the checkpoint, and a close that fails at its
129    /// checkpoint (a full disk, an archive directory it cannot write, a dirty
130    /// submodule) then leaves them running, with nothing to tell anyone.
131    ///
132    /// A child's durable output is the report it hands back, so a child is
133    /// stopped and removed the way a destroy removes one, with no checkpoint
134    /// of its own; its conversation is put into SessionWiki first, so it stays
135    /// searchable. The children are listed on the parent's record before any
136    /// of them is stopped, so the parent's model can be told about them when
137    /// the parent resumes, or at once when the close fails after this and
138    /// leaves the parent live.
139    ///
140    /// Only reading and recording that list can fail this. A child that
141    /// cannot be stopped is logged and its records are removed anyway: it
142    /// never fails the parent's suspend.
143    pub(super) async fn stop_subagents_for_suspend(
144        self: &Arc<Self>,
145        parent_session_id: &str,
146    ) -> Result<()> {
147        let stopped = blocking({
148            let parent_session_id = parent_session_id.to_owned();
149            move || {
150                let controller = Controller::load()?;
151                let stopped = active_child_session_ids(&controller.state, &parent_session_id)
152                    .iter()
153                    .map(|child_id| {
154                        crate::controller::stopped_subagent(&controller.state, child_id)
155                    })
156                    .collect::<Result<Vec<_>>>()?;
157                if !stopped.is_empty() {
158                    crate::database::record_stopped_subagents(&parent_session_id, &stopped)?;
159                }
160                Ok(stopped)
161            }
162        })
163        .await
164        .context("list the sub-agents this suspend stops")?;
165        if stopped.is_empty() {
166            return Ok(());
167        }
168        let not_handed_back = stopped.iter().filter(|child| !child.handed_back).count();
169        tracing::info!(
170            session_id = %parent_session_id,
171            stopped = stopped.len(),
172            not_handed_back,
173            "stopping sub-agents before suspending their parent"
174        );
175        // One pass indexes the parent's whole tree, so each child is
176        // destroyed below without indexing it again.
177        self.index_before_destroy(parent_session_id).await;
178        for child in stopped {
179            let child_id = child.child_session_id;
180            // A sub-agent borrows its parent's worker and owns no branch.
181            let Err(error) = Box::pin(self.force_destroy_indexed_session(
182                child_id.clone(),
183                BranchDisposition::Keep,
184                LifecycleKind::StopSubagent,
185            ))
186            .await
187            else {
188                continue;
189            };
190            tracing::warn!(
191                session_id = %parent_session_id,
192                %child_id,
193                error = format!("{error:#}"),
194                "could not stop a sub-agent for its parent's suspend; removing its records"
195            );
196            if let Err(error) = self.remove_subagent_records(&child_id).await {
197                tracing::warn!(
198                    session_id = %parent_session_id,
199                    %child_id,
200                    error = format!("{error:#}"),
201                    "could not remove the records of a sub-agent its parent's suspend stopped"
202                );
203            }
204        }
205        Ok(())
206    }
207
208    /// [`Self::stop_subagents_for_suspend`] as the step a close runs once the
209    /// parent's checkpoint is verified.
210    fn stop_subagents_before_close(self: &Arc<Self>, parent_session_id: &str) -> BeforeClose {
211        let state = Arc::clone(self);
212        let parent_session_id = parent_session_id.to_owned();
213        Box::pin(async move { state.stop_subagents_for_suspend(&parent_session_id).await })
214    }
215
216    /// Tell a parent that is still live which of its sub-agents a suspend or
217    /// a discard stopped before it failed.
218    ///
219    /// The list otherwise waits for the parent's next resume, and a parent
220    /// that never stopped may not resume for a long time; meanwhile its model
221    /// would wait on children that are gone, and the person would not know
222    /// why. So the relay takes the note for the parent's next prompt now, and
223    /// the conversation gets the line a resume records. A parent that is not
224    /// live keeps the list for its resume, and so does one whose relay
225    /// refuses the note.
226    pub(super) async fn tell_live_parent_about_stopped_subagents(
227        self: &Arc<Self>,
228        parent_session_id: &str,
229    ) {
230        let loaded = blocking({
231            let parent_session_id = parent_session_id.to_owned();
232            move || {
233                let live = Controller::load()?
234                    .state
235                    .sessions
236                    .get(&parent_session_id)
237                    .is_some_and(|session| {
238                        matches!(
239                            session.state,
240                            SessionState::Running | SessionState::Disconnected
241                        )
242                    });
243                if !live {
244                    return Ok(Vec::new());
245                }
246                crate::database::load_stopped_subagents(&parent_session_id)
247            }
248        })
249        .await;
250        let stopped = match loaded {
251            Ok(stopped) => stopped,
252            Err(error) => {
253                tracing::warn!(
254                    session_id = %parent_session_id,
255                    error = format!("{error:#}"),
256                    "could not read the sub-agents a failed suspend stopped"
257                );
258                return;
259            }
260        };
261        let Some(context) = mj_core::subagent::stopped_subagents_prompt_context(&stopped) else {
262            return;
263        };
264        let delivered = async {
265            let handle = self
266                .session_manager
267                .session(parent_session_id.to_owned())
268                .await?;
269            handle.install_prompt_context(context).await?;
270            if let Some(text) = mj_core::subagent::stopped_subagents_notice(&stopped) {
271                handle
272                    .submit(
273                        new_command_id("stopped-subagents")?,
274                        RelayCommand::RecordNotice { text },
275                    )
276                    .await?;
277            }
278            anyhow::Ok(())
279        }
280        .await;
281        if let Err(error) = delivered {
282            tracing::warn!(
283                session_id = %parent_session_id,
284                error = format!("{error:#}"),
285                "could not tell a live parent which sub-agents a failed suspend stopped; its next resume will"
286            );
287            return;
288        }
289        let delivered = stopped
290            .into_iter()
291            .map(|child| child.child_session_id)
292            .collect::<Vec<_>>();
293        // The relay owns the note now. Failing to forget the list only means
294        // a later resume tells the model again.
295        if let Err(error) = blocking({
296            let parent_session_id = parent_session_id.to_owned();
297            move || crate::database::clear_stopped_subagents(&parent_session_id, &delivered)
298        })
299        .await
300        {
301            tracing::warn!(
302                session_id = %parent_session_id,
303                error = format!("{error:#}"),
304                "could not clear the stopped sub-agents after telling the live parent"
305            );
306        }
307    }
308
309    /// Forget a sub-agent whose stop failed: its record, its relation to the
310    /// parent, its conversation and its attachments. Its worker lives on the
311    /// parent's target, which the parent's suspend releases next.
312    async fn remove_subagent_records(self: &Arc<Self>, child_id: &str) -> Result<()> {
313        self.preempt_active_lifecycle(child_id).await?;
314        blocking({
315            let child_id = child_id.to_owned();
316            move || {
317                if let Err(error) = mj_core::attachment::AttachmentStore::controller(&child_id)
318                    .and_then(|store| store.remove_session_data())
319                {
320                    tracing::warn!(
321                        %child_id,
322                        error = format!("{error:#}"),
323                        "could not remove a stopped sub-agent's attachments"
324                    );
325                }
326                crate::database::delete_session(&child_id)
327            }
328        })
329        .await?;
330        self.reload_controller().await?;
331        self.publish_revision();
332        Ok(())
333    }
334
335    pub(super) async fn wait_before_close(self: &Arc<Self>, session_id: &str) -> Result<()> {
336        // Cancellation is a request: the old owner must actually finish before
337        // close acquires the target, including an irreversible create commit.
338        let pending = {
339            let operations = self
340                .lifecycle
341                .lock()
342                .unwrap_or_else(PoisonError::into_inner);
343            operations
344                .get(session_id)
345                .filter(|operation| {
346                    !matches!(
347                        operation.kind,
348                        LifecycleKind::Suspend | LifecycleKind::Cleanup
349                    )
350                })
351                .map(|operation| {
352                    operation.request_cancel();
353                    operation.result.clone()
354                })
355        };
356        if let Some(pending) = pending {
357            if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
358                tracing::debug!(%session_id, %error, "previous lifecycle ended before close");
359            }
360            self.remove_completed_lifecycle(&pending);
361        }
362
363        Ok(())
364    }
365
366    async fn close_requested_session_with_ack(
367        self: &Arc<Self>,
368        session_id: String,
369        acknowledge_unpublished_work: bool,
370    ) -> Result<()> {
371        self.wait_before_close(&session_id).await?;
372        let route = blocking({
373            let session_id = session_id.clone();
374            move || {
375                let controller = Controller::load()?;
376                Ok(close_route(controller.state.sessions.get(&session_id)))
377            }
378        })
379        .await?;
380        match route {
381            CloseRoute::Done | CloseRoute::DeferredCleanup => {
382                // The parent is closed already, so nothing can fail after
383                // its sub-agents stop.
384                self.stop_subagents_for_suspend(&session_id).await?;
385                if route == CloseRoute::DeferredCleanup {
386                    self.start_deferred_cleanup(session_id)?;
387                }
388                return Ok(());
389            }
390            CloseRoute::Graceful
391            | CloseRoute::RecoverInterrupted
392            | CloseRoute::SettleWithoutCheckpoint => {}
393        }
394        let operation_session_id = session_id.clone();
395        let result = self
396            .run_lifecycle(
397                operation_session_id,
398                LifecycleKind::Suspend,
399                move |state, session_id, cancelled| async move {
400                    let _recovery_reservation = tokio::task::spawn_blocking({
401                        let observer = state.recovery_observer.clone();
402                        let session_id = session_id.clone();
403                        let cancelled = cancelled.clone();
404                        move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
405                    })
406                    .await
407                    .context("reserve recovery for daemon close task")??;
408                    let mut controller = tokio::task::spawn_blocking(Controller::load)
409                        .await
410                        .context("load controller for daemon close task")??;
411                    let executor = DaemonStageReportingExecutor::new(
412                        CancellableProcessExecutor::new(cancelled),
413                        state.clone(),
414                        session_id.clone(),
415                    );
416                    let deferred = match route {
417                        // `prepare_suspension` marks a live session `Closing`,
418                        // so this is also the route of every live suspend.
419                        CloseRoute::RecoverInterrupted => {
420                            controller
421                                .recover_interrupted_close_managed(
422                                    &session_id,
423                                    &executor,
424                                    &state.session_manager,
425                                    acknowledge_unpublished_work,
426                                    Some(state.stop_subagents_before_close(&session_id)),
427                                )
428                                .await?
429                        }
430                        // Nothing to archive and no relay to latch, so this
431                        // close only tears down and settles. The route was
432                        // decided after `wait_before_close` let any live
433                        // create or resume finish, so a session that is still
434                        // genuinely provisioning is not caught here.
435                        CloseRoute::SettleWithoutCheckpoint => {
436                            // No checkpoint can fail after the sub-agents stop.
437                            state.stop_subagents_for_suspend(&session_id).await?;
438                            controller.suspend_session_without_checkpoint(&session_id, &executor)?
439                        }
440                        _ => {
441                            controller
442                                .suspend_session_managed_controlled(
443                                    &session_id,
444                                    &executor,
445                                    &state.session_manager,
446                                    acknowledge_unpublished_work,
447                                    Some(state.stop_subagents_before_close(&session_id)),
448                                )
449                                .await?
450                        }
451                    };
452                    Ok(if deferred {
453                        DaemonLifecycleResult::DeferredCleanup
454                    } else {
455                        DaemonLifecycleResult::Done
456                    })
457                },
458            )
459            .await?;
460        let _ = result; // Deferred cleanup is handed off by the daemon-owned supervisor.
461        Ok(())
462    }
463
464    pub(super) fn start_deferred_cleanup(
465        self: &Arc<Self>,
466        session_id: String,
467    ) -> Result<LifecycleWatch> {
468        let result = self.start_or_join_lifecycle(
469            session_id.clone(),
470            LifecycleKind::Cleanup,
471            |state, session_id, cancelled| async move {
472                blocking(move || {
473                    let mut controller = Controller::load()?;
474                    let executor = DaemonStageReportingExecutor::new(
475                        CancellableProcessExecutor::new(cancelled),
476                        state,
477                        session_id.clone(),
478                    );
479                    controller.cleanup_stopped_target(&session_id, &executor)?;
480                    Ok(DaemonLifecycleResult::Done)
481                })
482                .await
483            },
484        )?;
485        let caller_result = result.clone();
486        let channel = result.clone();
487        let state = Arc::clone(self);
488        tokio::spawn(async move {
489            if let Err(error) = Self::wait_lifecycle_result(result).await {
490                tracing::warn!(%session_id, error = format!("{error:#}"), "deferred Podman cleanup failed");
491                state.push_notice(
492                    &session_id,
493                    "Container storage cleanup failed; the stopped session retains its target for retry.",
494                );
495            }
496            state.remove_completed_lifecycle(&channel);
497        });
498        Ok(caller_result)
499    }
500
501    pub(super) fn resume_retained_cleanups(self: &Arc<Self>) {
502        let session_ids = self
503            .controller
504            .lock()
505            .unwrap_or_else(PoisonError::into_inner)
506            .state
507            .sessions
508            .iter()
509            .filter(|(_, session)| {
510                session.state == SessionState::Stopped && session.target.is_some()
511            })
512            .map(|(session_id, _)| session_id.clone())
513            .collect::<Vec<_>>();
514        for session_id in session_ids {
515            if crate::controller::move_session::move_owns_session(&session_id) {
516                continue;
517            }
518            if let Err(error) = self.start_deferred_cleanup(session_id.clone()) {
519                tracing::warn!(%session_id, error = format!("{error:#}"), "could not resume deferred Podman cleanup");
520                self.push_notice(
521                    &session_id,
522                    format!("Could not resume container storage cleanup: {error:#}"),
523                );
524            }
525        }
526    }
527
528    pub(super) async fn wait_for_deferred_cleanup(
529        self: &Arc<Self>,
530        session_id: &str,
531    ) -> Result<()> {
532        let existing = {
533            let lifecycle = self
534                .lifecycle
535                .lock()
536                .unwrap_or_else(PoisonError::into_inner);
537            lifecycle.get(session_id).and_then(|active| {
538                (active.kind == LifecycleKind::Cleanup).then(|| active.result.clone())
539            })
540        };
541        let result = match existing {
542            Some(result) => result,
543            None => {
544                let needs_cleanup = blocking({
545                    let session_id = session_id.to_owned();
546                    move || {
547                        let controller = Controller::load()?;
548                        Ok(controller
549                            .state
550                            .sessions
551                            .get(&session_id)
552                            .is_some_and(|session| {
553                                session.state == SessionState::Stopped && session.target.is_some()
554                            }))
555                    }
556                })
557                .await?;
558                if !needs_cleanup {
559                    return Ok(());
560                }
561                self.start_deferred_cleanup(session_id.to_owned())?
562            }
563        };
564        let channel = result.clone();
565        let outcome = Self::wait_lifecycle_result(result).await;
566        self.remove_completed_lifecycle(&channel);
567        match outcome? {
568            DaemonLifecycleResult::Done => Ok(()),
569            DaemonLifecycleResult::Move(_) => unreachable!("cleanup cannot return a move outcome"),
570            DaemonLifecycleResult::DeferredCleanup => {
571                unreachable!("cleanup cannot schedule another cleanup")
572            }
573        }
574    }
575}