Skip to main content

mj_controller/worker_client/
reviewer.rs

1use super::*;
2
3impl RelayClient {
4    /// Start the second-opinion reviewer beside this session, or report the
5    /// running one when it already matches `config`.
6    ///
7    /// The reviewer's profile must already be staged on the target. Starting
8    /// can take as long as opening any harness session, so this uses the
9    /// handshake deadline rather than the bookkeeping one.
10    pub async fn start_reviewer(
11        &mut self,
12        role: Option<&str>,
13        config: ReviewerLaunchConfig,
14    ) -> Result<StartedReviewer> {
15        let request = self.reviewer_request(
16            role,
17            ReviewerRequest::Start {
18                config: Box::new(config),
19            },
20        )?;
21        match self
22            .call_with_timeout(request, RELAY_HANDSHAKE_TIMEOUT)
23            .await?
24        {
25            RelayResponsePayload::ReviewerStarted {
26                native_session_id,
27                config_options,
28                reused,
29                state,
30            } => Ok(StartedReviewer {
31                native_session_id,
32                config_options,
33                reused,
34                state: *state,
35            }),
36            _ => bail!("relay returned an unexpected reviewer start response"),
37        }
38    }
39
40    /// Replay the reviewer's journal from a cursor, exactly as [`Self::attach`]
41    /// does for the primary.
42    pub async fn attach_reviewer(
43        &mut self,
44        role: Option<&str>,
45        after_ordinal: u64,
46        after_digest: impl Into<String>,
47    ) -> Result<RelayAttachment> {
48        let after_digest = after_digest.into();
49        let request = self.reviewer_request(
50            role,
51            ReviewerRequest::Attach {
52                after_ordinal,
53                after_digest: after_digest.clone(),
54            },
55        )?;
56        let payload = self
57            .call_with_timeout(request, RELAY_HISTORY_TIMEOUT)
58            .await?;
59        let RelayResponsePayload::Attached {
60            state,
61            events,
62            through_ordinal,
63            through_digest,
64        } = payload
65        else {
66            bail!("relay returned an unexpected reviewer attach response");
67        };
68        // The reviewer's journal is verified the same way the primary's is: a
69        // sidecar's history is not exempt from the chain check.
70        let mut cursor = RelayCursor {
71            ordinal: after_ordinal,
72            digest: after_digest,
73        };
74        for event in &events {
75            validate_relay_event(cursor.ordinal, &cursor.digest, event)
76                .context("verify reviewer attachment event chain")?;
77            cursor.ordinal = event.ordinal;
78            cursor.digest.clone_from(&event.digest);
79        }
80        if cursor.ordinal != through_ordinal || cursor.digest != through_digest {
81            bail!("reviewer attachment frontier does not match its event chain");
82        }
83        Ok(RelayAttachment {
84            state,
85            events,
86            through_ordinal,
87            through_digest,
88        })
89    }
90
91    /// Advance the reviewer's acknowledged frontier so its journal can be
92    /// pruned once the controller has the events durably.
93    pub async fn acknowledge_reviewer(
94        &mut self,
95        role: Option<&str>,
96        through_ordinal: u64,
97        through_digest: impl Into<String>,
98    ) -> Result<RelayCursor> {
99        let request = self.reviewer_request(
100            role,
101            ReviewerRequest::Acknowledge {
102                through_ordinal,
103                through_digest: through_digest.into(),
104            },
105        )?;
106        match self
107            .call_with_timeout(request, RELAY_ACKNOWLEDGE_TIMEOUT)
108            .await?
109        {
110            RelayResponsePayload::Acknowledged {
111                through_ordinal,
112                through_digest,
113            } => Ok(RelayCursor {
114                ordinal: through_ordinal,
115                digest: through_digest,
116            }),
117            _ => bail!("relay returned an unexpected reviewer acknowledgement response"),
118        }
119    }
120
121    /// Queue one command on the reviewer's own relay.
122    pub async fn submit_to_reviewer(
123        &mut self,
124        role: Option<&str>,
125        command_id: impl Into<String>,
126        command: RelayCommand,
127    ) -> Result<u64> {
128        let command_id = command_id.into();
129        let request = self.reviewer_request(
130            role,
131            ReviewerRequest::Submit {
132                command_id: command_id.clone(),
133                command,
134            },
135        )?;
136        match self.call(request).await? {
137            RelayResponsePayload::Accepted {
138                command_id: accepted_id,
139                ordinal,
140            } if accepted_id == command_id => Ok(ordinal),
141            RelayResponsePayload::Accepted {
142                command_id: accepted_id,
143                ..
144            } => bail!("reviewer accepted command under ID {accepted_id}, expected {command_id}"),
145            _ => bail!("relay returned an unexpected reviewer command response"),
146        }
147    }
148
149    pub async fn reviewer_status(&mut self, role: Option<&str>) -> Result<RelayOperationalState> {
150        let request = self.reviewer_request(role, ReviewerRequest::Status)?;
151        match self.call(request).await? {
152            RelayResponsePayload::Status(status) => Ok(status),
153            _ => bail!("relay returned an unexpected reviewer status response"),
154        }
155    }
156
157    /// Answer a form the reviewer's harness is waiting on.
158    pub async fn respond_to_reviewer(
159        &mut self,
160        role: Option<&str>,
161        elicitation_id: String,
162        response: ElicitationResponse,
163    ) -> Result<()> {
164        let request = self.reviewer_request(
165            role,
166            ReviewerRequest::RespondElicitation {
167                elicitation_id: elicitation_id.clone(),
168                response,
169            },
170        )?;
171        match self.call(request).await? {
172            RelayResponsePayload::ElicitationResolved {
173                elicitation_id: resolved,
174            } if resolved == elicitation_id => Ok(()),
175            RelayResponsePayload::ElicitationResolved {
176                elicitation_id: resolved,
177            } => bail!("reviewer resolved elicitation {resolved:?}, expected {elicitation_id:?}"),
178            _ => bail!("relay returned an unexpected reviewer elicitation response"),
179        }
180    }
181
182    /// Cancel any reviewer turn in flight and stop its process group, keeping
183    /// its staged profile, native session and journal for the next review.
184    pub async fn pause_reviewer(&mut self, role: Option<&str>) -> Result<()> {
185        let request = self.reviewer_request(role, ReviewerRequest::Pause)?;
186        match self
187            .call_with_timeout(request, RELAY_ACKNOWLEDGE_TIMEOUT)
188            .await?
189        {
190            RelayResponsePayload::ReviewerPaused => Ok(()),
191            _ => bail!("relay returned an unexpected reviewer pause response"),
192        }
193    }
194
195    pub async fn pause_reviewer_generation(
196        &mut self,
197        role: Option<&str>,
198        generation: u64,
199    ) -> Result<()> {
200        let request =
201            self.reviewer_request(role, ReviewerRequest::PauseGeneration { generation })?;
202        match self
203            .call_with_timeout(request, RELAY_ACKNOWLEDGE_TIMEOUT)
204            .await?
205        {
206            RelayResponsePayload::ReviewerPaused => Ok(()),
207            _ => bail!("relay returned an unexpected reviewer pause response"),
208        }
209    }
210
211    /// Report what every workspace repository changed since the review
212    /// baselines the controller holds.
213    pub async fn capture_review_delta(
214        &mut self,
215        role: Option<&str>,
216        baselines: std::collections::BTreeMap<std::path::PathBuf, String>,
217    ) -> Result<Vec<mj_core::relay::RepoDelta>> {
218        let request = self.reviewer_request(role, ReviewerRequest::CaptureDelta { baselines })?;
219        match self
220            .call_with_timeout(request, REVIEW_CAPTURE_TIMEOUT)
221            .await?
222        {
223            RelayResponsePayload::ReviewDelta { repositories } => Ok(repositories),
224            _ => bail!("relay returned an unexpected review capture response"),
225        }
226    }
227
228    /// Record the trees a completed review reviewed through, so the next
229    /// review starts from them.
230    pub async fn advance_review_baseline(
231        &mut self,
232        role: Option<&str>,
233        trees: std::collections::BTreeMap<std::path::PathBuf, String>,
234    ) -> Result<()> {
235        let request = self.reviewer_request(role, ReviewerRequest::AdvanceBaseline { trees })?;
236        match self
237            .call_with_timeout(request, REVIEW_CAPTURE_TIMEOUT)
238            .await?
239        {
240            RelayResponsePayload::ReviewBaselineAdvanced => Ok(()),
241            _ => bail!("relay returned an unexpected review baseline response"),
242        }
243    }
244
245    /// Collect the specialist lanes the review supervisor asked for since the
246    /// last call.
247    pub async fn take_lane_dispatches(
248        &mut self,
249    ) -> Result<Vec<mj_core::review::lanes::ReviewSubagentRequest>> {
250        let request = self.reviewer_request(None, ReviewerRequest::TakeLaneDispatches)?;
251        match self.call(request).await? {
252            RelayResponsePayload::LaneDispatches { requests } => Ok(requests),
253            _ => bail!("relay returned an unexpected lane dispatch response"),
254        }
255    }
256
257    /// Wraps a reviewer action, refusing it on a worker too old to know what a
258    /// reviewer is rather than sending a method it would reject as unknown.
259    pub(super) fn reviewer_request(
260        &self,
261        role: Option<&str>,
262        request: ReviewerRequest,
263    ) -> Result<RelayRequest> {
264        let request = RelayRequest::Reviewer {
265            role: role.map(str::to_owned),
266            request,
267        };
268        Ok(request)
269    }
270
271    /// Answer an ACP form over the live relay connection. User-entered content
272    /// is intentionally excluded from the relay's durable command path.
273    pub async fn respond_elicitation(
274        &mut self,
275        elicitation_id: String,
276        response: ElicitationResponse,
277    ) -> Result<()> {
278        let request = RelayRequest::RespondElicitation {
279            elicitation_id: elicitation_id.clone(),
280            response,
281        };
282        match self.call(request).await? {
283            RelayResponsePayload::ElicitationResolved {
284                elicitation_id: resolved,
285            } if resolved == elicitation_id => Ok(()),
286            RelayResponsePayload::ElicitationResolved {
287                elicitation_id: resolved,
288            } => bail!("relay resolved elicitation {resolved:?}, expected {elicitation_id:?}"),
289            _ => bail!("relay returned an unexpected elicitation response"),
290        }
291    }
292
293    /// Ask the live worker to stop one process-local background task.
294    pub async fn stop_background_task(&mut self, background_task_id: String) -> Result<()> {
295        let request = RelayRequest::StopBackgroundTask {
296            background_task_id: background_task_id.clone(),
297        };
298        match self.call(request).await? {
299            RelayResponsePayload::BackgroundTaskStopRequested {
300                background_task_id: stopped,
301            } if stopped == background_task_id => Ok(()),
302            RelayResponsePayload::BackgroundTaskStopRequested {
303                background_task_id: stopped,
304            } => {
305                bail!("relay stopped background task {stopped:?}, expected {background_task_id:?}")
306            }
307            _ => bail!("relay returned an unexpected background task stop response"),
308        }
309    }
310
311    pub async fn subagent_requests(
312        &mut self,
313    ) -> Result<(
314        Vec<mj_core::subagent::SubagentToolRequest>,
315        Vec<mj_core::subagent::SubagentToolResult>,
316    )> {
317        let request = RelayRequest::SubagentRequests;
318        if !request.supported_at(self.protocol_version) {
319            return Ok((Vec::new(), Vec::new()));
320        }
321        match self.call(request).await? {
322            RelayResponsePayload::SubagentRequests { requests, results } => Ok((requests, results)),
323            _ => bail!("relay returned an unexpected sub-agent request response"),
324        }
325    }
326
327    pub async fn complete_subagent_request(
328        &mut self,
329        result: mj_core::subagent::SubagentToolResult,
330    ) -> Result<()> {
331        let request = RelayRequest::CompleteSubagentRequest { result };
332        match self.call(request).await? {
333            RelayResponsePayload::SubagentRequestCompleted => Ok(()),
334            _ => bail!("relay returned an unexpected sub-agent completion response"),
335        }
336    }
337
338    pub async fn detach(mut self) -> Result<()> {
339        self.input
340            .take()
341            .expect("connected relay owns proxy stdin")
342            .shutdown()
343            .await
344            .context("close relay proxy stdin")?;
345        let mut child = self.child.take().expect("connected relay owns proxy child");
346        match tokio::time::timeout(RELAY_PROXY_DETACH_GRACE, child.wait()).await {
347            Ok(status) => {
348                status.context("wait for relay proxy")?;
349            }
350            Err(_) => {
351                if let Err(error) = child.start_kill().context("stop relay proxy") {
352                    tracing::warn!(
353                        session_id = %self.session_id,
354                        operation = "detach",
355                        %error,
356                        "could not stop relay proxy after detach timeout"
357                    );
358                    return Err(error);
359                }
360                if let Err(error) = child.wait().await {
361                    tracing::warn!(
362                        session_id = %self.session_id,
363                        operation = "detach",
364                        %error,
365                        "could not reap relay proxy after stopping it"
366                    );
367                }
368            }
369        }
370        Ok(())
371    }
372}