1use std::fmt;
4use std::future::Future;
5use std::pin::Pin;
6use std::sync::Arc;
7use std::time::Duration;
8
9use agent_client_protocol::schema::v1::SessionConfigOption;
10use anyhow::{Context, Result, ensure};
11use mj_core::config::Config;
12use mj_core::elicitation::ElicitationResponse;
13use mj_core::state::{ManagedSessionSnapshot, SessionRecord};
14
15use mj_core::relay::{
16 AnalyzeDeltaRepository, RelayCommand, RelayCursor, RelayEvent, RelayOperationalState, RepoDelta,
17};
18use mj_core::worker_launch::ReviewerLaunchConfig;
19
20pub type BoxFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
21
22#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
23#[serde(tag = "kind", content = "detail", rename_all = "snake_case")]
24pub enum ViewError {
25 Unreachable(String),
26 TargetMissing(String),
27 ProjectionIntegrity(String),
28}
29
30impl ViewError {
31 pub fn detail(&self) -> &str {
32 match self {
33 Self::Unreachable(detail)
34 | Self::TargetMissing(detail)
35 | Self::ProjectionIntegrity(detail) => detail,
36 }
37 }
38}
39
40#[derive(Debug, Clone, PartialEq, Default)]
41pub struct ManagedSessionView {
42 pub snapshot: Option<ManagedSessionSnapshot>,
43 pub connected: bool,
44 pub error: Option<ViewError>,
45}
46
47#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
48pub struct RelayAttachment {
49 pub state: RelayOperationalState,
50 pub events: Vec<RelayEvent>,
51 pub through_ordinal: u64,
52 pub through_digest: String,
53}
54
55#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
56pub struct StartedReviewer {
57 pub native_session_id: Option<String>,
58 pub config_options: Vec<SessionConfigOption>,
59 pub reused: bool,
60 pub state: RelayOperationalState,
61}
62
63#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
64#[serde(rename_all = "snake_case")]
65pub enum ReviewerAction {
66 Start {
67 config: Box<ReviewerLaunchConfig>,
68 },
69 Submit {
70 command_id: String,
71 command: RelayCommand,
72 },
73 Attach {
74 after_ordinal: u64,
75 after_digest: String,
76 },
77 Acknowledge {
78 through_ordinal: u64,
79 through_digest: String,
80 },
81 Status,
82 RespondElicitation {
83 elicitation_id: String,
84 response: ElicitationResponse,
85 },
86 Pause,
87 PauseGeneration {
89 generation: u64,
90 },
91 CaptureDelta {
92 baselines: std::collections::BTreeMap<std::path::PathBuf, String>,
93 },
94 AdvanceBaseline {
95 trees: std::collections::BTreeMap<std::path::PathBuf, String>,
96 },
97 AnalyzeDelta {
98 repositories: Vec<AnalyzeDeltaRepository>,
99 },
100 TakeLaneDispatches,
101}
102
103impl ReviewerAction {
104 pub const fn operation_name(&self) -> &'static str {
105 match self {
106 Self::Start { .. } => "reviewer_start",
107 Self::Submit { .. } => "reviewer_submit",
108 Self::Attach { .. } => "reviewer_attach",
109 Self::Acknowledge { .. } => "reviewer_acknowledge",
110 Self::Status => "reviewer_status",
111 Self::RespondElicitation { .. } => "reviewer_respond_elicitation",
112 Self::Pause | Self::PauseGeneration { .. } => "reviewer_pause",
113 Self::CaptureDelta { .. } => "reviewer_capture_delta",
114 Self::AdvanceBaseline { .. } => "reviewer_advance_baseline",
115 Self::AnalyzeDelta { .. } => "reviewer_analyze_delta",
116 Self::TakeLaneDispatches => "reviewer_take_lane_dispatches",
117 }
118 }
119}
120
121#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
122#[serde(rename_all = "snake_case")]
123pub enum ReviewerOutcome {
124 Started(Box<StartedReviewer>),
125 Accepted {
126 ordinal: u64,
127 },
128 Attached(Box<RelayAttachment>),
129 Acknowledged(RelayCursor),
130 Status(Box<RelayOperationalState>),
131 ElicitationResolved,
132 Paused,
133 Delta {
134 repositories: Vec<RepoDelta>,
135 },
136 BaselineAdvanced,
137 ChangedFunctions {
138 packet: String,
139 },
140 LaneDispatches {
141 requests: Vec<mj_core::review::lanes::ReviewSubagentRequest>,
142 },
143}
144
145#[derive(Debug)]
147pub struct SubmitFailure {
148 pub message: String,
149 pub unconfirmed: bool,
150}
151impl From<String> for SubmitFailure {
152 fn from(message: String) -> Self {
153 Self {
154 message,
155 unconfirmed: false,
156 }
157 }
158}
159impl From<&str> for SubmitFailure {
160 fn from(message: &str) -> Self {
161 message.to_owned().into()
162 }
163}
164
165#[derive(Debug)]
167pub struct DeliveryUnconfirmed;
168impl std::fmt::Display for DeliveryUnconfirmed {
169 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
170 f.write_str("delivery unconfirmed")
171 }
172}
173impl std::error::Error for DeliveryUnconfirmed {}
174
175pub struct PendingRelaySubmit {
176 completion: BoxFuture<'static, Result<u64>>,
177}
178
179impl PendingRelaySubmit {
180 pub fn new(completion: BoxFuture<'static, Result<u64>>) -> Self {
181 Self { completion }
182 }
183
184 pub async fn wait(self) -> Result<u64> {
185 self.completion.await
186 }
187}
188
189pub struct PendingRelaySync {
190 completion: BoxFuture<'static, Result<()>>,
191}
192
193impl PendingRelaySync {
194 pub fn new(completion: BoxFuture<'static, Result<()>>) -> Self {
195 Self { completion }
196 }
197
198 pub async fn wait(self) -> Result<()> {
199 self.completion.await
200 }
201}
202
203#[derive(Debug, Default)]
204pub struct ReviewState {
205 pub review: Option<mj_core::storage::StoredReview>,
206}
207
208pub trait SessionHandleBackend: Send + Sync {
209 fn search_prompts(
210 &self,
211 bundle_id: String,
212 scope: mj_core::storage::HistoryScope,
213 query: String,
214 ) -> BoxFuture<'_, Result<Vec<mj_core::storage::PromptHistoryEntry>>>;
215 fn review_state(&self) -> BoxFuture<'_, Result<ReviewState>>;
216 fn resolve_review_settings(
217 &self,
218 cancelled: Arc<std::sync::atomic::AtomicBool>,
219 ) -> BoxFuture<'_, Result<mj_core::review::settings::ResolvedReviewSettings>> {
220 let _ = cancelled;
221 Box::pin(async { anyhow::bail!("review settings resolution is unavailable") })
222 }
223
224 fn config_result(&self, command_id: String) -> BoxFuture<'_, Result<Option<Option<String>>>>;
225
226 fn clone_box(&self) -> Box<dyn SessionHandleBackend>;
227 fn session_id(&self) -> &str;
228 fn view(&self) -> ManagedSessionView;
229 fn is_stopped(&self) -> bool;
230 fn has_changed(&self) -> Result<bool>;
231 fn changed(&mut self) -> BoxFuture<'_, Result<ManagedSessionView>>;
232 fn enqueue_submit(
233 &self,
234 command_id: String,
235 command: RelayCommand,
236 ) -> BoxFuture<'_, Result<PendingRelaySubmit>>;
237 fn enqueue_sync(&self) -> BoxFuture<'_, Result<PendingRelaySync>>;
238 fn respond_elicitation(
239 &self,
240 elicitation_id: String,
241 response: ElicitationResponse,
242 ) -> BoxFuture<'_, Result<()>>;
243 fn stop_background_task(&self, background_task_id: String) -> BoxFuture<'_, Result<()>>;
244 fn reviewer(
245 &self,
246 role: Option<String>,
247 action: ReviewerAction,
248 ) -> BoxFuture<'_, Result<ReviewerOutcome>>;
249}
250
251pub struct SessionHandle {
252 backend: Box<dyn SessionHandleBackend>,
253}
254
255impl SessionHandle {
256 pub async fn search_prompts(
257 &self,
258 bundle_id: String,
259 scope: mj_core::storage::HistoryScope,
260 query: String,
261 ) -> Result<Vec<mj_core::storage::PromptHistoryEntry>> {
262 self.backend.search_prompts(bundle_id, scope, query).await
263 }
264 pub async fn review_state(&self) -> Result<ReviewState> {
265 self.backend.review_state().await
266 }
267
268 pub async fn resolve_review_settings(
269 &self,
270 cancelled: Arc<std::sync::atomic::AtomicBool>,
271 ) -> Result<mj_core::review::settings::ResolvedReviewSettings> {
272 self.backend.resolve_review_settings(cancelled).await
273 }
274
275 pub fn new(backend: impl SessionHandleBackend + 'static) -> Self {
276 Self {
277 backend: Box::new(backend),
278 }
279 }
280
281 pub fn session_id(&self) -> &str {
282 self.backend.session_id()
283 }
284
285 pub fn view(&self) -> ManagedSessionView {
286 self.backend.view()
287 }
288
289 pub fn is_stopped(&self) -> bool {
290 self.backend.is_stopped()
291 }
292
293 pub fn has_changed(&self) -> Result<bool> {
294 self.backend.has_changed()
295 }
296
297 pub async fn changed(&mut self) -> Result<ManagedSessionView> {
298 self.backend.changed().await
299 }
300
301 pub async fn submit(&self, command_id: String, command: RelayCommand) -> Result<u64> {
302 self.enqueue_submit(command_id, command).await?.wait().await
303 }
304
305 pub async fn set_config(&self, key: String, value: String) -> Result<()> {
307 self.set_config_with_id(new_command_id("set-config")?, key, value)
308 .await
309 }
310
311 pub async fn set_config_with_id(
312 &self,
313 command_id: String,
314 key: String,
315 value: String,
316 ) -> Result<()> {
317 self.submit(command_id.clone(), RelayCommand::SetConfig { key, value })
318 .await?;
319 tokio::time::timeout(Duration::from_secs(60), async {
320 loop {
321 if let Some(error) = self.backend.config_result(command_id.clone()).await? {
322 if let Some(error) = error {
323 anyhow::bail!("{error}");
324 }
325 self.sync_now().await?;
326 return Ok(());
327 }
328 ensure!(
329 !self.is_stopped(),
330 "session stopped while applying configuration"
331 );
332 if let Some(error) = self.view().error {
333 anyhow::bail!("configuration connection failed: {}", error.detail());
334 }
335 tokio::time::sleep(Duration::from_millis(50)).await;
336 }
337 })
338 .await
339 .context("configuration command did not complete within 60 seconds")?
340 }
341
342 pub async fn apply_plan_control(
344 &self,
345 command_id: String,
346 control: mj_core::acp::PlanControl,
347 ) -> Result<()> {
348 match control {
349 mj_core::acp::PlanControl::SetConfig { key, value } => {
350 self.set_config_with_id(command_id, key, value).await
351 }
352 mj_core::acp::PlanControl::SetSessionMode { mode_id } => {
353 self.submit(
354 command_id,
355 RelayCommand::SetSessionMode {
356 mode_id: mode_id.clone(),
357 },
358 )
359 .await?;
360 tokio::time::timeout(Duration::from_secs(60), async {
361 loop {
362 self.sync_now().await?;
363 if self
364 .view()
365 .snapshot
366 .as_ref()
367 .and_then(|snapshot| snapshot.operational.modes.as_ref())
368 .is_some_and(|modes| modes.current_mode_id.to_string() == mode_id)
369 {
370 return Ok(());
371 }
372 ensure!(!self.is_stopped(), "session stopped while changing mode");
373 tokio::time::sleep(Duration::from_millis(50)).await;
374 }
375 })
376 .await
377 .context("session mode change did not complete within 60 seconds")?
378 }
379 }
380 }
381
382 pub async fn enqueue_submit(
383 &self,
384 command_id: String,
385 command: RelayCommand,
386 ) -> Result<PendingRelaySubmit> {
387 self.backend.enqueue_submit(command_id, command).await
388 }
389
390 pub async fn sync_now(&self) -> Result<()> {
391 self.enqueue_sync().await?.wait().await
392 }
393
394 pub async fn enqueue_sync(&self) -> Result<PendingRelaySync> {
395 self.backend.enqueue_sync().await
396 }
397
398 pub async fn respond_elicitation(
399 &self,
400 elicitation_id: String,
401 response: ElicitationResponse,
402 ) -> Result<()> {
403 self.backend
404 .respond_elicitation(elicitation_id, response)
405 .await
406 }
407
408 pub async fn stop_background_task(&self, background_task_id: String) -> Result<()> {
409 self.backend.stop_background_task(background_task_id).await
410 }
411
412 pub async fn reviewer(&self, action: ReviewerAction) -> Result<ReviewerOutcome> {
413 self.reviewer_as(None, action).await
414 }
415
416 pub async fn reviewer_as(
417 &self,
418 role: Option<String>,
419 action: ReviewerAction,
420 ) -> Result<ReviewerOutcome> {
421 self.backend.reviewer(role, action).await
422 }
423}
424
425impl Clone for SessionHandle {
426 fn clone(&self) -> Self {
427 Self {
428 backend: self.backend.clone_box(),
429 }
430 }
431}
432
433impl fmt::Debug for SessionHandle {
434 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
435 formatter
436 .debug_struct("SessionHandle")
437 .field("session_id", &self.session_id())
438 .finish_non_exhaustive()
439 }
440}
441
442pub trait SessionControlBackend: Send + Sync {
443 fn session(&self, session_id: String) -> BoxFuture<'_, Result<SessionHandle>>;
444}
445
446#[derive(Clone)]
447pub struct SessionControl {
448 backend: Arc<dyn SessionControlBackend>,
449}
450
451impl SessionControl {
452 pub fn new(backend: impl SessionControlBackend + 'static) -> Self {
453 Self {
454 backend: Arc::new(backend),
455 }
456 }
457
458 pub async fn session(&self, session_id: impl Into<String>) -> Result<SessionHandle> {
459 self.backend.session(session_id.into()).await
460 }
461
462 pub async fn wait_for_session(
463 &self,
464 session_id: &str,
465 timeout: Duration,
466 ) -> Result<SessionHandle> {
467 tokio::time::timeout(timeout, async {
468 loop {
469 match self.session(session_id.to_owned()).await {
470 Ok(handle) => return Ok(handle),
471 Err(error) => {
472 tracing::trace!(session_id, "waiting for session actor: {error:#}");
473 tokio::time::sleep(Duration::from_millis(25)).await;
474 }
475 }
476 }
477 })
478 .await
479 .with_context(|| {
480 format!(
481 "session {session_id} did not become available within {} seconds",
482 timeout.as_secs()
483 )
484 })?
485 }
486}
487
488impl fmt::Debug for SessionControl {
489 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
490 formatter.write_str("SessionControl(..)")
491 }
492}
493
494pub trait ReviewerStagerBackend: Send + Sync {
495 fn stage(
496 &self,
497 config: Config,
498 session: SessionRecord,
499 profile_id: String,
500 generation: u64,
501 cancelled: Arc<std::sync::atomic::AtomicBool>,
502 ) -> Result<ReviewerLaunchConfig>;
503}
504
505#[derive(Clone)]
506pub struct ReviewerStager {
507 backend: Arc<dyn ReviewerStagerBackend>,
508}
509
510impl ReviewerStager {
511 pub fn new(backend: impl ReviewerStagerBackend + 'static) -> Self {
512 Self {
513 backend: Arc::new(backend),
514 }
515 }
516
517 pub fn stage(
518 &self,
519 config: Config,
520 session: SessionRecord,
521 profile_id: String,
522 generation: u64,
523 cancelled: Arc<std::sync::atomic::AtomicBool>,
524 ) -> Result<ReviewerLaunchConfig> {
525 self.backend
526 .stage(config, session, profile_id, generation, cancelled)
527 }
528
529 #[doc(hidden)]
530 pub fn unavailable(message: impl Into<String>) -> Self {
531 Self::new(UnavailableReviewerStager(message.into()))
532 }
533}
534
535impl fmt::Debug for ReviewerStager {
536 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
537 formatter.write_str("ReviewerStager(..)")
538 }
539}
540
541struct UnavailableReviewerStager(String);
542
543impl ReviewerStagerBackend for UnavailableReviewerStager {
544 fn stage(
545 &self,
546 _config: Config,
547 _session: SessionRecord,
548 _profile_id: String,
549 _generation: u64,
550 _cancelled: Arc<std::sync::atomic::AtomicBool>,
551 ) -> Result<ReviewerLaunchConfig> {
552 anyhow::bail!(self.0.clone())
553 }
554}
555
556pub fn new_command_id(prefix: &str) -> Result<String> {
557 ensure!(!prefix.trim().is_empty(), "command ID prefix is required");
558 let mut random = [0_u8; 16];
559 getrandom::fill(&mut random)
560 .map_err(|error| anyhow::anyhow!("generate command ID: {error}"))?;
561 Ok(format!("{prefix}-{}", mj_core::hex::lower_hex(random)))
562}
563
564#[doc(hidden)]
569pub struct ReplacementSessionTestFixture {
570 pub stopped: SessionHandle,
571 pub control: SessionControl,
572 pub submitted: tokio::sync::mpsc::UnboundedReceiver<RelayCommand>,
573}
574
575#[derive(Clone)]
576struct ReplacementTestSession {
577 #[cfg(test)]
578 history: Option<tokio::sync::mpsc::UnboundedSender<HistoryTestRequest>>,
579 session_id: String,
580 stopped: bool,
581 accepted_ordinal: u64,
582 submitted: Option<tokio::sync::mpsc::UnboundedSender<RelayCommand>>,
583 view: tokio::sync::watch::Receiver<ManagedSessionView>,
584 _view_guard: Option<Arc<tokio::sync::watch::Sender<ManagedSessionView>>>,
585}
586
587impl SessionHandleBackend for ReplacementTestSession {
588 fn search_prompts(
589 &self,
590 _bundle_id: String,
591 _scope: mj_core::storage::HistoryScope,
592 _query: String,
593 ) -> BoxFuture<'_, Result<Vec<mj_core::storage::PromptHistoryEntry>>> {
594 #[cfg(test)]
595 if let Some(history) = &self.history {
596 let (response, result) = tokio::sync::oneshot::channel();
597 let sent = history.send(HistoryTestRequest {
598 bundle_id: _bundle_id,
599 scope: _scope,
600 query: _query,
601 response,
602 });
603 return Box::pin(async move {
604 sent.map_err(|_| anyhow::anyhow!("history backend closed"))?;
605 result
606 .await
607 .map_err(|_| anyhow::anyhow!("history response dropped"))?
608 });
609 }
610 Box::pin(async { Ok(Vec::new()) })
611 }
612 fn review_state(&self) -> BoxFuture<'_, Result<ReviewState>> {
613 Box::pin(async { Ok(ReviewState::default()) })
614 }
615
616 fn config_result(&self, _command_id: String) -> BoxFuture<'_, Result<Option<Option<String>>>> {
617 Box::pin(async { Ok(None) })
618 }
619 fn clone_box(&self) -> Box<dyn SessionHandleBackend> {
620 Box::new(self.clone())
621 }
622
623 fn session_id(&self) -> &str {
624 &self.session_id
625 }
626
627 fn view(&self) -> ManagedSessionView {
628 self.view.borrow().clone()
629 }
630
631 fn is_stopped(&self) -> bool {
632 self.stopped
633 }
634
635 fn has_changed(&self) -> Result<bool> {
636 self.view.has_changed().context("session manager stopped")
637 }
638
639 fn changed(&mut self) -> BoxFuture<'_, Result<ManagedSessionView>> {
640 Box::pin(async move {
641 self.view
642 .changed()
643 .await
644 .context("session manager stopped")?;
645 Ok(self.view())
646 })
647 }
648
649 fn enqueue_submit(
650 &self,
651 _command_id: String,
652 command: RelayCommand,
653 ) -> BoxFuture<'_, Result<PendingRelaySubmit>> {
654 let submitted = self.submitted.clone();
655 let stopped = self.stopped;
656 let accepted_ordinal = self.accepted_ordinal;
657 Box::pin(async move {
658 ensure!(!stopped, "session manager stopped");
659 let submitted = submitted.context("unsupported test operation")?;
660 submitted
661 .send(command)
662 .context("test submit observer stopped")?;
663 Ok(PendingRelaySubmit::new(Box::pin(async move {
664 Ok(accepted_ordinal)
665 })))
666 })
667 }
668
669 fn enqueue_sync(&self) -> BoxFuture<'_, Result<PendingRelaySync>> {
670 let stopped = self.stopped;
671 Box::pin(async move {
672 ensure!(!stopped, "session manager stopped");
673 Ok(PendingRelaySync::new(Box::pin(async { Ok(()) })))
674 })
675 }
676
677 fn respond_elicitation(
678 &self,
679 _elicitation_id: String,
680 _response: ElicitationResponse,
681 ) -> BoxFuture<'_, Result<()>> {
682 Box::pin(async { anyhow::bail!("unsupported test operation") })
683 }
684
685 fn stop_background_task(&self, _background_task_id: String) -> BoxFuture<'_, Result<()>> {
686 Box::pin(async { anyhow::bail!("unsupported test operation") })
687 }
688
689 fn reviewer(
690 &self,
691 _role: Option<String>,
692 _action: ReviewerAction,
693 ) -> BoxFuture<'_, Result<ReviewerOutcome>> {
694 Box::pin(async { anyhow::bail!("unsupported test operation") })
695 }
696}
697
698struct ReplacementTestControl {
699 session_id: String,
700 replacement: SessionHandle,
701}
702
703impl SessionControlBackend for ReplacementTestControl {
704 fn session(&self, session_id: String) -> BoxFuture<'_, Result<SessionHandle>> {
705 Box::pin(async move {
706 ensure!(
707 session_id == self.session_id,
708 "session {session_id} is not managed"
709 );
710 Ok(self.replacement.clone())
711 })
712 }
713}
714
715#[doc(hidden)]
716pub fn replacement_session_test_fixture(
717 session_id: &str,
718 accepted_ordinal: u64,
719) -> ReplacementSessionTestFixture {
720 let (stopped_view_tx, stopped_view) =
721 tokio::sync::watch::channel(ManagedSessionView::default());
722 drop(stopped_view_tx);
723 let stopped = SessionHandle::new(ReplacementTestSession {
724 #[cfg(test)]
725 history: None,
726 session_id: session_id.to_owned(),
727 stopped: true,
728 accepted_ordinal,
729 submitted: None,
730 view: stopped_view,
731 _view_guard: None,
732 });
733
734 let (view_tx, view) = tokio::sync::watch::channel(ManagedSessionView::default());
735 let (submitted_tx, submitted) = tokio::sync::mpsc::unbounded_channel();
736 let replacement = SessionHandle::new(ReplacementTestSession {
737 #[cfg(test)]
738 history: None,
739 session_id: session_id.to_owned(),
740 stopped: false,
741 accepted_ordinal,
742 submitted: Some(submitted_tx),
743 view,
744 _view_guard: Some(Arc::new(view_tx)),
745 });
746 let control = SessionControl::new(ReplacementTestControl {
747 session_id: session_id.to_owned(),
748 replacement,
749 });
750 ReplacementSessionTestFixture {
751 stopped,
752 control,
753 submitted,
754 }
755}
756
757#[cfg(test)]
758struct HistoryTestRequest {
759 bundle_id: String,
760 scope: mj_core::storage::HistoryScope,
761 query: String,
762 response: tokio::sync::oneshot::Sender<Result<Vec<mj_core::storage::PromptHistoryEntry>>>,
763}
764
765#[cfg(test)]
766mod storage_tests {
767 use super::*;
768 use mj_core::storage::{HistoryScope, PromptHistoryEntry};
769
770 #[tokio::test]
771 async fn history_search_yields_until_backend_responds_and_propagates_failures() {
772 let (history, mut requests) = tokio::sync::mpsc::unbounded_channel();
773 let (view_guard, view) = tokio::sync::watch::channel(ManagedSessionView::default());
774 let session = SessionHandle::new(ReplacementTestSession {
775 history: Some(history),
776 session_id: "session".into(),
777 stopped: false,
778 accepted_ordinal: 0,
779 submitted: None,
780 view,
781 _view_guard: Some(Arc::new(view_guard)),
782 });
783 let search =
784 session.search_prompts("bundle".into(), HistoryScope::Project, "needle".into());
785 tokio::pin!(search);
786 let request = tokio::select! {
787 biased;
788 result = &mut search => panic!("search completed before storage replied: {result:?}"),
789 request = requests.recv() => request.unwrap(),
790 };
791 assert_eq!(request.bundle_id, "bundle");
792 assert_eq!(request.scope, HistoryScope::Project);
793 assert_eq!(request.query, "needle");
794 request
795 .response
796 .send(Err(anyhow::anyhow!("storage unavailable")))
797 .unwrap();
798 assert!(
799 search
800 .await
801 .unwrap_err()
802 .to_string()
803 .contains("storage unavailable")
804 );
805
806 let search = session.search_prompts("bundle".into(), HistoryScope::Project, "retry".into());
807 tokio::pin!(search);
808 let request = tokio::select! {
809 biased;
810 result = &mut search => panic!("retry completed before storage replied: {result:?}"),
811 request = requests.recv() => request.unwrap(),
812 };
813 request
814 .response
815 .send(Ok(vec![PromptHistoryEntry {
816 id: 1,
817 session_id: "session".into(),
818 text: "retry works".into(),
819 }]))
820 .unwrap();
821 assert_eq!(search.await.unwrap()[0].text, "retry works");
822 }
823}