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        let upgrade_work = crate::upgrade::activity("startup prompt delivery")?;
666        queues.insert(
667            session_id.to_owned(),
668            StartupQueue {
669                pending: VecDeque::from([step]),
670                in_flight: false,
671                cancel: cancel.clone(),
672                task: None,
673            },
674        );
675        let runtime = Arc::clone(self);
676        let drain_session = session_id.to_owned();
677        // Outer task supervises inner task: a panic in the drain becomes a
678        // reported failure that restores the text, not a queue nobody drains.
679        let task = tokio::spawn(async move {
680            let _upgrade_work = upgrade_work;
681            let supervised = {
682                let runtime = Arc::clone(&runtime);
683                let session_id = drain_session.clone();
684                let cancel = cancel.clone();
685                tokio::spawn(async move { runtime.drain_startup_queue(&session_id, &cancel).await })
686            };
687            if let Err(error) = supervised.await {
688                runtime
689                    .fail_startup_queue(
690                        &drain_session,
691                        None,
692                        &format!("the daemon's delivery task failed: {error}"),
693                    )
694                    .await;
695            }
696        });
697        if let Some(queue) = queues.get_mut(session_id) {
698            queue.task = Some(task);
699        }
700        Ok(())
701    }
702
703    /// Wait for the session's harness, then carry out its queued steps in
704    /// order. Every failure path ends in [`Self::fail_startup_queue`], which
705    /// is what puts the text back where the person can see it.
706    async fn drain_startup_queue(self: Arc<Self>, session_id: &str, cancel: &CancellationToken) {
707        let handle = tokio::select! {
708            () = cancel.cancelled() => {
709                self.fail_startup_queue(
710                    session_id,
711                    None,
712                    "the daemon stopped before the session was ready",
713                )
714                .await;
715                return;
716            }
717            ready = self.wait_for_ready_session(session_id) => match ready {
718                Ok(handle) => handle,
719                Err(error) => {
720                    self.fail_startup_queue(session_id, None, &format!("{error:#}"))
721                        .await;
722                    return;
723                }
724            },
725        };
726        loop {
727            let step = {
728                let mut queues = self
729                    .startup_prompts
730                    .lock()
731                    .unwrap_or_else(PoisonError::into_inner);
732                let Some(queue) = queues.get_mut(session_id) else {
733                    return;
734                };
735                match queue.pending.pop_front() {
736                    Some(step) => {
737                        queue.in_flight = true;
738                        step
739                    }
740                    None => {
741                        queues.remove(session_id);
742                        return;
743                    }
744                }
745            };
746            let outcome = tokio::select! {
747                () = cancel.cancelled() => {
748                    Err(anyhow!("the daemon stopped before the prompt was sent"))
749                }
750                result = self.run_startup_step(session_id, &handle, &step) => result,
751            };
752            if let Err(error) = outcome {
753                self.fail_startup_queue(session_id, Some(step), &format!("{error:#}"))
754                    .await;
755                return;
756            }
757            let mut queues = self
758                .startup_prompts
759                .lock()
760                .unwrap_or_else(PoisonError::into_inner);
761            let Some(queue) = queues.get_mut(session_id) else {
762                return;
763            };
764            queue.in_flight = false;
765            if queue.pending.is_empty() {
766                queues.remove(session_id);
767                return;
768            }
769        }
770    }
771
772    async fn run_startup_step(
773        &self,
774        session_id: &str,
775        handle: &crate::session_manager::ManagedSessionHandle,
776        step: &StartupStep,
777    ) -> Result<()> {
778        match step {
779            StartupStep::InstallHandoff(snapshot) => {
780                self.install_archive_handoff(session_id, handle, snapshot)
781                    .await
782            }
783            StartupStep::Prompt {
784                text,
785                inherited_draft,
786            } => {
787                self.submit_startup_prompt(session_id, handle, text, inherited_draft.as_deref())
788                    .await
789            }
790        }
791    }
792
793    /// Submit one queued prompt and give it the history and draft handling a
794    /// prompt submitted from a live composer gets.
795    async fn submit_startup_prompt(
796        &self,
797        session_id: &str,
798        handle: &crate::session_manager::ManagedSessionHandle,
799        text: &str,
800        inherited_draft: Option<&str>,
801    ) -> Result<()> {
802        let bundle_id = self
803            .controller
804            .lock()
805            .unwrap_or_else(PoisonError::into_inner)
806            .state
807            .sessions
808            .get(session_id)
809            .map(|record| record.bundle_id.clone());
810        let ordinal = handle
811            .submit(
812                new_command_id("startup")?,
813                RelayCommand::Prompt {
814                    prompt: vec![ContentBlock::Text(TextContent::new(text.to_owned()))],
815                },
816            )
817            .await?;
818        if let Some(expected) = inherited_draft {
819            let persisted_id = session_id.to_owned();
820            let persisted_expected = expected.to_owned();
821            if let Err(error) = blocking(move || {
822                crate::database::clear_session_draft_input_if_matches(
823                    &persisted_id,
824                    &persisted_expected,
825                )
826            })
827            .await
828            {
829                tracing::warn!(
830                    session_id,
831                    error = format!("{error:#}"),
832                    "the delivered prompt's draft could not be cleared"
833                );
834            }
835            if let Some(record) = self
836                .controller
837                .lock()
838                .unwrap_or_else(PoisonError::into_inner)
839                .state
840                .sessions
841                .get_mut(session_id)
842                && record.draft_input == expected
843            {
844                record.draft_input.clear();
845            }
846            self.publish_revision();
847        }
848        if let Some(bundle_id) = bundle_id {
849            let history_id = session_id.to_owned();
850            let history_text = text.to_owned();
851            if let Err(error) = blocking(move || {
852                crate::database::record_prompt(
853                    &history_id,
854                    &bundle_id,
855                    ordinal,
856                    None,
857                    &history_text,
858                )
859            })
860            .await
861            {
862                tracing::warn!(
863                    session_id,
864                    error = format!("{error:#}"),
865                    "the queued prompt was accepted but its history could not be stored"
866                );
867            }
868        }
869        Ok(())
870    }
871
872    /// Give up on a session's queue: nothing typed is lost, so every prompt
873    /// still in it -- the one that failed and the ones behind it -- goes back
874    /// into the session's saved draft, with a notice saying why.
875    async fn fail_startup_queue(
876        &self,
877        session_id: &str,
878        failed: Option<StartupStep>,
879        reason: &str,
880    ) {
881        let remaining = self
882            .startup_prompts
883            .lock()
884            .unwrap_or_else(PoisonError::into_inner)
885            .remove(session_id)
886            .map(|queue| queue.pending)
887            .unwrap_or_default();
888        let mut texts = Vec::new();
889        let mut dropped_handoff = false;
890        for step in failed.into_iter().chain(remaining) {
891            match step {
892                StartupStep::Prompt { text, .. } => texts.push(text),
893                StartupStep::InstallHandoff(_) => dropped_handoff = true,
894            }
895        }
896        if dropped_handoff {
897            tracing::warn!(
898                session_id,
899                reason,
900                "could not install the restored archive's hand-off"
901            );
902            self.push_notice(
903                session_id,
904                format!("The restored session started without its archived hand-off: {reason}"),
905            );
906        }
907        if texts.is_empty() {
908            return;
909        }
910        let restored = texts.join("\n\n");
911        if let Err(error) = self.append_draft_input(session_id, &restored).await {
912            tracing::warn!(
913                session_id,
914                error = format!("{error:#}"),
915                "a queued prompt could not be saved back into the session's draft"
916            );
917        }
918        self.push_notice(
919            session_id,
920            format!(
921                "Your prompt could not be sent to session {} ({reason}); it is back in the composer draft.",
922                mj_core::state::short_id(session_id)
923            ),
924        );
925        tracing::warn!(
926            session_id,
927            reason,
928            "a queued startup prompt could not be delivered"
929        );
930    }
931
932    /// Put text back into the session's saved composer draft, after whatever
933    /// is already there. The database is the source of truth, because the
934    /// target refresher reloads the controller from disk regularly; the
935    /// in-memory record is updated too so the change shows up at once.
936    pub(super) async fn append_draft_input(&self, session_id: &str, text: &str) -> Result<()> {
937        let existing = self
938            .session_record(session_id)
939            .map(|record| record.draft_input)
940            .unwrap_or_default();
941        let combined = [existing.as_str(), text]
942            .into_iter()
943            .filter(|part| !part.is_empty())
944            .collect::<Vec<_>>()
945            .join("\n\n");
946        let persisted_id = session_id.to_owned();
947        let persisted = combined.clone();
948        let stored =
949            blocking(move || crate::database::set_session_draft_input(&persisted_id, &persisted))
950                .await;
951        if let Some(record) = self
952            .controller
953            .lock()
954            .unwrap_or_else(PoisonError::into_inner)
955            .state
956            .sessions
957            .get_mut(session_id)
958        {
959            record.draft_input = combined;
960        }
961        self.publish_revision();
962        stored
963    }
964
965    /// Stop every startup queue and wait for its drain to report, so the text
966    /// it holds reaches the database while the writer is still running.
967    pub(crate) async fn cancel_and_join_startup_prompts(&self) -> Result<()> {
968        let tasks = {
969            let mut queues = self
970                .startup_prompts
971                .lock()
972                .unwrap_or_else(PoisonError::into_inner);
973            queues
974                .values_mut()
975                .filter_map(|queue| {
976                    queue.cancel.cancel();
977                    queue.task.take()
978                })
979                .collect::<Vec<_>>()
980        };
981        let deadline = tokio::time::Instant::now() + Duration::from_secs(1);
982        let mut outcome = Ok(());
983        for task in tasks {
984            let joined = match tokio::time::timeout_at(deadline, task).await {
985                Ok(Ok(())) => Ok(()),
986                Ok(Err(error)) => Err(anyhow!("startup prompt delivery task failed: {error}")),
987                Err(_) => Err(anyhow!(
988                    "a startup prompt delivery task did not stop within 1s"
989                )),
990            };
991            if outcome.is_ok() {
992                outcome = joined;
993            } else if let Err(error) = joined {
994                tracing::warn!(%error, "another startup prompt drain did not stop cleanly");
995            }
996        }
997        outcome
998    }
999}