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