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 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 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 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 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 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}