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    /// Report what every workspace repository changed since the review
196    /// baselines the controller holds.
197    pub async fn capture_review_delta(
198        &mut self,
199        role: Option<&str>,
200        baselines: std::collections::BTreeMap<std::path::PathBuf, String>,
201    ) -> Result<Vec<mj_core::relay::RepoDelta>> {
202        let request = self.reviewer_request(role, ReviewerRequest::CaptureDelta { baselines })?;
203        match self
204            .call_with_timeout(request, REVIEW_CAPTURE_TIMEOUT)
205            .await?
206        {
207            RelayResponsePayload::ReviewDelta { repositories } => Ok(repositories),
208            _ => bail!("relay returned an unexpected review capture response"),
209        }
210    }
211
212    /// Record the trees a completed review reviewed through, so the next
213    /// review starts from them.
214    pub async fn advance_review_baseline(
215        &mut self,
216        role: Option<&str>,
217        trees: std::collections::BTreeMap<std::path::PathBuf, String>,
218    ) -> Result<()> {
219        let request = self.reviewer_request(role, ReviewerRequest::AdvanceBaseline { trees })?;
220        match self
221            .call_with_timeout(request, REVIEW_CAPTURE_TIMEOUT)
222            .await?
223        {
224            RelayResponsePayload::ReviewBaselineAdvanced => Ok(()),
225            _ => bail!("relay returned an unexpected review baseline response"),
226        }
227    }
228
229    /// Run Bifrost's semantic diff analysis over the captured trees. It can
230    /// take minutes on a large changeset, so it carries its own budget.
231    pub async fn analyze_review_delta(
232        &mut self,
233        role: Option<&str>,
234        repositories: Vec<mj_core::relay::AnalyzeDeltaRepository>,
235    ) -> Result<String> {
236        let request =
237            self.reviewer_request(role, ReviewerRequest::AnalyzeDelta { repositories })?;
238        match self
239            .call_with_timeout(request, REVIEW_ANALYSIS_TIMEOUT)
240            .await?
241        {
242            RelayResponsePayload::ReviewChangedFunctions { packet } => Ok(packet),
243            _ => bail!("relay returned an unexpected review analysis response"),
244        }
245    }
246
247    /// Collect the specialist lanes the review supervisor asked for since the
248    /// last call.
249    pub async fn take_lane_dispatches(
250        &mut self,
251    ) -> Result<Vec<mj_core::review::lanes::ReviewSubagentRequest>> {
252        let request = self.reviewer_request(None, ReviewerRequest::TakeLaneDispatches)?;
253        match self.call(request).await? {
254            RelayResponsePayload::LaneDispatches { requests } => Ok(requests),
255            _ => bail!("relay returned an unexpected lane dispatch response"),
256        }
257    }
258
259    /// Wraps a reviewer action, refusing it on a worker too old to know what a
260    /// reviewer is rather than sending a method it would reject as unknown.
261    pub(super) fn reviewer_request(
262        &self,
263        role: Option<&str>,
264        request: ReviewerRequest,
265    ) -> Result<RelayRequest> {
266        let request = RelayRequest::Reviewer {
267            role: role.map(str::to_owned),
268            request,
269        };
270        Ok(request)
271    }
272
273    /// Answer an ACP form over the live relay connection. User-entered content
274    /// is intentionally excluded from the relay's durable command path.
275    pub async fn respond_elicitation(
276        &mut self,
277        elicitation_id: String,
278        response: ElicitationResponse,
279    ) -> Result<()> {
280        let request = RelayRequest::RespondElicitation {
281            elicitation_id: elicitation_id.clone(),
282            response,
283        };
284        match self.call(request).await? {
285            RelayResponsePayload::ElicitationResolved {
286                elicitation_id: resolved,
287            } if resolved == elicitation_id => Ok(()),
288            RelayResponsePayload::ElicitationResolved {
289                elicitation_id: resolved,
290            } => bail!("relay resolved elicitation {resolved:?}, expected {elicitation_id:?}"),
291            _ => bail!("relay returned an unexpected elicitation response"),
292        }
293    }
294
295    /// Ask the live worker to stop one process-local background task.
296    pub async fn stop_background_task(&mut self, background_task_id: String) -> Result<()> {
297        let request = RelayRequest::StopBackgroundTask {
298            background_task_id: background_task_id.clone(),
299        };
300        match self.call(request).await? {
301            RelayResponsePayload::BackgroundTaskStopRequested {
302                background_task_id: stopped,
303            } if stopped == background_task_id => Ok(()),
304            RelayResponsePayload::BackgroundTaskStopRequested {
305                background_task_id: stopped,
306            } => {
307                bail!("relay stopped background task {stopped:?}, expected {background_task_id:?}")
308            }
309            _ => bail!("relay returned an unexpected background task stop response"),
310        }
311    }
312
313    pub async fn subagent_requests(
314        &mut self,
315    ) -> Result<(
316        Vec<mj_core::subagent::SubagentToolRequest>,
317        Vec<mj_core::subagent::SubagentToolResult>,
318    )> {
319        let request = RelayRequest::SubagentRequests;
320        if !request.supported_at(self.protocol_version) {
321            return Ok((Vec::new(), Vec::new()));
322        }
323        match self.call(request).await? {
324            RelayResponsePayload::SubagentRequests { requests, results } => Ok((requests, results)),
325            _ => bail!("relay returned an unexpected sub-agent request response"),
326        }
327    }
328
329    pub async fn complete_subagent_request(
330        &mut self,
331        result: mj_core::subagent::SubagentToolResult,
332    ) -> Result<()> {
333        let request = RelayRequest::CompleteSubagentRequest { result };
334        match self.call(request).await? {
335            RelayResponsePayload::SubagentRequestCompleted => Ok(()),
336            _ => bail!("relay returned an unexpected sub-agent completion response"),
337        }
338    }
339
340    pub async fn detach(mut self) -> Result<()> {
341        self.input
342            .take()
343            .expect("connected relay owns proxy stdin")
344            .shutdown()
345            .await
346            .context("close relay proxy stdin")?;
347        let mut child = self.child.take().expect("connected relay owns proxy child");
348        match tokio::time::timeout(RELAY_PROXY_DETACH_GRACE, child.wait()).await {
349            Ok(status) => {
350                status.context("wait for relay proxy")?;
351            }
352            Err(_) => {
353                if let Err(error) = child.start_kill().context("stop relay proxy") {
354                    tracing::warn!(
355                        session_id = %self.session_id,
356                        operation = "detach",
357                        %error,
358                        "could not stop relay proxy after detach timeout"
359                    );
360                    return Err(error);
361                }
362                if let Err(error) = child.wait().await {
363                    tracing::warn!(
364                        session_id = %self.session_id,
365                        operation = "detach",
366                        %error,
367                        "could not reap relay proxy after stopping it"
368                    );
369                }
370            }
371        }
372        Ok(())
373    }
374}