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 let upgrade_work = crate::upgrade::activity("startup prompt delivery")?;
666 queues.insert(
667 session_id.to_owned(),
668 StartupQueue {
669 pending: VecDeque::from([step]),
670 in_flight: false,
671 cancel: cancel.clone(),
672 task: None,
673 },
674 );
675 let runtime = Arc::clone(self);
676 let drain_session = session_id.to_owned();
677 let task = tokio::spawn(async move {
680 let _upgrade_work = upgrade_work;
681 let supervised = {
682 let runtime = Arc::clone(&runtime);
683 let session_id = drain_session.clone();
684 let cancel = cancel.clone();
685 tokio::spawn(async move { runtime.drain_startup_queue(&session_id, &cancel).await })
686 };
687 if let Err(error) = supervised.await {
688 runtime
689 .fail_startup_queue(
690 &drain_session,
691 None,
692 &format!("the daemon's delivery task failed: {error}"),
693 )
694 .await;
695 }
696 });
697 if let Some(queue) = queues.get_mut(session_id) {
698 queue.task = Some(task);
699 }
700 Ok(())
701 }
702
703 async fn drain_startup_queue(self: Arc<Self>, session_id: &str, cancel: &CancellationToken) {
707 let handle = tokio::select! {
708 () = cancel.cancelled() => {
709 self.fail_startup_queue(
710 session_id,
711 None,
712 "the daemon stopped before the session was ready",
713 )
714 .await;
715 return;
716 }
717 ready = self.wait_for_ready_session(session_id) => match ready {
718 Ok(handle) => handle,
719 Err(error) => {
720 self.fail_startup_queue(session_id, None, &format!("{error:#}"))
721 .await;
722 return;
723 }
724 },
725 };
726 loop {
727 let step = {
728 let mut queues = self
729 .startup_prompts
730 .lock()
731 .unwrap_or_else(PoisonError::into_inner);
732 let Some(queue) = queues.get_mut(session_id) else {
733 return;
734 };
735 match queue.pending.pop_front() {
736 Some(step) => {
737 queue.in_flight = true;
738 step
739 }
740 None => {
741 queues.remove(session_id);
742 return;
743 }
744 }
745 };
746 let outcome = tokio::select! {
747 () = cancel.cancelled() => {
748 Err(anyhow!("the daemon stopped before the prompt was sent"))
749 }
750 result = self.run_startup_step(session_id, &handle, &step) => result,
751 };
752 if let Err(error) = outcome {
753 self.fail_startup_queue(session_id, Some(step), &format!("{error:#}"))
754 .await;
755 return;
756 }
757 let mut queues = self
758 .startup_prompts
759 .lock()
760 .unwrap_or_else(PoisonError::into_inner);
761 let Some(queue) = queues.get_mut(session_id) else {
762 return;
763 };
764 queue.in_flight = false;
765 if queue.pending.is_empty() {
766 queues.remove(session_id);
767 return;
768 }
769 }
770 }
771
772 async fn run_startup_step(
773 &self,
774 session_id: &str,
775 handle: &crate::session_manager::ManagedSessionHandle,
776 step: &StartupStep,
777 ) -> Result<()> {
778 match step {
779 StartupStep::InstallHandoff(snapshot) => {
780 self.install_archive_handoff(session_id, handle, snapshot)
781 .await
782 }
783 StartupStep::Prompt {
784 text,
785 inherited_draft,
786 } => {
787 self.submit_startup_prompt(session_id, handle, text, inherited_draft.as_deref())
788 .await
789 }
790 }
791 }
792
793 async fn submit_startup_prompt(
796 &self,
797 session_id: &str,
798 handle: &crate::session_manager::ManagedSessionHandle,
799 text: &str,
800 inherited_draft: Option<&str>,
801 ) -> Result<()> {
802 let bundle_id = self
803 .controller
804 .lock()
805 .unwrap_or_else(PoisonError::into_inner)
806 .state
807 .sessions
808 .get(session_id)
809 .map(|record| record.bundle_id.clone());
810 let ordinal = handle
811 .submit(
812 new_command_id("startup")?,
813 RelayCommand::Prompt {
814 prompt: vec![ContentBlock::Text(TextContent::new(text.to_owned()))],
815 },
816 )
817 .await?;
818 if let Some(expected) = inherited_draft {
819 let persisted_id = session_id.to_owned();
820 let persisted_expected = expected.to_owned();
821 if let Err(error) = blocking(move || {
822 crate::database::clear_session_draft_input_if_matches(
823 &persisted_id,
824 &persisted_expected,
825 )
826 })
827 .await
828 {
829 tracing::warn!(
830 session_id,
831 error = format!("{error:#}"),
832 "the delivered prompt's draft could not be cleared"
833 );
834 }
835 if let Some(record) = self
836 .controller
837 .lock()
838 .unwrap_or_else(PoisonError::into_inner)
839 .state
840 .sessions
841 .get_mut(session_id)
842 && record.draft_input == expected
843 {
844 record.draft_input.clear();
845 }
846 self.publish_revision();
847 }
848 if let Some(bundle_id) = bundle_id {
849 let history_id = session_id.to_owned();
850 let history_text = text.to_owned();
851 if let Err(error) = blocking(move || {
852 crate::database::record_prompt(
853 &history_id,
854 &bundle_id,
855 ordinal,
856 None,
857 &history_text,
858 )
859 })
860 .await
861 {
862 tracing::warn!(
863 session_id,
864 error = format!("{error:#}"),
865 "the queued prompt was accepted but its history could not be stored"
866 );
867 }
868 }
869 Ok(())
870 }
871
872 async fn fail_startup_queue(
876 &self,
877 session_id: &str,
878 failed: Option<StartupStep>,
879 reason: &str,
880 ) {
881 let remaining = self
882 .startup_prompts
883 .lock()
884 .unwrap_or_else(PoisonError::into_inner)
885 .remove(session_id)
886 .map(|queue| queue.pending)
887 .unwrap_or_default();
888 let mut texts = Vec::new();
889 let mut dropped_handoff = false;
890 for step in failed.into_iter().chain(remaining) {
891 match step {
892 StartupStep::Prompt { text, .. } => texts.push(text),
893 StartupStep::InstallHandoff(_) => dropped_handoff = true,
894 }
895 }
896 if dropped_handoff {
897 tracing::warn!(
898 session_id,
899 reason,
900 "could not install the restored archive's hand-off"
901 );
902 self.push_notice(
903 session_id,
904 format!("The restored session started without its archived hand-off: {reason}"),
905 );
906 }
907 if texts.is_empty() {
908 return;
909 }
910 let restored = texts.join("\n\n");
911 if let Err(error) = self.append_draft_input(session_id, &restored).await {
912 tracing::warn!(
913 session_id,
914 error = format!("{error:#}"),
915 "a queued prompt could not be saved back into the session's draft"
916 );
917 }
918 self.push_notice(
919 session_id,
920 format!(
921 "Your prompt could not be sent to session {} ({reason}); it is back in the composer draft.",
922 mj_core::state::short_id(session_id)
923 ),
924 );
925 tracing::warn!(
926 session_id,
927 reason,
928 "a queued startup prompt could not be delivered"
929 );
930 }
931
932 pub(super) async fn append_draft_input(&self, session_id: &str, text: &str) -> Result<()> {
937 let existing = self
938 .session_record(session_id)
939 .map(|record| record.draft_input)
940 .unwrap_or_default();
941 let combined = [existing.as_str(), text]
942 .into_iter()
943 .filter(|part| !part.is_empty())
944 .collect::<Vec<_>>()
945 .join("\n\n");
946 let persisted_id = session_id.to_owned();
947 let persisted = combined.clone();
948 let stored =
949 blocking(move || crate::database::set_session_draft_input(&persisted_id, &persisted))
950 .await;
951 if let Some(record) = self
952 .controller
953 .lock()
954 .unwrap_or_else(PoisonError::into_inner)
955 .state
956 .sessions
957 .get_mut(session_id)
958 {
959 record.draft_input = combined;
960 }
961 self.publish_revision();
962 stored
963 }
964
965 pub(crate) async fn cancel_and_join_startup_prompts(&self) -> Result<()> {
968 let tasks = {
969 let mut queues = self
970 .startup_prompts
971 .lock()
972 .unwrap_or_else(PoisonError::into_inner);
973 queues
974 .values_mut()
975 .filter_map(|queue| {
976 queue.cancel.cancel();
977 queue.task.take()
978 })
979 .collect::<Vec<_>>()
980 };
981 let deadline = tokio::time::Instant::now() + Duration::from_secs(1);
982 let mut outcome = Ok(());
983 for task in tasks {
984 let joined = match tokio::time::timeout_at(deadline, task).await {
985 Ok(Ok(())) => Ok(()),
986 Ok(Err(error)) => Err(anyhow!("startup prompt delivery task failed: {error}")),
987 Err(_) => Err(anyhow!(
988 "a startup prompt delivery task did not stop within 1s"
989 )),
990 };
991 if outcome.is_ok() {
992 outcome = joined;
993 } else if let Err(error) = joined {
994 tracing::warn!(%error, "another startup prompt drain did not stop cleanly");
995 }
996 }
997 outcome
998 }
999}