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 Some(checkout) = controller
228                    .state
229                    .sessions
230                    .get(&session_id)
231                    .and_then(|session| session.managed_worktree.as_ref())
232                    .filter(|checkout| checkout.kind == mj_core::state::ManagedCheckoutKind::Clone)
233                else {
234                    return Ok(None);
235                };
236                if crate::controller::path_exists_on_managed_target(
237                    &crate::targets::ProcessExecutor,
238                    &checkout.target,
239                    &checkout.worktree_root,
240                )? {
241                    Ok(Some(checkout.worktree_root.clone()))
242                } else {
243                    Ok(None)
244                }
245            }
246        })
247        .await;
248        match owned_clone {
249            Ok(Some(path)) => {
250                self.push_notice(&session_id, format!(
251                    "Session {short} lost its target, but its owned clone remains at {}; keeping the record for inspection",
252                    path.display(),
253                ));
254                return;
255            }
256            Err(error) => {
257                tracing::warn!(%session_id, error = %format!("{error:#}"), "could not inspect owned clone before lost-session cleanup");
258                return;
259            }
260            Ok(None) => {}
261        }
262        match self
263            .tear_down_stopped_session(
264                session_id.clone(),
265                LifecycleKind::DestroyStopped,
266                BranchDisposition::Keep,
267                CheckoutDisposition::KeepWhenDirty,
268            )
269            .await
270        {
271            Ok(retained) => {
272                let mut text = format!(
273                    "Session {short} was lost because its managed target no longer exists; its record was removed."
274                );
275                if let Some(path) = retained {
276                    text.push_str(&format!(
277                        " Its checkout has uncommitted changes, so it was kept at {}.",
278                        path.display()
279                    ));
280                }
281                self.push_notice(&session_id, text);
282            }
283            Err(error) => {
284                tracing::warn!(
285                    %session_id,
286                    error = format!("{error:#}"),
287                    "could not discard the record of a lost session"
288                );
289                self.push_notice(
290                    &session_id,
291                    format!(
292                        "Session {short} was lost, but its record could not be removed: {error:#}"
293                    ),
294                );
295            }
296        }
297    }
298
299    /// Archive every stopped session older than `older_than_days` whose
300    /// conversation SessionWiki holds. Answers with how many were archived.
301    ///
302    /// The whole job belongs in a background task: it runs a full index sync,
303    /// which walks every tool's store, and then one lifecycle per session.
304    pub(crate) async fn archive_aged_sessions(
305        self: &Arc<Self>,
306        older_than_days: u32,
307    ) -> Result<usize> {
308        self.wiki()
309            .sync_now(true)
310            .await
311            .context("sync SessionWiki before archiving stopped sessions")?;
312        // The sync is restartable and holds no admission; the archive changes
313        // sessions, so a handoff waits for it from here. It does not start
314        // while a handoff is waiting: the next daemon's tick runs it.
315        let Ok(_work) = crate::upgrade::activity_unless_draining("SessionWiki archive") else {
316            tracing::debug!(
317                "a daemon upgrade is waiting; the SessionWiki archive pass waits for the next daemon"
318            );
319            return Ok(0);
320        };
321        // Recheck a small bounded batch of older clone checkpoints. A remote
322        // that was offline during suspension may now prove its saved refs.
323        let refreshable = blocking(move || {
324            let controller = Controller::load()?;
325            let cutoff = chrono::Utc::now() - chrono::Duration::days(i64::from(older_than_days));
326            let mut sessions = controller
327                .state
328                .sessions
329                .values()
330                .filter(|session| session.state == SessionState::Stopped)
331                .filter(|session| {
332                    chrono::DateTime::parse_from_rfc3339(&session.updated_at)
333                        .is_ok_and(|time| time.with_timezone(&chrono::Utc) <= cutoff)
334                })
335                .filter(|session| {
336                    session.managed_worktree.as_ref().is_some_and(|owned| {
337                        owned.kind == mj_core::state::ManagedCheckoutKind::Clone
338                    })
339                })
340                .filter(|session| {
341                    session.publication.as_ref().is_none_or(|evidence| {
342                        evidence.state != mj_core::state::PublicationState::Published
343                            && !evidence.dirty
344                            && !evidence.stashed
345                    })
346                })
347                .cloned()
348                .collect::<Vec<_>>();
349            sessions.sort_by(|a, b| {
350                a.publication
351                    .as_ref()
352                    .map(|e| &e.checked_at)
353                    .cmp(&b.publication.as_ref().map(|e| &e.checked_at))
354            });
355            sessions.truncate(4);
356            Ok(sessions)
357        })
358        .await?;
359        let mut checks = tokio::task::JoinSet::new();
360        for session in refreshable {
361            checks.spawn_blocking(move || {
362                let assessment =
363                    crate::controller::publication::refresh_stopped_clone_publication(&session);
364                (session.id, assessment)
365            });
366        }
367        let mut refreshed = false;
368        while let Some(done) = checks.join_next().await {
369            let (id, assessment) = done.context("publication refresh task failed")?;
370            if let Some(assessment) = assessment {
371                refreshed |= blocking(move || {
372                    crate::database::set_publication_assessment_if_current(&id, &assessment)
373                })
374                .await?;
375            }
376        }
377        if refreshed {
378            self.reload_controller().await?;
379            self.publish_revision();
380        }
381        let (candidates, children) = blocking(move || {
382            let controller = Controller::load()?;
383            let candidates = crate::sessionwiki::sessions_ready_to_archive(
384                &controller.state.sessions,
385                &controller.state.subagents,
386                chrono::Utc::now(),
387                older_than_days,
388            );
389            let children: std::collections::BTreeSet<String> = candidates
390                .iter()
391                .filter(|id| controller.state.is_subagent_session(id))
392                .cloned()
393                .collect();
394            Ok((candidates, children))
395        })
396        .await
397        .context("select the stopped sessions old enough to archive")?;
398        if candidates.is_empty() {
399            return Ok(0);
400        }
401        let indexed = blocking({
402            let candidates = candidates
403                .iter()
404                .filter(|id| !children.contains(*id))
405                .cloned()
406                .collect::<Vec<_>>();
407            move || crate::sessionwiki::indexed_with_messages(&candidates)
408        })
409        .await
410        .context("check the SessionWiki index before archiving")?;
411        let mut archived = 0;
412        for session_id in candidates {
413            if !children.contains(&session_id) && !indexed.contains(&session_id) {
414                tracing::warn!(
415                    %session_id,
416                    "SessionWiki holds no conversation for this stopped session; keeping it"
417                );
418                continue;
419            }
420            match self.archive_stopped_session(session_id.clone()).await {
421                Ok(()) => {
422                    archived += 1;
423                    tracing::info!(
424                        %session_id,
425                        older_than_days,
426                        child = children.contains(&session_id),
427                        "archived a stopped session; top-level conversations are retained in SessionWiki"
428                    );
429                }
430                Err(error) => tracing::warn!(
431                    %session_id,
432                    error = %format!("{error:#}"),
433                    "could not archive a stopped session"
434                ),
435            }
436        }
437        if archived > 0 {
438            // Flip the rows the job just emptied to archived.
439            self.wiki().request_sync(false);
440        }
441        Ok(archived)
442    }
443
444    /// Destroy a stopped session the way the archive job wants: the record,
445    /// the checkpoint, and the attachments go, and the session's git branch
446    /// goes only when another branch already contains all of its commits.
447    /// A top-level conversation stays searchable and restorable through
448    /// SessionWiki; child conversations are excluded.
449    async fn archive_stopped_session(self: &Arc<Self>, session_id: String) -> Result<()> {
450        self.tear_down_stopped_session(
451            session_id,
452            LifecycleKind::ArchiveStopped,
453            BranchDisposition::DeleteIfMerged,
454            CheckoutDisposition::Remove,
455        )
456        .await
457        .map(|_| ())
458    }
459
460    /// Answers with the managed checkout the teardown kept, if it kept one.
461    async fn tear_down_stopped_session(
462        self: &Arc<Self>,
463        session_id: String,
464        kind: LifecycleKind,
465        branch: BranchDisposition,
466        checkout: CheckoutDisposition,
467    ) -> Result<Option<PathBuf>> {
468        // Preserve the root conversation before tearing down its children.
469        // Children are intentionally excluded from SessionWiki.
470        self.index_before_destroy(&session_id).await;
471        let children = blocking({
472            let session_id = session_id.clone();
473            move || {
474                Ok(crate::database::list_subagents(&session_id)?
475                    .into_iter()
476                    .map(|child| child.child_session_id)
477                    .collect::<Vec<_>>())
478            }
479        })
480        .await?;
481        for child_id in children {
482            // A sub-agent borrows its parent's worker and never owns a managed
483            // worktree, so it has no branch of its own to keep.
484            Box::pin(self.force_destroy_indexed_session(
485                child_id.clone(),
486                BranchDisposition::Keep,
487                LifecycleKind::ForceDestroy,
488            ))
489            .await
490            .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
491        }
492        self.wait_for_deferred_cleanup(&session_id).await?;
493        let exists = blocking({
494            let session_id = session_id.clone();
495            move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
496        })
497        .await?;
498        if !exists {
499            return Ok(None);
500        }
501        // The retained path is produced inside the lifecycle task, which can
502        // only answer with a `DaemonLifecycleResult`, so it comes back here.
503        let retained = Arc::new(Mutex::new(None));
504        self.run_lifecycle(session_id, kind, {
505            let retained = retained.clone();
506            move |state, session_id, cancelled| async move {
507                blocking(move || {
508                    let mut controller = Controller::load()?;
509                    let executor = DaemonStageReportingExecutor::new(
510                        CancellableProcessExecutor::new(cancelled),
511                        state,
512                        session_id.clone(),
513                    );
514                    let kept = controller.destroy_session_controlled_with_checkout(
515                        &session_id,
516                        &executor,
517                        branch,
518                        checkout,
519                    )?;
520                    *retained.lock().unwrap_or_else(PoisonError::into_inner) = kept;
521                    Ok(DaemonLifecycleResult::Done)
522                })
523                .await
524            }
525        })
526        .await?;
527        let kept = retained
528            .lock()
529            .unwrap_or_else(PoisonError::into_inner)
530            .clone();
531        Ok(kept)
532    }
533
534    /// Preserve a top-level session before destroy so its conversation remains
535    /// searchable afterwards. Child sessions are excluded from the index.
536    ///
537    /// Waits at most [`crate::sessionwiki::DESTROY_SYNC_WAIT`] for a sync
538    /// pass, then indexes the sessions on their own. A destroy is never
539    /// refused for this: when the index cannot take the sessions, the log
540    /// says why and the destroy goes ahead.
541    pub(super) async fn index_before_destroy(self: &Arc<Self>, session_id: &str) {
542        use crate::sessionwiki::IndexedBeforeDestroy;
543        let outcome = self
544            .wiki()
545            .index_before_destroy(session_id, crate::sessionwiki::DESTROY_SYNC_WAIT)
546            .await;
547        match outcome {
548            IndexedBeforeDestroy::Unavailable(reason) => tracing::info!(
549                %session_id,
550                reason,
551                "destroying a session without indexing it in SessionWiki"
552            ),
553            IndexedBeforeDestroy::Failed(reason) => tracing::warn!(
554                %session_id,
555                %reason,
556                "could not index a session in SessionWiki before destroying it"
557            ),
558            outcome => tracing::debug!(
559                %session_id,
560                ?outcome,
561                "indexed a session in SessionWiki before destroying it"
562            ),
563        }
564    }
565
566    /// Cancel any in-flight lifecycle for `session_id` and wait for it to
567    /// finish.
568    ///
569    /// Force destruction is the escape hatch for a wedged operation, so it
570    /// takes over rather than queueing behind one — but only after the running
571    /// task has stopped, because a cancelled create or close re-persists its
572    /// record as it unwinds and would otherwise resurrect the row this
573    /// operation deletes. A lifecycle that ignores cancellation for longer
574    /// than [`FORCE_DESTROY_PREEMPT_TIMEOUT`] is reported instead of destroyed
575    /// under.
576    pub(super) async fn preempt_active_lifecycle(self: &Arc<Self>, session_id: &str) -> Result<()> {
577        let (mut result, tearing_down, cancelled_restart) = {
578            let mut lifecycle_owner = self.owner();
579            let lifecycle = &mut lifecycle_owner.lifecycle;
580            let Some(active) = lifecycle.get_mut(session_id) else {
581                return Ok(());
582            };
583            if !active.is_running() {
584                return Ok(());
585            }
586            let tearing_down = active.kind.is_teardown();
587            // A teardown already under way is the fact this destroy wants, so
588            // it is waited for, not cancelled half way: cancelling it is what
589            // left a sub-agent's destroy "cancelled before its parent" when
590            // the parent's destroy and the person's own destroy of the child
591            // overlapped. Only one that outlives the wait is cancelled below.
592            if !tearing_down {
593                active.request_cancel();
594            }
595            (
596                active.result.clone(),
597                tearing_down,
598                active.kind == LifecycleKind::Restart && !tearing_down,
599            )
600        };
601        if cancelled_restart {
602            blocking({
603                let session_id = session_id.to_owned();
604                move || crate::database::cancel_session_restart(&session_id)
605            })
606            .await?;
607        }
608        let mut finished = wait_for_lifecycle_result(&mut result).await;
609        if finished.is_err() && tearing_down {
610            {
611                let mut lifecycle_owner = self.owner();
612                if let Some(active) = lifecycle_owner.lifecycle.get_mut(session_id) {
613                    active.request_cancel();
614                }
615            }
616            finished = wait_for_lifecycle_result(&mut result).await;
617        }
618        match finished {
619            // The wait only returns once the watch holds a result or its
620            // sender died; distinguish those two, and the timeout separately.
621            Ok(Ok(())) => Ok(()),
622            Ok(Err(())) => bail!(
623                "daemon lifecycle operation stopped without a result for session {session_id}"
624            ),
625            Err(_) => bail!(
626                "session {session_id} still has an operation that did not stop after cancellation; try again"
627            ),
628        }
629    }
630
631    /// Permanently destroy a session from any state, cancelling whatever
632    /// lifecycle operation holds it first. Data loss is the caller's confirmed
633    /// decision; see [`Controller::force_destroy_session`]. The session's git
634    /// branch survives unless `branch` says to delete it.
635    pub async fn force_destroy_session(
636        self: &Arc<Self>,
637        session_id: String,
638        branch: BranchDisposition,
639    ) -> Result<()> {
640        // The destroy holds this from the acknowledgement to its last step, so
641        // a graceful stop or handoff waits for it, sub-agents first (#1191).
642        let _upgrade_work = crate::upgrade::destroy_activity(&session_id)?;
643        self.index_before_destroy(&session_id).await;
644        self.force_destroy_indexed_session(session_id, branch, LifecycleKind::ForceDestroy)
645            .await
646    }
647
648    /// [`Self::force_destroy_session`] once the session and its sub-agents
649    /// have been indexed. `kind` is `ForceDestroy`, or `StopSubagent` for a
650    /// sub-agent its parent's suspend stops; its own sub-agents go the same
651    /// way.
652    pub(super) async fn force_destroy_indexed_session(
653        self: &Arc<Self>,
654        session_id: String,
655        branch: BranchDisposition,
656        kind: LifecycleKind,
657    ) -> Result<()> {
658        let children = blocking({
659            let session_id = session_id.clone();
660            move || {
661                Ok(crate::database::list_subagents(&session_id)?
662                    .into_iter()
663                    .map(|child| child.child_session_id)
664                    .collect::<Vec<_>>())
665            }
666        })
667        .await?;
668        for child_id in children {
669            // Sub-agents borrow their parent's worker and own no branch.
670            Box::pin(self.force_destroy_indexed_session(
671                child_id.clone(),
672                BranchDisposition::Keep,
673                kind,
674            ))
675            .await
676            .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
677        }
678        self.preempt_active_lifecycle(&session_id).await?;
679        let exists = blocking({
680            let session_id = session_id.clone();
681            move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
682        })
683        .await?;
684        if !exists {
685            return Ok(());
686        }
687        self.run_lifecycle(
688            session_id,
689            kind,
690            move |state, session_id, cancelled| async move {
691                let _recovery_reservation = tokio::task::spawn_blocking({
692                    let observer = state.recovery_observer.clone();
693                    let session_id = session_id.clone();
694                    let cancelled = cancelled.clone();
695                    move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
696                })
697                .await
698                .context("reserve recovery for daemon force-destroy task")??;
699                blocking({
700                    let session_id = session_id.clone();
701                    move || {
702                        let mut controller = Controller::load()?;
703                        let executor = DaemonStageReportingExecutor::new(
704                            CancellableProcessExecutor::new(cancelled),
705                            state,
706                            session_id.clone(),
707                        );
708                        controller.force_destroy_session(&session_id, &executor, branch)?;
709                        crate::controller::move_session::release_move_queue_hold(&session_id);
710                        Ok(DaemonLifecycleResult::Done)
711                    }
712                })
713                .await
714            },
715        )
716        .await?;
717        Ok(())
718    }
719}
720
721/// Wait up to [`FORCE_DESTROY_PREEMPT_TIMEOUT`] for a lifecycle to publish its
722/// result. `Ok(Err(()))` means the task ended without one.
723async fn wait_for_lifecycle_result(
724    result: &mut LifecycleWatch,
725) -> std::result::Result<std::result::Result<(), ()>, tokio::time::error::Elapsed> {
726    tokio::time::timeout(FORCE_DESTROY_PREEMPT_TIMEOUT, async {
727        loop {
728            if result.borrow().is_some() {
729                return Ok(());
730            }
731            if result.changed().await.is_err() {
732                return Err(());
733            }
734        }
735    })
736    .await
737}