mj_controller/worker_client/
reviewer.rs1use super::*;
2
3impl RelayClient {
4 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 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 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 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 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 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 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 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 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 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 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 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 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 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}