mj_controller/session_manager/
handle.rs1use super::*;
2
3#[derive(Clone)]
4pub struct SessionManagerControl {
5 pub(super) commands: mpsc::Sender<ManagerCommand>,
6}
7
8#[derive(Clone, Debug)]
9pub struct ManagedSessionHandle {
10 pub(super) session_id: String,
11 pub(super) commands: mpsc::Sender<ActorCommand>,
12 pub(super) releases: mpsc::UnboundedSender<ReturnedConnection>,
13 pub(super) view: watch::Receiver<ManagedSessionView>,
14}
15
16#[derive(Debug, Clone, PartialEq, Eq)]
20pub struct ReviewDeliveryAdmission {
21 pub(super) session_id: String,
22 pub(super) epoch: u64,
23 pub(super) command_id: String,
24}
25
26impl ReviewDeliveryAdmission {
27 pub(crate) fn new(session_id: String, epoch: u64, command_id: String) -> Self {
28 Self {
29 session_id,
30 epoch,
31 command_id,
32 }
33 }
34
35 pub(crate) fn session_id(&self) -> &str {
36 &self.session_id
37 }
38
39 pub(crate) const fn epoch(&self) -> u64 {
40 self.epoch
41 }
42
43 pub(crate) fn command_id(&self) -> &str {
44 &self.command_id
45 }
46}
47
48pub struct ManagedSessionLease {
58 pub(super) session_id: String,
59 pub(super) lease_id: Option<u64>,
60 pub(super) connection: Option<StandaloneSession>,
61 pub(super) releases: mpsc::UnboundedSender<ReturnedConnection>,
62}
63
64impl ManagedSessionLease {
65 pub fn connection_mut(&mut self) -> &mut StandaloneSession {
66 self.connection
67 .as_mut()
68 .expect("managed session lease has already been released")
69 }
70
71 pub fn replace_connection(&mut self, connection: StandaloneSession) {
74 drop(self.connection.take());
75 self.connection = Some(connection);
76 }
77
78 pub fn release(mut self) {
79 let lease_id = self
80 .lease_id
81 .take()
82 .expect("managed session lease has already been released");
83 let connection = self.connection.take();
84 if let Err(error) = self.releases.send(ReturnedConnection {
85 lease_id,
86 connection,
87 }) {
88 tracing::warn!(
89 session_id = %self.session_id,
90 operation = "lease_release",
91 %error,
92 "session actor stopped before receiving released relay connection"
93 );
94 }
95 }
96}
97
98impl Drop for ManagedSessionLease {
99 fn drop(&mut self) {
100 let Some(lease_id) = self.lease_id.take() else {
101 return;
102 };
103 drop(self.connection.take());
106 if let Err(error) = self.releases.send(ReturnedConnection {
107 lease_id,
108 connection: None,
109 }) {
110 tracing::warn!(
111 session_id = %self.session_id,
112 operation = "lease_drop",
113 %error,
114 "session actor stopped before receiving dropped relay lease"
115 );
116 }
117 }
118}
119
120impl ManagedSessionHandle {
121 pub fn client(&self) -> mj_client::session::SessionHandle {
123 mj_client::session::SessionHandle::new(ClientSessionHandle(self.clone()))
124 }
125
126 pub fn session_id(&self) -> &str {
127 &self.session_id
128 }
129
130 pub fn view(&self) -> ManagedSessionView {
131 self.view.borrow().clone()
132 }
133
134 pub fn is_stopped(&self) -> bool {
138 self.commands.is_closed()
139 }
140
141 pub fn has_changed(&self) -> Result<bool> {
142 self.view.has_changed().context("session manager stopped")
143 }
144
145 pub async fn changed(&mut self) -> Result<ManagedSessionView> {
146 self.view
147 .changed()
148 .await
149 .context("session manager stopped")?;
150 Ok(self.view())
151 }
152
153 pub async fn submit(&self, command_id: String, command: RelayCommand) -> Result<u64> {
154 self.enqueue_submit(command_id, command).await?.wait().await
155 }
156
157 pub(crate) async fn submit_review_delivery(
161 &self,
162 admission: ReviewDeliveryAdmission,
163 command: RelayCommand,
164 ) -> Result<u64> {
165 let command_id = admission.command_id.clone();
166 self.enqueue_submit_with_admission(command_id, command, Some(admission))
167 .await?
168 .wait()
169 .await
170 }
171
172 pub async fn enqueue_submit(
173 &self,
174 command_id: String,
175 command: RelayCommand,
176 ) -> Result<PendingRelaySubmit> {
177 self.enqueue_submit_with_admission(command_id, command, None)
178 .await
179 }
180
181 pub(super) async fn enqueue_submit_with_admission(
182 &self,
183 command_id: String,
184 command: RelayCommand,
185 admission: Option<ReviewDeliveryAdmission>,
186 ) -> Result<PendingRelaySubmit> {
187 let (reply, response) = oneshot::channel();
188 self.commands
189 .send(ActorCommand::Submit {
190 queued_at: Instant::now(),
191 command_id,
192 command,
193 admission,
194 reply,
195 })
196 .await
197 .context("session manager stopped")?;
198 Ok(PendingRelaySubmit { response })
199 }
200
201 pub async fn sync_now(&self) -> Result<()> {
202 self.enqueue_sync().await?.wait().await
203 }
204
205 pub async fn respond_elicitation(
206 &self,
207 elicitation_id: String,
208 response: ElicitationResponse,
209 ) -> Result<()> {
210 let (reply, result) = oneshot::channel();
211 self.commands
212 .send(ActorCommand::RespondElicitation {
213 elicitation_id,
214 response,
215 reply,
216 })
217 .await
218 .context("session manager stopped")?;
219 result
220 .await
221 .context("session manager stopped")?
222 .map_err(anyhow::Error::msg)
223 }
224
225 pub async fn stop_background_task(&self, background_task_id: String) -> Result<()> {
226 let (reply, result) = oneshot::channel();
227 self.commands
228 .send(ActorCommand::StopBackgroundTask {
229 background_task_id,
230 reply,
231 })
232 .await
233 .context("session manager stopped")?;
234 result
235 .await
236 .context("session manager stopped")?
237 .map_err(anyhow::Error::msg)
238 }
239
240 pub async fn install_prompt_context(&self, text: String) -> Result<()> {
243 let (reply, result) = oneshot::channel();
244 self.commands
245 .send(ActorCommand::InstallPromptContext { text, reply })
246 .await
247 .context("session manager stopped")?;
248 result
249 .await
250 .context("session manager stopped")?
251 .map_err(anyhow::Error::msg)
252 }
253
254 pub async fn reviewer(&self, action: ReviewerAction) -> Result<ReviewerOutcome> {
260 self.reviewer_as(None, action).await
261 }
262
263 pub async fn reviewer_as(
267 &self,
268 role: Option<String>,
269 action: ReviewerAction,
270 ) -> Result<ReviewerOutcome> {
271 let (reply, result) = oneshot::channel();
272 self.commands
273 .send(ActorCommand::Reviewer {
274 role,
275 action,
276 reply,
277 })
278 .await
279 .context("session manager stopped")?;
280 result
281 .await
282 .context("session manager stopped")?
283 .map_err(anyhow::Error::msg)
284 }
285
286 pub async fn enqueue_sync(&self) -> Result<PendingRelaySync> {
287 let (reply, response) = oneshot::channel();
288 self.commands
289 .send(ActorCommand::Sync { reply })
290 .await
291 .context("session manager stopped")?;
292 Ok(PendingRelaySync { response })
293 }
294
295 pub async fn lease_connection(&self) -> Result<ManagedSessionLease> {
296 let (reply, response) = oneshot::channel();
297 self.commands
298 .send(ActorCommand::Lease { reply })
299 .await
300 .context("session manager stopped")?;
301 let (lease_id, connection) = response.await.context("session manager stopped")??;
302 Ok(ManagedSessionLease {
303 session_id: self.session_id.clone(),
304 lease_id: Some(lease_id),
305 connection: Some(connection),
306 releases: self.releases.clone(),
307 })
308 }
309}
310
311pub struct PendingRelaySubmit {
312 pub(super) response:
313 oneshot::Receiver<std::result::Result<u64, mj_client::session::SubmitFailure>>,
314}
315
316impl PendingRelaySubmit {
317 pub async fn wait(self) -> Result<u64> {
318 self.response
319 .await
320 .map_err(|error| {
321 anyhow::Error::new(error).context(mj_client::session::DeliveryUnconfirmed)
322 })?
323 .map_err(|failure| {
324 let error = anyhow::Error::msg(failure.message);
325 if failure.unconfirmed {
326 error.context(mj_client::session::DeliveryUnconfirmed)
327 } else {
328 error
329 }
330 })
331 }
332}
333
334pub struct PendingRelaySync {
335 pub(super) response: oneshot::Receiver<std::result::Result<(), String>>,
336}
337
338impl PendingRelaySync {
339 pub async fn wait(self) -> Result<()> {
340 self.response
341 .await
342 .context("session manager stopped")?
343 .map_err(anyhow::Error::msg)
344 }
345}