Skip to main content

mj_controller/session_manager/
handle.rs

1use 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/// A one-command capability issued by the review host while its prompt hold
17/// is open. It is intentionally opaque to callers: the session actor checks
18/// it against the host's live hold registry before bypassing prompt refusal.
19#[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
48/// Exclusive ownership of a session actor's existing relay connection.
49///
50/// Lifecycle operations use this instead of opening a competing projection
51/// client. Dropping an unreleased lease drops the proxy connection, which in
52/// turn cancels any ordinary relay checkpoint barrier.
53///
54/// Prompt submissions that arrive while the lease is active are not rejected.
55/// The actor queues them and forwards them in arrival order once the lease is
56/// released or dropped.
57pub 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    /// Swap the leased proxy after the worker process behind it was replaced.
72    /// The actor stays leased, so queued prompts cannot race the new latch.
73    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 the proxy before telling the actor to reconnect so the relay
104        // observes EOF and releases any abandoned checkpoint barrier first.
105        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    /// Narrow this controller-owned handle to the operations a control surface uses.
122    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    /// Whether the per-session actor behind this handle has retired. The
135    /// manager itself may still be alive with a replacement actor, so callers
136    /// holding long-lived handles use this to reacquire the current one.
137    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    /// Submit the review's corrective prompt through the one admission that
158    /// corresponds to its live prompt hold. Generic submissions continue to
159    /// use [`Self::submit`] and remain subject to review refusal.
160    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    /// Install background text the harness reads with the next real prompt.
241    /// It creates no transcript turn, so the user never sees it.
242    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    /// Drive the session's second-opinion reviewer.
255    ///
256    /// The reviewer shares this session's relay connection, so its actions
257    /// queue behind the session's own and are refused while a lifecycle
258    /// operation holds the connection.
259    pub async fn reviewer(&self, action: ReviewerAction) -> Result<ReviewerOutcome> {
260        self.reviewer_as(None, action).await
261    }
262
263    /// Drive one reviewing role. `None` is the default role, which is the one
264    /// plan review uses; a turn review in the extended tier names its
265    /// supervisor, its intent analyst, and each specialist lane.
266    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        self.lease_connection_for(None).await
297    }
298
299    /// Background replacement must never take a busy session's control channel.
300    pub async fn lease_idle_connection(
301        &self,
302        harness: mj_core::config::HarnessKind,
303    ) -> Result<Option<ManagedSessionLease>> {
304        match self.lease_connection_for(Some(harness)).await {
305            Ok(lease) => Ok(Some(lease)),
306            Err(error) if error.is::<SessionNotIdle>() => Ok(None),
307            Err(error) => Err(error),
308        }
309    }
310
311    async fn lease_connection_for(
312        &self,
313        idle_harness: Option<mj_core::config::HarnessKind>,
314    ) -> Result<ManagedSessionLease> {
315        let (reply, response) = oneshot::channel();
316        self.commands
317            .send(ActorCommand::Lease {
318                idle_harness,
319                reply,
320            })
321            .await
322            .context("session manager stopped")?;
323        let (lease_id, connection) = response.await.context("session manager stopped")??;
324        Ok(ManagedSessionLease {
325            session_id: self.session_id.clone(),
326            lease_id: Some(lease_id),
327            connection: Some(connection),
328            releases: self.releases.clone(),
329        })
330    }
331}
332
333#[derive(Debug)]
334pub(super) struct SessionNotIdle;
335
336impl std::fmt::Display for SessionNotIdle {
337    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
338        f.write_str("session still has work in flight; upgrade deferred")
339    }
340}
341
342impl std::error::Error for SessionNotIdle {}
343
344pub struct PendingRelaySubmit {
345    pub(super) response:
346        oneshot::Receiver<std::result::Result<u64, mj_client::session::SubmitFailure>>,
347}
348
349impl PendingRelaySubmit {
350    pub async fn wait(self) -> Result<u64> {
351        self.response
352            .await
353            .map_err(|error| {
354                anyhow::Error::new(error).context(mj_client::session::DeliveryUnconfirmed)
355            })?
356            .map_err(|failure| {
357                let error = anyhow::Error::msg(failure.message);
358                if failure.unconfirmed {
359                    error.context(mj_client::session::DeliveryUnconfirmed)
360                } else {
361                    error
362                }
363            })
364    }
365}
366
367pub struct PendingRelaySync {
368    pub(super) response: oneshot::Receiver<std::result::Result<(), String>>,
369}
370
371impl PendingRelaySync {
372    pub async fn wait(self) -> Result<()> {
373        self.response
374            .await
375            .context("session manager stopped")?
376            .map_err(anyhow::Error::msg)
377    }
378}