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