Skip to main content

mj_controller/session_manager/
client_backend.rs

1use super::*;
2
3#[derive(Clone)]
4pub(super) struct ClientSessionHandle(pub(super) ManagedSessionHandle);
5
6impl mj_client::session::SessionHandleBackend for ClientSessionHandle {
7    fn transcript_history(
8        &self,
9        before: Option<mj_core::storage::TranscriptCursor>,
10    ) -> mj_client::session::BoxFuture<'_, Result<mj_core::storage::TranscriptHistoryPage>> {
11        let session_id = self.0.session_id().to_owned();
12        Box::pin(async move {
13            tokio::task::spawn_blocking(move || {
14                crate::database::load_transcript_history(&session_id, before.as_ref(), 128)?
15                    .context("session history is unavailable")
16            })
17            .await
18            .context("earlier conversation history task failed")?
19        })
20    }
21
22    fn search_prompts(
23        &self,
24        bundle_id: String,
25        scope: mj_core::storage::HistoryScope,
26        query: String,
27    ) -> mj_client::session::BoxFuture<'_, Result<Vec<mj_core::storage::PromptHistoryEntry>>> {
28        let session_id = self.0.session_id().to_owned();
29        Box::pin(async move {
30            tokio::task::spawn_blocking(move || {
31                crate::database::search_prompts(&session_id, &bundle_id, scope, &query)
32            })
33            .await
34            .context("history search task")?
35        })
36    }
37    fn review_state(
38        &self,
39    ) -> mj_client::session::BoxFuture<'_, Result<mj_client::session::ReviewState>> {
40        let session_id = self.0.session_id().to_owned();
41        Box::pin(async move {
42            tokio::task::spawn_blocking(move || {
43                Ok(mj_client::session::ReviewState {
44                    review: crate::database::active_review(&session_id)?,
45                })
46            })
47            .await
48            .context("review restoration task")?
49        })
50    }
51
52    fn resolve_review_settings(
53        &self,
54        cancelled: std::sync::Arc<std::sync::atomic::AtomicBool>,
55    ) -> mj_client::session::BoxFuture<'_, Result<mj_core::review::settings::ResolvedReviewSettings>>
56    {
57        Box::pin(crate::review_selection::resolve(
58            self.0.clone(),
59            None,
60            false,
61            cancelled,
62        ))
63    }
64
65    fn config_result(
66        &self,
67        command_id: String,
68    ) -> mj_client::session::BoxFuture<'_, Result<Option<Option<String>>>> {
69        let session_id = self.session_id().to_owned();
70        Box::pin(async move {
71            tokio::task::spawn_blocking(move || {
72                crate::database::load_config_result(&session_id, &command_id)
73            })
74            .await
75            .context("read configuration completion task")?
76        })
77    }
78
79    fn clone_box(&self) -> Box<dyn mj_client::session::SessionHandleBackend> {
80        Box::new(self.clone())
81    }
82
83    fn session_id(&self) -> &str {
84        self.0.session_id()
85    }
86
87    fn view(&self) -> ManagedSessionView {
88        self.0.view()
89    }
90
91    fn is_stopped(&self) -> bool {
92        self.0.is_stopped()
93    }
94
95    fn has_changed(&self) -> Result<bool> {
96        self.0.has_changed()
97    }
98
99    fn changed(&mut self) -> mj_client::session::BoxFuture<'_, Result<ManagedSessionView>> {
100        Box::pin(self.0.changed())
101    }
102
103    fn enqueue_submit(
104        &self,
105        command_id: String,
106        command: RelayCommand,
107    ) -> mj_client::session::BoxFuture<'_, Result<mj_client::session::PendingRelaySubmit>> {
108        Box::pin(async move {
109            let pending = self.0.enqueue_submit(command_id, command).await?;
110            Ok(mj_client::session::PendingRelaySubmit::new(Box::pin(
111                pending.wait(),
112            )))
113        })
114    }
115
116    fn enqueue_sync(
117        &self,
118    ) -> mj_client::session::BoxFuture<'_, Result<mj_client::session::PendingRelaySync>> {
119        Box::pin(async move {
120            let pending = self.0.enqueue_sync().await?;
121            Ok(mj_client::session::PendingRelaySync::new(Box::pin(
122                pending.wait(),
123            )))
124        })
125    }
126
127    fn respond_elicitation(
128        &self,
129        elicitation_id: String,
130        response: ElicitationResponse,
131    ) -> mj_client::session::BoxFuture<'_, Result<()>> {
132        Box::pin(self.0.respond_elicitation(elicitation_id, response))
133    }
134
135    fn stop_background_task(
136        &self,
137        background_task_id: String,
138    ) -> mj_client::session::BoxFuture<'_, Result<()>> {
139        Box::pin(self.0.stop_background_task(background_task_id))
140    }
141
142    fn reviewer(
143        &self,
144        role: Option<String>,
145        action: ReviewerAction,
146    ) -> mj_client::session::BoxFuture<'_, Result<ReviewerOutcome>> {
147        Box::pin(self.0.reviewer_as(role, action))
148    }
149}
150
151#[derive(Clone)]
152pub(super) struct ClientSessionControl(pub(super) SessionManagerControl);
153
154impl mj_client::session::SessionControlBackend for ClientSessionControl {
155    fn session(
156        &self,
157        session_id: String,
158    ) -> mj_client::session::BoxFuture<'_, Result<mj_client::session::SessionHandle>> {
159        Box::pin(async move { Ok(self.0.session(session_id).await?.client()) })
160    }
161}
162
163impl SessionManagerControl {
164    /// Narrow this controller-owned manager to session lookup for a control surface.
165    pub fn client(&self) -> mj_client::session::SessionControl {
166        mj_client::session::SessionControl::new(ClientSessionControl(self.clone()))
167    }
168
169    pub async fn session(&self, session_id: impl Into<String>) -> Result<ManagedSessionHandle> {
170        let session_id = session_id.into();
171        let (reply, response) = oneshot::channel();
172        self.commands
173            .send(ManagerCommand::Session {
174                session_id: session_id.clone(),
175                reply,
176            })
177            .await
178            .context("session manager stopped")?;
179        response
180            .await
181            .context("session manager stopped")?
182            .with_context(|| format!("session {session_id} is not managed"))
183    }
184
185    pub async fn wait_for_session(
186        &self,
187        session_id: &str,
188        timeout: Duration,
189    ) -> Result<ManagedSessionHandle> {
190        tokio::time::timeout(timeout, async {
191            loop {
192                match self.session(session_id.to_owned()).await {
193                    Ok(handle) => return Ok(handle),
194                    Err(error) => {
195                        tracing::trace!(session_id, "waiting for session actor: {error:#}");
196                        tokio::time::sleep(Duration::from_millis(25)).await;
197                    }
198                }
199            }
200        })
201        .await
202        .with_context(|| {
203            format!(
204                "session {session_id} did not become available within {} seconds",
205                timeout.as_secs()
206            )
207        })?
208    }
209}