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 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 pub fn review_host(&self) -> &TurnReviewHost {
115 &self.review_host
116 }
117
118 pub fn wiki(&self) -> &crate::sessionwiki::WikiIndexer {
120 &self.wiki
121 }
122
123 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 let rows =
131 blocking(move || crate::sessionwiki::query_rows(&query, limit, &live, false)).await?;
132 Ok(WikiSearchPage {
135 rows,
136 status: self.wiki.status(),
137 })
138 }
139
140 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 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 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 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 pub async fn wiki_session(
188 &self,
189 wiki_id: String,
190 ) -> Result<Option<mj_client::daemon::WikiSessionInfo>> {
191 let known = self.live_session_ids();
194 blocking(move || crate::sessionwiki::wiki_session(&wiki_id, &known)).await
195 }
196
197 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 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 session_title_override: Some(archived.title.clone()),
254 })
255 .await?;
256 let session_id = registered.session.id.clone();
257 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 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 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 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 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 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 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 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 pub fn with_config<T>(&self, read: impl FnOnce(&Config) -> T) -> T {
618 read(&self.owner().controller().config)
619 }
620
621 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 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 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 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 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 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 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 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 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 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 async fn drain_startup_queue(self: Arc<Self>, session_id: &str, cancel: &CancellationToken) {
879 loop {
880 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 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 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 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 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 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 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 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 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 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 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 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 #[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 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}