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 wait_prompt_backend: std::sync::OnceLock::new(),
91 wait_prompt_retries: Mutex::new(BTreeMap::new()),
92 credential_targets: Arc::new(tokio::sync::watch::channel(Vec::new()).0),
93 feed: Mutex::new(feed::RuntimeHistory::default()),
94 committed: None,
95 workspace_closes: Mutex::new(BTreeMap::new()),
96 workspace_resume_admission: Mutex::new(BTreeMap::new()),
97 harness_readiness: Mutex::new(HarnessReadinessWatch::default()),
98 startup_prompts: Mutex::new(BTreeMap::new()),
99 startup_enqueue: tokio::sync::Mutex::new(()),
100 controller_loader,
101 config_mutation: tokio::sync::Mutex::new(()),
102 projects: Arc::new(crate::project_catalog::Catalog::default()),
103 profile_catalog,
104 recovery_observer,
105 worker_upgrade_observer,
106 notices: Mutex::new(VecDeque::new()),
107 next_notice_id: AtomicU64::new(1),
108 quota: Mutex::new(QuotaBoard::default()),
109 capacity: std::sync::OnceLock::new(),
110 review_config,
111 review_host,
112 wiki: crate::sessionwiki::WikiIndexer::spawn(),
113 }
114 }
115
116 pub(crate) fn projects(&self) -> Arc<crate::project_catalog::Catalog> {
117 self.projects.clone()
118 }
119
120 pub async fn project_catalog(
121 &self,
122 refresh: bool,
123 retry: bool,
124 ) -> Result<mj_core::project_catalog::ProjectCatalogView> {
125 if refresh {
126 self.projects.request(retry);
127 }
128 let catalog = self.projects.clone();
129 blocking(move || catalog.view()).await
130 }
131
132 pub fn review_host(&self) -> &TurnReviewHost {
134 &self.review_host
135 }
136
137 pub fn wiki(&self) -> &crate::sessionwiki::WikiIndexer {
139 &self.wiki
140 }
141
142 pub async fn wiki_search(&self, query: String, limit: usize) -> Result<WikiSearchPage> {
145 self.request_wiki_sync_if_stale();
146 let live = self.live_session_ids();
147 let rows =
150 blocking(move || crate::sessionwiki::query_rows(&query, limit, &live, false)).await?;
151 Ok(WikiSearchPage {
154 rows,
155 status: self.wiki.status(),
156 })
157 }
158
159 fn request_wiki_sync_if_stale(&self) {
164 if crate::sessionwiki::sync_is_stale(self.wiki.last_success()) {
165 self.wiki.request_sync(false);
166 }
167 }
168
169 pub async fn session_text_search(&self, query: String) -> Result<Vec<SessionTextMatch>> {
173 self.request_wiki_sync_if_stale();
174 let live = self.live_session_ids();
175 blocking(move || crate::sessionwiki::session_text_matches(&query, &live)).await
176 }
177
178 pub async fn wiki_brief(&self, wiki_id: String, max_chars: usize) -> Result<Option<String>> {
181 blocking(move || crate::sessionwiki::brief(&wiki_id, max_chars)).await
182 }
183
184 pub async fn wiki_hits(
187 &self,
188 wiki_id: String,
189 query: String,
190 context_messages: usize,
191 per_message_chars: usize,
192 ) -> Result<Option<WikiHitTranscript>> {
193 blocking(move || {
194 crate::sessionwiki::transcript_hits(
195 &wiki_id,
196 &query,
197 context_messages,
198 per_message_chars,
199 )
200 })
201 .await
202 }
203
204 pub async fn wiki_session(
207 &self,
208 wiki_id: String,
209 ) -> Result<Option<mj_client::daemon::WikiSessionInfo>> {
210 let known = self.live_session_ids();
213 blocking(move || crate::sessionwiki::wiki_session(&wiki_id, &known)).await
214 }
215
216 pub async fn restore_wiki_session(
224 self: &Arc<Self>,
225 request: WikiRestoreRequest,
226 cancellation: &CancellationToken,
227 ) -> Result<Option<RegisteredSession>> {
228 let wiki_id = request.wiki_id.clone();
229 let Some(archived) =
230 blocking(move || crate::sessionwiki::archived_session(&wiki_id)).await?
231 else {
232 return Ok(None);
233 };
234 let project_directory = request
235 .project_directory
236 .clone()
237 .or_else(|| archived.project_directory.clone())
238 .context(
239 "name a project directory: the archived session's own project is no longer on this machine",
240 )?;
241 let source = project_directory.display().to_string();
242 let bundle_id = blocking(move || {
243 crate::controller::create_bundle_from_sources(&[source])
244 .map(|created| created.bundle_id)
245 .map_err(anyhow::Error::new)
246 })
247 .await
248 .context("find or create a bundle for the restored session's project")?;
249 let _startup_admission = self.startup_enqueue.lock().await;
252 let registered = self
253 .start_create_session(CreateSessionRequest {
254 at: None,
255 branch: None,
256 base: None,
257 create_managed_worktree: None,
258 subagents: None,
259 review: None,
260 initial_prompt: None,
261 workspace_id: request.workspace_id,
262 profile_id: request.profile_id,
263 bundle_id,
264 project_directory: Some(project_directory),
265 target_template_id: request.target_template_id,
266 additional_mounts: request.additional_mounts,
267 resource_allocation: request.resource_allocation,
268 title: archived.title.clone(),
269 session_title_override: Some(archived.title.clone()),
274 })
275 .await?;
276 let session_id = registered.session.id.clone();
277 self.queue_startup_steps_admitted(
281 &session_id,
282 vec![(
283 new_command_id("startup")?,
284 StartupStep::InstallHandoff(Box::new(archived.snapshot)),
285 )],
286 None,
287 cancellation,
288 )
289 .await?;
290 Ok(Some(registered))
291 }
292
293 pub(super) fn live_session_ids(&self) -> BTreeSet<String> {
294 self.owner()
295 .controller()
296 .state
297 .sessions
298 .keys()
299 .cloned()
300 .collect()
301 }
302
303 async fn prepare_archive_handoff(
307 &self,
308 session_id: &str,
309 snapshot: &mj_core::archive::CanonicalSessionSnapshot,
310 ) -> Result<String> {
311 let (config, profile_id) = {
312 let controller_owner = self.owner();
313 let controller = controller_owner.controller();
314 let profile_id = controller
315 .state
316 .sessions
317 .get(session_id)
318 .map(|record| record.last_profile.clone());
319 (controller.config.clone(), profile_id)
320 };
321 let context_bytes = crate::handoff::profile_handoff_bytes(
322 profile_id.and_then(|id| config.profiles.get(&id)),
323 );
324 let cancel = CancellationToken::new();
325 let handoff = crate::handoff::build_handoff_context(
326 session_id,
327 &config,
328 snapshot,
329 context_bytes,
330 &cancel,
331 )
332 .await
333 .context("compact the archived transcript")?;
334 Ok(format!(
335 "{} {handoff}",
336 crate::compaction::ARCHIVE_HANDOFF_PREAMBLE
337 ))
338 }
339
340 pub(super) async fn wait_for_ready_session(
342 &self,
343 session_id: &str,
344 ) -> Result<crate::session_manager::ManagedSessionHandle> {
345 const POLL: Duration = Duration::from_millis(250);
346 let deadline = tokio::time::Instant::now() + Duration::from_secs(30 * 60);
347 loop {
348 match self.session_state(session_id) {
352 Some(
353 SessionState::Provisioning
354 | SessionState::Running
355 | SessionState::Disconnected
356 | SessionState::Checkpointing,
357 ) => {}
358 Some(state) => bail!("session {session_id} is {state:?} before its hand-off"),
359 None => bail!("session {session_id} disappeared before its hand-off"),
360 }
361 if let Ok(handle) = self.session_manager.session(session_id).await {
362 let view = handle.view();
363 if let Some(ViewError::TargetMissing(detail)) = &view.error {
366 bail!("session {session_id} lost its target: {detail}");
367 }
368 if view.connected
369 && view
370 .snapshot
371 .is_some_and(|snapshot| snapshot.operational.native_session_is_ready())
372 {
373 return Ok(handle);
374 }
375 }
376 ensure!(
377 tokio::time::Instant::now() < deadline,
378 "session {session_id} was not ready for its hand-off within 30 minutes"
379 );
380 tokio::time::sleep(POLL).await;
381 }
382 }
383
384 pub(super) async fn fail_unfinished_provisioning(
399 self: &Arc<Self>,
400 session_id: &str,
401 error: &str,
402 ) {
403 let provisioning = {
404 let controller_owner = self.owner();
405 let controller = controller_owner.controller();
406 durable_session_state(controller, session_id) == Some(SessionState::Provisioning)
407 };
408 if !provisioning {
409 return;
410 }
411 let cause = format!("session provisioning ended without finishing: {error}");
412 let applied = blocking({
413 let session_id = session_id.to_owned();
414 move || {
415 let mut controller = Controller::load()?;
416 controller.fail_interrupted_lifecycle(&session_id, &cause)
417 }
418 })
419 .await;
420 match applied {
421 Ok(true) => {
422 if let Err(error) = self.reload_controller().await {
423 tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after recording a failed provision");
424 }
425 }
426 Ok(false) => {}
427 Err(error) => tracing::warn!(
428 %session_id,
429 error = format!("{error:#}"),
430 "could not record that provisioning ended without finishing"
431 ),
432 }
433 }
434
435 pub(super) async fn record_failed_close(
443 self: &Arc<Self>,
444 session_id: &str,
445 reference: &str,
446 failure: &LifecycleFailure,
447 ) {
448 self.record_lifecycle_failure(
449 session_id,
450 reference,
451 failure,
452 mj_core::state::CLOSE_FAILURE_PREFIX,
453 )
454 .await;
455 }
456
457 pub(crate) async fn record_lifecycle_failure(
458 self: &Arc<Self>,
459 session_id: &str,
460 reference: &str,
461 failure: &LifecycleFailure,
462 prefix: &str,
463 ) {
464 let cause = match &failure.refusal {
465 Some(refusal) => format!("{prefix}: {refusal}"),
466 None => {
467 format!("{prefix}; the daemon log records the reason under reference {reference}")
468 }
469 };
470 let applied = blocking({
471 let session_id = session_id.to_owned();
472 let cause = cause.clone();
473 move || {
474 let mut controller = Controller::load()?;
475 controller.record_failed_close(&session_id, &cause)
476 }
477 })
478 .await;
479 match applied {
480 Ok(true) => {
481 if let Err(error) = self.reload_controller().await {
482 tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after recording a failed close");
483 }
484 self.publish_revision();
485 }
486 Ok(false) => {}
487 Err(error) => tracing::warn!(
488 %session_id,
489 error = format!("{error:#}"),
490 "could not record why a close failed"
491 ),
492 }
493 }
494
495 pub async fn clear_recorded_close_failure(self: &Arc<Self>, session_id: &str) {
503 let recorded = self
504 .owner()
505 .controller()
506 .state
507 .sessions
508 .get(session_id)
509 .is_some_and(|record| record.public_error().is_some());
510 if !recorded {
511 return;
512 }
513 let cleared = blocking({
514 let session_id = session_id.to_owned();
515 move || {
516 let mut controller = Controller::load()?;
517 controller.clear_recorded_close_failure(&session_id)
518 }
519 })
520 .await;
521 match cleared {
522 Ok(true) => {
523 if let Err(error) = self.reload_controller().await {
524 tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after clearing a recorded close failure");
525 }
526 self.publish_revision();
527 }
528 Ok(false) => {}
529 Err(error) => tracing::warn!(
530 %session_id,
531 error = format!("{error:#}"),
532 "could not clear a recorded close failure"
533 ),
534 }
535 }
536
537 pub(super) fn note_lifecycle_outcome(&self, session_id: &str) {
538 let stopped = {
539 let controller_owner = self.owner();
540 let controller = controller_owner.controller();
541 durable_session_state(controller, session_id) == Some(SessionState::Stopped)
542 };
543 if stopped {
544 self.wiki.request_sync(false);
545 }
546 }
547
548 pub fn allocate_revision(&self) -> u64 {
549 self.revisions.allocate()
550 }
551
552 pub(super) fn publish_revision(&self) -> u64 {
553 self.revisions.publish()
554 }
555
556 pub(super) fn attachments(&self) -> std::sync::MutexGuard<'_, BTreeMap<String, Attachment>> {
557 self.attachments
558 .lock()
559 .unwrap_or_else(PoisonError::into_inner)
560 }
561
562 pub(super) fn prune_dead_clients(&self) {
563 self.attachments()
564 .retain(|_, attachment| process_is_alive(attachment.pid));
565 }
566
567 pub(super) fn workspace_has_active_resume(&self, workspace_id: &str) -> bool {
568 self.owner().lifecycle.values().any(|active| {
569 active.is_running() && active.resume_workspace_id.as_deref() == Some(workspace_id)
570 })
571 }
572
573 pub fn publish_web_access(&self, access: crate::server::WebViewerAccess) {
574 use crate::server::WebViewerAccess;
575 let status = match &access {
576 WebViewerAccess::Starting => WebViewerStatus::Starting,
577 WebViewerAccess::Ready {
578 viewer_url,
579 viewer_code,
580 qr_login_url,
581 fallback_reason,
582 ..
583 } => {
584 tracing::info!("the web viewer and API are ready");
587 WebViewerStatus::Ready {
588 viewer_url: viewer_url.clone(),
589 viewer_code: viewer_code.clone(),
590 qr_login_url: qr_login_url.clone(),
591 fallback_reason: fallback_reason.clone(),
592 }
593 }
594 WebViewerAccess::Failed {
595 address, message, ..
596 } => WebViewerStatus::Error {
597 message: format!("{message} Address: {address}"),
598 },
599 WebViewerAccess::Unavailable(message) => WebViewerStatus::Error {
600 message: message.clone(),
601 },
602 };
603 self.web_viewer.publish(access);
604 self.set_phone_status(status);
605 }
606
607 pub(super) fn set_phone_status(&self, status: WebViewerStatus) {
608 *self
609 .phone_status
610 .lock()
611 .unwrap_or_else(PoisonError::into_inner) = status;
612 }
613
614 pub(super) fn phone_status(&self) -> WebViewerStatus {
615 self.phone_status
616 .lock()
617 .unwrap_or_else(PoisonError::into_inner)
618 .clone()
619 }
620
621 pub(super) fn workspaces(&self) -> tokio::sync::watch::Receiver<Vec<WorkspaceRecord>> {
622 self.workspaces_tx.subscribe()
623 }
624
625 #[cfg(test)]
626 pub(super) fn worker_poll_exclusion_session_ids(&self) -> BTreeSet<String> {
627 let owner = self.owner();
628 owner
629 .lifecycle
630 .keys()
631 .filter(|id| owner.worker_is_owned(id))
632 .cloned()
633 .collect()
634 }
635
636 pub fn revisions(&self) -> tokio::sync::watch::Receiver<u64> {
637 self.revisions.subscribe()
638 }
639
640 pub fn with_config<T>(&self, read: impl FnOnce(&Config) -> T) -> T {
643 read(&self.owner().controller().config)
644 }
645
646 pub async fn create_quick_bundle(
651 &self,
652 source: String,
653 ) -> std::result::Result<
654 crate::controller::QuickBundleCreation,
655 crate::controller::QuickBundleFailure,
656 > {
657 let _mutation = self.config_mutation.lock().await;
658 tokio::task::spawn_blocking(move || crate::controller::create_quick_bundle(&source))
659 .await
660 .map_err(|error| {
661 crate::controller::QuickBundleFailure::Persistence(anyhow!(
662 "bundle creation task panicked: {error}"
663 ))
664 })?
665 }
666
667 pub async fn create_bundle_from_sources(
670 &self,
671 sources: Vec<String>,
672 ) -> std::result::Result<
673 crate::controller::QuickBundleCreation,
674 crate::controller::QuickBundleFailure,
675 > {
676 let _mutation = self.config_mutation.lock().await;
677 tokio::task::spawn_blocking(move || crate::controller::create_bundle_from_sources(&sources))
678 .await
679 .map_err(|error| {
680 crate::controller::QuickBundleFailure::Persistence(anyhow!(
681 "bundle creation task panicked: {error}"
682 ))
683 })?
684 }
685
686 pub(crate) fn publish_workspaces(&self, workspaces: Vec<WorkspaceRecord>) {
688 self.workspaces_tx.send_replace(workspaces);
689 self.publish_revision();
690 }
691
692 pub(crate) async fn refresh_workspaces(&self) -> Result<()> {
693 let _refresh = self.workspace_refresh.lock().await;
696 let workspaces = tokio::task::spawn_blocking(crate::database::list_workspaces)
697 .await
698 .context("daemon workspace refresh task panicked")??;
699 self.publish_workspaces(workspaces);
700 Ok(())
701 }
702
703 pub(crate) async fn queue_startup_step(
710 self: &Arc<Self>,
711 session_id: &str,
712 step: StartupStep,
713 cancellation: &CancellationToken,
714 ) -> Result<()> {
715 self.queue_startup_steps_with_ids(
716 session_id,
717 vec![(new_command_id("startup")?, step)],
718 None,
719 cancellation,
720 )
721 .await
722 }
723
724 pub(crate) async fn queue_startup_steps_with_ids(
725 self: &Arc<Self>,
726 session_id: &str,
727 steps: Vec<(String, StartupStep)>,
728 group_id: Option<String>,
729 cancellation: &CancellationToken,
730 ) -> Result<()> {
731 let _admission = self.startup_enqueue.lock().await;
732 self.queue_startup_steps_admitted(session_id, steps, group_id, cancellation)
733 .await
734 }
735
736 async fn queue_startup_steps_admitted(
738 self: &Arc<Self>,
739 session_id: &str,
740 steps: Vec<(String, StartupStep)>,
741 group_id: Option<String>,
742 cancellation: &CancellationToken,
743 ) -> Result<()> {
744 match self.session_state(session_id) {
747 Some(
748 SessionState::Provisioning
749 | SessionState::Running
750 | SessionState::Disconnected
751 | SessionState::Checkpointing,
752 ) => {}
753 Some(state) => {
754 bail!("session {session_id} is {state:?}; it cannot take a queued prompt")
755 }
756 None => bail!("unknown session {session_id}"),
757 }
758 let deliveries = steps
759 .into_iter()
760 .map(|(command_id, step)| {
761 Ok(crate::database::StartupDelivery {
762 session_id: session_id.to_owned(),
763 command_id,
764 step_json: serde_json::to_string(&step)?,
765 phase: "pending".into(),
766 group_id: group_id.clone(),
767 accepted_ordinal: None,
768 error: None,
769 })
770 })
771 .collect::<Result<Vec<_>>>()?;
772 let inserted =
773 blocking(move || crate::database::enqueue_startup_deliveries(deliveries)).await?;
774 for delivery in inserted {
775 self.start_persisted_startup_delivery(delivery, cancellation);
776 }
777 Ok(())
778 }
779
780 pub(crate) async fn withdraw_startup_prompt(
784 &self,
785 session_id: &str,
786 text: &str,
787 ) -> Result<bool> {
788 let (id, prompt) = (session_id.to_owned(), text.to_owned());
789 blocking(move || crate::database::withdraw_startup_prompt(&id, &prompt)).await
790 }
791
792 pub(crate) async fn restore_startup_deliveries(
793 self: &Arc<Self>,
794 cancellation: &CancellationToken,
795 ) -> Result<()> {
796 let pruned = blocking(crate::database::prune_settled_startup_deliveries).await?;
797 if pruned > 0 {
798 tracing::info!(rows = pruned, "pruned settled startup steps");
799 }
800 let deliveries = blocking(crate::database::load_startup_deliveries).await?;
801 for delivery in deliveries {
802 self.start_persisted_startup_delivery(delivery, cancellation);
803 }
804 Ok(())
805 }
806
807 pub(super) fn start_persisted_startup_delivery(
808 self: &Arc<Self>,
809 delivery: crate::database::StartupDelivery,
810 cancellation: &CancellationToken,
811 ) {
812 let session_id = delivery.session_id.clone();
813 let mut queues = self
814 .startup_prompts
815 .lock()
816 .unwrap_or_else(PoisonError::into_inner);
817 if queues.contains_key(&session_id) {
818 return;
819 }
820 let cancel = cancellation.child_token();
821 let identity = Arc::new(());
822 queues.insert(
823 session_id.to_owned(),
824 StartupQueue {
825 identity: Arc::clone(&identity),
826 last_error: None,
827 cancel: cancel.clone(),
828 task: None,
829 },
830 );
831 let runtime = Arc::clone(self);
832 let drain_session = session_id.to_owned();
833 let task = tokio::spawn(async move {
836 use futures::FutureExt;
837 let mut backoff = Duration::from_secs(1);
838 let mut failed_rounds: u32 = 0;
839 loop {
840 let result = std::panic::AssertUnwindSafe(
841 Arc::clone(&runtime).drain_startup_queue(&drain_session, &cancel),
842 )
843 .catch_unwind()
844 .await;
845 if result.is_err() {
846 runtime
847 .fail_startup_queue(&drain_session, "the startup delivery task panicked")
848 .await;
849 }
850 if cancel.is_cancelled() {
851 runtime.retire_startup_drain(&drain_session, &identity);
852 return;
853 }
854 if !runtime
855 .startup_prompts
856 .lock()
857 .unwrap_or_else(PoisonError::into_inner)
858 .get(&drain_session)
859 .is_some_and(|queue| Arc::ptr_eq(&queue.identity, &identity))
860 {
861 return;
862 }
863 failed_rounds += 1;
867 if failed_rounds >= STARTUP_STEP_ATTEMPTS {
868 runtime.abandon_startup_step(&drain_session).await;
869 failed_rounds = 0;
870 backoff = Duration::from_secs(1);
871 }
872 tokio::select! {
873 () = cancel.cancelled() => {
874 runtime.retire_startup_drain(&drain_session, &identity);
875 return;
876 }
877 () = tokio::time::sleep(backoff) => {}
878 }
879 backoff = (backoff * 2).min(Duration::from_secs(30));
880 }
881 });
882 if let Some(queue) = queues.get_mut(&session_id) {
883 queue.task = Some(task);
884 }
885 }
886
887 fn retire_startup_drain(&self, session_id: &str, identity: &Arc<()>) {
888 let mut queues = self
889 .startup_prompts
890 .lock()
891 .unwrap_or_else(PoisonError::into_inner);
892 if queues
893 .get(session_id)
894 .is_some_and(|queue| Arc::ptr_eq(&queue.identity, identity))
895 {
896 queues.remove(session_id);
897 }
898 }
899
900 async fn drain_startup_queue(self: Arc<Self>, session_id: &str, cancel: &CancellationToken) {
904 loop {
905 let admission = tokio::select! {
908 () = cancel.cancelled() => return,
909 admission = self.startup_enqueue.lock() => admission,
910 };
911 let lookup_id = session_id.to_owned();
912 let step =
913 match blocking(move || crate::database::next_startup_delivery(&lookup_id)).await {
914 Ok(Some(step)) => step,
915 Ok(None) => {
916 self.startup_prompts
917 .lock()
918 .unwrap_or_else(PoisonError::into_inner)
919 .remove(session_id);
920 return;
921 }
922 Err(error) => {
923 drop(admission);
924 self.fail_startup_queue(session_id, &format!("{error:#}"))
925 .await;
926 return;
927 }
928 };
929 drop(admission);
930 let delivers = !matches!(step.phase.as_str(), "accepted" | "cancelling" | "rejecting");
933 let handle = if !delivers {
934 None
935 } else {
936 tokio::select! {
937 () = cancel.cancelled() => {
938 self.fail_startup_queue(session_id, "the daemon stopped before the session was ready",
939 )
940 .await;
941 return;
942 }
943 ready = self.wait_for_ready_session(session_id) => match ready {
944 Ok(handle) => Some(handle),
945 Err(error) => {
946 let id = session_id.to_owned();
947 let reason = format!("{error:#}");
948 if let Err(persistence) = blocking(move || crate::database::fail_unavailable_startup_groups(&id, &reason)).await {
949 tracing::error!(session_id, %persistence, "could not record unavailable startup session");
950 }
951 self.fail_startup_queue(session_id, &format!("{error:#}"))
952 .await;
953 return;
954 }
955 },
956 }
957 };
958 let command_id = step.command_id.clone();
959 let mut decoded: StartupStep = match serde_json::from_str(&step.step_json) {
960 Ok(step) => step,
961 Err(error) => {
962 self.fail_startup_queue(
963 session_id,
964 &format!("invalid durable startup step: {error}"),
965 )
966 .await;
967 return;
968 }
969 };
970 if !matches!(step.phase.as_str(), "cancelling" | "rejecting")
971 && let StartupStep::InstallHandoff(snapshot) = &decoded
972 {
973 let prepared = tokio::select! {
974 () = cancel.cancelled() => return,
975 prepared = self.prepare_archive_handoff(session_id, snapshot) => prepared,
976 };
977 let text = match prepared {
978 Ok(text) => text,
979 Err(error) => {
980 self.fail_startup_queue(session_id, &format!("{error:#}"))
981 .await;
982 return;
983 }
984 };
985 decoded = StartupStep::PreparedHandoff { text };
986 let prepared_json = match serde_json::to_string(&decoded) {
987 Ok(json) => json,
988 Err(error) => {
989 self.fail_startup_queue(session_id, &format!("{error:#}"))
990 .await;
991 return;
992 }
993 };
994 let prepared_id = command_id.clone();
995 if let Err(error) = blocking(move || {
996 crate::database::prepare_startup_delivery(&prepared_id, prepared_json)
997 })
998 .await
999 {
1000 self.fail_startup_queue(session_id, &format!("{error:#}"))
1001 .await;
1002 return;
1003 }
1004 }
1005 if let Some(handle) = &handle {
1006 let persisted_id = command_id.clone();
1007 match blocking(move || crate::database::claim_startup_delivery(&persisted_id)).await
1008 {
1009 Ok(true) => {}
1010 Ok(false) => continue,
1013 Err(error) => {
1014 self.fail_startup_queue(session_id, &format!("{error:#}"))
1015 .await;
1016 return;
1017 }
1018 }
1019 let outcome = tokio::select! {
1020 () = cancel.cancelled() => Err(anyhow!("the daemon stopped before startup delivery settled")),
1021 result = self.run_startup_step(session_id, handle, &decoded, &command_id) => result,
1022 };
1023 let ordinal = match outcome {
1024 Ok(ordinal) => ordinal,
1025 Err(error) => {
1026 if error.is::<super::startup_followup::StartupRejected>()
1027 && let Some(group_id) = step.group_id.clone()
1028 {
1029 let id = session_id.to_owned();
1030 let reason = format!("{error:#}");
1031 let persisted = reason.clone();
1032 match blocking(move || {
1033 crate::database::fail_startup_group(&id, &group_id, &persisted)
1034 })
1035 .await
1036 {
1037 Ok(()) => {
1040 self.reject_startup(session_id, &reason).await;
1041 continue;
1042 }
1043 Err(persistence) => {
1044 tracing::error!(session_id, %persistence, "could not persist startup rejection");
1045 }
1046 }
1047 }
1048 let refused = error
1049 .downcast_ref::<mj_client::session::Refused>()
1050 .is_some();
1051 if refused
1052 && step.group_id.is_none()
1053 && let StartupStep::Prompt { text, .. } = &decoded
1054 {
1055 self.return_startup_prompt_to_draft(
1059 session_id,
1060 &command_id,
1061 text,
1062 &format!("{error:#}"),
1063 )
1064 .await;
1065 continue;
1066 }
1067 self.fail_startup_queue(session_id, &format!("{error:#}"))
1068 .await;
1069 return;
1070 }
1071 };
1072 let settled_id = command_id.clone();
1073 if let Err(error) = blocking(move || {
1074 crate::database::set_startup_delivery_accepted(&settled_id, ordinal)
1075 })
1076 .await
1077 {
1078 self.fail_startup_queue(
1079 session_id,
1080 &format!("delivery accepted but settlement failed: {error:#}"),
1081 )
1082 .await;
1083 return;
1084 }
1085 }
1086 let settled_id = command_id;
1087 let final_phase = match step.phase.as_str() {
1088 "cancelling" => "dismissed",
1089 "rejecting" => "failed",
1090 _ => "done",
1091 };
1092 let final_error = step.error.clone();
1093 if let Err(error) = blocking(move || {
1094 crate::database::set_startup_delivery_phase(
1095 &settled_id,
1096 final_phase,
1097 final_error.as_deref(),
1098 )
1099 })
1100 .await
1101 {
1102 self.fail_startup_queue(
1103 session_id,
1104 &format!("startup settlement failed: {error:#}"),
1105 )
1106 .await;
1107 return;
1108 }
1109 }
1110 }
1111
1112 async fn run_startup_step(
1113 &self,
1114 session_id: &str,
1115 handle: &crate::session_manager::ManagedSessionHandle,
1116 step: &StartupStep,
1117 command_id: &str,
1118 ) -> Result<Option<u64>> {
1119 match step {
1120 StartupStep::InstallHandoff(_) => bail!("archive startup step was not prepared"),
1121 StartupStep::PreparedHandoff { text } => handle
1122 .submit(
1123 command_id.to_owned(),
1124 RelayCommand::InstallPromptContext { text: text.clone() },
1125 )
1126 .await
1127 .map(Some),
1128 StartupStep::Prompt {
1129 text,
1130 inherited_draft,
1131 } => self
1132 .submit_startup_prompt(
1133 session_id,
1134 handle,
1135 text,
1136 inherited_draft.as_deref(),
1137 command_id,
1138 )
1139 .await
1140 .map(Some),
1141 StartupStep::Configure {
1142 key,
1143 value,
1144 optional,
1145 } => {
1146 super::startup_followup::configure_startup(
1147 &handle.client(),
1148 command_id,
1149 key,
1150 value,
1151 *optional,
1152 )
1153 .await
1154 }
1155 StartupStep::ApiPrompt { text } => {
1156 let ordinal = self
1157 .submit_startup_prompt(session_id, handle, text, None, command_id)
1158 .await?;
1159 let id = session_id.to_owned();
1160 blocking(move || crate::database::record_subagent_prompt(&id, ordinal)).await?;
1161 Ok(Some(ordinal))
1162 }
1163 }
1164 }
1165
1166 async fn submit_startup_prompt(
1169 &self,
1170 session_id: &str,
1171 handle: &crate::session_manager::ManagedSessionHandle,
1172 text: &str,
1173 inherited_draft: Option<&str>,
1174 command_id: &str,
1175 ) -> Result<u64> {
1176 let bundle_id = self
1177 .owner()
1178 .controller()
1179 .state
1180 .sessions
1181 .get(session_id)
1182 .map(|record| record.bundle_id.clone());
1183 let ordinal = handle
1184 .submit(
1185 command_id.to_owned(),
1186 RelayCommand::Prompt {
1187 prompt: vec![ContentBlock::Text(TextContent::new(text.to_owned()))],
1188 },
1189 )
1190 .await?;
1191 if let Some(expected) = inherited_draft {
1192 let persisted_id = session_id.to_owned();
1193 let persisted_expected = expected.to_owned();
1194 if let Err(error) = blocking(move || {
1195 crate::database::clear_session_draft_input_if_matches(
1196 &persisted_id,
1197 &persisted_expected,
1198 )
1199 })
1200 .await
1201 {
1202 tracing::warn!(
1203 session_id,
1204 error = format!("{error:#}"),
1205 "the delivered prompt's draft could not be cleared"
1206 );
1207 }
1208 self.publish_revision();
1209 }
1210 if let Some(bundle_id) = bundle_id {
1211 let history_id = session_id.to_owned();
1212 let history_text = text.to_owned();
1213 if let Err(error) = blocking(move || {
1214 crate::database::record_prompt(
1215 &history_id,
1216 &bundle_id,
1217 ordinal,
1218 None,
1219 &history_text,
1220 )
1221 })
1222 .await
1223 {
1224 tracing::warn!(
1225 session_id,
1226 error = format!("{error:#}"),
1227 "the queued prompt was accepted but its history could not be stored"
1228 );
1229 }
1230 }
1231 Ok(ordinal)
1232 }
1233
1234 async fn abandon_startup_step(self: &Arc<Self>, session_id: &str) {
1239 let lookup_id = session_id.to_owned();
1240 let step = match blocking(move || crate::database::next_startup_delivery(&lookup_id)).await
1241 {
1242 Ok(Some(step)) => step,
1243 Ok(None) => return,
1244 Err(error) => {
1245 tracing::error!(session_id, %error, "could not read the startup step to abandon");
1246 return;
1247 }
1248 };
1249 let reason = format!(
1250 "startup delivery gave up after {STARTUP_STEP_ATTEMPTS} attempts: {}",
1251 self.startup_prompts
1252 .lock()
1253 .unwrap_or_else(PoisonError::into_inner)
1254 .get(session_id)
1255 .and_then(|queue| queue.last_error.clone())
1256 .unwrap_or_else(|| "delivery kept failing".to_owned())
1257 );
1258 if let Some(group_id) = step.group_id.clone() {
1259 let id = session_id.to_owned();
1260 let persisted = reason.clone();
1261 if let Err(error) =
1262 blocking(move || crate::database::fail_startup_group(&id, &group_id, &persisted))
1263 .await
1264 {
1265 tracing::error!(session_id, %error, "could not fail an abandoned startup group");
1266 }
1267 self.reject_startup(session_id, &reason).await;
1268 return;
1269 }
1270 let decoded: Option<StartupStep> = serde_json::from_str(&step.step_json).ok();
1271 if let Some(StartupStep::Prompt { text, .. }) = &decoded {
1272 self.return_startup_prompt_to_draft(session_id, &step.command_id, text, &reason)
1273 .await;
1274 return;
1275 }
1276 let command_id = step.command_id.clone();
1277 let persisted = reason.clone();
1278 if let Err(error) = blocking(move || {
1279 crate::database::set_startup_delivery_phase(&command_id, "failed", Some(&persisted))
1280 })
1281 .await
1282 {
1283 tracing::error!(session_id, %error, "could not mark an abandoned startup step failed");
1284 }
1285 self.push_notice(session_id, reason);
1286 }
1287
1288 async fn reject_startup(self: &Arc<Self>, session_id: &str, reason: &str) {
1295 tracing::warn!(session_id, reason, "session startup failed");
1296 self.push_notice(session_id, format!("Session startup failed: {reason}"));
1297 self.fail_subagent_start(session_id, reason).await;
1298 }
1299
1300 async fn return_startup_prompt_to_draft(
1303 &self,
1304 session_id: &str,
1305 command_id: &str,
1306 text: &str,
1307 reason: &str,
1308 ) {
1309 let settled_id = command_id.to_owned();
1310 let persisted = reason.to_owned();
1311 if let Err(error) = blocking(move || {
1312 crate::database::set_startup_delivery_phase(&settled_id, "failed", Some(&persisted))
1313 })
1314 .await
1315 {
1316 tracing::error!(session_id, %error, "could not mark a refused startup prompt failed");
1317 return;
1318 }
1319 if let Err(error) = self.append_draft_input(session_id, text).await {
1320 tracing::error!(session_id, %error, "could not return a refused startup prompt to the draft");
1321 self.push_notice(
1322 session_id,
1323 format!("A queued prompt was refused and could not be saved as a draft: {reason}"),
1324 );
1325 return;
1326 }
1327 self.push_notice(
1328 session_id,
1329 format!("A queued prompt was refused and returned to your draft: {reason}"),
1330 );
1331 }
1332
1333 async fn fail_startup_queue(&self, session_id: &str, reason: &str) {
1335 let changed = {
1336 let mut queues = self
1337 .startup_prompts
1338 .lock()
1339 .unwrap_or_else(PoisonError::into_inner);
1340 let Some(queue) = queues.get_mut(session_id) else {
1341 return;
1342 };
1343 if queue.last_error.as_deref() == Some(reason) {
1344 false
1345 } else {
1346 queue.last_error = Some(reason.to_owned());
1347 true
1348 }
1349 };
1350 if !changed {
1351 return;
1352 }
1353 self.push_notice(session_id, format!("Startup work remains saved for this session: {reason}. It has not been restored as an unsent draft because delivery may have been accepted."));
1354 tracing::warn!(session_id, reason, "durable startup delivery paused");
1355 }
1356
1357 pub(super) async fn append_draft_input(&self, session_id: &str, text: &str) -> Result<()> {
1359 let session_id = session_id.to_owned();
1360 let text = text.to_owned();
1361 blocking(move || crate::database::append_session_draft_input(&session_id, &text)).await?;
1362 self.publish_revision();
1363 Ok(())
1364 }
1365
1366 pub(crate) async fn cancel_and_join_startup_prompts(&self) -> Result<()> {
1369 let tasks = {
1370 let mut queues = self
1371 .startup_prompts
1372 .lock()
1373 .unwrap_or_else(PoisonError::into_inner);
1374 queues
1375 .values_mut()
1376 .filter_map(|queue| {
1377 queue.cancel.cancel();
1378 queue.task.take()
1379 })
1380 .collect::<Vec<_>>()
1381 };
1382 let deadline = tokio::time::Instant::now() + Duration::from_secs(1);
1383 let mut outcome = Ok(());
1384 for mut task in tasks {
1385 let joined = match tokio::time::timeout_at(deadline, &mut task).await {
1386 Ok(Ok(())) => Ok(()),
1387 Ok(Err(error)) => Err(anyhow!("startup prompt delivery task failed: {error}")),
1388 Err(_) => {
1389 task.abort();
1390 let _ = task.await;
1391 Err(anyhow!(
1392 "a startup prompt delivery task did not stop within 1s; durable work retained"
1393 ))
1394 }
1395 };
1396 if outcome.is_ok() {
1397 outcome = joined;
1398 } else if let Err(error) = joined {
1399 tracing::warn!(%error, "another startup prompt drain did not stop cleanly");
1400 }
1401 }
1402 outcome
1403 }
1404}
1405
1406#[cfg(test)]
1407mod tests {
1408 use super::super::tests::test_runtime_state;
1409 use std::sync::Arc;
1410
1411 #[tokio::test]
1416 async fn a_session_text_search_asks_for_a_sync_like_the_resume_search() {
1417 let mut state = test_runtime_state();
1418 Arc::get_mut(&mut state).expect("the only handle").wiki =
1419 crate::sessionwiki::WikiIndexer::inert();
1420 assert!(!state.wiki().sync_requested());
1421
1422 state.session_text_search(" ".into()).await.unwrap();
1424 assert!(
1425 state.wiki().sync_requested(),
1426 "a stale index is synced for the next search"
1427 );
1428 }
1429}