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 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 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 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 pub(super) fn reviewer_request(
248 &self,
249 role: Option<&str>,
250 request: ReviewerRequest,
251 ) -> Result<RelayRequest> {
252 let request = RelayRequest::Reviewer {
253 role: role.map(str::to_owned),
254 request,
255 };
256 Ok(request)
257 }
258
259 pub async fn respond_elicitation(
262 &mut self,
263 elicitation_id: String,
264 response: ElicitationResponse,
265 ) -> Result<()> {
266 let request = RelayRequest::RespondElicitation {
267 elicitation_id: elicitation_id.clone(),
268 response,
269 };
270 match self.call(request).await? {
271 RelayResponsePayload::ElicitationResolved {
272 elicitation_id: resolved,
273 } if resolved == elicitation_id => Ok(()),
274 RelayResponsePayload::ElicitationResolved {
275 elicitation_id: resolved,
276 } => bail!("relay resolved elicitation {resolved:?}, expected {elicitation_id:?}"),
277 _ => bail!("relay returned an unexpected elicitation response"),
278 }
279 }
280
281 pub async fn stop_background_task(&mut self, background_task_id: String) -> Result<()> {
283 let request = RelayRequest::StopBackgroundTask {
284 background_task_id: background_task_id.clone(),
285 };
286 match self.call(request).await? {
287 RelayResponsePayload::BackgroundTaskStopRequested {
288 background_task_id: stopped,
289 } if stopped == background_task_id => Ok(()),
290 RelayResponsePayload::BackgroundTaskStopRequested {
291 background_task_id: stopped,
292 } => {
293 bail!("relay stopped background task {stopped:?}, expected {background_task_id:?}")
294 }
295 _ => bail!("relay returned an unexpected background task stop response"),
296 }
297 }
298
299 pub async fn subagent_requests(
300 &mut self,
301 ) -> Result<(
302 Vec<mj_core::subagent::SubagentToolRequest>,
303 Vec<mj_core::subagent::SubagentToolResult>,
304 )> {
305 let request = RelayRequest::SubagentRequests;
306 if !request.supported_at(self.protocol_version) {
307 return Ok((Vec::new(), Vec::new()));
308 }
309 match self.call(request).await? {
310 RelayResponsePayload::SubagentRequests { requests, results } => Ok((requests, results)),
311 _ => bail!("relay returned an unexpected sub-agent request response"),
312 }
313 }
314
315 pub async fn set_subagent_admission(&mut self, open: bool) -> Result<()> {
316 let request = RelayRequest::SetSubagentAdmission { open };
317 if !request.supported_at(self.protocol_version) {
318 bail!(
319 "source worker does not support safe in-place sub-agent draining; update or restart the session worker, then retry Move"
320 );
321 }
322 match self.call(request).await? {
323 RelayResponsePayload::SubagentAdmissionChanged { open: returned }
324 if returned == open =>
325 {
326 Ok(())
327 }
328 RelayResponsePayload::SubagentAdmissionChanged { open: returned } => {
329 bail!("worker set sub-agent admission to {returned}, expected {open}")
330 }
331 _ => bail!("relay returned an unexpected sub-agent admission response"),
332 }
333 }
334
335 pub async fn complete_subagent_request(
336 &mut self,
337 result: mj_core::subagent::SubagentToolResult,
338 ) -> Result<()> {
339 let request = RelayRequest::CompleteSubagentRequest { result };
340 match self.call(request).await? {
341 RelayResponsePayload::SubagentRequestCompleted => Ok(()),
342 _ => bail!("relay returned an unexpected sub-agent completion response"),
343 }
344 }
345
346 pub async fn detach(mut self) -> Result<()> {
347 self.input
348 .take()
349 .expect("connected relay owns proxy stdin")
350 .shutdown()
351 .await
352 .context("close relay proxy stdin")?;
353 let mut child = self.child.take().expect("connected relay owns proxy child");
354 match tokio::time::timeout(RELAY_PROXY_DETACH_GRACE, child.wait()).await {
355 Ok(status) => {
356 status.context("wait for relay proxy")?;
357 }
358 Err(_) => {
359 if let Err(error) = child.start_kill().context("stop relay proxy") {
360 tracing::warn!(
361 session_id = %self.session_id,
362 operation = "detach",
363 %error,
364 "could not stop relay proxy after detach timeout"
365 );
366 return Err(error);
367 }
368 if let Err(error) = child.wait().await {
369 tracing::warn!(
370 session_id = %self.session_id,
371 operation = "detach",
372 %error,
373 "could not reap relay proxy after stopping it"
374 );
375 }
376 }
377 }
378 Ok(())
379 }
380}