Skip to main content

mj_controller/daemon/
state.rs

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