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 install_github_token(&self, token: String) -> Result<()> {
324 let (reply, response) = oneshot::channel();
325 self.run_on_connection(Box::new(InstallGithubTokenJob { token, reply }))
326 .await;
327 response
328 .await
329 .context("GitHub token refresh did not receive a relay result")?
330 .map_err(anyhow::Error::msg)
331 }
332
333 pub async fn lease_connection(&self) -> Result<ManagedSessionLease> {
334 self.lease_connection_for(None).await
335 }
336
337 pub async fn lease_idle_connection(
339 &self,
340 harness: mj_core::config::HarnessKind,
341 ) -> Result<Option<ManagedSessionLease>> {
342 match self.lease_connection_for(Some(harness)).await {
343 Ok(lease) => Ok(Some(lease)),
344 Err(error) if error.is::<SessionNotIdle>() => Ok(None),
345 Err(error) => Err(error),
346 }
347 }
348
349 async fn lease_connection_for(
350 &self,
351 idle_harness: Option<mj_core::config::HarnessKind>,
352 ) -> Result<ManagedSessionLease> {
353 let (reply, response) = oneshot::channel();
354 self.commands
355 .send(ActorCommand::Lease {
356 idle_harness,
357 reply,
358 })
359 .await
360 .context("session manager stopped")?;
361 let (lease_id, connection) = response.await.context("session manager stopped")??;
362 Ok(ManagedSessionLease {
363 session_id: self.session_id.clone(),
364 lease_id: Some(lease_id),
365 connection: Some(connection),
366 releases: self.releases.clone(),
367 })
368 }
369}
370
371struct InstallGithubTokenJob {
372 token: String,
373 reply: oneshot::Sender<std::result::Result<(), String>>,
374}
375
376impl RelayConnectionJob for InstallGithubTokenJob {
377 fn run<'a>(self: Box<Self>, client: &'a mut RelayClient) -> futures::future::BoxFuture<'a, ()> {
378 Box::pin(async move {
379 let expected = mj_core::credentials::GithubTokenSnapshot::of(&self.token);
380 let result = client
381 .install_github_token(&self.token)
382 .await
383 .and_then(|installed| {
384 anyhow::ensure!(
385 installed == expected,
386 "worker GitHub token fingerprint did not match after install"
387 );
388 Ok(())
389 })
390 .map_err(|error| format!("{error:#}"));
391 if self.reply.send(result).is_err() {
392 tracing::debug!("GitHub token refresh result receiver was dropped");
393 }
394 })
395 }
396
397 fn refuse(self: Box<Self>, error: anyhow::Error) {
398 let _ = self.reply.send(Err(format!("{error:#}")));
399 }
400}
401
402#[derive(Debug)]
403pub(super) struct SessionNotIdle;
404
405impl std::fmt::Display for SessionNotIdle {
406 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
407 f.write_str("session still has work in flight; upgrade deferred")
408 }
409}
410
411impl std::error::Error for SessionNotIdle {}
412
413pub struct PendingRelaySubmit {
414 pub(super) response:
415 oneshot::Receiver<std::result::Result<u64, mj_client::session::SubmitFailure>>,
416}
417
418impl PendingRelaySubmit {
419 pub async fn wait(self) -> Result<u64> {
420 self.response
421 .await
422 .map_err(|error| {
423 anyhow::Error::new(error).context(mj_client::session::DeliveryUnconfirmed)
424 })?
425 .map_err(|failure| {
426 let error = anyhow::Error::msg(failure.message);
427 if failure.unconfirmed {
428 error.context(mj_client::session::DeliveryUnconfirmed)
429 } else if failure.refused {
430 error.context(mj_client::session::Refused)
431 } else {
432 error
433 }
434 })
435 }
436}
437
438pub struct PendingRelaySync {
439 pub(super) response: oneshot::Receiver<std::result::Result<(), String>>,
440}
441
442impl PendingRelaySync {
443 pub async fn wait(self) -> Result<()> {
444 self.response
445 .await
446 .context("session manager stopped")?
447 .map_err(anyhow::Error::msg)
448 }
449}