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    /// Run Bifrost's semantic diff analysis over the captured trees. It can
246    /// take minutes on a large changeset, so it carries its own budget.
247    pub async fn analyze_review_delta(
248        &mut self,
249        role: Option<&str>,
250        repositories: Vec<mj_core::relay::AnalyzeDeltaRepository>,
251    ) -> Result<String> {
252        let request =
253            self.reviewer_request(role, ReviewerRequest::AnalyzeDelta { repositories })?;
254        match self
255            .call_with_timeout(request, REVIEW_ANALYSIS_TIMEOUT)
256            .await?
257        {
258            RelayResponsePayload::ReviewChangedFunctions { packet } => Ok(packet),
259            _ => bail!("relay returned an unexpected review analysis response"),
260        }
261    }
262
263    /// Collect the specialist lanes the review supervisor asked for since the
264    /// last call.
265    pub async fn take_lane_dispatches(
266        &mut self,
267    ) -> Result<Vec<mj_core::review::lanes::ReviewSubagentRequest>> {
268        let request = self.reviewer_request(None, ReviewerRequest::TakeLaneDispatches)?;
269        match self.call(request).await? {
270            RelayResponsePayload::LaneDispatches { requests } => Ok(requests),
271            _ => bail!("relay returned an unexpected lane dispatch response"),
272        }
273    }
274
275    /// Wraps a reviewer action, refusing it on a worker too old to know what a
276    /// reviewer is rather than sending a method it would reject as unknown.
277    pub(super) fn reviewer_request(
278        &self,
279        role: Option<&str>,
280        request: ReviewerRequest,
281    ) -> Result<RelayRequest> {
282        let request = RelayRequest::Reviewer {
283            role: role.map(str::to_owned),
284            request,
285        };
286        Ok(request)
287    }
288
289    /// Answer an ACP form over the live relay connection. User-entered content
290    /// is intentionally excluded from the relay's durable command path.
291    pub async fn respond_elicitation(
292        &mut self,
293        elicitation_id: String,
294        response: ElicitationResponse,
295    ) -> Result<()> {
296        let request = RelayRequest::RespondElicitation {
297            elicitation_id: elicitation_id.clone(),
298            response,
299        };
300        match self.call(request).await? {
301            RelayResponsePayload::ElicitationResolved {
302                elicitation_id: resolved,
303            } if resolved == elicitation_id => Ok(()),
304            RelayResponsePayload::ElicitationResolved {
305                elicitation_id: resolved,
306            } => bail!("relay resolved elicitation {resolved:?}, expected {elicitation_id:?}"),
307            _ => bail!("relay returned an unexpected elicitation response"),
308        }
309    }
310
311    /// Ask the live worker to stop one process-local background task.
312    pub async fn stop_background_task(&mut self, background_task_id: String) -> Result<()> {
313        let request = RelayRequest::StopBackgroundTask {
314            background_task_id: background_task_id.clone(),
315        };
316        match self.call(request).await? {
317            RelayResponsePayload::BackgroundTaskStopRequested {
318                background_task_id: stopped,
319            } if stopped == background_task_id => Ok(()),
320            RelayResponsePayload::BackgroundTaskStopRequested {
321                background_task_id: stopped,
322            } => {
323                bail!("relay stopped background task {stopped:?}, expected {background_task_id:?}")
324            }
325            _ => bail!("relay returned an unexpected background task stop response"),
326        }
327    }
328
329    pub async fn subagent_requests(
330        &mut self,
331    ) -> Result<(
332        Vec<mj_core::subagent::SubagentToolRequest>,
333        Vec<mj_core::subagent::SubagentToolResult>,
334    )> {
335        let request = RelayRequest::SubagentRequests;
336        if !request.supported_at(self.protocol_version) {
337            return Ok((Vec::new(), Vec::new()));
338        }
339        match self.call(request).await? {
340            RelayResponsePayload::SubagentRequests { requests, results } => Ok((requests, results)),
341            _ => bail!("relay returned an unexpected sub-agent request response"),
342        }
343    }
344
345    pub async fn complete_subagent_request(
346        &mut self,
347        result: mj_core::subagent::SubagentToolResult,
348    ) -> Result<()> {
349        let request = RelayRequest::CompleteSubagentRequest { result };
350        match self.call(request).await? {
351            RelayResponsePayload::SubagentRequestCompleted => Ok(()),
352            _ => bail!("relay returned an unexpected sub-agent completion response"),
353        }
354    }
355
356    pub async fn detach(mut self) -> Result<()> {
357        self.input
358            .take()
359            .expect("connected relay owns proxy stdin")
360            .shutdown()
361            .await
362            .context("close relay proxy stdin")?;
363        let mut child = self.child.take().expect("connected relay owns proxy child");
364        match tokio::time::timeout(RELAY_PROXY_DETACH_GRACE, child.wait()).await {
365            Ok(status) => {
366                status.context("wait for relay proxy")?;
367            }
368            Err(_) => {
369                if let Err(error) = child.start_kill().context("stop relay proxy") {
370                    tracing::warn!(
371                        session_id = %self.session_id,
372                        operation = "detach",
373                        %error,
374                        "could not stop relay proxy after detach timeout"
375                    );
376                    return Err(error);
377                }
378                if let Err(error) = child.wait().await {
379                    tracing::warn!(
380                        session_id = %self.session_id,
381                        operation = "detach",
382                        %error,
383                        "could not reap relay proxy after stopping it"
384                    );
385                }
386            }
387        }
388        Ok(())
389    }
390}