Skip to main content

mj_client/
session.rs

1//! The session operations used by interactive control surfaces.
2
3use 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    /// Stop only the preparation that owns this generation; stale cleanup must not stop its replacement.
88    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/// Preserve whether a submit failed before delivery or lost its acknowledgement.
146#[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/// Submission lost its acknowledgement; callers must reconcile before retrying.
166#[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    /// Apply a setting and wait for its durable success or rejection.
306    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    /// A follow-up prompt must wait for the mode change, not just admission.
343    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/// A stopped session and a manager that resolves its live replacement.
565///
566/// Chat's cross-crate tests use this hand-written client fake to verify actor
567/// replacement without depending on the controller implementation crate.
568#[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}