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