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