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