Skip to main content

mj_controller/session_manager/
handle.rs

1use 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/// A one-command capability issued by the review host while its prompt hold
18/// is open. It is intentionally opaque to callers: the session actor checks
19/// it against the host's live hold registry before bypassing prompt refusal.
20#[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
49/// Exclusive ownership of a session actor's existing relay connection.
50///
51/// Lifecycle operations use this instead of opening a competing projection
52/// client. Dropping an unreleased lease drops the proxy connection, which in
53/// turn cancels any ordinary relay checkpoint barrier.
54///
55/// Prompt submissions that arrive while the lease is active are not rejected.
56/// The actor queues them and forwards them in arrival order once the lease is
57/// released or dropped.
58pub 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    /// Swap the leased proxy after the worker process behind it was replaced.
73    /// The actor stays leased, so queued prompts cannot race the new latch.
74    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 the proxy before telling the actor to reconnect so the relay
105        // observes EOF and releases any abandoned checkpoint barrier first.
106        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    /// Narrow this controller-owned handle to the operations a control surface uses.
123    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    /// Whether the per-session actor behind this handle has retired. The
136    /// manager itself may still be alive with a replacement actor, so callers
137    /// holding long-lived handles use this to reacquire the current one.
138    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    /// Submit the review's corrective prompt through the one admission that
159    /// corresponds to its live prompt hold. Generic submissions continue to
160    /// use [`Self::submit`] and remain subject to review refusal.
161    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        // Every surface answers through here. A reviewing role's form goes
212        // back to that role's harness in the worker, as the review pane's
213        // own answer does.
214        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    /// Install background text the harness reads with the next real prompt.
258    /// It creates no transcript turn, so the user never sees it.
259    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    /// Drive the session's second-opinion reviewer.
272    ///
273    /// The reviewer shares this session's relay connection, so its actions
274    /// queue behind the session's own and are refused while a lifecycle
275    /// operation holds the connection.
276    pub async fn reviewer(&self, action: ReviewerAction) -> Result<ReviewerOutcome> {
277        self.reviewer_as(None, action).await
278    }
279
280    /// Drive one isolated reviewer role. `None` selects the default role
281    /// shared by plan and turn review; settings discovery uses a named role.
282    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    /// Run `job` on this session actor's own relay connection. The job is
312    /// refused, never dropped, when the actor cannot run it.
313    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    /// Replace the worker's token on its actor-owned relay connection. This
322    /// keeps export from opening a second connection while a turn is active.
323    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    /// Background replacement must never take a busy session's control channel.
338    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}