Skip to main content

mj_controller/daemon/
state.rs

1use super::*;
2
3impl RuntimeState {
4    pub(super) fn new(
5        session_manager: SessionManagerControl,
6        controller: Controller,
7        recovery_observer: RecoveryObserver,
8        worker_upgrade_observer: WorkerUpgradeObserver,
9        workspaces: Vec<WorkspaceRecord>,
10    ) -> Self {
11        Self::new_with_controller_loader(
12            session_manager,
13            controller,
14            recovery_observer,
15            worker_upgrade_observer,
16            workspaces,
17            Controller::load,
18        )
19    }
20
21    pub(super) fn new_with_controller_loader(
22        session_manager: SessionManagerControl,
23        controller: Controller,
24        recovery_observer: RecoveryObserver,
25        worker_upgrade_observer: WorkerUpgradeObserver,
26        workspaces: Vec<WorkspaceRecord>,
27        controller_loader: fn() -> Result<Controller>,
28    ) -> Self {
29        // Revisions are opaque cursors, so give every daemon incarnation a
30        // fresh high-water mark. Clients that survive a daemon restart must
31        // never wait on, or render, a cursor from the previous process as if
32        // it belonged to the new feed.
33        let initial_revision = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap_or(1);
34        let revisions = RuntimeRevisions::new(initial_revision);
35        let (workspaces_tx, _) = tokio::sync::watch::channel(workspaces);
36        // The host reads `[review]` at each trigger decision. The target
37        // refresher already reloads config.toml every 500 ms and installs the
38        // result here, so arming needs no reload machinery of its own.
39        let review_config = Arc::new(Mutex::new(controller.config.review.clone()));
40        let review_host = TurnReviewHost::spawn_notifying(
41            session_manager.clone(),
42            {
43                let installed = review_config.clone();
44                Arc::new(move || {
45                    installed
46                        .lock()
47                        .unwrap_or_else(PoisonError::into_inner)
48                        .clone()
49                })
50            },
51            revisions.notifier(),
52            Some(recovery_observer.gate.clone()),
53        );
54        Self {
55            attachments: Mutex::new(BTreeMap::new()),
56            phone_status: Mutex::new(WebViewerStatus::Starting),
57            web_viewer: crate::web_viewer::ViewerControl::new(),
58            ever_attached: AtomicBool::new(false),
59            sessions: Mutex::new(BTreeMap::new()),
60            background_policies: Mutex::new(BTreeMap::new()),
61            revisions,
62            workspaces_tx,
63            workspace_refresh: tokio::sync::Mutex::new(()),
64            session_manager,
65            lifecycle: Mutex::new(BTreeMap::new()),
66            workspace_closes: Mutex::new(BTreeMap::new()),
67            workspace_resume_admission: Mutex::new(BTreeMap::new()),
68            harness_readiness: Mutex::new(HarnessReadinessWatch::default()),
69            startup_prompts: Mutex::new(BTreeMap::new()),
70            close_requested: Mutex::new(BTreeSet::new()),
71            controller: Mutex::new(controller),
72            controller_loader,
73            config_mutation: tokio::sync::Mutex::new(()),
74            recovery_observer,
75            worker_upgrade_observer,
76            notices: Mutex::new(VecDeque::new()),
77            next_notice_id: AtomicU64::new(1),
78            review_config,
79            review_host,
80            wiki: crate::sessionwiki::WikiIndexer::spawn(),
81        }
82    }
83
84    /// The review host, for the surfaces that project and resolve reviews.
85    pub fn review_host(&self) -> &TurnReviewHost {
86        &self.review_host
87    }
88
89    /// The SessionWiki indexer, for the surfaces and jobs that trigger a sync.
90    pub fn wiki(&self) -> &crate::sessionwiki::WikiIndexer {
91        &self.wiki
92    }
93
94    /// Search the user's SessionWiki index, with this daemon's own live
95    /// sessions marked so a surface can resume them instead of restoring them.
96    pub async fn wiki_search(&self, query: String, limit: usize) -> Result<WikiSearchPage> {
97        if crate::sessionwiki::sync_is_stale(self.wiki.last_success()) {
98            // Fresh enough matters less than answering now: the sync runs in
99            // the background and the next keystroke sees its result.
100            self.wiki.request_sync(false);
101        }
102        let live = self.live_session_ids();
103        // Every caller is a resume list, and a sub-agent is never resumed on
104        // its own.
105        let rows =
106            blocking(move || crate::sessionwiki::query_rows(&query, limit, &live, false)).await?;
107        // The status is read after the rows, so a sync that finished while the
108        // query ran is reported as finished.
109        Ok(WikiSearchPage {
110            rows,
111            status: self.wiki.status(),
112        })
113    }
114
115    /// The markdown briefing for one indexed session, or `None` when the index
116    /// holds no session with that id.
117    pub async fn wiki_brief(&self, wiki_id: String, max_chars: usize) -> Result<Option<String>> {
118        blocking(move || crate::sessionwiki::brief(&wiki_id, max_chars)).await
119    }
120
121    /// The passages of one indexed session that match a query, or `None` when
122    /// the index holds no session with that id.
123    pub async fn wiki_hits(
124        &self,
125        wiki_id: String,
126        query: String,
127        context_messages: usize,
128        per_message_chars: usize,
129    ) -> Result<Option<WikiHitTranscript>> {
130        blocking(move || {
131            crate::sessionwiki::transcript_hits(
132                &wiki_id,
133                &query,
134                context_messages,
135                per_message_chars,
136            )
137        })
138        .await
139    }
140
141    /// What one indexed session is, and what continuing it would mean, or
142    /// `None` when the index holds no session with that id.
143    pub async fn wiki_session(
144        &self,
145        wiki_id: String,
146    ) -> Result<Option<mj_client::daemon::WikiSessionInfo>> {
147        // A session with a record here is one `mj resume` can take; anything
148        // else has to be restored or imported first.
149        let known = self.live_session_ids();
150        blocking(move || crate::sessionwiki::wiki_session(&wiki_id, &known)).await
151    }
152
153    /// Start a new session carrying a hand-off compacted from an archived one.
154    ///
155    /// The session starts like any other; the hand-off is installed in the
156    /// background once the harness is ready, because building it can take
157    /// several summarizer requests and the caller should not hold a socket
158    /// open for them.
159    /// `None` means the index holds no session with that id.
160    pub async fn restore_wiki_session(
161        self: &Arc<Self>,
162        request: WikiRestoreRequest,
163        cancellation: &CancellationToken,
164    ) -> Result<Option<RegisteredSession>> {
165        let wiki_id = request.wiki_id.clone();
166        let Some(archived) =
167            blocking(move || crate::sessionwiki::archived_session(&wiki_id)).await?
168        else {
169            return Ok(None);
170        };
171        let project_directory = request
172            .project_directory
173            .clone()
174            .or_else(|| archived.project_directory.clone())
175            .context(
176                "name a project directory: the archived session's own project is no longer on this machine",
177            )?;
178        let source = project_directory.display().to_string();
179        let bundle_id = blocking(move || {
180            crate::controller::create_bundle_from_sources(&[source])
181                .map(|created| created.bundle_id)
182                .map_err(anyhow::Error::new)
183        })
184        .await
185        .context("find or create a bundle for the restored session's project")?;
186        let registered = self
187            .start_create_session(CreateSessionRequest {
188                launch_base: None,
189                launch_branch: None,
190                create_managed_worktree: None,
191                mjolnir_subagents: None,
192                initial_prompt: None,
193                workspace_id: request.workspace_id,
194                profile_id: request.profile_id,
195                bundle_id,
196                project_directory: Some(project_directory),
197                target_template_id: request.target_template_id,
198                additional_mounts: request.additional_mounts,
199                resource_allocation: request.resource_allocation,
200                title: archived.title.clone(),
201                // The harness names a session after its first message, and the
202                // first message here carries the hidden hand-off. Pinning the
203                // archived session's own title keeps that text out of every
204                // list the session appears in.
205                session_title_override: Some(archived.title.clone()),
206            })
207            .await?;
208        let session_id = registered.session.id.clone();
209        // The hand-off rides the session's startup queue so that a prompt
210        // typed while the session starts is submitted after the hand-off it
211        // is supposed to read, not before it.
212        self.queue_startup_step(
213            &session_id,
214            StartupStep::InstallHandoff(Box::new(archived.snapshot)),
215            cancellation,
216        )?;
217        Ok(Some(registered))
218    }
219
220    pub(super) fn live_session_ids(&self) -> BTreeSet<String> {
221        self.controller
222            .lock()
223            .unwrap_or_else(PoisonError::into_inner)
224            .state
225            .sessions
226            .keys()
227            .cloned()
228            .collect()
229    }
230
231    /// Compact an archived transcript and hand it to the new session's harness
232    /// as hidden context for its first prompt, which is what the cross-harness
233    /// resume does with a checkpoint.
234    async fn install_archive_handoff(
235        &self,
236        session_id: &str,
237        handle: &crate::session_manager::ManagedSessionHandle,
238        snapshot: &mj_core::archive::CanonicalSessionSnapshot,
239    ) -> Result<()> {
240        let (config, profile_id) = {
241            let controller = self
242                .controller
243                .lock()
244                .unwrap_or_else(PoisonError::into_inner);
245            let profile_id = controller
246                .state
247                .sessions
248                .get(session_id)
249                .map(|record| record.last_profile.clone());
250            (controller.config.clone(), profile_id)
251        };
252        let context_bytes = crate::handoff::profile_handoff_bytes(
253            profile_id.and_then(|id| config.profiles.get(&id)),
254        );
255        let cancel = CancellationToken::new();
256        let handoff = crate::handoff::build_handoff_context(
257            session_id,
258            &config,
259            snapshot,
260            context_bytes,
261            &cancel,
262        )
263        .await
264        .context("compact the archived transcript")?;
265        handle
266            .install_prompt_context(format!(
267                "{} {handoff}",
268                crate::compaction::ARCHIVE_HANDOFF_PREAMBLE
269            ))
270            .await
271            .context("install the archived hand-off")?;
272        tracing::info!(
273            session_id,
274            bytes = handoff.len(),
275            "installed the restored archive's hand-off"
276        );
277        Ok(())
278    }
279
280    /// Wait until a just-created session has a harness that can be handed to.
281    pub(super) async fn wait_for_ready_session(
282        &self,
283        session_id: &str,
284    ) -> Result<crate::session_manager::ManagedSessionHandle> {
285        const POLL: Duration = Duration::from_millis(250);
286        let deadline = tokio::time::Instant::now() + Duration::from_secs(30 * 60);
287        loop {
288            // A session that is coming up moves through Disconnected and
289            // Checkpointing on its way; only a state it cannot leave ends the
290            // wait. This is the same set the API's first-prompt wait accepts.
291            match self.session_state(session_id) {
292                Some(
293                    SessionState::Provisioning
294                    | SessionState::Running
295                    | SessionState::Disconnected
296                    | SessionState::Checkpointing,
297                ) => {}
298                Some(state) => bail!("session {session_id} is {state:?} before its hand-off"),
299                None => bail!("session {session_id} disappeared before its hand-off"),
300            }
301            if let Ok(handle) = self.session_manager.session(session_id).await {
302                let view = handle.view();
303                // A target that is gone never becomes ready. Reporting it now
304                // beats holding the queued work for the full deadline.
305                if let Some(ViewError::TargetMissing(detail)) = &view.error {
306                    bail!("session {session_id} lost its target: {detail}");
307                }
308                if view.connected
309                    && view
310                        .snapshot
311                        .is_some_and(|snapshot| snapshot.operational.native_session_is_ready())
312                {
313                    return Ok(handle);
314                }
315            }
316            ensure!(
317                tokio::time::Instant::now() < deadline,
318                "session {session_id} was not ready for its hand-off within 30 minutes"
319            );
320            tokio::time::sleep(POLL).await;
321        }
322    }
323
324    /// React to the durable outcome of one lifecycle operation.
325    ///
326    /// A session that has just reached `Stopped` is checkpointed and torn
327    /// down, so its transcript is complete and ready to index. This is the one
328    /// place the daemon sees every operation's reloaded durable state.
329    /// Apply a failed create or resume to a record the operation left in
330    /// `Provisioning`.
331    ///
332    /// Both of those operations roll their own record back when they return an
333    /// error, but a task that panics, or one dropped with its runtime, never
334    /// reaches that rollback. The stored result is then the only evidence the
335    /// operation ended, and nothing else owns a `Provisioning` record, so the
336    /// session waits for a provision that will never resume. The operation's
337    /// owner applies the failure here instead.
338    pub(super) async fn fail_unfinished_provisioning(
339        self: &Arc<Self>,
340        session_id: &str,
341        error: &str,
342    ) {
343        let provisioning = {
344            let controller = self
345                .controller
346                .lock()
347                .unwrap_or_else(PoisonError::into_inner);
348            durable_session_state(&controller, session_id) == Some(SessionState::Provisioning)
349        };
350        if !provisioning {
351            return;
352        }
353        let cause = format!("session provisioning ended without finishing: {error}");
354        let applied = blocking({
355            let session_id = session_id.to_owned();
356            move || {
357                let mut controller = Controller::load()?;
358                controller.fail_interrupted_lifecycle(&session_id, &cause)
359            }
360        })
361        .await;
362        match applied {
363            Ok(true) => {
364                if let Err(error) = self.reload_controller().await {
365                    tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after recording a failed provision");
366                }
367            }
368            Ok(false) => {}
369            Err(error) => tracing::warn!(
370                %session_id,
371                error = format!("{error:#}"),
372                "could not record that provisioning ended without finishing"
373            ),
374        }
375    }
376
377    /// Record why a close failed, on the session it was for.
378    ///
379    /// The sentence is written for the person, not copied from the error
380    /// chain: `last_error` is published, and a close failure's chain names
381    /// project paths and SSH hosts. A failure that said what the caller can do
382    /// about it supplies that sentence; every other one points at the daemon
383    /// log entry that carries the whole reason.
384    pub(super) async fn record_failed_close(
385        self: &Arc<Self>,
386        session_id: &str,
387        reference: &str,
388        failure: &LifecycleFailure,
389    ) {
390        self.record_lifecycle_failure(
391            session_id,
392            reference,
393            failure,
394            mj_core::state::CLOSE_FAILURE_PREFIX,
395        )
396        .await;
397    }
398
399    pub(crate) async fn record_lifecycle_failure(
400        self: &Arc<Self>,
401        session_id: &str,
402        reference: &str,
403        failure: &LifecycleFailure,
404        prefix: &str,
405    ) {
406        let cause = match &failure.refusal {
407            Some(refusal) => format!("{prefix}: {refusal}"),
408            None => {
409                format!("{prefix}; the daemon log records the reason under reference {reference}")
410            }
411        };
412        let applied = blocking({
413            let session_id = session_id.to_owned();
414            let cause = cause.clone();
415            move || {
416                let mut controller = Controller::load()?;
417                controller.record_failed_close(&session_id, &cause)
418            }
419        })
420        .await;
421        match applied {
422            Ok(true) => {
423                if let Err(error) = self.reload_controller().await {
424                    tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after recording a failed close");
425                }
426                self.publish_revision();
427            }
428            Ok(false) => {}
429            Err(error) => tracing::warn!(
430                %session_id,
431                error = format!("{error:#}"),
432                "could not record why a close failed"
433            ),
434        }
435    }
436
437    /// Retire a recorded close failure once something for the session has
438    /// succeeded.
439    ///
440    /// The reason is published for a session that is alive, so it has to stop
441    /// being published for one that is working again; a lifecycle transition
442    /// clears `last_error` on its own, and this covers the ordinary actions,
443    /// such as a prompt, that do not.
444    pub async fn clear_recorded_close_failure(self: &Arc<Self>, session_id: &str) {
445        let recorded = self
446            .controller
447            .lock()
448            .unwrap_or_else(PoisonError::into_inner)
449            .state
450            .sessions
451            .get(session_id)
452            .is_some_and(|record| record.public_error().is_some());
453        if !recorded {
454            return;
455        }
456        let cleared = blocking({
457            let session_id = session_id.to_owned();
458            move || {
459                let mut controller = Controller::load()?;
460                controller.clear_recorded_close_failure(&session_id)
461            }
462        })
463        .await;
464        match cleared {
465            Ok(true) => {
466                if let Err(error) = self.reload_controller().await {
467                    tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after clearing a recorded close failure");
468                }
469                self.publish_revision();
470            }
471            Ok(false) => {}
472            Err(error) => tracing::warn!(
473                %session_id,
474                error = format!("{error:#}"),
475                "could not clear a recorded close failure"
476            ),
477        }
478    }
479
480    pub(super) fn note_lifecycle_outcome(&self, session_id: &str) {
481        let stopped = {
482            let controller = self
483                .controller
484                .lock()
485                .unwrap_or_else(PoisonError::into_inner);
486            durable_session_state(&controller, session_id) == Some(SessionState::Stopped)
487        };
488        if stopped {
489            self.wiki.request_sync(false);
490        }
491    }
492
493    pub fn allocate_revision(&self) -> u64 {
494        self.revisions.allocate()
495    }
496
497    pub(super) fn publish_revision(&self) -> u64 {
498        self.revisions.publish()
499    }
500
501    pub(super) fn attachments(&self) -> std::sync::MutexGuard<'_, BTreeMap<String, Attachment>> {
502        self.attachments
503            .lock()
504            .unwrap_or_else(PoisonError::into_inner)
505    }
506
507    pub(super) fn prune_dead_clients(&self) {
508        self.attachments()
509            .retain(|_, attachment| process_is_alive(attachment.pid));
510    }
511
512    pub(super) fn workspace_has_active_resume(&self, workspace_id: &str) -> bool {
513        self.lifecycle
514            .lock()
515            .unwrap_or_else(PoisonError::into_inner)
516            .values()
517            .any(|active| {
518                active.result.borrow().is_none()
519                    && active.resume_workspace_id.as_deref() == Some(workspace_id)
520            })
521    }
522
523    pub fn publish_web_access(&self, access: crate::server::WebViewerAccess) {
524        use crate::server::WebViewerAccess;
525        let status = match &access {
526            WebViewerAccess::Starting => WebViewerStatus::Starting,
527            WebViewerAccess::Ready {
528                viewer_url,
529                viewer_code,
530                qr_login_url,
531                fallback_reason,
532                ..
533            } => WebViewerStatus::Ready {
534                viewer_url: viewer_url.clone(),
535                viewer_code: viewer_code.clone(),
536                qr_login_url: qr_login_url.clone(),
537                fallback_reason: fallback_reason.clone(),
538            },
539            WebViewerAccess::Failed {
540                address, message, ..
541            } => WebViewerStatus::Error {
542                message: format!("{message} Address: {address}"),
543            },
544            WebViewerAccess::Unavailable(message) => WebViewerStatus::Error {
545                message: message.clone(),
546            },
547        };
548        self.web_viewer.publish(access);
549        self.set_phone_status(status);
550    }
551
552    pub(super) fn set_phone_status(&self, status: WebViewerStatus) {
553        *self
554            .phone_status
555            .lock()
556            .unwrap_or_else(PoisonError::into_inner) = status;
557    }
558
559    pub(super) fn phone_status(&self) -> WebViewerStatus {
560        self.phone_status
561            .lock()
562            .unwrap_or_else(PoisonError::into_inner)
563            .clone()
564    }
565
566    pub(super) fn workspaces(&self) -> tokio::sync::watch::Receiver<Vec<WorkspaceRecord>> {
567        self.workspaces_tx.subscribe()
568    }
569
570    pub(super) fn worker_poll_exclusion_session_ids(
571        &self,
572        controller: &Controller,
573    ) -> BTreeSet<String> {
574        self.lifecycle
575            .lock()
576            .unwrap_or_else(PoisonError::into_inner)
577            .iter()
578            .filter(|(session_id, active)| {
579                active.result.borrow().is_none()
580                    && (active.move_source_closed
581                        || lifecycle_owns_worker_target(
582                            active.kind,
583                            controller
584                                .state
585                                .sessions
586                                .get(*session_id)
587                                .map(|session| session.state),
588                        ))
589            })
590            .map(|(session_id, _)| session_id.clone())
591            .collect()
592    }
593
594    pub fn revisions(&self) -> tokio::sync::watch::Receiver<u64> {
595        self.revisions.subscribe()
596    }
597
598    /// Read the config the daemon serves right now. A task on a schedule reads
599    /// it again on every tick, so a reload reaches it without a restart.
600    pub fn with_config<T>(&self, read: impl FnOnce(&Config) -> T) -> T {
601        read(
602            &self
603                .controller
604                .lock()
605                .unwrap_or_else(PoisonError::into_inner)
606                .config,
607        )
608    }
609
610    /// Create a bundle under the daemon's config-mutation coordinator. The
611    /// controller helper also takes the cross-process config lock, so a TUI
612    /// transaction cannot race this one while the daemon's other config
613    /// writers are excluded by this mutex.
614    pub async fn create_quick_bundle(
615        &self,
616        source: String,
617    ) -> std::result::Result<
618        crate::controller::QuickBundleCreation,
619        crate::controller::QuickBundleFailure,
620    > {
621        let _mutation = self.config_mutation.lock().await;
622        tokio::task::spawn_blocking(move || crate::controller::create_quick_bundle(&source))
623            .await
624            .map_err(|error| {
625                crate::controller::QuickBundleFailure::Persistence(anyhow!(
626                    "bundle creation task panicked: {error}"
627                ))
628            })?
629    }
630
631    /// Persist exactly the picker-selected repository set under the same
632    /// config-mutation coordinator as legacy quick-bundle creation.
633    pub async fn create_bundle_from_sources(
634        &self,
635        sources: Vec<String>,
636    ) -> std::result::Result<
637        crate::controller::QuickBundleCreation,
638        crate::controller::QuickBundleFailure,
639    > {
640        let _mutation = self.config_mutation.lock().await;
641        tokio::task::spawn_blocking(move || crate::controller::create_bundle_from_sources(&sources))
642            .await
643            .map_err(|error| {
644                crate::controller::QuickBundleFailure::Persistence(anyhow!(
645                    "bundle creation task panicked: {error}"
646                ))
647            })?
648    }
649
650    /// Hand a changed workspace list to the terminal clients and the web viewer.
651    pub(crate) fn publish_workspaces(&self, workspaces: Vec<WorkspaceRecord>) {
652        self.workspaces_tx.send_replace(workspaces);
653        self.publish_revision();
654    }
655
656    pub(crate) async fn refresh_workspaces(&self) -> Result<()> {
657        // Keep the read and publication together: a delayed read must not
658        // publish an older list after a newer removal has been published.
659        let _refresh = self.workspace_refresh.lock().await;
660        let workspaces = tokio::task::spawn_blocking(crate::database::list_workspaces)
661            .await
662            .context("daemon workspace refresh task panicked")??;
663        self.publish_workspaces(workspaces);
664        Ok(())
665    }
666
667    /// Queue one piece of startup work for a session, starting the drain task
668    /// when this is the session's first step.
669    ///
670    /// Steps are carried out in the order they arrive. The entry in the map
671    /// exists only while a drain owns it, so "no entry" and "no live task"
672    /// are the same condition and a second call never starts a second drain.
673    pub(crate) fn queue_startup_step(
674        self: &Arc<Self>,
675        session_id: &str,
676        step: StartupStep,
677        cancellation: &CancellationToken,
678    ) -> Result<()> {
679        // The same set `wait_for_ready_session` accepts: anything else will
680        // never become ready, so queueing would only lose the text later.
681        match self.session_state(session_id) {
682            Some(
683                SessionState::Provisioning
684                | SessionState::Running
685                | SessionState::Disconnected
686                | SessionState::Checkpointing,
687            ) => {}
688            Some(state) => {
689                bail!("session {session_id} is {state:?}; it cannot take a queued prompt")
690            }
691            None => bail!("unknown session {session_id}"),
692        }
693        let mut queues = self
694            .startup_prompts
695            .lock()
696            .unwrap_or_else(PoisonError::into_inner);
697        if let Some(queue) = queues.get_mut(session_id) {
698            queue.pending.push_back(step);
699            return Ok(());
700        }
701        let cancel = cancellation.child_token();
702        let upgrade_work = crate::upgrade::activity("startup prompt delivery")?;
703        queues.insert(
704            session_id.to_owned(),
705            StartupQueue {
706                pending: VecDeque::from([step]),
707                in_flight: false,
708                cancel: cancel.clone(),
709                task: None,
710            },
711        );
712        let runtime = Arc::clone(self);
713        let drain_session = session_id.to_owned();
714        // Outer task supervises inner task: a panic in the drain becomes a
715        // reported failure that restores the text, not a queue nobody drains.
716        let task = tokio::spawn(async move {
717            let _upgrade_work = upgrade_work;
718            let supervised = {
719                let runtime = Arc::clone(&runtime);
720                let session_id = drain_session.clone();
721                let cancel = cancel.clone();
722                tokio::spawn(async move { runtime.drain_startup_queue(&session_id, &cancel).await })
723            };
724            if let Err(error) = supervised.await {
725                runtime
726                    .fail_startup_queue(
727                        &drain_session,
728                        None,
729                        &format!("the daemon's delivery task failed: {error}"),
730                    )
731                    .await;
732            }
733        });
734        if let Some(queue) = queues.get_mut(session_id) {
735            queue.task = Some(task);
736        }
737        Ok(())
738    }
739
740    /// Wait for the session's harness, then carry out its queued steps in
741    /// order. Every failure path ends in [`Self::fail_startup_queue`], which
742    /// is what puts the text back where the person can see it.
743    async fn drain_startup_queue(self: Arc<Self>, session_id: &str, cancel: &CancellationToken) {
744        let handle = tokio::select! {
745            () = cancel.cancelled() => {
746                self.fail_startup_queue(
747                    session_id,
748                    None,
749                    "the daemon stopped before the session was ready",
750                )
751                .await;
752                return;
753            }
754            ready = self.wait_for_ready_session(session_id) => match ready {
755                Ok(handle) => handle,
756                Err(error) => {
757                    self.fail_startup_queue(session_id, None, &format!("{error:#}"))
758                        .await;
759                    return;
760                }
761            },
762        };
763        loop {
764            let step = {
765                let mut queues = self
766                    .startup_prompts
767                    .lock()
768                    .unwrap_or_else(PoisonError::into_inner);
769                let Some(queue) = queues.get_mut(session_id) else {
770                    return;
771                };
772                match queue.pending.pop_front() {
773                    Some(step) => {
774                        queue.in_flight = true;
775                        step
776                    }
777                    None => {
778                        queues.remove(session_id);
779                        return;
780                    }
781                }
782            };
783            let outcome = tokio::select! {
784                () = cancel.cancelled() => {
785                    Err(anyhow!("the daemon stopped before the prompt was sent"))
786                }
787                result = self.run_startup_step(session_id, &handle, &step) => result,
788            };
789            if let Err(error) = outcome {
790                self.fail_startup_queue(session_id, Some(step), &format!("{error:#}"))
791                    .await;
792                return;
793            }
794            let mut queues = self
795                .startup_prompts
796                .lock()
797                .unwrap_or_else(PoisonError::into_inner);
798            let Some(queue) = queues.get_mut(session_id) else {
799                return;
800            };
801            queue.in_flight = false;
802            if queue.pending.is_empty() {
803                queues.remove(session_id);
804                return;
805            }
806        }
807    }
808
809    async fn run_startup_step(
810        &self,
811        session_id: &str,
812        handle: &crate::session_manager::ManagedSessionHandle,
813        step: &StartupStep,
814    ) -> Result<()> {
815        match step {
816            StartupStep::InstallHandoff(snapshot) => {
817                self.install_archive_handoff(session_id, handle, snapshot)
818                    .await
819            }
820            StartupStep::Prompt {
821                text,
822                inherited_draft,
823            } => {
824                self.submit_startup_prompt(session_id, handle, text, inherited_draft.as_deref())
825                    .await
826            }
827        }
828    }
829
830    /// Submit one queued prompt and give it the history and draft handling a
831    /// prompt submitted from a live composer gets.
832    async fn submit_startup_prompt(
833        &self,
834        session_id: &str,
835        handle: &crate::session_manager::ManagedSessionHandle,
836        text: &str,
837        inherited_draft: Option<&str>,
838    ) -> Result<()> {
839        let bundle_id = self
840            .controller
841            .lock()
842            .unwrap_or_else(PoisonError::into_inner)
843            .state
844            .sessions
845            .get(session_id)
846            .map(|record| record.bundle_id.clone());
847        let ordinal = handle
848            .submit(
849                new_command_id("startup")?,
850                RelayCommand::Prompt {
851                    prompt: vec![ContentBlock::Text(TextContent::new(text.to_owned()))],
852                },
853            )
854            .await?;
855        if let Some(expected) = inherited_draft {
856            let persisted_id = session_id.to_owned();
857            let persisted_expected = expected.to_owned();
858            if let Err(error) = blocking(move || {
859                crate::database::clear_session_draft_input_if_matches(
860                    &persisted_id,
861                    &persisted_expected,
862                )
863            })
864            .await
865            {
866                tracing::warn!(
867                    session_id,
868                    error = format!("{error:#}"),
869                    "the delivered prompt's draft could not be cleared"
870                );
871            }
872            if let Some(record) = self
873                .controller
874                .lock()
875                .unwrap_or_else(PoisonError::into_inner)
876                .state
877                .sessions
878                .get_mut(session_id)
879                && record.draft_input == expected
880            {
881                record.draft_input.clear();
882            }
883            self.publish_revision();
884        }
885        if let Some(bundle_id) = bundle_id {
886            let history_id = session_id.to_owned();
887            let history_text = text.to_owned();
888            if let Err(error) = blocking(move || {
889                crate::database::record_prompt(
890                    &history_id,
891                    &bundle_id,
892                    ordinal,
893                    None,
894                    &history_text,
895                )
896            })
897            .await
898            {
899                tracing::warn!(
900                    session_id,
901                    error = format!("{error:#}"),
902                    "the queued prompt was accepted but its history could not be stored"
903                );
904            }
905        }
906        Ok(())
907    }
908
909    /// Give up on a session's queue: nothing typed is lost, so every prompt
910    /// still in it -- the one that failed and the ones behind it -- goes back
911    /// into the session's saved draft, with a notice saying why.
912    async fn fail_startup_queue(
913        &self,
914        session_id: &str,
915        failed: Option<StartupStep>,
916        reason: &str,
917    ) {
918        let remaining = self
919            .startup_prompts
920            .lock()
921            .unwrap_or_else(PoisonError::into_inner)
922            .remove(session_id)
923            .map(|queue| queue.pending)
924            .unwrap_or_default();
925        let mut texts = Vec::new();
926        let mut dropped_handoff = false;
927        for step in failed.into_iter().chain(remaining) {
928            match step {
929                StartupStep::Prompt { text, .. } => texts.push(text),
930                StartupStep::InstallHandoff(_) => dropped_handoff = true,
931            }
932        }
933        if dropped_handoff {
934            tracing::warn!(
935                session_id,
936                reason,
937                "could not install the restored archive's hand-off"
938            );
939            self.push_notice(
940                session_id,
941                format!("The restored session started without its archived hand-off: {reason}"),
942            );
943        }
944        if texts.is_empty() {
945            return;
946        }
947        let restored = texts.join("\n\n");
948        if let Err(error) = self.append_draft_input(session_id, &restored).await {
949            tracing::warn!(
950                session_id,
951                error = format!("{error:#}"),
952                "a queued prompt could not be saved back into the session's draft"
953            );
954        }
955        self.push_notice(
956            session_id,
957            format!(
958                "Your prompt could not be sent to session {} ({reason}); it is back in the composer draft.",
959                mj_core::state::short_id(session_id)
960            ),
961        );
962        tracing::warn!(
963            session_id,
964            reason,
965            "a queued startup prompt could not be delivered"
966        );
967    }
968
969    /// Put text back into the session's saved composer draft, after whatever
970    /// is already there. The database is the source of truth, because the
971    /// target refresher reloads the controller from disk regularly; the
972    /// in-memory record is updated too so the change shows up at once.
973    pub(super) async fn append_draft_input(&self, session_id: &str, text: &str) -> Result<()> {
974        let existing = self
975            .session_record(session_id)
976            .map(|record| record.draft_input)
977            .unwrap_or_default();
978        let combined = [existing.as_str(), text]
979            .into_iter()
980            .filter(|part| !part.is_empty())
981            .collect::<Vec<_>>()
982            .join("\n\n");
983        let persisted_id = session_id.to_owned();
984        let persisted = combined.clone();
985        let stored =
986            blocking(move || crate::database::set_session_draft_input(&persisted_id, &persisted))
987                .await;
988        if let Some(record) = self
989            .controller
990            .lock()
991            .unwrap_or_else(PoisonError::into_inner)
992            .state
993            .sessions
994            .get_mut(session_id)
995        {
996            record.draft_input = combined;
997        }
998        self.publish_revision();
999        stored
1000    }
1001
1002    /// Stop every startup queue and wait for its drain to report, so the text
1003    /// it holds reaches the database while the writer is still running.
1004    pub(crate) async fn cancel_and_join_startup_prompts(&self) -> Result<()> {
1005        let tasks = {
1006            let mut queues = self
1007                .startup_prompts
1008                .lock()
1009                .unwrap_or_else(PoisonError::into_inner);
1010            queues
1011                .values_mut()
1012                .filter_map(|queue| {
1013                    queue.cancel.cancel();
1014                    queue.task.take()
1015                })
1016                .collect::<Vec<_>>()
1017        };
1018        let deadline = tokio::time::Instant::now() + Duration::from_secs(1);
1019        let mut outcome = Ok(());
1020        for task in tasks {
1021            let joined = match tokio::time::timeout_at(deadline, task).await {
1022                Ok(Ok(())) => Ok(()),
1023                Ok(Err(error)) => Err(anyhow!("startup prompt delivery task failed: {error}")),
1024                Err(_) => Err(anyhow!(
1025                    "a startup prompt delivery task did not stop within 1s"
1026                )),
1027            };
1028            if outcome.is_ok() {
1029                outcome = joined;
1030            } else if let Err(error) = joined {
1031                tracing::warn!(%error, "another startup prompt drain did not stop cleanly");
1032            }
1033        }
1034        outcome
1035    }
1036}