mj_controller/session_manager/
handle.rs1use super::*;
2
3#[derive(Clone)]
4pub struct SessionManagerControl {
5 pub(crate) session_cpu: watch::Receiver<SessionCpuTable>,
6 pub(super) commands: mpsc::Sender<ManagerCommand>,
7}
8
9#[derive(Clone, Debug)]
10pub struct ManagedSessionHandle {
11 pub(super) session_id: String,
12 pub(super) commands: mpsc::Sender<ActorCommand>,
13 pub(super) releases: mpsc::UnboundedSender<ReturnedConnection>,
14 pub(super) view: watch::Receiver<ManagedSessionView>,
15}
16
17#[derive(Debug, Clone, PartialEq, Eq)]
21pub struct ReviewDeliveryAdmission {
22 pub(super) session_id: String,
23 pub(super) epoch: u64,
24 pub(super) command_id: String,
25}
26
27impl ReviewDeliveryAdmission {
28 pub(crate) fn new(session_id: String, epoch: u64, command_id: String) -> Self {
29 Self {
30 session_id,
31 epoch,
32 command_id,
33 }
34 }
35
36 pub(crate) fn session_id(&self) -> &str {
37 &self.session_id
38 }
39
40 pub(crate) const fn epoch(&self) -> u64 {
41 self.epoch
42 }
43
44 pub(crate) fn command_id(&self) -> &str {
45 &self.command_id
46 }
47}
48
49pub struct ManagedSessionLease {
59 pub(super) session_id: String,
60 pub(super) lease_id: Option<u64>,
61 pub(super) connection: Option<StandaloneSession>,
62 pub(super) releases: mpsc::UnboundedSender<ReturnedConnection>,
63}
64
65impl ManagedSessionLease {
66 pub fn connection_mut(&mut self) -> &mut StandaloneSession {
67 self.connection
68 .as_mut()
69 .expect("managed session lease has already been released")
70 }
71
72 pub fn replace_connection(&mut self, connection: StandaloneSession) {
75 drop(self.connection.take());
76 self.connection = Some(connection);
77 }
78
79 pub fn release(mut self) {
80 let lease_id = self
81 .lease_id
82 .take()
83 .expect("managed session lease has already been released");
84 let connection = self.connection.take();
85 if let Err(error) = self.releases.send(ReturnedConnection {
86 lease_id,
87 connection,
88 }) {
89 tracing::warn!(
90 session_id = %self.session_id,
91 operation = "lease_release",
92 %error,
93 "session actor stopped before receiving released relay connection"
94 );
95 }
96 }
97}
98
99impl Drop for ManagedSessionLease {
100 fn drop(&mut self) {
101 let Some(lease_id) = self.lease_id.take() else {
102 return;
103 };
104 drop(self.connection.take());
107 if let Err(error) = self.releases.send(ReturnedConnection {
108 lease_id,
109 connection: None,
110 }) {
111 tracing::warn!(
112 session_id = %self.session_id,
113 operation = "lease_drop",
114 %error,
115 "session actor stopped before receiving dropped relay lease"
116 );
117 }
118 }
119}
120
121impl ManagedSessionHandle {
122 pub fn client(&self) -> mj_client::session::SessionHandle {
124 mj_client::session::SessionHandle::new(ClientSessionHandle(self.clone()))
125 }
126
127 pub fn session_id(&self) -> &str {
128 &self.session_id
129 }
130
131 pub fn view(&self) -> ManagedSessionView {
132 self.view.borrow().clone()
133 }
134
135 pub fn is_stopped(&self) -> bool {
139 self.commands.is_closed()
140 }
141
142 pub fn has_changed(&self) -> Result<bool> {
143 self.view.has_changed().context("session manager stopped")
144 }
145
146 pub async fn changed(&mut self) -> Result<ManagedSessionView> {
147 self.view
148 .changed()
149 .await
150 .context("session manager stopped")?;
151 Ok(self.view())
152 }
153
154 pub async fn submit(&self, command_id: String, command: RelayCommand) -> Result<u64> {
155 self.enqueue_submit(command_id, command).await?.wait().await
156 }
157
158 pub(crate) async fn submit_review_delivery(
162 &self,
163 admission: ReviewDeliveryAdmission,
164 command: RelayCommand,
165 ) -> Result<u64> {
166 let command_id = admission.command_id.clone();
167 self.enqueue_submit_with_admission(command_id, command, Some(admission))
168 .await?
169 .wait()
170 .await
171 }
172
173 pub async fn enqueue_submit(
174 &self,
175 command_id: String,
176 command: RelayCommand,
177 ) -> Result<PendingRelaySubmit> {
178 self.enqueue_submit_with_admission(command_id, command, None)
179 .await
180 }
181
182 pub(super) async fn enqueue_submit_with_admission(
183 &self,
184 command_id: String,
185 command: RelayCommand,
186 admission: Option<ReviewDeliveryAdmission>,
187 ) -> Result<PendingRelaySubmit> {
188 let (reply, response) = oneshot::channel();
189 self.commands
190 .send(ActorCommand::Submit {
191 queued_at: Instant::now(),
192 command_id,
193 command,
194 admission,
195 reply,
196 })
197 .await
198 .context("session manager stopped")?;
199 Ok(PendingRelaySubmit { response })
200 }
201
202 pub async fn sync_now(&self) -> Result<()> {
203 self.enqueue_sync().await?.wait().await
204 }
205
206 pub async fn respond_elicitation(
207 &self,
208 elicitation_id: String,
209 response: ElicitationResponse,
210 ) -> Result<()> {
211 if let Some((role, elicitation_id)) =
215 mj_client::review::parse_reviewer_question_id(&elicitation_id)
216 {
217 self.reviewer_as(
218 Some(role.to_owned()),
219 ReviewerAction::RespondElicitation {
220 elicitation_id: elicitation_id.to_owned(),
221 response,
222 },
223 )
224 .await?;
225 return Ok(());
226 }
227 let (reply, result) = oneshot::channel();
228 self.commands
229 .send(ActorCommand::RespondElicitation {
230 elicitation_id,
231 response,
232 reply,
233 })
234 .await
235 .context("session manager stopped")?;
236 result
237 .await
238 .context("session manager stopped")?
239 .map_err(anyhow::Error::msg)
240 }
241
242 pub async fn stop_background_task(&self, background_task_id: String) -> Result<()> {
243 let (reply, result) = oneshot::channel();
244 self.commands
245 .send(ActorCommand::StopBackgroundTask {
246 background_task_id,
247 reply,
248 })
249 .await
250 .context("session manager stopped")?;
251 result
252 .await
253 .context("session manager stopped")?
254 .map_err(anyhow::Error::msg)
255 }
256
257 pub async fn install_prompt_context(&self, text: String) -> Result<()> {
260 let (reply, result) = oneshot::channel();
261 self.commands
262 .send(ActorCommand::InstallPromptContext { text, reply })
263 .await
264 .context("session manager stopped")?;
265 result
266 .await
267 .context("session manager stopped")?
268 .map_err(anyhow::Error::msg)
269 }
270
271 pub async fn reviewer(&self, action: ReviewerAction) -> Result<ReviewerOutcome> {
277 self.reviewer_as(None, action).await
278 }
279
280 pub async fn reviewer_as(
283 &self,
284 role: Option<String>,
285 action: ReviewerAction,
286 ) -> Result<ReviewerOutcome> {
287 let (reply, result) = oneshot::channel();
288 self.commands
289 .send(ActorCommand::Reviewer {
290 role,
291 action,
292 reply,
293 })
294 .await
295 .context("session manager stopped")?;
296 result
297 .await
298 .context("session manager stopped")?
299 .map_err(anyhow::Error::msg)
300 }
301
302 pub async fn enqueue_sync(&self) -> Result<PendingRelaySync> {
303 let (reply, response) = oneshot::channel();
304 self.commands
305 .send(ActorCommand::Sync { reply })
306 .await
307 .context("session manager stopped")?;
308 Ok(PendingRelaySync { response })
309 }
310
311 pub async fn run_on_connection(&self, job: Box<dyn RelayConnectionJob>) {
314 if let Err(error) = self.commands.send(ActorCommand::RelayJob { job }).await
315 && let ActorCommand::RelayJob { job } = error.0
316 {
317 job.refuse(anyhow::anyhow!("session actor stopped"));
318 }
319 }
320
321 pub async fn lease_connection(&self) -> Result<ManagedSessionLease> {
322 self.lease_connection_for(None).await
323 }
324
325 pub async fn lease_idle_connection(
327 &self,
328 harness: mj_core::config::HarnessKind,
329 ) -> Result<Option<ManagedSessionLease>> {
330 match self.lease_connection_for(Some(harness)).await {
331 Ok(lease) => Ok(Some(lease)),
332 Err(error) if error.is::<SessionNotIdle>() => Ok(None),
333 Err(error) => Err(error),
334 }
335 }
336
337 async fn lease_connection_for(
338 &self,
339 idle_harness: Option<mj_core::config::HarnessKind>,
340 ) -> Result<ManagedSessionLease> {
341 let (reply, response) = oneshot::channel();
342 self.commands
343 .send(ActorCommand::Lease {
344 idle_harness,
345 reply,
346 })
347 .await
348 .context("session manager stopped")?;
349 let (lease_id, connection) = response.await.context("session manager stopped")??;
350 Ok(ManagedSessionLease {
351 session_id: self.session_id.clone(),
352 lease_id: Some(lease_id),
353 connection: Some(connection),
354 releases: self.releases.clone(),
355 })
356 }
357}
358
359#[derive(Debug)]
360pub(super) struct SessionNotIdle;
361
362impl std::fmt::Display for SessionNotIdle {
363 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
364 f.write_str("session still has work in flight; upgrade deferred")
365 }
366}
367
368impl std::error::Error for SessionNotIdle {}
369
370pub struct PendingRelaySubmit {
371 pub(super) response:
372 oneshot::Receiver<std::result::Result<u64, mj_client::session::SubmitFailure>>,
373}
374
375impl PendingRelaySubmit {
376 pub async fn wait(self) -> Result<u64> {
377 self.response
378 .await
379 .map_err(|error| {
380 anyhow::Error::new(error).context(mj_client::session::DeliveryUnconfirmed)
381 })?
382 .map_err(|failure| {
383 let error = anyhow::Error::msg(failure.message);
384 if failure.unconfirmed {
385 error.context(mj_client::session::DeliveryUnconfirmed)
386 } else if failure.refused {
387 error.context(mj_client::session::Refused)
388 } else {
389 error
390 }
391 })
392 }
393}
394
395pub struct PendingRelaySync {
396 pub(super) response: oneshot::Receiver<std::result::Result<(), String>>,
397}
398
399impl PendingRelaySync {
400 pub async fn wait(self) -> Result<()> {
401 self.response
402 .await
403 .context("session manager stopped")?
404 .map_err(anyhow::Error::msg)
405 }
406}