Skip to main content

mj_controller/daemon/
resume.rs

1use super::*;
2
3impl RuntimeState {
4    /// Resume a session, and return nothing.
5    ///
6    /// This used to answer with the whole `MaterializedSession`. That reply
7    /// travels as one JSON frame against `MAX_FRAME_BYTES`, so a session whose
8    /// projection outgrew 8 MiB could not be resumed at all — it built a
9    /// several-hundred-megabyte buffer and then refused to send it. The
10    /// projection is already durable; a viewer reads it from the store.
11    pub async fn resume_session(self: &Arc<Self>, request: ResumeSessionRequest) -> Result<()> {
12        let _admission = self
13            .workspace_resume_gate(&request.workspace_id)
14            .read_owned()
15            .await;
16        let session_id = request.session_id.clone();
17        self.wait_for_deferred_cleanup(&session_id).await?;
18        // Whether it is already running is a boolean. Answering it used to
19        // load the entire projection so it could be handed back as the reply.
20        let already_running = blocking({
21            let session_id = session_id.clone();
22            move || {
23                let controller = Controller::load()?;
24                Ok(controller
25                    .state
26                    .sessions
27                    .get(&session_id)
28                    .is_some_and(|session| session.state == SessionState::Running))
29            }
30        })
31        .await?;
32        if already_running {
33            return Ok(());
34        }
35        let profile_id = request.profile_id.clone();
36        let target_template_id = request.target_template_id.clone();
37        let workspace_id = request.workspace_id.clone();
38        let rebind_workspace_id = workspace_id.clone();
39        let operation_session_id = session_id.clone();
40        let result = self.start_or_join_lifecycle_for_workspace(
41            session_id,
42            LifecycleKind::Resume,
43            Some(workspace_id),
44            move |state, session_id, cancelled| async move {
45                let _recovery_reservation = tokio::task::spawn_blocking({
46                    let observer = state.recovery_observer.clone();
47                    let session_id = session_id.clone();
48                    let cancelled = cancelled.clone();
49                    move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
50                })
51                .await
52                .context("reserve recovery for daemon resume task")??;
53                blocking({
54                    let session_id = session_id.clone();
55                    move || {
56                        crate::database::reassign_resumable_session_workspace(
57                            &session_id,
58                            &rebind_workspace_id,
59                        )
60                    }
61                })
62                .await?;
63                let restore_request = request.clone();
64                let mut controller = tokio::task::spawn_blocking(move || {
65                    session_move::load_controller_for_resume(&restore_request)
66                })
67                .await
68                .context("load controller for daemon resume task")??;
69                let executor = DaemonStageReportingExecutor::new(
70                    CancellableProcessExecutor::new(cancelled),
71                    state.clone(),
72                    session_id.clone(),
73                );
74                let materialized = controller
75                    .resume_session_controlled_with_repository_preflight(
76                        &session_id,
77                        &request.profile_id,
78                        &request.target_template_id,
79                        SessionResumeOptions {
80                            additional_mounts: request.additional_mounts,
81                            resource_allocation: request.resource_allocation,
82                            discard_queue: request.discard_queue,
83                        },
84                        request.repository_preflight,
85                        &executor,
86                    )
87                    .await?;
88                // The projection stays where it was written. A viewer reads
89                // it from the store; shipping it back through the daemon
90                // reply put a whole transcript in one IPC frame.
91                let _ = materialized;
92                Ok(DaemonLifecycleResult::Done)
93            },
94        )?;
95        self.set_lifecycle_resume_destination(
96            &operation_session_id,
97            profile_id,
98            target_template_id,
99        );
100        let channel = result.clone();
101        let result = Self::wait_lifecycle_result(result).await;
102        self.remove_completed_lifecycle(&channel);
103        match result? {
104            DaemonLifecycleResult::Done => {}
105            DaemonLifecycleResult::Move(_) | DaemonLifecycleResult::Park(_) => {
106                unreachable!("resume cannot return a move outcome")
107            }
108            DaemonLifecycleResult::DeferredCleanup => {
109                unreachable!("session resume cannot schedule target cleanup")
110            }
111        }
112        blocking(move || {
113            if let Some(mut operation) =
114                crate::database::load_move_operation(&operation_session_id)?
115                && !operation.queue_admission_started
116            {
117                operation.phase = mj_core::state::MovePhase::Cancelled;
118                operation.queue_admission_finished = true;
119                operation.updated_at = chrono::Utc::now().to_rfc3339();
120                operation.error = Some("Recovered through an explicit Resume operation".into());
121                crate::database::save_move_operation(&operation)?;
122            }
123            Ok(())
124        })
125        .await?;
126        Ok(())
127    }
128
129    pub(super) async fn discard_since_checkpoint(
130        self: &Arc<Self>,
131        session_id: String,
132        checkpoint: mj_core::state::CheckpointMetadata,
133    ) -> Result<()> {
134        let operation_session_id = session_id.clone();
135        let result = self
136            .run_lifecycle(
137                operation_session_id,
138                LifecycleKind::ForceStop,
139                move |state, session_id, cancelled| async move {
140                    // Refuse before anything stops when the copy changed.
141                    blocking({
142                        let session_id = session_id.clone();
143                        let checkpoint = checkpoint.clone();
144                        move || {
145                            ensure!(
146                                Controller::load()?
147                                    .state
148                                    .sessions
149                                    .get(&session_id)
150                                    .and_then(|s| s.checkpoint.as_ref())
151                                    == Some(&checkpoint),
152                                "the recovery copy changed; review it before discarding changes"
153                            );
154                            Ok(())
155                        }
156                    })
157                    .await?;
158                    // The parent is about to go back to an older recovery
159                    // copy, and its sub-agents stop exactly as they do when
160                    // it is suspended.
161                    state.stop_subagents_for_suspend(&session_id).await?;
162                    blocking(move || {
163                        let mut controller = Controller::load()?;
164                        let executor = DaemonStageReportingExecutor::new(
165                            CancellableProcessExecutor::new(cancelled),
166                            state,
167                            session_id.clone(),
168                        );
169                        ensure!(
170                            controller
171                                .state
172                                .sessions
173                                .get(&session_id)
174                                .and_then(|s| s.checkpoint.as_ref())
175                                == Some(&checkpoint),
176                            "the recovery copy changed; review it before discarding changes"
177                        );
178                        let deferred = controller.force_stop(&session_id, &executor)?;
179                        Ok(if deferred {
180                            DaemonLifecycleResult::DeferredCleanup
181                        } else {
182                            DaemonLifecycleResult::Done
183                        })
184                    })
185                    .await
186                },
187            )
188            .await;
189        if result.is_err() {
190            self.tell_live_parent_about_stopped_subagents(&session_id)
191                .await;
192        }
193        let _ = result?; // The lifecycle supervisor owns the cleanup handoff.
194        Ok(())
195    }
196
197    pub(super) async fn destroy_stopped_session(
198        self: &Arc<Self>,
199        session_id: String,
200        branch: BranchDisposition,
201    ) -> Result<()> {
202        self.tear_down_stopped_session(
203            session_id,
204            LifecycleKind::DestroyStopped,
205            branch,
206            CheckoutDisposition::Remove,
207        )
208        .await
209        .map(|_| ())
210    }
211
212    /// Discard the record of a session that ended as `Lost`: its managed
213    /// target is gone, it has no checkpoint to resume from, and the only
214    /// action it still offers is Destroy. Keeping it would leave a tombstone
215    /// the person has to find and delete by hand.
216    ///
217    /// The branch stays, and so does a checkout that holds uncommitted
218    /// changes: nothing here was confirmed by a person, so nothing here may
219    /// discard work. A teardown that fails leaves the row where it is, with
220    /// its cause, and says so.
221    pub(crate) async fn discard_lost_session(self: &Arc<Self>, session_id: String) {
222        let short = mj_core::state::short_id(&session_id).to_owned();
223        let owned_clone = blocking({
224            let session_id = session_id.clone();
225            move || {
226                let controller = Controller::load()?;
227                let checkout = match controller.state.checkout(&session_id)? {
228                    mj_core::state::Checkout::ManagedWorktree { worktree, .. }
229                        if worktree.kind == mj_core::state::ManagedCheckoutKind::Clone =>
230                    {
231                        worktree
232                    }
233                    _ => return Ok(None),
234                };
235                if !controller.state.sessions.contains_key(&session_id) {
236                    return Ok(None);
237                }
238                if crate::controller::path_exists_on_managed_target(
239                    &crate::targets::ProcessExecutor,
240                    &checkout.target,
241                    &checkout.worktree_root,
242                )? {
243                    Ok(Some(checkout.worktree_root.clone()))
244                } else {
245                    Ok(None)
246                }
247            }
248        })
249        .await;
250        match owned_clone {
251            Ok(Some(path)) => {
252                self.push_notice(&session_id, format!(
253                    "Session {short} lost its target, but its owned clone remains at {}; keeping the record for inspection",
254                    path.display(),
255                ));
256                return;
257            }
258            Err(error) => {
259                tracing::warn!(%session_id, error = %format!("{error:#}"), "could not inspect owned clone before lost-session cleanup");
260                return;
261            }
262            Ok(None) => {}
263        }
264        match self
265            .tear_down_stopped_session(
266                session_id.clone(),
267                LifecycleKind::DestroyStopped,
268                BranchDisposition::Keep,
269                CheckoutDisposition::KeepWhenDirty,
270            )
271            .await
272        {
273            Ok(retained) => {
274                let mut text = format!(
275                    "Session {short} was lost because its managed target no longer exists; its record was removed."
276                );
277                if let Some(path) = retained {
278                    text.push_str(&format!(
279                        " Its checkout has uncommitted changes, so it was kept at {}.",
280                        path.display()
281                    ));
282                }
283                self.push_notice(&session_id, text);
284            }
285            Err(error) => {
286                tracing::warn!(
287                    %session_id,
288                    error = format!("{error:#}"),
289                    "could not discard the record of a lost session"
290                );
291                self.push_notice(
292                    &session_id,
293                    format!(
294                        "Session {short} was lost, but its record could not be removed: {error:#}"
295                    ),
296                );
297            }
298        }
299    }
300
301    /// Archive every stopped session older than `older_than_days` whose
302    /// conversation SessionWiki holds. Answers with how many were archived.
303    ///
304    /// The whole job belongs in a background task: it runs a full index sync,
305    /// which walks every tool's store, and then one lifecycle per session.
306    pub(crate) async fn archive_aged_sessions(
307        self: &Arc<Self>,
308        older_than_days: u32,
309    ) -> Result<usize> {
310        self.wiki()
311            .sync_now(true)
312            .await
313            .context("sync SessionWiki before archiving stopped sessions")?;
314        // The sync is restartable and holds no admission; the archive changes
315        // sessions, so a handoff waits for it from here. It does not start
316        // while a handoff is waiting: the next daemon's tick runs it.
317        let Ok(_work) = crate::upgrade::activity_unless_draining("SessionWiki archive") else {
318            tracing::debug!(
319                "a daemon upgrade is waiting; the SessionWiki archive pass waits for the next daemon"
320            );
321            return Ok(0);
322        };
323        // Recheck a small bounded batch of older clone checkpoints. A remote
324        // that was offline during suspension may now prove its saved refs.
325        let refreshable = blocking(move || {
326            let controller = Controller::load()?;
327            let cutoff = chrono::Utc::now() - chrono::Duration::days(i64::from(older_than_days));
328            let mut sessions = controller
329                .state
330                .sessions
331                .values()
332                .filter(|session| session.state == SessionState::Stopped)
333                .filter(|session| {
334                    chrono::DateTime::parse_from_rfc3339(&session.updated_at)
335                        .is_ok_and(|time| time.with_timezone(&chrono::Utc) <= cutoff)
336                })
337                .filter(|session| {
338                    session.publication.as_ref().is_none_or(|evidence| {
339                        evidence.state != mj_core::state::PublicationState::Published
340                            && !evidence.dirty
341                            && !evidence.stashed
342                    })
343                })
344                .filter_map(|session| {
345                    let mj_core::state::Checkout::ManagedWorktree { worktree, .. } =
346                        controller.state.checkout(&session.id).ok()?
347                    else {
348                        return None;
349                    };
350                    (worktree.kind == mj_core::state::ManagedCheckoutKind::Clone)
351                        .then(|| (session.clone(), worktree.clone()))
352                })
353                .collect::<Vec<_>>();
354            sessions.sort_by(|(a, _), (b, _)| {
355                a.publication
356                    .as_ref()
357                    .map(|e| &e.checked_at)
358                    .cmp(&b.publication.as_ref().map(|e| &e.checked_at))
359            });
360            sessions.truncate(4);
361            Ok(sessions)
362        })
363        .await?;
364        let mut checks = tokio::task::JoinSet::new();
365        for (session, checkout) in refreshable {
366            checks.spawn_blocking(move || {
367                let assessment =
368                    crate::controller::publication::refresh_stopped_clone_publication_for_worktree(
369                        &session, &checkout,
370                    );
371                (session.id, assessment)
372            });
373        }
374        let mut refreshed = false;
375        while let Some(done) = checks.join_next().await {
376            let (id, assessment) = done.context("publication refresh task failed")?;
377            if let Some(assessment) = assessment {
378                refreshed |= blocking(move || {
379                    crate::database::set_publication_assessment_if_current(&id, &assessment)
380                })
381                .await?;
382            }
383        }
384        if refreshed {
385            self.reload_controller().await?;
386            self.publish_revision();
387        }
388        let (candidates, children) = blocking(move || {
389            let controller = Controller::load()?;
390            let candidates = crate::sessionwiki::sessions_ready_to_archive_from_state(
391                &controller.state,
392                chrono::Utc::now(),
393                older_than_days,
394            );
395            let children: std::collections::BTreeSet<String> = candidates
396                .iter()
397                .filter(|id| controller.state.is_subagent_session(id))
398                .cloned()
399                .collect();
400            Ok((candidates, children))
401        })
402        .await
403        .context("select the stopped sessions old enough to archive")?;
404        if candidates.is_empty() {
405            return Ok(0);
406        }
407        let indexed = blocking({
408            let candidates = candidates
409                .iter()
410                .filter(|id| !children.contains(*id))
411                .cloned()
412                .collect::<Vec<_>>();
413            move || crate::sessionwiki::indexed_with_messages(&candidates)
414        })
415        .await
416        .context("check the SessionWiki index before archiving")?;
417        let mut archived = 0;
418        for session_id in candidates {
419            if !children.contains(&session_id) && !indexed.contains(&session_id) {
420                tracing::warn!(
421                    %session_id,
422                    "SessionWiki holds no conversation for this stopped session; keeping it"
423                );
424                continue;
425            }
426            match self.archive_stopped_session(session_id.clone()).await {
427                Ok(()) => {
428                    archived += 1;
429                    tracing::info!(
430                        %session_id,
431                        older_than_days,
432                        child = children.contains(&session_id),
433                        "archived a stopped session; top-level conversations are retained in SessionWiki"
434                    );
435                }
436                Err(error) => tracing::warn!(
437                    %session_id,
438                    error = %format!("{error:#}"),
439                    "could not archive a stopped session"
440                ),
441            }
442        }
443        if archived > 0 {
444            // Flip the rows the job just emptied to archived.
445            self.wiki().request_sync(false);
446        }
447        Ok(archived)
448    }
449
450    /// Destroy a stopped session the way the archive job wants: the record,
451    /// the checkpoint, and the attachments go, and the session's git branch
452    /// goes only when another branch already contains all of its commits.
453    /// A top-level conversation stays searchable and restorable through
454    /// SessionWiki; child conversations are excluded.
455    async fn archive_stopped_session(self: &Arc<Self>, session_id: String) -> Result<()> {
456        self.tear_down_stopped_session(
457            session_id,
458            LifecycleKind::ArchiveStopped,
459            BranchDisposition::DeleteIfMerged,
460            CheckoutDisposition::Remove,
461        )
462        .await
463        .map(|_| ())
464    }
465
466    /// Answers with the managed checkout the teardown kept, if it kept one.
467    async fn tear_down_stopped_session(
468        self: &Arc<Self>,
469        session_id: String,
470        kind: LifecycleKind,
471        branch: BranchDisposition,
472        checkout: CheckoutDisposition,
473    ) -> Result<Option<PathBuf>> {
474        // Preserve the root conversation before tearing down its children.
475        // Children are intentionally excluded from SessionWiki.
476        self.index_before_destroy(&session_id).await;
477        let children = blocking({
478            let session_id = session_id.clone();
479            move || {
480                Ok(crate::database::list_subagents(&session_id)?
481                    .into_iter()
482                    .map(|child| child.child_session_id)
483                    .collect::<Vec<_>>())
484            }
485        })
486        .await?;
487        for child_id in children {
488            // A sub-agent borrows its parent's worker and never owns a managed
489            // worktree, so it has no branch of its own to keep.
490            Box::pin(self.force_destroy_indexed_session(
491                child_id.clone(),
492                BranchDisposition::Keep,
493                LifecycleKind::ForceDestroy,
494            ))
495            .await
496            .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
497        }
498        self.wait_for_deferred_cleanup(&session_id).await?;
499        let exists = blocking({
500            let session_id = session_id.clone();
501            move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
502        })
503        .await?;
504        if !exists {
505            return Ok(None);
506        }
507        // The retained path is produced inside the lifecycle task, which can
508        // only answer with a `DaemonLifecycleResult`, so it comes back here.
509        let retained = Arc::new(Mutex::new(None));
510        self.run_lifecycle(session_id, kind, {
511            let retained = retained.clone();
512            move |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                    let kept = controller.destroy_session_controlled_with_checkout(
521                        &session_id,
522                        &executor,
523                        branch,
524                        checkout,
525                    )?;
526                    *retained.lock().unwrap_or_else(PoisonError::into_inner) = kept;
527                    Ok(DaemonLifecycleResult::Done)
528                })
529                .await
530            }
531        })
532        .await?;
533        let kept = retained
534            .lock()
535            .unwrap_or_else(PoisonError::into_inner)
536            .clone();
537        Ok(kept)
538    }
539
540    /// Preserve a top-level session before destroy so its conversation remains
541    /// searchable afterwards. Child sessions are excluded from the index.
542    ///
543    /// Waits at most [`crate::sessionwiki::DESTROY_SYNC_WAIT`] for a sync
544    /// pass, then indexes the sessions on their own. A destroy is never
545    /// refused for this: when the index cannot take the sessions, the log
546    /// says why and the destroy goes ahead.
547    pub(super) async fn index_before_destroy(self: &Arc<Self>, session_id: &str) {
548        use crate::sessionwiki::IndexedBeforeDestroy;
549        let outcome = self
550            .wiki()
551            .index_before_destroy(session_id, crate::sessionwiki::DESTROY_SYNC_WAIT)
552            .await;
553        match outcome {
554            IndexedBeforeDestroy::Unavailable(reason) => tracing::info!(
555                %session_id,
556                reason,
557                "destroying a session without indexing it in SessionWiki"
558            ),
559            IndexedBeforeDestroy::Failed(reason) => tracing::warn!(
560                %session_id,
561                %reason,
562                "could not index a session in SessionWiki before destroying it"
563            ),
564            outcome => tracing::debug!(
565                %session_id,
566                ?outcome,
567                "indexed a session in SessionWiki before destroying it"
568            ),
569        }
570    }
571
572    /// Cancel any in-flight lifecycle for `session_id` and wait for it to
573    /// finish.
574    ///
575    /// Force destruction is the escape hatch for a wedged operation, so it
576    /// takes over rather than queueing behind one — but only after the running
577    /// task has stopped, because a cancelled create or close re-persists its
578    /// record as it unwinds and would otherwise resurrect the row this
579    /// operation deletes. A lifecycle that ignores cancellation for longer
580    /// than [`FORCE_DESTROY_PREEMPT_TIMEOUT`] is reported instead of destroyed
581    /// under.
582    pub(super) async fn preempt_active_lifecycle(self: &Arc<Self>, session_id: &str) -> Result<()> {
583        let (mut result, tearing_down, cancelled_restart) = {
584            let mut lifecycle_owner = self.owner();
585            let lifecycle = &mut lifecycle_owner.lifecycle;
586            let Some(active) = lifecycle.get_mut(session_id) else {
587                return Ok(());
588            };
589            if !active.is_running() {
590                return Ok(());
591            }
592            let tearing_down = active.kind.is_teardown();
593            // A teardown already under way is the fact this destroy wants, so
594            // it is waited for, not cancelled half way: cancelling it is what
595            // left a sub-agent's destroy "cancelled before its parent" when
596            // the parent's destroy and the person's own destroy of the child
597            // overlapped. Only one that outlives the wait is cancelled below.
598            if !tearing_down {
599                active.request_cancel();
600            }
601            (
602                active.result.clone(),
603                tearing_down,
604                active.kind == LifecycleKind::Restart && !tearing_down,
605            )
606        };
607        if cancelled_restart {
608            blocking({
609                let session_id = session_id.to_owned();
610                move || crate::database::cancel_session_restart(&session_id)
611            })
612            .await?;
613        }
614        let mut finished = wait_for_lifecycle_result(&mut result).await;
615        if finished.is_err() && tearing_down {
616            {
617                let mut lifecycle_owner = self.owner();
618                if let Some(active) = lifecycle_owner.lifecycle.get_mut(session_id) {
619                    active.request_cancel();
620                }
621            }
622            finished = wait_for_lifecycle_result(&mut result).await;
623        }
624        match finished {
625            // The wait only returns once the watch holds a result or its
626            // sender died; distinguish those two, and the timeout separately.
627            Ok(Ok(())) => Ok(()),
628            Ok(Err(())) => bail!(
629                "daemon lifecycle operation stopped without a result for session {session_id}"
630            ),
631            Err(_) => bail!(
632                "session {session_id} still has an operation that did not stop after cancellation; try again"
633            ),
634        }
635    }
636
637    /// Permanently destroy a session from any state, cancelling whatever
638    /// lifecycle operation holds it first. Data loss is the caller's confirmed
639    /// decision; see [`Controller::force_destroy_session`]. The session's git
640    /// branch survives unless `branch` says to delete it.
641    pub async fn force_destroy_session(
642        self: &Arc<Self>,
643        session_id: String,
644        branch: BranchDisposition,
645    ) -> Result<()> {
646        // The destroy holds this from the acknowledgement to its last step, so
647        // a graceful stop or handoff waits for it, sub-agents first (#1191).
648        let _upgrade_work = crate::upgrade::destroy_activity(&session_id)?;
649        self.index_before_destroy(&session_id).await;
650        self.force_destroy_indexed_session(session_id, branch, LifecycleKind::ForceDestroy)
651            .await
652    }
653
654    /// [`Self::force_destroy_session`] once the session and its sub-agents
655    /// have been indexed. `kind` is `ForceDestroy`, or `StopSubagent` for a
656    /// sub-agent its parent's suspend stops; its own sub-agents go the same
657    /// way.
658    pub(super) async fn force_destroy_indexed_session(
659        self: &Arc<Self>,
660        session_id: String,
661        branch: BranchDisposition,
662        kind: LifecycleKind,
663    ) -> Result<()> {
664        let children = blocking({
665            let session_id = session_id.clone();
666            move || {
667                Ok(crate::database::list_subagents(&session_id)?
668                    .into_iter()
669                    .map(|child| child.child_session_id)
670                    .collect::<Vec<_>>())
671            }
672        })
673        .await?;
674        for child_id in children {
675            // Sub-agents borrow their parent's worker and own no branch.
676            Box::pin(self.force_destroy_indexed_session(
677                child_id.clone(),
678                BranchDisposition::Keep,
679                kind,
680            ))
681            .await
682            .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
683        }
684        self.preempt_active_lifecycle(&session_id).await?;
685        let exists = blocking({
686            let session_id = session_id.clone();
687            move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
688        })
689        .await?;
690        if !exists {
691            return Ok(());
692        }
693        self.run_lifecycle(
694            session_id,
695            kind,
696            move |state, session_id, cancelled| async move {
697                let _recovery_reservation = tokio::task::spawn_blocking({
698                    let observer = state.recovery_observer.clone();
699                    let session_id = session_id.clone();
700                    let cancelled = cancelled.clone();
701                    move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
702                })
703                .await
704                .context("reserve recovery for daemon force-destroy task")??;
705                blocking({
706                    let session_id = session_id.clone();
707                    move || {
708                        let mut controller = Controller::load()?;
709                        let executor = DaemonStageReportingExecutor::new(
710                            CancellableProcessExecutor::new(cancelled),
711                            state,
712                            session_id.clone(),
713                        );
714                        controller.force_destroy_session(&session_id, &executor, branch)?;
715                        crate::controller::move_session::release_move_queue_hold(&session_id);
716                        Ok(DaemonLifecycleResult::Done)
717                    }
718                })
719                .await
720            },
721        )
722        .await?;
723        Ok(())
724    }
725}
726
727/// Wait up to [`FORCE_DESTROY_PREEMPT_TIMEOUT`] for a lifecycle to publish its
728/// result. `Ok(Err(()))` means the task ended without one.
729async fn wait_for_lifecycle_result(
730    result: &mut LifecycleWatch,
731) -> std::result::Result<std::result::Result<(), ()>, tokio::time::error::Elapsed> {
732    tokio::time::timeout(FORCE_DESTROY_PREEMPT_TIMEOUT, async {
733        loop {
734            if result.borrow().is_some() {
735                return Ok(());
736            }
737            if result.changed().await.is_err() {
738                return Err(());
739            }
740        }
741    })
742    .await
743}