1use super::*;
2
3impl RuntimeState {
4 pub(super) fn new(
5 session_manager: SessionManagerControl,
6 controller: Controller,
7 recovery_observer: RecoveryObserver,
8 worker_upgrade_observer: WorkerUpgradeObserver,
9 workspaces: Vec<WorkspaceRecord>,
10 ) -> Self {
11 Self::new_with_controller_loader(
12 session_manager,
13 controller,
14 recovery_observer,
15 worker_upgrade_observer,
16 workspaces,
17 Controller::load,
18 )
19 }
20
21 pub(super) fn new_with_controller_loader(
22 session_manager: SessionManagerControl,
23 controller: Controller,
24 recovery_observer: RecoveryObserver,
25 worker_upgrade_observer: WorkerUpgradeObserver,
26 workspaces: Vec<WorkspaceRecord>,
27 controller_loader: fn() -> Result<Controller>,
28 ) -> Self {
29 let initial_revision = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap_or(1);
34 let revisions = RuntimeRevisions::new(initial_revision);
35 let (workspaces_tx, _) = tokio::sync::watch::channel(workspaces);
36 let review_config = Arc::new(Mutex::new(controller.config.review.clone()));
40 let review_host = TurnReviewHost::spawn_notifying(
41 session_manager.clone(),
42 {
43 let installed = review_config.clone();
44 Arc::new(move || {
45 installed
46 .lock()
47 .unwrap_or_else(PoisonError::into_inner)
48 .clone()
49 })
50 },
51 revisions.notifier(),
52 Some(recovery_observer.gate.clone()),
53 );
54 Self {
55 attachments: Mutex::new(BTreeMap::new()),
56 phone_status: Mutex::new(WebViewerStatus::Starting),
57 web_viewer: crate::web_viewer::ViewerControl::new(),
58 ever_attached: AtomicBool::new(false),
59 sessions: Mutex::new(BTreeMap::new()),
60 background_policies: Mutex::new(BTreeMap::new()),
61 revisions,
62 workspaces_tx,
63 workspace_refresh: tokio::sync::Mutex::new(()),
64 session_manager,
65 lifecycle: Mutex::new(BTreeMap::new()),
66 workspace_closes: Mutex::new(BTreeMap::new()),
67 workspace_resume_admission: Mutex::new(BTreeMap::new()),
68 harness_readiness: Mutex::new(HarnessReadinessWatch::default()),
69 startup_prompts: Mutex::new(BTreeMap::new()),
70 close_requested: Mutex::new(BTreeSet::new()),
71 controller: Mutex::new(controller),
72 controller_loader,
73 config_mutation: tokio::sync::Mutex::new(()),
74 recovery_observer,
75 worker_upgrade_observer,
76 notices: Mutex::new(VecDeque::new()),
77 next_notice_id: AtomicU64::new(1),
78 review_config,
79 review_host,
80 wiki: crate::sessionwiki::WikiIndexer::spawn(),
81 }
82 }
83
84 pub fn review_host(&self) -> &TurnReviewHost {
86 &self.review_host
87 }
88
89 pub fn wiki(&self) -> &crate::sessionwiki::WikiIndexer {
91 &self.wiki
92 }
93
94 pub async fn wiki_search(&self, query: String, limit: usize) -> Result<WikiSearchPage> {
97 if crate::sessionwiki::sync_is_stale(self.wiki.last_success()) {
98 self.wiki.request_sync(false);
101 }
102 let live = self.live_session_ids();
103 let rows =
106 blocking(move || crate::sessionwiki::query_rows(&query, limit, &live, false)).await?;
107 Ok(WikiSearchPage {
110 rows,
111 status: self.wiki.status(),
112 })
113 }
114
115 pub async fn wiki_brief(&self, wiki_id: String, max_chars: usize) -> Result<Option<String>> {
118 blocking(move || crate::sessionwiki::brief(&wiki_id, max_chars)).await
119 }
120
121 pub async fn wiki_hits(
124 &self,
125 wiki_id: String,
126 query: String,
127 context_messages: usize,
128 per_message_chars: usize,
129 ) -> Result<Option<WikiHitTranscript>> {
130 blocking(move || {
131 crate::sessionwiki::transcript_hits(
132 &wiki_id,
133 &query,
134 context_messages,
135 per_message_chars,
136 )
137 })
138 .await
139 }
140
141 pub async fn wiki_session(
144 &self,
145 wiki_id: String,
146 ) -> Result<Option<mj_client::daemon::WikiSessionInfo>> {
147 let known = self.live_session_ids();
150 blocking(move || crate::sessionwiki::wiki_session(&wiki_id, &known)).await
151 }
152
153 pub async fn restore_wiki_session(
161 self: &Arc<Self>,
162 request: WikiRestoreRequest,
163 cancellation: &CancellationToken,
164 ) -> Result<Option<RegisteredSession>> {
165 let wiki_id = request.wiki_id.clone();
166 let Some(archived) =
167 blocking(move || crate::sessionwiki::archived_session(&wiki_id)).await?
168 else {
169 return Ok(None);
170 };
171 let project_directory = request
172 .project_directory
173 .clone()
174 .or_else(|| archived.project_directory.clone())
175 .context(
176 "name a project directory: the archived session's own project is no longer on this machine",
177 )?;
178 let source = project_directory.display().to_string();
179 let bundle_id = blocking(move || {
180 crate::controller::create_bundle_from_sources(&[source])
181 .map(|created| created.bundle_id)
182 .map_err(anyhow::Error::new)
183 })
184 .await
185 .context("find or create a bundle for the restored session's project")?;
186 let registered = self
187 .start_create_session(CreateSessionRequest {
188 launch_base: None,
189 launch_branch: None,
190 create_managed_worktree: None,
191 mjolnir_subagents: None,
192 initial_prompt: None,
193 workspace_id: request.workspace_id,
194 profile_id: request.profile_id,
195 bundle_id,
196 project_directory: Some(project_directory),
197 target_template_id: request.target_template_id,
198 additional_mounts: request.additional_mounts,
199 resource_allocation: request.resource_allocation,
200 title: archived.title.clone(),
201 session_title_override: Some(archived.title.clone()),
206 })
207 .await?;
208 let session_id = registered.session.id.clone();
209 self.queue_startup_step(
213 &session_id,
214 StartupStep::InstallHandoff(Box::new(archived.snapshot)),
215 cancellation,
216 )?;
217 Ok(Some(registered))
218 }
219
220 pub(super) fn live_session_ids(&self) -> BTreeSet<String> {
221 self.controller
222 .lock()
223 .unwrap_or_else(PoisonError::into_inner)
224 .state
225 .sessions
226 .keys()
227 .cloned()
228 .collect()
229 }
230
231 async fn install_archive_handoff(
235 &self,
236 session_id: &str,
237 handle: &crate::session_manager::ManagedSessionHandle,
238 snapshot: &mj_core::archive::CanonicalSessionSnapshot,
239 ) -> Result<()> {
240 let (config, profile_id) = {
241 let controller = self
242 .controller
243 .lock()
244 .unwrap_or_else(PoisonError::into_inner);
245 let profile_id = controller
246 .state
247 .sessions
248 .get(session_id)
249 .map(|record| record.last_profile.clone());
250 (controller.config.clone(), profile_id)
251 };
252 let context_bytes = crate::handoff::profile_handoff_bytes(
253 profile_id.and_then(|id| config.profiles.get(&id)),
254 );
255 let cancel = CancellationToken::new();
256 let handoff = crate::handoff::build_handoff_context(
257 session_id,
258 &config,
259 snapshot,
260 context_bytes,
261 &cancel,
262 )
263 .await
264 .context("compact the archived transcript")?;
265 handle
266 .install_prompt_context(format!(
267 "{} {handoff}",
268 crate::compaction::ARCHIVE_HANDOFF_PREAMBLE
269 ))
270 .await
271 .context("install the archived hand-off")?;
272 tracing::info!(
273 session_id,
274 bytes = handoff.len(),
275 "installed the restored archive's hand-off"
276 );
277 Ok(())
278 }
279
280 pub(super) async fn wait_for_ready_session(
282 &self,
283 session_id: &str,
284 ) -> Result<crate::session_manager::ManagedSessionHandle> {
285 const POLL: Duration = Duration::from_millis(250);
286 let deadline = tokio::time::Instant::now() + Duration::from_secs(30 * 60);
287 loop {
288 match self.session_state(session_id) {
292 Some(
293 SessionState::Provisioning
294 | SessionState::Running
295 | SessionState::Disconnected
296 | SessionState::Checkpointing,
297 ) => {}
298 Some(state) => bail!("session {session_id} is {state:?} before its hand-off"),
299 None => bail!("session {session_id} disappeared before its hand-off"),
300 }
301 if let Ok(handle) = self.session_manager.session(session_id).await {
302 let view = handle.view();
303 if let Some(ViewError::TargetMissing(detail)) = &view.error {
306 bail!("session {session_id} lost its target: {detail}");
307 }
308 if view.connected
309 && view
310 .snapshot
311 .is_some_and(|snapshot| snapshot.operational.native_session_is_ready())
312 {
313 return Ok(handle);
314 }
315 }
316 ensure!(
317 tokio::time::Instant::now() < deadline,
318 "session {session_id} was not ready for its hand-off within 30 minutes"
319 );
320 tokio::time::sleep(POLL).await;
321 }
322 }
323
324 pub(super) async fn fail_unfinished_provisioning(
339 self: &Arc<Self>,
340 session_id: &str,
341 error: &str,
342 ) {
343 let provisioning = {
344 let controller = self
345 .controller
346 .lock()
347 .unwrap_or_else(PoisonError::into_inner);
348 durable_session_state(&controller, session_id) == Some(SessionState::Provisioning)
349 };
350 if !provisioning {
351 return;
352 }
353 let cause = format!("session provisioning ended without finishing: {error}");
354 let applied = blocking({
355 let session_id = session_id.to_owned();
356 move || {
357 let mut controller = Controller::load()?;
358 controller.fail_interrupted_lifecycle(&session_id, &cause)
359 }
360 })
361 .await;
362 match applied {
363 Ok(true) => {
364 if let Err(error) = self.reload_controller().await {
365 tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after recording a failed provision");
366 }
367 }
368 Ok(false) => {}
369 Err(error) => tracing::warn!(
370 %session_id,
371 error = format!("{error:#}"),
372 "could not record that provisioning ended without finishing"
373 ),
374 }
375 }
376
377 pub(super) async fn record_failed_close(
385 self: &Arc<Self>,
386 session_id: &str,
387 reference: &str,
388 failure: &LifecycleFailure,
389 ) {
390 self.record_lifecycle_failure(
391 session_id,
392 reference,
393 failure,
394 mj_core::state::CLOSE_FAILURE_PREFIX,
395 )
396 .await;
397 }
398
399 pub(crate) async fn record_lifecycle_failure(
400 self: &Arc<Self>,
401 session_id: &str,
402 reference: &str,
403 failure: &LifecycleFailure,
404 prefix: &str,
405 ) {
406 let cause = match &failure.refusal {
407 Some(refusal) => format!("{prefix}: {refusal}"),
408 None => {
409 format!("{prefix}; the daemon log records the reason under reference {reference}")
410 }
411 };
412 let applied = blocking({
413 let session_id = session_id.to_owned();
414 let cause = cause.clone();
415 move || {
416 let mut controller = Controller::load()?;
417 controller.record_failed_close(&session_id, &cause)
418 }
419 })
420 .await;
421 match applied {
422 Ok(true) => {
423 if let Err(error) = self.reload_controller().await {
424 tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after recording a failed close");
425 }
426 self.publish_revision();
427 }
428 Ok(false) => {}
429 Err(error) => tracing::warn!(
430 %session_id,
431 error = format!("{error:#}"),
432 "could not record why a close failed"
433 ),
434 }
435 }
436
437 pub async fn clear_recorded_close_failure(self: &Arc<Self>, session_id: &str) {
445 let recorded = self
446 .controller
447 .lock()
448 .unwrap_or_else(PoisonError::into_inner)
449 .state
450 .sessions
451 .get(session_id)
452 .is_some_and(|record| record.public_error().is_some());
453 if !recorded {
454 return;
455 }
456 let cleared = blocking({
457 let session_id = session_id.to_owned();
458 move || {
459 let mut controller = Controller::load()?;
460 controller.clear_recorded_close_failure(&session_id)
461 }
462 })
463 .await;
464 match cleared {
465 Ok(true) => {
466 if let Err(error) = self.reload_controller().await {
467 tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after clearing a recorded close failure");
468 }
469 self.publish_revision();
470 }
471 Ok(false) => {}
472 Err(error) => tracing::warn!(
473 %session_id,
474 error = format!("{error:#}"),
475 "could not clear a recorded close failure"
476 ),
477 }
478 }
479
480 pub(super) fn note_lifecycle_outcome(&self, session_id: &str) {
481 let stopped = {
482 let controller = self
483 .controller
484 .lock()
485 .unwrap_or_else(PoisonError::into_inner);
486 durable_session_state(&controller, session_id) == Some(SessionState::Stopped)
487 };
488 if stopped {
489 self.wiki.request_sync(false);
490 }
491 }
492
493 pub fn allocate_revision(&self) -> u64 {
494 self.revisions.allocate()
495 }
496
497 pub(super) fn publish_revision(&self) -> u64 {
498 self.revisions.publish()
499 }
500
501 pub(super) fn attachments(&self) -> std::sync::MutexGuard<'_, BTreeMap<String, Attachment>> {
502 self.attachments
503 .lock()
504 .unwrap_or_else(PoisonError::into_inner)
505 }
506
507 pub(super) fn prune_dead_clients(&self) {
508 self.attachments()
509 .retain(|_, attachment| process_is_alive(attachment.pid));
510 }
511
512 pub(super) fn workspace_has_active_resume(&self, workspace_id: &str) -> bool {
513 self.lifecycle
514 .lock()
515 .unwrap_or_else(PoisonError::into_inner)
516 .values()
517 .any(|active| {
518 active.result.borrow().is_none()
519 && active.resume_workspace_id.as_deref() == Some(workspace_id)
520 })
521 }
522
523 pub fn publish_web_access(&self, access: crate::server::WebViewerAccess) {
524 use crate::server::WebViewerAccess;
525 let status = match &access {
526 WebViewerAccess::Starting => WebViewerStatus::Starting,
527 WebViewerAccess::Ready {
528 viewer_url,
529 viewer_code,
530 qr_login_url,
531 fallback_reason,
532 ..
533 } => WebViewerStatus::Ready {
534 viewer_url: viewer_url.clone(),
535 viewer_code: viewer_code.clone(),
536 qr_login_url: qr_login_url.clone(),
537 fallback_reason: fallback_reason.clone(),
538 },
539 WebViewerAccess::Failed {
540 address, message, ..
541 } => WebViewerStatus::Error {
542 message: format!("{message} Address: {address}"),
543 },
544 WebViewerAccess::Unavailable(message) => WebViewerStatus::Error {
545 message: message.clone(),
546 },
547 };
548 self.web_viewer.publish(access);
549 self.set_phone_status(status);
550 }
551
552 pub(super) fn set_phone_status(&self, status: WebViewerStatus) {
553 *self
554 .phone_status
555 .lock()
556 .unwrap_or_else(PoisonError::into_inner) = status;
557 }
558
559 pub(super) fn phone_status(&self) -> WebViewerStatus {
560 self.phone_status
561 .lock()
562 .unwrap_or_else(PoisonError::into_inner)
563 .clone()
564 }
565
566 pub(super) fn workspaces(&self) -> tokio::sync::watch::Receiver<Vec<WorkspaceRecord>> {
567 self.workspaces_tx.subscribe()
568 }
569
570 pub(super) fn worker_poll_exclusion_session_ids(
571 &self,
572 controller: &Controller,
573 ) -> BTreeSet<String> {
574 self.lifecycle
575 .lock()
576 .unwrap_or_else(PoisonError::into_inner)
577 .iter()
578 .filter(|(session_id, active)| {
579 active.result.borrow().is_none()
580 && (active.move_source_closed
581 || lifecycle_owns_worker_target(
582 active.kind,
583 controller
584 .state
585 .sessions
586 .get(*session_id)
587 .map(|session| session.state),
588 ))
589 })
590 .map(|(session_id, _)| session_id.clone())
591 .collect()
592 }
593
594 pub fn revisions(&self) -> tokio::sync::watch::Receiver<u64> {
595 self.revisions.subscribe()
596 }
597
598 pub fn with_config<T>(&self, read: impl FnOnce(&Config) -> T) -> T {
601 read(
602 &self
603 .controller
604 .lock()
605 .unwrap_or_else(PoisonError::into_inner)
606 .config,
607 )
608 }
609
610 pub async fn create_quick_bundle(
615 &self,
616 source: String,
617 ) -> std::result::Result<
618 crate::controller::QuickBundleCreation,
619 crate::controller::QuickBundleFailure,
620 > {
621 let _mutation = self.config_mutation.lock().await;
622 tokio::task::spawn_blocking(move || crate::controller::create_quick_bundle(&source))
623 .await
624 .map_err(|error| {
625 crate::controller::QuickBundleFailure::Persistence(anyhow!(
626 "bundle creation task panicked: {error}"
627 ))
628 })?
629 }
630
631 pub async fn create_bundle_from_sources(
634 &self,
635 sources: Vec<String>,
636 ) -> std::result::Result<
637 crate::controller::QuickBundleCreation,
638 crate::controller::QuickBundleFailure,
639 > {
640 let _mutation = self.config_mutation.lock().await;
641 tokio::task::spawn_blocking(move || crate::controller::create_bundle_from_sources(&sources))
642 .await
643 .map_err(|error| {
644 crate::controller::QuickBundleFailure::Persistence(anyhow!(
645 "bundle creation task panicked: {error}"
646 ))
647 })?
648 }
649
650 pub(crate) fn publish_workspaces(&self, workspaces: Vec<WorkspaceRecord>) {
652 self.workspaces_tx.send_replace(workspaces);
653 self.publish_revision();
654 }
655
656 pub(crate) async fn refresh_workspaces(&self) -> Result<()> {
657 let _refresh = self.workspace_refresh.lock().await;
660 let workspaces = tokio::task::spawn_blocking(crate::database::list_workspaces)
661 .await
662 .context("daemon workspace refresh task panicked")??;
663 self.publish_workspaces(workspaces);
664 Ok(())
665 }
666
667 pub(crate) fn queue_startup_step(
674 self: &Arc<Self>,
675 session_id: &str,
676 step: StartupStep,
677 cancellation: &CancellationToken,
678 ) -> Result<()> {
679 match self.session_state(session_id) {
682 Some(
683 SessionState::Provisioning
684 | SessionState::Running
685 | SessionState::Disconnected
686 | SessionState::Checkpointing,
687 ) => {}
688 Some(state) => {
689 bail!("session {session_id} is {state:?}; it cannot take a queued prompt")
690 }
691 None => bail!("unknown session {session_id}"),
692 }
693 let mut queues = self
694 .startup_prompts
695 .lock()
696 .unwrap_or_else(PoisonError::into_inner);
697 if let Some(queue) = queues.get_mut(session_id) {
698 queue.pending.push_back(step);
699 return Ok(());
700 }
701 let cancel = cancellation.child_token();
702 let upgrade_work = crate::upgrade::activity("startup prompt delivery")?;
703 queues.insert(
704 session_id.to_owned(),
705 StartupQueue {
706 pending: VecDeque::from([step]),
707 in_flight: false,
708 cancel: cancel.clone(),
709 task: None,
710 },
711 );
712 let runtime = Arc::clone(self);
713 let drain_session = session_id.to_owned();
714 let task = tokio::spawn(async move {
717 let _upgrade_work = upgrade_work;
718 let supervised = {
719 let runtime = Arc::clone(&runtime);
720 let session_id = drain_session.clone();
721 let cancel = cancel.clone();
722 tokio::spawn(async move { runtime.drain_startup_queue(&session_id, &cancel).await })
723 };
724 if let Err(error) = supervised.await {
725 runtime
726 .fail_startup_queue(
727 &drain_session,
728 None,
729 &format!("the daemon's delivery task failed: {error}"),
730 )
731 .await;
732 }
733 });
734 if let Some(queue) = queues.get_mut(session_id) {
735 queue.task = Some(task);
736 }
737 Ok(())
738 }
739
740 async fn drain_startup_queue(self: Arc<Self>, session_id: &str, cancel: &CancellationToken) {
744 let handle = tokio::select! {
745 () = cancel.cancelled() => {
746 self.fail_startup_queue(
747 session_id,
748 None,
749 "the daemon stopped before the session was ready",
750 )
751 .await;
752 return;
753 }
754 ready = self.wait_for_ready_session(session_id) => match ready {
755 Ok(handle) => handle,
756 Err(error) => {
757 self.fail_startup_queue(session_id, None, &format!("{error:#}"))
758 .await;
759 return;
760 }
761 },
762 };
763 loop {
764 let step = {
765 let mut queues = self
766 .startup_prompts
767 .lock()
768 .unwrap_or_else(PoisonError::into_inner);
769 let Some(queue) = queues.get_mut(session_id) else {
770 return;
771 };
772 match queue.pending.pop_front() {
773 Some(step) => {
774 queue.in_flight = true;
775 step
776 }
777 None => {
778 queues.remove(session_id);
779 return;
780 }
781 }
782 };
783 let outcome = tokio::select! {
784 () = cancel.cancelled() => {
785 Err(anyhow!("the daemon stopped before the prompt was sent"))
786 }
787 result = self.run_startup_step(session_id, &handle, &step) => result,
788 };
789 if let Err(error) = outcome {
790 self.fail_startup_queue(session_id, Some(step), &format!("{error:#}"))
791 .await;
792 return;
793 }
794 let mut queues = self
795 .startup_prompts
796 .lock()
797 .unwrap_or_else(PoisonError::into_inner);
798 let Some(queue) = queues.get_mut(session_id) else {
799 return;
800 };
801 queue.in_flight = false;
802 if queue.pending.is_empty() {
803 queues.remove(session_id);
804 return;
805 }
806 }
807 }
808
809 async fn run_startup_step(
810 &self,
811 session_id: &str,
812 handle: &crate::session_manager::ManagedSessionHandle,
813 step: &StartupStep,
814 ) -> Result<()> {
815 match step {
816 StartupStep::InstallHandoff(snapshot) => {
817 self.install_archive_handoff(session_id, handle, snapshot)
818 .await
819 }
820 StartupStep::Prompt {
821 text,
822 inherited_draft,
823 } => {
824 self.submit_startup_prompt(session_id, handle, text, inherited_draft.as_deref())
825 .await
826 }
827 }
828 }
829
830 async fn submit_startup_prompt(
833 &self,
834 session_id: &str,
835 handle: &crate::session_manager::ManagedSessionHandle,
836 text: &str,
837 inherited_draft: Option<&str>,
838 ) -> Result<()> {
839 let bundle_id = self
840 .controller
841 .lock()
842 .unwrap_or_else(PoisonError::into_inner)
843 .state
844 .sessions
845 .get(session_id)
846 .map(|record| record.bundle_id.clone());
847 let ordinal = handle
848 .submit(
849 new_command_id("startup")?,
850 RelayCommand::Prompt {
851 prompt: vec![ContentBlock::Text(TextContent::new(text.to_owned()))],
852 },
853 )
854 .await?;
855 if let Some(expected) = inherited_draft {
856 let persisted_id = session_id.to_owned();
857 let persisted_expected = expected.to_owned();
858 if let Err(error) = blocking(move || {
859 crate::database::clear_session_draft_input_if_matches(
860 &persisted_id,
861 &persisted_expected,
862 )
863 })
864 .await
865 {
866 tracing::warn!(
867 session_id,
868 error = format!("{error:#}"),
869 "the delivered prompt's draft could not be cleared"
870 );
871 }
872 if let Some(record) = self
873 .controller
874 .lock()
875 .unwrap_or_else(PoisonError::into_inner)
876 .state
877 .sessions
878 .get_mut(session_id)
879 && record.draft_input == expected
880 {
881 record.draft_input.clear();
882 }
883 self.publish_revision();
884 }
885 if let Some(bundle_id) = bundle_id {
886 let history_id = session_id.to_owned();
887 let history_text = text.to_owned();
888 if let Err(error) = blocking(move || {
889 crate::database::record_prompt(
890 &history_id,
891 &bundle_id,
892 ordinal,
893 None,
894 &history_text,
895 )
896 })
897 .await
898 {
899 tracing::warn!(
900 session_id,
901 error = format!("{error:#}"),
902 "the queued prompt was accepted but its history could not be stored"
903 );
904 }
905 }
906 Ok(())
907 }
908
909 async fn fail_startup_queue(
913 &self,
914 session_id: &str,
915 failed: Option<StartupStep>,
916 reason: &str,
917 ) {
918 let remaining = self
919 .startup_prompts
920 .lock()
921 .unwrap_or_else(PoisonError::into_inner)
922 .remove(session_id)
923 .map(|queue| queue.pending)
924 .unwrap_or_default();
925 let mut texts = Vec::new();
926 let mut dropped_handoff = false;
927 for step in failed.into_iter().chain(remaining) {
928 match step {
929 StartupStep::Prompt { text, .. } => texts.push(text),
930 StartupStep::InstallHandoff(_) => dropped_handoff = true,
931 }
932 }
933 if dropped_handoff {
934 tracing::warn!(
935 session_id,
936 reason,
937 "could not install the restored archive's hand-off"
938 );
939 self.push_notice(
940 session_id,
941 format!("The restored session started without its archived hand-off: {reason}"),
942 );
943 }
944 if texts.is_empty() {
945 return;
946 }
947 let restored = texts.join("\n\n");
948 if let Err(error) = self.append_draft_input(session_id, &restored).await {
949 tracing::warn!(
950 session_id,
951 error = format!("{error:#}"),
952 "a queued prompt could not be saved back into the session's draft"
953 );
954 }
955 self.push_notice(
956 session_id,
957 format!(
958 "Your prompt could not be sent to session {} ({reason}); it is back in the composer draft.",
959 mj_core::state::short_id(session_id)
960 ),
961 );
962 tracing::warn!(
963 session_id,
964 reason,
965 "a queued startup prompt could not be delivered"
966 );
967 }
968
969 pub(super) async fn append_draft_input(&self, session_id: &str, text: &str) -> Result<()> {
974 let existing = self
975 .session_record(session_id)
976 .map(|record| record.draft_input)
977 .unwrap_or_default();
978 let combined = [existing.as_str(), text]
979 .into_iter()
980 .filter(|part| !part.is_empty())
981 .collect::<Vec<_>>()
982 .join("\n\n");
983 let persisted_id = session_id.to_owned();
984 let persisted = combined.clone();
985 let stored =
986 blocking(move || crate::database::set_session_draft_input(&persisted_id, &persisted))
987 .await;
988 if let Some(record) = self
989 .controller
990 .lock()
991 .unwrap_or_else(PoisonError::into_inner)
992 .state
993 .sessions
994 .get_mut(session_id)
995 {
996 record.draft_input = combined;
997 }
998 self.publish_revision();
999 stored
1000 }
1001
1002 pub(crate) async fn cancel_and_join_startup_prompts(&self) -> Result<()> {
1005 let tasks = {
1006 let mut queues = self
1007 .startup_prompts
1008 .lock()
1009 .unwrap_or_else(PoisonError::into_inner);
1010 queues
1011 .values_mut()
1012 .filter_map(|queue| {
1013 queue.cancel.cancel();
1014 queue.task.take()
1015 })
1016 .collect::<Vec<_>>()
1017 };
1018 let deadline = tokio::time::Instant::now() + Duration::from_secs(1);
1019 let mut outcome = Ok(());
1020 for task in tasks {
1021 let joined = match tokio::time::timeout_at(deadline, task).await {
1022 Ok(Ok(())) => Ok(()),
1023 Ok(Err(error)) => Err(anyhow!("startup prompt delivery task failed: {error}")),
1024 Err(_) => Err(anyhow!(
1025 "a startup prompt delivery task did not stop within 1s"
1026 )),
1027 };
1028 if outcome.is_ok() {
1029 outcome = joined;
1030 } else if let Err(error) = joined {
1031 tracing::warn!(%error, "another startup prompt drain did not stop cleanly");
1032 }
1033 }
1034 outcome
1035 }
1036}