Skip to main content

mj_controller/daemon/
state.rs

1/// Consecutive failed delivery rounds before a startup step is given up.
2const STARTUP_STEP_ATTEMPTS: u32 = 5;
3
4use super::*;
5
6impl RuntimeState {
7    pub(crate) fn worker_background_gate(&self) -> Arc<crate::recovery_gate::RecoveryGate> {
8        self.recovery_observer.gate.clone()
9    }
10
11    pub(crate) fn new(
12        session_manager: SessionManagerControl,
13        controller: Controller,
14        recovery_observer: RecoveryObserver,
15        worker_upgrade_observer: WorkerUpgradeObserver,
16        workspaces: Vec<WorkspaceRecord>,
17    ) -> Self {
18        let mut state = Self::new_with_controller_loader(
19            session_manager,
20            controller,
21            recovery_observer,
22            worker_upgrade_observer,
23            workspaces,
24            Controller::load,
25        );
26        state.committed = crate::database::database_writer_installed().then(|| {
27            crate::database::subscribe_committed_state().expect("installed database writer")
28        });
29        state
30    }
31
32    pub(super) fn new_with_controller_loader(
33        session_manager: SessionManagerControl,
34        controller: Controller,
35        recovery_observer: RecoveryObserver,
36        worker_upgrade_observer: WorkerUpgradeObserver,
37        workspaces: Vec<WorkspaceRecord>,
38        controller_loader: fn() -> Result<Controller>,
39    ) -> Self {
40        // Revisions are opaque cursors, so give every daemon incarnation a
41        // fresh high-water mark. Clients that survive a daemon restart must
42        // never wait on, or render, a cursor from the previous process as if
43        // it belonged to the new feed.
44        let initial_revision = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap_or(1);
45        let revisions = RuntimeRevisions::new(initial_revision);
46        let (workspaces_tx, _) = tokio::sync::watch::channel(workspaces);
47        // The host reads `[review]` at each trigger decision. The target
48        // refresher already reloads config.toml every 500 ms and installs the
49        // result here, so arming needs no reload machinery of its own.
50        let review_config = Arc::new(Mutex::new(controller.config.review.clone()));
51        // A session's own review choice is read from the latest committed
52        // records, which the writer publishes without the owner's lock.
53        let committed_sessions = crate::database::database_writer_installed()
54            .then(|| crate::database::subscribe_committed_state().ok())
55            .flatten();
56        let profile_catalog = crate::review_host::SharedProfileCatalog::default();
57        let review_host = TurnReviewHost::spawn_notifying(
58            session_manager.clone(),
59            {
60                let installed = review_config.clone();
61                Arc::new(move |session_id: &str| {
62                    let session = committed_sessions.as_ref().and_then(|committed| {
63                        committed
64                            .borrow()
65                            .as_ref()
66                            .ok()
67                            .and_then(|committed| committed.state.sessions.get(session_id))
68                            .and_then(|session| session.review.clone())
69                    });
70                    installed
71                        .lock()
72                        .unwrap_or_else(PoisonError::into_inner)
73                        .for_session(session.as_ref())
74                })
75            },
76            revisions.notifier(),
77            Some(recovery_observer.gate.clone()),
78            Some(profile_catalog.clone()),
79        );
80        Self {
81            attachments: Mutex::new(BTreeMap::new()),
82            phone_status: Mutex::new(WebViewerStatus::Starting),
83            web_viewer: crate::web_viewer::ViewerControl::new(),
84            ever_attached: AtomicBool::new(false),
85            revisions,
86            workspaces_tx,
87            workspace_refresh: tokio::sync::Mutex::new(()),
88            session_manager,
89            owner: Mutex::new(RuntimeStateOwner::new(controller)),
90            wait_prompt_backend: std::sync::OnceLock::new(),
91            wait_prompt_retries: Mutex::new(BTreeMap::new()),
92            credential_targets: Arc::new(tokio::sync::watch::channel(Vec::new()).0),
93            feed: Mutex::new(feed::RuntimeHistory::default()),
94            committed: None,
95            workspace_closes: Mutex::new(BTreeMap::new()),
96            workspace_resume_admission: Mutex::new(BTreeMap::new()),
97            harness_readiness: Mutex::new(HarnessReadinessWatch::default()),
98            startup_prompts: Mutex::new(BTreeMap::new()),
99            startup_enqueue: tokio::sync::Mutex::new(()),
100            controller_loader,
101            config_mutation: tokio::sync::Mutex::new(()),
102            projects: Arc::new(crate::project_catalog::Catalog::default()),
103            profile_catalog,
104            recovery_observer,
105            worker_upgrade_observer,
106            notices: Mutex::new(VecDeque::new()),
107            next_notice_id: AtomicU64::new(1),
108            quota: Mutex::new(QuotaBoard::default()),
109            capacity: std::sync::OnceLock::new(),
110            review_config,
111            review_host,
112            wiki: crate::sessionwiki::WikiIndexer::spawn(),
113        }
114    }
115
116    pub(crate) fn projects(&self) -> Arc<crate::project_catalog::Catalog> {
117        self.projects.clone()
118    }
119
120    pub async fn project_catalog(
121        &self,
122        refresh: bool,
123        retry: bool,
124    ) -> Result<mj_core::project_catalog::ProjectCatalogView> {
125        if refresh {
126            self.projects.request(retry);
127        }
128        let catalog = self.projects.clone();
129        blocking(move || catalog.view()).await
130    }
131
132    /// The review host, for the surfaces that project and resolve reviews.
133    pub fn review_host(&self) -> &TurnReviewHost {
134        &self.review_host
135    }
136
137    /// The SessionWiki indexer, for the surfaces and jobs that trigger a sync.
138    pub fn wiki(&self) -> &crate::sessionwiki::WikiIndexer {
139        &self.wiki
140    }
141
142    /// Search the user's SessionWiki index, with this daemon's own live
143    /// sessions marked so a surface can resume them instead of restoring them.
144    pub async fn wiki_search(&self, query: String, limit: usize) -> Result<WikiSearchPage> {
145        self.request_wiki_sync_if_stale();
146        let live = self.live_session_ids();
147        // Every caller is a resume list, and a sub-agent is never resumed on
148        // its own.
149        let rows =
150            blocking(move || crate::sessionwiki::query_rows(&query, limit, &live, false)).await?;
151        // The status is read after the rows, so a sync that finished while the
152        // query ran is reported as finished.
153        Ok(WikiSearchPage {
154            rows,
155            status: self.wiki.status(),
156        })
157    }
158
159    /// Ask for a background sync when the index is stale. Every search asks
160    /// here, so this is the one place that decides when a search syncs.
161    /// Answering now matters more than answering fresh: the sync runs in the
162    /// background and the next keystroke sees its result.
163    fn request_wiki_sync_if_stale(&self) {
164        if crate::sessionwiki::sync_is_stale(self.wiki.last_success()) {
165            self.wiki.request_sync(false);
166        }
167    }
168
169    /// The live sessions whose user or agent messages contain `query`. Like
170    /// [`Self::wiki_search`], it answers from the index as it is and asks for
171    /// a sync for the next search.
172    pub async fn session_text_search(&self, query: String) -> Result<Vec<SessionTextMatch>> {
173        self.request_wiki_sync_if_stale();
174        let live = self.live_session_ids();
175        blocking(move || crate::sessionwiki::session_text_matches(&query, &live)).await
176    }
177
178    /// The markdown briefing for one indexed session, or `None` when the index
179    /// holds no session with that id.
180    pub async fn wiki_brief(&self, wiki_id: String, max_chars: usize) -> Result<Option<String>> {
181        blocking(move || crate::sessionwiki::brief(&wiki_id, max_chars)).await
182    }
183
184    /// The passages of one indexed session that match a query, or `None` when
185    /// the index holds no session with that id.
186    pub async fn wiki_hits(
187        &self,
188        wiki_id: String,
189        query: String,
190        context_messages: usize,
191        per_message_chars: usize,
192    ) -> Result<Option<WikiHitTranscript>> {
193        blocking(move || {
194            crate::sessionwiki::transcript_hits(
195                &wiki_id,
196                &query,
197                context_messages,
198                per_message_chars,
199            )
200        })
201        .await
202    }
203
204    /// What one indexed session is, and what continuing it would mean, or
205    /// `None` when the index holds no session with that id.
206    pub async fn wiki_session(
207        &self,
208        wiki_id: String,
209    ) -> Result<Option<mj_client::daemon::WikiSessionInfo>> {
210        // A session with a record here is one `mj resume` can take; anything
211        // else has to be restored or imported first.
212        let known = self.live_session_ids();
213        blocking(move || crate::sessionwiki::wiki_session(&wiki_id, &known)).await
214    }
215
216    /// Start a new session carrying a hand-off compacted from an archived one.
217    ///
218    /// The session starts like any other; the hand-off is installed in the
219    /// background once the harness is ready, because building it can take
220    /// several summarizer requests and the caller should not hold a socket
221    /// open for them.
222    /// `None` means the index holds no session with that id.
223    pub async fn restore_wiki_session(
224        self: &Arc<Self>,
225        request: WikiRestoreRequest,
226        cancellation: &CancellationToken,
227    ) -> Result<Option<RegisteredSession>> {
228        let wiki_id = request.wiki_id.clone();
229        let Some(archived) =
230            blocking(move || crate::sessionwiki::archived_session(&wiki_id)).await?
231        else {
232            return Ok(None);
233        };
234        let project_directory = request
235            .project_directory
236            .clone()
237            .or_else(|| archived.project_directory.clone())
238            .context(
239                "name a project directory: the archived session's own project is no longer on this machine",
240            )?;
241        let source = project_directory.display().to_string();
242        let bundle_id = blocking(move || {
243            crate::controller::create_bundle_from_sources(&[source])
244                .map(|created| created.bundle_id)
245                .map_err(anyhow::Error::new)
246        })
247        .await
248        .context("find or create a bundle for the restored session's project")?;
249        // Keep queued user prompts behind the archive context even if another
250        // client observes the newly registered session before this call returns.
251        let _startup_admission = self.startup_enqueue.lock().await;
252        let registered = self
253            .start_create_session(CreateSessionRequest {
254                at: None,
255                branch: None,
256                base: None,
257                create_managed_worktree: None,
258                subagents: None,
259                review: None,
260                initial_prompt: None,
261                workspace_id: request.workspace_id,
262                profile_id: request.profile_id,
263                bundle_id,
264                project_directory: Some(project_directory),
265                target_template_id: request.target_template_id,
266                additional_mounts: request.additional_mounts,
267                resource_allocation: request.resource_allocation,
268                title: archived.title.clone(),
269                // The harness names a session after its first message, and the
270                // first message here carries the hidden hand-off. Pinning the
271                // archived session's own title keeps that text out of every
272                // list the session appears in.
273                session_title_override: Some(archived.title.clone()),
274            })
275            .await?;
276        let session_id = registered.session.id.clone();
277        // The hand-off rides the session's startup queue so that a prompt
278        // typed while the session starts is submitted after the hand-off it
279        // is supposed to read, not before it.
280        self.queue_startup_steps_admitted(
281            &session_id,
282            vec![(
283                new_command_id("startup")?,
284                StartupStep::InstallHandoff(Box::new(archived.snapshot)),
285            )],
286            None,
287            cancellation,
288        )
289        .await?;
290        Ok(Some(registered))
291    }
292
293    pub(super) fn live_session_ids(&self) -> BTreeSet<String> {
294        self.owner()
295            .controller()
296            .state
297            .sessions
298            .keys()
299            .cloned()
300            .collect()
301    }
302
303    /// Compact an archived transcript and hand it to the new session's harness
304    /// as hidden context for its first prompt, which is what the cross-harness
305    /// resume does with a checkpoint.
306    async fn prepare_archive_handoff(
307        &self,
308        session_id: &str,
309        snapshot: &mj_core::archive::CanonicalSessionSnapshot,
310    ) -> Result<String> {
311        let (config, profile_id) = {
312            let controller_owner = self.owner();
313            let controller = controller_owner.controller();
314            let profile_id = controller
315                .state
316                .sessions
317                .get(session_id)
318                .map(|record| record.last_profile.clone());
319            (controller.config.clone(), profile_id)
320        };
321        let context_bytes = crate::handoff::profile_handoff_bytes(
322            profile_id.and_then(|id| config.profiles.get(&id)),
323        );
324        let cancel = CancellationToken::new();
325        let handoff = crate::handoff::build_handoff_context(
326            session_id,
327            &config,
328            snapshot,
329            context_bytes,
330            &cancel,
331        )
332        .await
333        .context("compact the archived transcript")?;
334        Ok(format!(
335            "{} {handoff}",
336            crate::compaction::ARCHIVE_HANDOFF_PREAMBLE
337        ))
338    }
339
340    /// Wait until a just-created session has a harness that can be handed to.
341    pub(super) async fn wait_for_ready_session(
342        &self,
343        session_id: &str,
344    ) -> Result<crate::session_manager::ManagedSessionHandle> {
345        const POLL: Duration = Duration::from_millis(250);
346        let deadline = tokio::time::Instant::now() + Duration::from_secs(30 * 60);
347        loop {
348            // A session that is coming up moves through Disconnected and
349            // Checkpointing on its way; only a state it cannot leave ends the
350            // wait. This is the same set the API's first-prompt wait accepts.
351            match self.session_state(session_id) {
352                Some(
353                    SessionState::Provisioning
354                    | SessionState::Running
355                    | SessionState::Disconnected
356                    | SessionState::Checkpointing,
357                ) => {}
358                Some(state) => bail!("session {session_id} is {state:?} before its hand-off"),
359                None => bail!("session {session_id} disappeared before its hand-off"),
360            }
361            if let Ok(handle) = self.session_manager.session(session_id).await {
362                let view = handle.view();
363                // A target that is gone never becomes ready. Reporting it now
364                // beats holding the queued work for the full deadline.
365                if let Some(ViewError::TargetMissing(detail)) = &view.error {
366                    bail!("session {session_id} lost its target: {detail}");
367                }
368                if view.connected
369                    && view
370                        .snapshot
371                        .is_some_and(|snapshot| snapshot.operational.native_session_is_ready())
372                {
373                    return Ok(handle);
374                }
375            }
376            ensure!(
377                tokio::time::Instant::now() < deadline,
378                "session {session_id} was not ready for its hand-off within 30 minutes"
379            );
380            tokio::time::sleep(POLL).await;
381        }
382    }
383
384    /// React to the durable outcome of one lifecycle operation.
385    ///
386    /// A session that has just reached `Stopped` is checkpointed and torn
387    /// down, so its transcript is complete and ready to index. This is the one
388    /// place the daemon sees every operation's reloaded durable state.
389    /// Apply a failed create or resume to a record the operation left in
390    /// `Provisioning`.
391    ///
392    /// Both of those operations roll their own record back when they return an
393    /// error, but a task that panics, or one dropped with its runtime, never
394    /// reaches that rollback. The stored result is then the only evidence the
395    /// operation ended, and nothing else owns a `Provisioning` record, so the
396    /// session waits for a provision that will never resume. The operation's
397    /// owner applies the failure here instead.
398    pub(super) async fn fail_unfinished_provisioning(
399        self: &Arc<Self>,
400        session_id: &str,
401        error: &str,
402    ) {
403        let provisioning = {
404            let controller_owner = self.owner();
405            let controller = controller_owner.controller();
406            durable_session_state(controller, session_id) == Some(SessionState::Provisioning)
407        };
408        if !provisioning {
409            return;
410        }
411        let cause = format!("session provisioning ended without finishing: {error}");
412        let applied = blocking({
413            let session_id = session_id.to_owned();
414            move || {
415                let mut controller = Controller::load()?;
416                controller.fail_interrupted_lifecycle(&session_id, &cause)
417            }
418        })
419        .await;
420        match applied {
421            Ok(true) => {
422                if let Err(error) = self.reload_controller().await {
423                    tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after recording a failed provision");
424                }
425            }
426            Ok(false) => {}
427            Err(error) => tracing::warn!(
428                %session_id,
429                error = format!("{error:#}"),
430                "could not record that provisioning ended without finishing"
431            ),
432        }
433    }
434
435    /// Record why a close failed, on the session it was for.
436    ///
437    /// The sentence is written for the person, not copied from the error
438    /// chain: `last_error` is published, and a close failure's chain names
439    /// project paths and SSH hosts. A failure that said what the caller can do
440    /// about it supplies that sentence; every other one points at the daemon
441    /// log entry that carries the whole reason.
442    pub(super) async fn record_failed_close(
443        self: &Arc<Self>,
444        session_id: &str,
445        reference: &str,
446        failure: &LifecycleFailure,
447    ) {
448        self.record_lifecycle_failure(
449            session_id,
450            reference,
451            failure,
452            mj_core::state::CLOSE_FAILURE_PREFIX,
453        )
454        .await;
455    }
456
457    pub(crate) async fn record_lifecycle_failure(
458        self: &Arc<Self>,
459        session_id: &str,
460        reference: &str,
461        failure: &LifecycleFailure,
462        prefix: &str,
463    ) {
464        let cause = match &failure.refusal {
465            Some(refusal) => format!("{prefix}: {refusal}"),
466            None => {
467                format!("{prefix}; the daemon log records the reason under reference {reference}")
468            }
469        };
470        let applied = blocking({
471            let session_id = session_id.to_owned();
472            let cause = cause.clone();
473            move || {
474                let mut controller = Controller::load()?;
475                controller.record_failed_close(&session_id, &cause)
476            }
477        })
478        .await;
479        match applied {
480            Ok(true) => {
481                if let Err(error) = self.reload_controller().await {
482                    tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after recording a failed close");
483                }
484                self.publish_revision();
485            }
486            Ok(false) => {}
487            Err(error) => tracing::warn!(
488                %session_id,
489                error = format!("{error:#}"),
490                "could not record why a close failed"
491            ),
492        }
493    }
494
495    /// Retire a recorded close failure once something for the session has
496    /// succeeded.
497    ///
498    /// The reason is published for a session that is alive, so it has to stop
499    /// being published for one that is working again; a lifecycle transition
500    /// clears `last_error` on its own, and this covers the ordinary actions,
501    /// such as a prompt, that do not.
502    pub async fn clear_recorded_close_failure(self: &Arc<Self>, session_id: &str) {
503        let recorded = self
504            .owner()
505            .controller()
506            .state
507            .sessions
508            .get(session_id)
509            .is_some_and(|record| record.public_error().is_some());
510        if !recorded {
511            return;
512        }
513        let cleared = blocking({
514            let session_id = session_id.to_owned();
515            move || {
516                let mut controller = Controller::load()?;
517                controller.clear_recorded_close_failure(&session_id)
518            }
519        })
520        .await;
521        match cleared {
522            Ok(true) => {
523                if let Err(error) = self.reload_controller().await {
524                    tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after clearing a recorded close failure");
525                }
526                self.publish_revision();
527            }
528            Ok(false) => {}
529            Err(error) => tracing::warn!(
530                %session_id,
531                error = format!("{error:#}"),
532                "could not clear a recorded close failure"
533            ),
534        }
535    }
536
537    pub(super) fn note_lifecycle_outcome(&self, session_id: &str) {
538        let stopped = {
539            let controller_owner = self.owner();
540            let controller = controller_owner.controller();
541            durable_session_state(controller, session_id) == Some(SessionState::Stopped)
542        };
543        if stopped {
544            self.wiki.request_sync(false);
545        }
546    }
547
548    pub fn allocate_revision(&self) -> u64 {
549        self.revisions.allocate()
550    }
551
552    pub(super) fn publish_revision(&self) -> u64 {
553        self.revisions.publish()
554    }
555
556    pub(super) fn attachments(&self) -> std::sync::MutexGuard<'_, BTreeMap<String, Attachment>> {
557        self.attachments
558            .lock()
559            .unwrap_or_else(PoisonError::into_inner)
560    }
561
562    pub(super) fn prune_dead_clients(&self) {
563        self.attachments()
564            .retain(|_, attachment| process_is_alive(attachment.pid));
565    }
566
567    pub(super) fn workspace_has_active_resume(&self, workspace_id: &str) -> bool {
568        self.owner().lifecycle.values().any(|active| {
569            active.is_running() && active.resume_workspace_id.as_deref() == Some(workspace_id)
570        })
571    }
572
573    pub fn publish_web_access(&self, access: crate::server::WebViewerAccess) {
574        use crate::server::WebViewerAccess;
575        let status = match &access {
576            WebViewerAccess::Starting => WebViewerStatus::Starting,
577            WebViewerAccess::Ready {
578                viewer_url,
579                viewer_code,
580                qr_login_url,
581                fallback_reason,
582                ..
583            } => {
584                // API clients wait for this after the daemon answers, so its
585                // time after startup is part of every command's.
586                tracing::info!("the web viewer and API are ready");
587                WebViewerStatus::Ready {
588                    viewer_url: viewer_url.clone(),
589                    viewer_code: viewer_code.clone(),
590                    qr_login_url: qr_login_url.clone(),
591                    fallback_reason: fallback_reason.clone(),
592                }
593            }
594            WebViewerAccess::Failed {
595                address, message, ..
596            } => WebViewerStatus::Error {
597                message: format!("{message} Address: {address}"),
598            },
599            WebViewerAccess::Unavailable(message) => WebViewerStatus::Error {
600                message: message.clone(),
601            },
602        };
603        self.web_viewer.publish(access);
604        self.set_phone_status(status);
605    }
606
607    pub(super) fn set_phone_status(&self, status: WebViewerStatus) {
608        *self
609            .phone_status
610            .lock()
611            .unwrap_or_else(PoisonError::into_inner) = status;
612    }
613
614    pub(super) fn phone_status(&self) -> WebViewerStatus {
615        self.phone_status
616            .lock()
617            .unwrap_or_else(PoisonError::into_inner)
618            .clone()
619    }
620
621    pub(super) fn workspaces(&self) -> tokio::sync::watch::Receiver<Vec<WorkspaceRecord>> {
622        self.workspaces_tx.subscribe()
623    }
624
625    #[cfg(test)]
626    pub(super) fn worker_poll_exclusion_session_ids(&self) -> BTreeSet<String> {
627        let owner = self.owner();
628        owner
629            .lifecycle
630            .keys()
631            .filter(|id| owner.worker_is_owned(id))
632            .cloned()
633            .collect()
634    }
635
636    pub fn revisions(&self) -> tokio::sync::watch::Receiver<u64> {
637        self.revisions.subscribe()
638    }
639
640    /// Read the config the daemon serves right now. A task on a schedule reads
641    /// it again on every tick, so a reload reaches it without a restart.
642    pub fn with_config<T>(&self, read: impl FnOnce(&Config) -> T) -> T {
643        read(&self.owner().controller().config)
644    }
645
646    /// Create a bundle under the daemon's config-mutation coordinator. The
647    /// controller helper also takes the cross-process config lock, so a TUI
648    /// transaction cannot race this one while the daemon's other config
649    /// writers are excluded by this mutex.
650    pub async fn create_quick_bundle(
651        &self,
652        source: String,
653    ) -> std::result::Result<
654        crate::controller::QuickBundleCreation,
655        crate::controller::QuickBundleFailure,
656    > {
657        let _mutation = self.config_mutation.lock().await;
658        tokio::task::spawn_blocking(move || crate::controller::create_quick_bundle(&source))
659            .await
660            .map_err(|error| {
661                crate::controller::QuickBundleFailure::Persistence(anyhow!(
662                    "bundle creation task panicked: {error}"
663                ))
664            })?
665    }
666
667    /// Persist exactly the picker-selected repository set under the same
668    /// config-mutation coordinator as legacy quick-bundle creation.
669    pub async fn create_bundle_from_sources(
670        &self,
671        sources: Vec<String>,
672    ) -> std::result::Result<
673        crate::controller::QuickBundleCreation,
674        crate::controller::QuickBundleFailure,
675    > {
676        let _mutation = self.config_mutation.lock().await;
677        tokio::task::spawn_blocking(move || crate::controller::create_bundle_from_sources(&sources))
678            .await
679            .map_err(|error| {
680                crate::controller::QuickBundleFailure::Persistence(anyhow!(
681                    "bundle creation task panicked: {error}"
682                ))
683            })?
684    }
685
686    /// Hand a changed workspace list to the terminal clients and the web viewer.
687    pub(crate) fn publish_workspaces(&self, workspaces: Vec<WorkspaceRecord>) {
688        self.workspaces_tx.send_replace(workspaces);
689        self.publish_revision();
690    }
691
692    pub(crate) async fn refresh_workspaces(&self) -> Result<()> {
693        // Keep the read and publication together: a delayed read must not
694        // publish an older list after a newer removal has been published.
695        let _refresh = self.workspace_refresh.lock().await;
696        let workspaces = tokio::task::spawn_blocking(crate::database::list_workspaces)
697            .await
698            .context("daemon workspace refresh task panicked")??;
699        self.publish_workspaces(workspaces);
700        Ok(())
701    }
702
703    /// Queue one piece of startup work for a session, starting the drain task
704    /// when this is the session's first step.
705    ///
706    /// Steps are carried out in the order they arrive. The entry in the map
707    /// exists only while a drain owns it, so "no entry" and "no live task"
708    /// are the same condition and a second call never starts a second drain.
709    pub(crate) async fn queue_startup_step(
710        self: &Arc<Self>,
711        session_id: &str,
712        step: StartupStep,
713        cancellation: &CancellationToken,
714    ) -> Result<()> {
715        self.queue_startup_steps_with_ids(
716            session_id,
717            vec![(new_command_id("startup")?, step)],
718            None,
719            cancellation,
720        )
721        .await
722    }
723
724    pub(crate) async fn queue_startup_steps_with_ids(
725        self: &Arc<Self>,
726        session_id: &str,
727        steps: Vec<(String, StartupStep)>,
728        group_id: Option<String>,
729        cancellation: &CancellationToken,
730    ) -> Result<()> {
731        let _admission = self.startup_enqueue.lock().await;
732        self.queue_startup_steps_admitted(session_id, steps, group_id, cancellation)
733            .await
734    }
735
736    /// Caller holds startup_enqueue across durable insertion and drain registration.
737    async fn queue_startup_steps_admitted(
738        self: &Arc<Self>,
739        session_id: &str,
740        steps: Vec<(String, StartupStep)>,
741        group_id: Option<String>,
742        cancellation: &CancellationToken,
743    ) -> Result<()> {
744        // The same set `wait_for_ready_session` accepts: anything else will
745        // never become ready, so queueing would only lose the text later.
746        match self.session_state(session_id) {
747            Some(
748                SessionState::Provisioning
749                | SessionState::Running
750                | SessionState::Disconnected
751                | SessionState::Checkpointing,
752            ) => {}
753            Some(state) => {
754                bail!("session {session_id} is {state:?}; it cannot take a queued prompt")
755            }
756            None => bail!("unknown session {session_id}"),
757        }
758        let deliveries = steps
759            .into_iter()
760            .map(|(command_id, step)| {
761                Ok(crate::database::StartupDelivery {
762                    session_id: session_id.to_owned(),
763                    command_id,
764                    step_json: serde_json::to_string(&step)?,
765                    phase: "pending".into(),
766                    group_id: group_id.clone(),
767                    accepted_ordinal: None,
768                    error: None,
769                })
770            })
771            .collect::<Result<Vec<_>>>()?;
772        let inserted =
773            blocking(move || crate::database::enqueue_startup_deliveries(deliveries)).await?;
774        for delivery in inserted {
775            self.start_persisted_startup_delivery(delivery, cancellation);
776        }
777        Ok(())
778    }
779
780    /// Take a queued prompt back before delivery starts. The caller becomes
781    /// its only owner. Once the drain has claimed the step, the prompt is on
782    /// its way to the session and this returns false.
783    pub(crate) async fn withdraw_startup_prompt(
784        &self,
785        session_id: &str,
786        text: &str,
787    ) -> Result<bool> {
788        let (id, prompt) = (session_id.to_owned(), text.to_owned());
789        blocking(move || crate::database::withdraw_startup_prompt(&id, &prompt)).await
790    }
791
792    pub(crate) async fn restore_startup_deliveries(
793        self: &Arc<Self>,
794        cancellation: &CancellationToken,
795    ) -> Result<()> {
796        let pruned = blocking(crate::database::prune_settled_startup_deliveries).await?;
797        if pruned > 0 {
798            tracing::info!(rows = pruned, "pruned settled startup steps");
799        }
800        let deliveries = blocking(crate::database::load_startup_deliveries).await?;
801        for delivery in deliveries {
802            self.start_persisted_startup_delivery(delivery, cancellation);
803        }
804        Ok(())
805    }
806
807    pub(super) fn start_persisted_startup_delivery(
808        self: &Arc<Self>,
809        delivery: crate::database::StartupDelivery,
810        cancellation: &CancellationToken,
811    ) {
812        let session_id = delivery.session_id.clone();
813        let mut queues = self
814            .startup_prompts
815            .lock()
816            .unwrap_or_else(PoisonError::into_inner);
817        if queues.contains_key(&session_id) {
818            return;
819        }
820        let cancel = cancellation.child_token();
821        let identity = Arc::new(());
822        queues.insert(
823            session_id.to_owned(),
824            StartupQueue {
825                identity: Arc::clone(&identity),
826                last_error: None,
827                cancel: cancel.clone(),
828                task: None,
829            },
830        );
831        let runtime = Arc::clone(self);
832        let drain_session = session_id.to_owned();
833        // One owned future: abort+join also settles the producer before a
834        // cancellation releases its durable command receipts.
835        let task = tokio::spawn(async move {
836            use futures::FutureExt;
837            let mut backoff = Duration::from_secs(1);
838            let mut failed_rounds: u32 = 0;
839            loop {
840                let result = std::panic::AssertUnwindSafe(
841                    Arc::clone(&runtime).drain_startup_queue(&drain_session, &cancel),
842                )
843                .catch_unwind()
844                .await;
845                if result.is_err() {
846                    runtime
847                        .fail_startup_queue(&drain_session, "the startup delivery task panicked")
848                        .await;
849                }
850                if cancel.is_cancelled() {
851                    runtime.retire_startup_drain(&drain_session, &identity);
852                    return;
853                }
854                if !runtime
855                    .startup_prompts
856                    .lock()
857                    .unwrap_or_else(PoisonError::into_inner)
858                    .get(&drain_session)
859                    .is_some_and(|queue| Arc::ptr_eq(&queue.identity, &identity))
860                {
861                    return;
862                }
863                // The queue is still registered, so this round ended in a
864                // failure. A step that keeps failing is abandoned rather than
865                // retried forever: a parent waiting on a child gets an answer.
866                failed_rounds += 1;
867                if failed_rounds >= STARTUP_STEP_ATTEMPTS {
868                    runtime.abandon_startup_step(&drain_session).await;
869                    failed_rounds = 0;
870                    backoff = Duration::from_secs(1);
871                }
872                tokio::select! {
873                    () = cancel.cancelled() => {
874                        runtime.retire_startup_drain(&drain_session, &identity);
875                        return;
876                    }
877                    () = tokio::time::sleep(backoff) => {}
878                }
879                backoff = (backoff * 2).min(Duration::from_secs(30));
880            }
881        });
882        if let Some(queue) = queues.get_mut(&session_id) {
883            queue.task = Some(task);
884        }
885    }
886
887    fn retire_startup_drain(&self, session_id: &str, identity: &Arc<()>) {
888        let mut queues = self
889            .startup_prompts
890            .lock()
891            .unwrap_or_else(PoisonError::into_inner);
892        if queues
893            .get(session_id)
894            .is_some_and(|queue| Arc::ptr_eq(&queue.identity, identity))
895        {
896            queues.remove(session_id);
897        }
898    }
899
900    /// Wait for the session's harness, then carry out its queued steps in
901    /// order. Failure leaves the durable queue intact; accepted commands use
902    /// retained worker receipts so restart can safely reconcile lost replies.
903    async fn drain_startup_queue(self: Arc<Self>, session_id: &str, cancel: &CancellationToken) {
904        loop {
905            // Admission orders durable insertion and the empty-queue decision.
906            // No pending payload list or notification can lose an earlier row.
907            let admission = tokio::select! {
908                () = cancel.cancelled() => return,
909                admission = self.startup_enqueue.lock() => admission,
910            };
911            let lookup_id = session_id.to_owned();
912            let step =
913                match blocking(move || crate::database::next_startup_delivery(&lookup_id)).await {
914                    Ok(Some(step)) => step,
915                    Ok(None) => {
916                        self.startup_prompts
917                            .lock()
918                            .unwrap_or_else(PoisonError::into_inner)
919                            .remove(session_id);
920                        return;
921                    }
922                    Err(error) => {
923                        drop(admission);
924                        self.fail_startup_queue(session_id, &format!("{error:#}"))
925                            .await;
926                        return;
927                    }
928                };
929            drop(admission);
930            // A step that is accepted, cancelling or rejecting only needs its
931            // final phase written, so it settles even after its session is gone.
932            let delivers = !matches!(step.phase.as_str(), "accepted" | "cancelling" | "rejecting");
933            let handle = if !delivers {
934                None
935            } else {
936                tokio::select! {
937                () = cancel.cancelled() => {
938                    self.fail_startup_queue(session_id, "the daemon stopped before the session was ready",
939                    )
940                    .await;
941                    return;
942                }
943                ready = self.wait_for_ready_session(session_id) => match ready {
944                    Ok(handle) => Some(handle),
945                    Err(error) => {
946                        let id = session_id.to_owned();
947                        let reason = format!("{error:#}");
948                        if let Err(persistence) = blocking(move || crate::database::fail_unavailable_startup_groups(&id, &reason)).await {
949                            tracing::error!(session_id, %persistence, "could not record unavailable startup session");
950                        }
951                        self.fail_startup_queue(session_id, &format!("{error:#}"))
952                            .await;
953                        return;
954                    }
955                },
956                }
957            };
958            let command_id = step.command_id.clone();
959            let mut decoded: StartupStep = match serde_json::from_str(&step.step_json) {
960                Ok(step) => step,
961                Err(error) => {
962                    self.fail_startup_queue(
963                        session_id,
964                        &format!("invalid durable startup step: {error}"),
965                    )
966                    .await;
967                    return;
968                }
969            };
970            if !matches!(step.phase.as_str(), "cancelling" | "rejecting")
971                && let StartupStep::InstallHandoff(snapshot) = &decoded
972            {
973                let prepared = tokio::select! {
974                    () = cancel.cancelled() => return,
975                    prepared = self.prepare_archive_handoff(session_id, snapshot) => prepared,
976                };
977                let text = match prepared {
978                    Ok(text) => text,
979                    Err(error) => {
980                        self.fail_startup_queue(session_id, &format!("{error:#}"))
981                            .await;
982                        return;
983                    }
984                };
985                decoded = StartupStep::PreparedHandoff { text };
986                let prepared_json = match serde_json::to_string(&decoded) {
987                    Ok(json) => json,
988                    Err(error) => {
989                        self.fail_startup_queue(session_id, &format!("{error:#}"))
990                            .await;
991                        return;
992                    }
993                };
994                let prepared_id = command_id.clone();
995                if let Err(error) = blocking(move || {
996                    crate::database::prepare_startup_delivery(&prepared_id, prepared_json)
997                })
998                .await
999                {
1000                    self.fail_startup_queue(session_id, &format!("{error:#}"))
1001                        .await;
1002                    return;
1003                }
1004            }
1005            if let Some(handle) = &handle {
1006                let persisted_id = command_id.clone();
1007                match blocking(move || crate::database::claim_startup_delivery(&persisted_id)).await
1008                {
1009                    Ok(true) => {}
1010                    // Withdrawn while this drain waited for the harness: the
1011                    // next read sees the step cancelling and settles it.
1012                    Ok(false) => continue,
1013                    Err(error) => {
1014                        self.fail_startup_queue(session_id, &format!("{error:#}"))
1015                            .await;
1016                        return;
1017                    }
1018                }
1019                let outcome = tokio::select! {
1020                    () = cancel.cancelled() => Err(anyhow!("the daemon stopped before startup delivery settled")),
1021                    result = self.run_startup_step(session_id, handle, &decoded, &command_id) => result,
1022                };
1023                let ordinal = match outcome {
1024                    Ok(ordinal) => ordinal,
1025                    Err(error) => {
1026                        if error.is::<super::startup_followup::StartupRejected>()
1027                            && let Some(group_id) = step.group_id.clone()
1028                        {
1029                            let id = session_id.to_owned();
1030                            let reason = format!("{error:#}");
1031                            let persisted = reason.clone();
1032                            match blocking(move || {
1033                                crate::database::fail_startup_group(&id, &group_id, &persisted)
1034                            })
1035                            .await
1036                            {
1037                                // The group failed as a whole, and was never
1038                                // accepted: say so, and settle its other rows.
1039                                Ok(()) => {
1040                                    self.reject_startup(session_id, &reason).await;
1041                                    continue;
1042                                }
1043                                Err(persistence) => {
1044                                    tracing::error!(session_id, %persistence, "could not persist startup rejection");
1045                                }
1046                            }
1047                        }
1048                        let refused = error
1049                            .downcast_ref::<mj_client::session::Refused>()
1050                            .is_some();
1051                        if refused
1052                            && step.group_id.is_none()
1053                            && let StartupStep::Prompt { text, .. } = &decoded
1054                        {
1055                            // The worker refused this prompt outright, so it was
1056                            // never accepted: give the text back rather than
1057                            // retrying something that will be refused again.
1058                            self.return_startup_prompt_to_draft(
1059                                session_id,
1060                                &command_id,
1061                                text,
1062                                &format!("{error:#}"),
1063                            )
1064                            .await;
1065                            continue;
1066                        }
1067                        self.fail_startup_queue(session_id, &format!("{error:#}"))
1068                            .await;
1069                        return;
1070                    }
1071                };
1072                let settled_id = command_id.clone();
1073                if let Err(error) = blocking(move || {
1074                    crate::database::set_startup_delivery_accepted(&settled_id, ordinal)
1075                })
1076                .await
1077                {
1078                    self.fail_startup_queue(
1079                        session_id,
1080                        &format!("delivery accepted but settlement failed: {error:#}"),
1081                    )
1082                    .await;
1083                    return;
1084                }
1085            }
1086            let settled_id = command_id;
1087            let final_phase = match step.phase.as_str() {
1088                "cancelling" => "dismissed",
1089                "rejecting" => "failed",
1090                _ => "done",
1091            };
1092            let final_error = step.error.clone();
1093            if let Err(error) = blocking(move || {
1094                crate::database::set_startup_delivery_phase(
1095                    &settled_id,
1096                    final_phase,
1097                    final_error.as_deref(),
1098                )
1099            })
1100            .await
1101            {
1102                self.fail_startup_queue(
1103                    session_id,
1104                    &format!("startup settlement failed: {error:#}"),
1105                )
1106                .await;
1107                return;
1108            }
1109        }
1110    }
1111
1112    async fn run_startup_step(
1113        &self,
1114        session_id: &str,
1115        handle: &crate::session_manager::ManagedSessionHandle,
1116        step: &StartupStep,
1117        command_id: &str,
1118    ) -> Result<Option<u64>> {
1119        match step {
1120            StartupStep::InstallHandoff(_) => bail!("archive startup step was not prepared"),
1121            StartupStep::PreparedHandoff { text } => handle
1122                .submit(
1123                    command_id.to_owned(),
1124                    RelayCommand::InstallPromptContext { text: text.clone() },
1125                )
1126                .await
1127                .map(Some),
1128            StartupStep::Prompt {
1129                text,
1130                inherited_draft,
1131            } => self
1132                .submit_startup_prompt(
1133                    session_id,
1134                    handle,
1135                    text,
1136                    inherited_draft.as_deref(),
1137                    command_id,
1138                )
1139                .await
1140                .map(Some),
1141            StartupStep::Configure {
1142                key,
1143                value,
1144                optional,
1145            } => {
1146                super::startup_followup::configure_startup(
1147                    &handle.client(),
1148                    command_id,
1149                    key,
1150                    value,
1151                    *optional,
1152                )
1153                .await
1154            }
1155            StartupStep::ApiPrompt { text } => {
1156                let ordinal = self
1157                    .submit_startup_prompt(session_id, handle, text, None, command_id)
1158                    .await?;
1159                let id = session_id.to_owned();
1160                blocking(move || crate::database::record_subagent_prompt(&id, ordinal)).await?;
1161                Ok(Some(ordinal))
1162            }
1163        }
1164    }
1165
1166    /// Submit one queued prompt and give it the history and draft handling a
1167    /// prompt submitted from a live composer gets.
1168    async fn submit_startup_prompt(
1169        &self,
1170        session_id: &str,
1171        handle: &crate::session_manager::ManagedSessionHandle,
1172        text: &str,
1173        inherited_draft: Option<&str>,
1174        command_id: &str,
1175    ) -> Result<u64> {
1176        let bundle_id = self
1177            .owner()
1178            .controller()
1179            .state
1180            .sessions
1181            .get(session_id)
1182            .map(|record| record.bundle_id.clone());
1183        let ordinal = handle
1184            .submit(
1185                command_id.to_owned(),
1186                RelayCommand::Prompt {
1187                    prompt: vec![ContentBlock::Text(TextContent::new(text.to_owned()))],
1188                },
1189            )
1190            .await?;
1191        if let Some(expected) = inherited_draft {
1192            let persisted_id = session_id.to_owned();
1193            let persisted_expected = expected.to_owned();
1194            if let Err(error) = blocking(move || {
1195                crate::database::clear_session_draft_input_if_matches(
1196                    &persisted_id,
1197                    &persisted_expected,
1198                )
1199            })
1200            .await
1201            {
1202                tracing::warn!(
1203                    session_id,
1204                    error = format!("{error:#}"),
1205                    "the delivered prompt's draft could not be cleared"
1206                );
1207            }
1208            self.publish_revision();
1209        }
1210        if let Some(bundle_id) = bundle_id {
1211            let history_id = session_id.to_owned();
1212            let history_text = text.to_owned();
1213            if let Err(error) = blocking(move || {
1214                crate::database::record_prompt(
1215                    &history_id,
1216                    &bundle_id,
1217                    ordinal,
1218                    None,
1219                    &history_text,
1220                )
1221            })
1222            .await
1223            {
1224                tracing::warn!(
1225                    session_id,
1226                    error = format!("{error:#}"),
1227                    "the queued prompt was accepted but its history could not be stored"
1228                );
1229            }
1230        }
1231        Ok(ordinal)
1232    }
1233
1234    /// A step that failed `STARTUP_STEP_ATTEMPTS` rounds in a row is given
1235    /// up: an API group fails so its client's wait answers, a user's prompt
1236    /// goes back to the draft, and anything else is marked failed. The drain
1237    /// then continues with the next step.
1238    async fn abandon_startup_step(self: &Arc<Self>, session_id: &str) {
1239        let lookup_id = session_id.to_owned();
1240        let step = match blocking(move || crate::database::next_startup_delivery(&lookup_id)).await
1241        {
1242            Ok(Some(step)) => step,
1243            Ok(None) => return,
1244            Err(error) => {
1245                tracing::error!(session_id, %error, "could not read the startup step to abandon");
1246                return;
1247            }
1248        };
1249        let reason = format!(
1250            "startup delivery gave up after {STARTUP_STEP_ATTEMPTS} attempts: {}",
1251            self.startup_prompts
1252                .lock()
1253                .unwrap_or_else(PoisonError::into_inner)
1254                .get(session_id)
1255                .and_then(|queue| queue.last_error.clone())
1256                .unwrap_or_else(|| "delivery kept failing".to_owned())
1257        );
1258        if let Some(group_id) = step.group_id.clone() {
1259            let id = session_id.to_owned();
1260            let persisted = reason.clone();
1261            if let Err(error) =
1262                blocking(move || crate::database::fail_startup_group(&id, &group_id, &persisted))
1263                    .await
1264            {
1265                tracing::error!(session_id, %error, "could not fail an abandoned startup group");
1266            }
1267            self.reject_startup(session_id, &reason).await;
1268            return;
1269        }
1270        let decoded: Option<StartupStep> = serde_json::from_str(&step.step_json).ok();
1271        if let Some(StartupStep::Prompt { text, .. }) = &decoded {
1272            self.return_startup_prompt_to_draft(session_id, &step.command_id, text, &reason)
1273                .await;
1274            return;
1275        }
1276        let command_id = step.command_id.clone();
1277        let persisted = reason.clone();
1278        if let Err(error) = blocking(move || {
1279            crate::database::set_startup_delivery_phase(&command_id, "failed", Some(&persisted))
1280        })
1281        .await
1282        {
1283            tracing::error!(session_id, %error, "could not mark an abandoned startup step failed");
1284        }
1285        self.push_notice(session_id, reason);
1286    }
1287
1288    /// An API startup group failed for good: its step was refused, or kept
1289    /// failing until it was given up. Nothing of it was accepted, so the
1290    /// notice says the startup failed, not that work remains saved. A
1291    /// sub-agent whose first prompt can therefore never run is recorded as
1292    /// failed and its worker stopped (I1-2); its parent reads the same cause
1293    /// from `wait` and `list_agents`.
1294    async fn reject_startup(self: &Arc<Self>, session_id: &str, reason: &str) {
1295        tracing::warn!(session_id, reason, "session startup failed");
1296        self.push_notice(session_id, format!("Session startup failed: {reason}"));
1297        self.fail_subagent_start(session_id, reason).await;
1298    }
1299
1300    /// A prompt the worker will not take is the user's text again, not a
1301    /// row that retries forever.
1302    async fn return_startup_prompt_to_draft(
1303        &self,
1304        session_id: &str,
1305        command_id: &str,
1306        text: &str,
1307        reason: &str,
1308    ) {
1309        let settled_id = command_id.to_owned();
1310        let persisted = reason.to_owned();
1311        if let Err(error) = blocking(move || {
1312            crate::database::set_startup_delivery_phase(&settled_id, "failed", Some(&persisted))
1313        })
1314        .await
1315        {
1316            tracing::error!(session_id, %error, "could not mark a refused startup prompt failed");
1317            return;
1318        }
1319        if let Err(error) = self.append_draft_input(session_id, text).await {
1320            tracing::error!(session_id, %error, "could not return a refused startup prompt to the draft");
1321            self.push_notice(
1322                session_id,
1323                format!("A queued prompt was refused and could not be saved as a draft: {reason}"),
1324            );
1325            return;
1326        }
1327        self.push_notice(
1328            session_id,
1329            format!("A queued prompt was refused and returned to your draft: {reason}"),
1330        );
1331    }
1332
1333    /// Stop this drain without turning uncertain accepted input into a new draft.
1334    async fn fail_startup_queue(&self, session_id: &str, reason: &str) {
1335        let changed = {
1336            let mut queues = self
1337                .startup_prompts
1338                .lock()
1339                .unwrap_or_else(PoisonError::into_inner);
1340            let Some(queue) = queues.get_mut(session_id) else {
1341                return;
1342            };
1343            if queue.last_error.as_deref() == Some(reason) {
1344                false
1345            } else {
1346                queue.last_error = Some(reason.to_owned());
1347                true
1348            }
1349        };
1350        if !changed {
1351            return;
1352        }
1353        self.push_notice(session_id, format!("Startup work remains saved for this session: {reason}. It has not been restored as an unsent draft because delivery may have been accepted."));
1354        tracing::warn!(session_id, reason, "durable startup delivery paused");
1355    }
1356
1357    /// Restore input through the writer; the owner observes only committed data.
1358    pub(super) async fn append_draft_input(&self, session_id: &str, text: &str) -> Result<()> {
1359        let session_id = session_id.to_owned();
1360        let text = text.to_owned();
1361        blocking(move || crate::database::append_session_draft_input(&session_id, &text)).await?;
1362        self.publish_revision();
1363        Ok(())
1364    }
1365
1366    /// Stop every startup queue and wait for its drain to report, so the text
1367    /// it holds reaches the database while the writer is still running.
1368    pub(crate) async fn cancel_and_join_startup_prompts(&self) -> Result<()> {
1369        let tasks = {
1370            let mut queues = self
1371                .startup_prompts
1372                .lock()
1373                .unwrap_or_else(PoisonError::into_inner);
1374            queues
1375                .values_mut()
1376                .filter_map(|queue| {
1377                    queue.cancel.cancel();
1378                    queue.task.take()
1379                })
1380                .collect::<Vec<_>>()
1381        };
1382        let deadline = tokio::time::Instant::now() + Duration::from_secs(1);
1383        let mut outcome = Ok(());
1384        for mut task in tasks {
1385            let joined = match tokio::time::timeout_at(deadline, &mut task).await {
1386                Ok(Ok(())) => Ok(()),
1387                Ok(Err(error)) => Err(anyhow!("startup prompt delivery task failed: {error}")),
1388                Err(_) => {
1389                    task.abort();
1390                    let _ = task.await;
1391                    Err(anyhow!(
1392                        "a startup prompt delivery task did not stop within 1s; durable work retained"
1393                    ))
1394                }
1395            };
1396            if outcome.is_ok() {
1397                outcome = joined;
1398            } else if let Err(error) = joined {
1399                tracing::warn!(%error, "another startup prompt drain did not stop cleanly");
1400            }
1401        }
1402        outcome
1403    }
1404}
1405
1406#[cfg(test)]
1407mod tests {
1408    use super::super::tests::test_runtime_state;
1409    use std::sync::Arc;
1410
1411    /// RCL-1 (2026-09-29): the Sessions filter's conversation search asks for
1412    /// the same non-forced sync the resume dialog's search asks for, so a
1413    /// message sent since the last sync is found by the next keystroke.
1414    // Hard-won: 13425c7c: session search missed messages newer than SessionWiki until it requested a sync.
1415    #[tokio::test]
1416    async fn a_session_text_search_asks_for_a_sync_like_the_resume_search() {
1417        let mut state = test_runtime_state();
1418        Arc::get_mut(&mut state).expect("the only handle").wiki =
1419            crate::sessionwiki::WikiIndexer::inert();
1420        assert!(!state.wiki().sync_requested());
1421
1422        // An empty query never reads the index, so the request is all it does.
1423        state.session_text_search("  ".into()).await.unwrap();
1424        assert!(
1425            state.wiki().sync_requested(),
1426            "a stale index is synced for the next search"
1427        );
1428    }
1429}