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