1use 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 Ok(RelayResponsePayload::SubagentRequests { requests, results }) => {
311 Ok((requests, results))
312 }
313 Ok(_) => bail!("relay returned an unexpected sub-agent request response"),
314 Err(error)
320 if error
321 .downcast_ref::<RelayRejected>()
322 .is_some_and(|rejected| rejected.0.code == RelayErrorCode::InvalidState) =>
323 {
324 Ok((Vec::new(), Vec::new()))
325 }
326 Err(error) => Err(error),
327 }
328 }
329
330 pub async fn set_subagent_admission(&mut self, open: bool) -> Result<()> {
331 let request = RelayRequest::SetSubagentAdmission { open };
332 if !request.supported_at(self.protocol_version) {
333 bail!(
334 "source worker does not support safe in-place sub-agent draining; update or restart the session worker, then retry Move"
335 );
336 }
337 match self.call(request).await? {
338 RelayResponsePayload::SubagentAdmissionChanged { open: returned }
339 if returned == open =>
340 {
341 Ok(())
342 }
343 RelayResponsePayload::SubagentAdmissionChanged { open: returned } => {
344 bail!("worker set sub-agent admission to {returned}, expected {open}")
345 }
346 _ => bail!("relay returned an unexpected sub-agent admission response"),
347 }
348 }
349
350 pub async fn complete_subagent_request(
351 &mut self,
352 result: mj_core::subagent::SubagentToolResult,
353 ) -> Result<bool> {
354 let request = RelayRequest::CompleteSubagentRequest { result };
355 subagent_completion_delivery(self.call(request).await?)
356 }
357
358 pub async fn detach(mut self) -> Result<()> {
359 self.input
360 .take()
361 .expect("connected relay owns proxy stdin")
362 .shutdown()
363 .await
364 .context("close relay proxy stdin")?;
365 let mut child = self.child.take().expect("connected relay owns proxy child");
366 match tokio::time::timeout(RELAY_PROXY_DETACH_GRACE, child.wait()).await {
367 Ok(status) => {
368 status.context("wait for relay proxy")?;
369 }
370 Err(_) => {
371 if let Err(error) = child.start_kill().context("stop relay proxy") {
372 tracing::warn!(
373 session_id = %self.session_id,
374 operation = "detach",
375 %error,
376 "could not stop relay proxy after detach timeout"
377 );
378 return Err(error);
379 }
380 if let Err(error) = child.wait().await {
381 tracing::warn!(
382 session_id = %self.session_id,
383 operation = "detach",
384 %error,
385 "could not reap relay proxy after stopping it"
386 );
387 }
388 }
389 }
390 Ok(())
391 }
392}
393
394fn subagent_completion_delivery(payload: RelayResponsePayload) -> Result<bool> {
395 match payload {
396 RelayResponsePayload::SubagentRequestCompleted => Ok(true),
399 RelayResponsePayload::SubagentRequestCompletedWithDelivery {
400 delivered_to_waiter,
401 } => Ok(delivered_to_waiter),
402 _ => bail!("relay returned an unexpected sub-agent completion response"),
403 }
404}
405
406#[cfg(test)]
407mod tests {
408 use super::*;
409
410 #[test]
411 fn old_worker_completion_without_delivery_status_counts_as_delivered() {
412 assert!(
413 subagent_completion_delivery(RelayResponsePayload::SubagentRequestCompleted).unwrap()
414 );
415 assert!(
416 !subagent_completion_delivery(
417 RelayResponsePayload::SubagentRequestCompletedWithDelivery {
418 delivered_to_waiter: false,
419 }
420 )
421 .unwrap()
422 );
423 }
424}