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