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