1const 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 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 let review_config = Arc::new(Mutex::new(controller.config.review.clone()));
51 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 pub fn review_host(&self) -> &TurnReviewHost {
132 &self.review_host
133 }
134
135 pub fn wiki(&self) -> &crate::sessionwiki::WikiIndexer {
137 &self.wiki
138 }
139
140 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 let rows =
148 blocking(move || crate::sessionwiki::query_rows(&query, limit, &live, false)).await?;
149 Ok(WikiSearchPage {
152 rows,
153 status: self.wiki.status(),
154 })
155 }
156
157 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 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 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 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 pub async fn wiki_session(
205 &self,
206 wiki_id: String,
207 ) -> Result<Option<mj_client::daemon::WikiSessionInfo>> {
208 let known = self.live_session_ids();
211 blocking(move || crate::sessionwiki::wiki_session(&wiki_id, &known)).await
212 }
213
214 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 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 session_title_override: Some(archived.title.clone()),
272 })
273 .await?;
274 let session_id = registered.session.id.clone();
275 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 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 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 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 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 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 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 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 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 pub fn with_config<T>(&self, read: impl FnOnce(&Config) -> T) -> T {
641 read(&self.owner().controller().config)
642 }
643
644 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 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 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 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 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 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 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 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 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 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 async fn drain_startup_queue(self: Arc<Self>, session_id: &str, cancel: &CancellationToken) {
902 loop {
903 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 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 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 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 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 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 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 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 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 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 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 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 #[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 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}