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