1use std::collections::{BTreeMap, BTreeSet, VecDeque};
4use std::path::Path;
5use std::process::Stdio;
6use std::sync::Arc;
7use std::time::{Duration, Instant};
8
9use anyhow::{Context, Result, anyhow, bail};
10use base64::Engine as _;
11use base64::engine::general_purpose::STANDARD as BASE64;
12use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWriteExt, BufReader};
13use tokio::process::{Child, ChildStdin, ChildStdout, Command};
14use tokio::sync::{mpsc, watch};
15
16use crate::targets::CommandSpec;
17use mj_core::config::harness_authentication_marker;
18use mj_core::credentials::{
19 CredentialSnapshot, CredentialSyncAction, CredentialSyncHandle, CredentialSyncOutcome,
20 CredentialSyncResult, CredentialSyncTarget, SYNC_INTERVAL, SyncAction, SyncTrigger, enqueue,
21 profiles_with_targets, read_credential_file, reconcile, validate_credential_payload,
22 write_credential_file,
23};
24use mj_core::elicitation::ElicitationResponse;
25use mj_core::relay::{
26 MAX_FRAME_BYTES, RELAY_EVENT_GENESIS_DIGEST, RELAY_MIN_PROTOCOL_VERSION,
27 RELAY_PROTOCOL_VERSION, RelayCommand, RelayCursor, RelayErrorCode, RelayEvent,
28 RelayOperationalState, RelayProtocolError, RelayRequest, RelayRequestEnvelope,
29 RelayResponseBody, RelayResponseEnvelope, RelayResponsePayload, RelayVersionRange,
30 ReviewerRequest, validate_relay_event,
31};
32
33pub use mj_client::session::{RelayAttachment, StartedReviewer};
34use mj_core::worker_launch::ReviewerLaunchConfig;
35
36const RELAY_RPC_TIMEOUT: Duration = Duration::from_secs(15);
37const RELAY_SLOW_OPERATION_WARNING: Duration = Duration::from_secs(5);
38const RELAY_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(300);
42const RELAY_HISTORY_TIMEOUT: Duration = Duration::from_secs(900);
46const RELAY_ACKNOWLEDGE_TIMEOUT: Duration = Duration::from_secs(300);
50const REVIEW_CAPTURE_TIMEOUT: Duration = Duration::from_secs(300);
53const REVIEW_ANALYSIS_TIMEOUT: Duration = Duration::from_secs(660);
57const RELAY_PROXY_DETACH_GRACE: Duration = Duration::from_millis(500);
58const RELAY_PROXY_REAP_POLL: Duration = Duration::from_millis(10);
59
60async fn drain_proxy_stderr(
64 errors: tokio::process::ChildStderr,
65 purpose: String,
66 session_id: String,
67) {
68 let mut lines = BufReader::new(errors).lines();
69 loop {
70 match lines.next_line().await {
71 Ok(Some(line)) if line.trim().is_empty() => continue,
72 Ok(Some(line)) => {
73 tracing::warn!(%session_id, %purpose, %line, "relay proxy stderr")
74 }
75 Ok(None) => return,
76 Err(error) => {
77 tracing::warn!(%session_id, %purpose, %error, "read relay proxy stderr");
78 return;
79 }
80 }
81 }
82}
83
84#[derive(Debug, Clone)]
88pub struct RelayEventPage {
89 pub events: Vec<RelayEvent>,
90 pub through_ordinal: u64,
91 pub through_digest: String,
92}
93
94#[derive(Debug, Clone)]
95pub struct RelayCatchUp {
96 pub state: RelayOperationalState,
97 pub frontier: RelayCursor,
98 pub first_page: RelayEventPage,
99}
100
101#[derive(Debug)]
102pub struct RelayRejected(pub RelayProtocolError);
103
104impl std::fmt::Display for RelayRejected {
105 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
106 write!(
107 formatter,
108 "relay rejected request ({:?}): {}",
109 self.0.code, self.0.message
110 )
111 }
112}
113
114impl std::error::Error for RelayRejected {}
115
116impl RelayRejected {
117 pub fn is_desynchronized(&self) -> bool {
118 self.0.code == RelayErrorCode::Desynchronized
119 }
120
121 pub fn is_retryable(&self) -> bool {
124 self.0.retryable
125 }
126}
127
128#[derive(Debug)]
135pub struct RelayTransportDead {
136 message: String,
137 handshake_failed: bool,
138}
139
140impl std::fmt::Display for RelayTransportDead {
141 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
142 formatter.write_str(&self.message)
143 }
144}
145
146impl std::error::Error for RelayTransportDead {}
147
148impl RelayTransportDead {
149 pub fn new(message: impl Into<String>) -> Self {
150 Self {
151 message: message.into(),
152 handshake_failed: false,
153 }
154 }
155
156 fn from_io(error: std::io::Error, kind: ExchangeKind) -> Self {
159 Self::during_exchange(error.to_string(), kind)
160 }
161
162 fn during_exchange(message: impl Into<String>, kind: ExchangeKind) -> Self {
163 Self {
164 message: message.into(),
165 handshake_failed: kind == ExchangeKind::Handshake,
166 }
167 }
168
169 pub fn marks(error: &anyhow::Error) -> bool {
171 error.downcast_ref::<Self>().is_some()
172 }
173
174 pub fn marks_failed_handshake(error: &anyhow::Error) -> bool {
179 error
180 .downcast_ref::<Self>()
181 .is_some_and(|failure| failure.handshake_failed)
182 }
183}
184
185#[derive(Clone, Copy, PartialEq, Eq)]
192enum ExchangeKind {
193 Handshake,
194 Call,
195}
196
197pub struct RelayClient {
203 child: Option<Child>,
204 input: Option<ChildStdin>,
205 output: BufReader<ChildStdout>,
206 request_timeout: Duration,
207 abandoned: Option<String>,
210 next_request: u64,
211 connection_nonce: u64,
212 protocol_version: u32,
213 session_id: String,
214 relay_version: String,
215 worker_build: Option<String>,
218 latest_ordinal: u64,
219 latest_digest: String,
220}
221
222impl RelayClient {
223 pub async fn connect(spec: &CommandSpec, expected_session_id: &str) -> Result<Self> {
224 Self::connect_with_timeouts(
225 spec,
226 expected_session_id,
227 RELAY_RPC_TIMEOUT,
228 RELAY_HANDSHAKE_TIMEOUT,
229 )
230 .await
231 }
232
233 #[cfg(all(test, unix))]
234 async fn connect_with_timeout(
235 spec: &CommandSpec,
236 expected_session_id: &str,
237 request_timeout: Duration,
238 ) -> Result<Self> {
239 Self::connect_with_timeouts(spec, expected_session_id, request_timeout, request_timeout)
240 .await
241 }
242
243 async fn connect_with_timeouts(
244 spec: &CommandSpec,
245 expected_session_id: &str,
246 request_timeout: Duration,
247 handshake_timeout: Duration,
248 ) -> Result<Self> {
249 let mut child = Command::new(&spec.program)
250 .args(&spec.args)
251 .envs(&spec.env)
252 .stdin(Stdio::piped())
253 .stdout(Stdio::piped())
254 .stderr(Stdio::piped())
258 .kill_on_drop(true)
259 .spawn()
260 .with_context(|| format!("start session relay proxy for {}", spec.purpose))
261 .map_err(|error| {
262 tracing::warn!(
263 session_id = %expected_session_id,
264 operation = "connect",
265 purpose = %spec.purpose,
266 error = %error,
267 "could not start relay proxy"
268 );
269 error
270 })?;
271 if let Some(errors) = child.stderr.take() {
272 let purpose = spec.purpose.clone();
273 let session_id = expected_session_id.to_owned();
274 tokio::spawn(drain_proxy_stderr(errors, purpose, session_id));
275 }
276 let input = child
277 .stdin
278 .take()
279 .context("relay proxy stdin unavailable")
280 .map_err(|error| {
281 tracing::warn!(
282 session_id = %expected_session_id,
283 operation = "connect",
284 purpose = %spec.purpose,
285 error = %error,
286 "relay proxy did not provide stdin"
287 );
288 error
289 })?;
290 let output = child
291 .stdout
292 .take()
293 .context("relay proxy stdout unavailable")
294 .map_err(|error| {
295 tracing::warn!(
296 session_id = %expected_session_id,
297 operation = "connect",
298 purpose = %spec.purpose,
299 error = %error,
300 "relay proxy did not provide stdout"
301 );
302 error
303 })?;
304 let mut nonce_bytes = [0_u8; 8];
305 getrandom::fill(&mut nonce_bytes).map_err(|error| {
306 let error = anyhow!("generate relay request nonce: {error}");
307 tracing::warn!(
308 session_id = %expected_session_id,
309 operation = "connect",
310 error = %error,
311 "could not initialize relay request nonce"
312 );
313 error
314 })?;
315 let mut client = Self {
316 child: Some(child),
317 input: Some(input),
318 output: BufReader::new(output),
319 request_timeout,
320 abandoned: None,
321 next_request: 1,
322 connection_nonce: u64::from_le_bytes(nonce_bytes),
323 protocol_version: RELAY_PROTOCOL_VERSION,
324 session_id: expected_session_id.to_owned(),
328 relay_version: String::new(),
329 worker_build: None,
330 latest_ordinal: 0,
331 latest_digest: RELAY_EVENT_GENESIS_DIGEST.to_owned(),
332 };
333 let response = client
334 .call_hello(
335 RelayRequest::Hello {
336 controller_version: env!("CARGO_PKG_VERSION").to_owned(),
337 supported: RelayVersionRange::CURRENT,
338 },
339 handshake_timeout,
340 )
341 .await?;
342 let RelayResponsePayload::Hello {
343 negotiated,
344 relay_version,
345 session_id,
346 worker_build,
347 } = response
348 else {
349 let error = anyhow!("relay returned an unexpected hello response");
350 log_relay_client_failure(&client, "hello", "relay-hello", &error);
351 return Err(error);
352 };
353 if session_id != expected_session_id {
354 let error = anyhow!("relay belongs to session {session_id}, not {expected_session_id}");
355 log_relay_client_failure(&client, "hello", "relay-hello", &error);
356 return Err(error);
357 }
358 if !RelayVersionRange::CURRENT.contains(negotiated) {
359 let error = anyhow!(
360 "relay negotiated unsupported protocol {negotiated}; this controller supports {}-{}",
361 RELAY_MIN_PROTOCOL_VERSION,
362 RELAY_PROTOCOL_VERSION
363 );
364 log_relay_client_failure(&client, "hello", "relay-hello", &error);
365 return Err(error);
366 }
367 client.protocol_version = negotiated;
368 client.session_id = session_id;
369 client.relay_version = relay_version;
370 client.worker_build = worker_build;
371 Ok(client)
372 }
373
374 pub fn session_id(&self) -> &str {
375 &self.session_id
376 }
377
378 pub const fn supports_project_memory_sync(&self) -> bool {
379 self.protocol_version >= 4
380 }
381
382 pub fn relay_version(&self) -> &str {
383 &self.relay_version
384 }
385
386 pub fn worker_build(&self) -> Option<&str> {
390 self.worker_build.as_deref()
391 }
392
393 pub fn protocol_version(&self) -> u32 {
394 self.protocol_version
395 }
396
397 pub fn latest_ordinal(&self) -> u64 {
398 self.latest_ordinal
399 }
400
401 pub fn latest_digest(&self) -> &str {
402 &self.latest_digest
403 }
404
405 pub async fn attach(
406 &mut self,
407 after_ordinal: u64,
408 after_digest: impl Into<String>,
409 ) -> Result<RelayAttachment> {
410 let after_digest = after_digest.into();
411 match self
412 .call_with_timeout(
413 RelayRequest::Attach {
414 after_ordinal,
415 after_digest: after_digest.clone(),
416 },
417 RELAY_HISTORY_TIMEOUT,
418 )
419 .await?
420 {
421 RelayResponsePayload::Attached {
422 state,
423 events,
424 through_ordinal,
425 through_digest,
426 } => {
427 let mut cursor = RelayCursor {
428 ordinal: after_ordinal,
429 digest: after_digest,
430 };
431 for event in &events {
432 validate_relay_event(cursor.ordinal, &cursor.digest, event)
433 .context("verify relay attachment event chain")?;
434 cursor.ordinal = event.ordinal;
435 cursor.digest.clone_from(&event.digest);
436 }
437 if cursor.ordinal != through_ordinal || cursor.digest != through_digest {
438 bail!("relay attachment frontier does not match its event chain");
439 }
440 self.latest_ordinal = state.latest_ordinal;
441 self.latest_digest = state.latest_digest.clone();
442 Ok(RelayAttachment {
443 state,
444 events,
445 through_ordinal,
446 through_digest,
447 })
448 }
449 _ => bail!("relay returned an unexpected attach response"),
450 }
451 }
452
453 pub async fn begin_catch_up(
458 &mut self,
459 after_ordinal: u64,
460 after_digest: impl Into<String>,
461 ) -> Result<RelayCatchUp> {
462 let after_digest = after_digest.into();
463 let first = self.attach(after_ordinal, after_digest.clone()).await?;
464 let frontier = RelayCursor {
465 ordinal: first.state.latest_ordinal,
466 digest: first.state.latest_digest.clone(),
467 };
468 let previous = RelayCursor {
469 ordinal: after_ordinal,
470 digest: after_digest,
471 };
472 let state = first.state.clone();
473 let first_page = clip_catch_up_page(first, &previous, &frontier)?;
474 Ok(RelayCatchUp {
475 state,
476 frontier,
477 first_page,
478 })
479 }
480
481 pub async fn next_catch_up_page(
485 &mut self,
486 previous: &RelayCursor,
487 frontier: &RelayCursor,
488 ) -> Result<RelayEventPage> {
489 if previous.ordinal >= frontier.ordinal {
490 bail!("relay catch-up is already at its fixed frontier");
491 }
492 let attachment = self
493 .attach(previous.ordinal, previous.digest.clone())
494 .await?;
495 clip_catch_up_page(attachment, previous, frontier)
496 }
497
498 pub async fn acknowledge(
499 &mut self,
500 through_ordinal: u64,
501 through_digest: impl Into<String>,
502 ) -> Result<RelayCursor> {
503 match self
504 .call_with_timeout(
505 RelayRequest::Acknowledge {
506 through_ordinal,
507 through_digest: through_digest.into(),
508 },
509 RELAY_ACKNOWLEDGE_TIMEOUT,
510 )
511 .await?
512 {
513 RelayResponsePayload::Acknowledged {
514 through_ordinal,
515 through_digest,
516 } => Ok(RelayCursor {
517 ordinal: through_ordinal,
518 digest: through_digest,
519 }),
520 _ => bail!("relay returned an unexpected acknowledgement response"),
521 }
522 }
523
524 pub async fn status(&mut self) -> Result<RelayOperationalState> {
525 match self.call(RelayRequest::Status).await? {
526 RelayResponsePayload::Status(status) => {
527 self.latest_ordinal = status.latest_ordinal;
528 self.latest_digest = status.latest_digest.clone();
529 Ok(status)
530 }
531 _ => bail!("relay returned an unexpected status response"),
532 }
533 }
534
535 pub async fn credential_state(&mut self) -> Result<CredentialSnapshot> {
538 credential_snapshot(self.call(RelayRequest::CredentialState).await?)
539 }
540
541 pub async fn read_credentials(&mut self) -> Result<Vec<u8>> {
544 match self.call(RelayRequest::ReadCredentials).await? {
545 RelayResponsePayload::Credentials { data } => BASE64
546 .decode(data.as_bytes())
547 .context("decode relay credential payload"),
548 _ => bail!("relay returned an unexpected credential response"),
549 }
550 }
551
552 pub async fn install_credentials(&mut self, bytes: &[u8]) -> Result<CredentialSnapshot> {
555 credential_snapshot(
556 self.call(RelayRequest::InstallCredentials {
557 data: BASE64.encode(bytes),
558 })
559 .await?,
560 )
561 }
562
563 pub async fn github_token_state(
564 &mut self,
565 ) -> Result<mj_core::credentials::GithubTokenSnapshot> {
566 github_token_snapshot(self.call(RelayRequest::GithubTokenState).await?)
567 }
568
569 pub async fn install_github_token(
570 &mut self,
571 token: &str,
572 ) -> Result<mj_core::credentials::GithubTokenSnapshot> {
573 github_token_snapshot(
574 self.call(RelayRequest::InstallGithubToken {
575 data: BASE64.encode(token.as_bytes()),
576 })
577 .await?,
578 )
579 }
580
581 pub async fn remove_github_token(
582 &mut self,
583 ) -> Result<mj_core::credentials::GithubTokenSnapshot> {
584 github_token_snapshot(self.call(RelayRequest::RemoveGithubToken).await?)
585 }
586
587 pub async fn skills_state(&mut self) -> Result<mj_core::skills::SkillsSyncState> {
590 skills_sync_state(self.call(RelayRequest::SkillsState).await?)
591 }
592
593 pub async fn install_prompt_context(&mut self, text: String) -> Result<()> {
596 let request = RelayRequest::InstallPromptContext { text };
597 if !request.supported_at(self.protocol_version) {
598 bail!(
599 "hidden prompt context requires relay protocol {}; this session negotiated {}",
600 request.minimum_protocol(),
601 self.protocol_version
602 );
603 }
604 match self.call(request).await? {
605 RelayResponsePayload::PromptContextInstalled => Ok(()),
606 _ => bail!("relay returned an unexpected prompt-context response"),
607 }
608 }
609
610 pub async fn project_memory_snapshot(
611 &mut self,
612 ) -> Result<(
613 mj_core::project_memory::ProjectMemorySnapshot,
614 mj_core::project_memory::ProjectMemorySnapshot,
615 )> {
616 let request = RelayRequest::ProjectMemorySnapshot;
617 if !request.supported_at(self.protocol_version) {
618 bail!(
619 "project memory synchronization requires relay protocol {}; this session negotiated {}",
620 request.minimum_protocol(),
621 self.protocol_version
622 );
623 }
624 match self.call(request).await? {
625 RelayResponsePayload::ProjectMemorySnapshot { baseline, replica } => {
626 Ok((baseline, replica))
627 }
628 _ => bail!("relay returned an unexpected project-memory response"),
629 }
630 }
631
632 pub async fn install_project_memory_snapshot(
633 &mut self,
634 snapshot: mj_core::project_memory::ProjectMemorySnapshot,
635 ) -> Result<()> {
636 let request = RelayRequest::InstallProjectMemorySnapshot { snapshot };
637 if !request.supported_at(self.protocol_version) {
638 bail!(
639 "project memory synchronization requires relay protocol {}; this session negotiated {}",
640 request.minimum_protocol(),
641 self.protocol_version
642 );
643 }
644 match self.call(request).await? {
645 RelayResponsePayload::ProjectMemorySnapshotInstalled => Ok(()),
646 _ => bail!("relay returned an unexpected project-memory install response"),
647 }
648 }
649
650 pub async fn install_skills(
654 &mut self,
655 archive_bytes: &[u8],
656 ) -> Result<mj_core::skills::SkillsSyncState> {
657 skills_sync_state(
658 self.call(RelayRequest::InstallSkills {
659 data: BASE64.encode(archive_bytes),
660 })
661 .await?,
662 )
663 }
664
665 pub async fn ensure_attachment(
667 &mut self,
668 reference: &mj_core::attachment::AttachmentRef,
669 ) -> Result<()> {
670 anyhow::ensure!(
671 self.protocol_version >= 8,
672 "photo attachments require an updated worker (protocol 8); upgrade the worker and retry"
673 );
674 match self
675 .call(RelayRequest::AttachmentPresent {
676 reference: reference.clone(),
677 })
678 .await?
679 {
680 RelayResponsePayload::AttachmentPresent { present: true } => return Ok(()),
681 RelayResponsePayload::AttachmentPresent { present: false } => {}
682 _ => bail!("unexpected image presence response"),
683 }
684 let store = mj_core::attachment::AttachmentStore::controller(&self.session_id)?;
685 let reference_copy = reference.clone();
686 let bytes = tokio::task::spawn_blocking(move || store.read(&reference_copy))
687 .await
688 .context("image loading task failed")??;
689 match self
690 .call(RelayRequest::InstallAttachment {
691 reference: reference.clone(),
692 data: BASE64.encode(bytes),
693 })
694 .await?
695 {
696 RelayResponsePayload::AttachmentInstalled => Ok(()),
697 _ => bail!("unexpected image upload response"),
698 }
699 }
700
701 pub async fn cache_attachment(
703 &mut self,
704 reference: &mj_core::attachment::AttachmentRef,
705 ) -> Result<()> {
706 let store = mj_core::attachment::AttachmentStore::controller(&self.session_id)?;
707 let local = store.clone();
708 let reference_copy = reference.clone();
709 if tokio::task::spawn_blocking(move || local.contains(&reference_copy))
710 .await
711 .context("image lookup task failed")??
712 {
713 return Ok(());
714 }
715 let RelayResponsePayload::AttachmentData { data } = self
716 .call(RelayRequest::ReadAttachment {
717 reference: reference.clone(),
718 })
719 .await?
720 else {
721 bail!("unexpected image download response")
722 };
723 let reference = reference.clone();
724 tokio::task::spawn_blocking(move || {
725 anyhow::ensure!(
726 data.len() <= mj_core::attachment::MAX_IMAGE_BYTES.div_ceil(3) * 4,
727 "image download is too large"
728 );
729 store.install(&reference, &BASE64.decode(data)?)
730 })
731 .await
732 .context("image caching task failed")?
733 }
734
735 pub async fn submit(
736 &mut self,
737 command_id: impl Into<String>,
738 command: RelayCommand,
739 ) -> Result<u64> {
740 let command_id = command_id.into();
741 if let RelayCommand::Prompt { prompt } = &command {
742 for reference in mj_core::attachment::references(prompt)? {
743 self.ensure_attachment(&reference).await?;
744 }
745 }
746 match self
747 .call(RelayRequest::Submit {
748 command_id: command_id.clone(),
749 command,
750 })
751 .await?
752 {
753 RelayResponsePayload::Accepted {
754 command_id: accepted_id,
755 ordinal,
756 } if accepted_id == command_id => Ok(ordinal),
757 RelayResponsePayload::Accepted {
758 command_id: accepted_id,
759 ..
760 } => bail!("relay accepted command under ID {accepted_id}, expected {command_id}"),
761 _ => bail!("relay returned an unexpected command response"),
762 }
763 }
764
765 pub async fn start_reviewer(
772 &mut self,
773 role: Option<&str>,
774 config: ReviewerLaunchConfig,
775 ) -> Result<StartedReviewer> {
776 let request = self.reviewer_request(
777 role,
778 ReviewerRequest::Start {
779 config: Box::new(config),
780 },
781 )?;
782 match self
783 .call_with_timeout(request, RELAY_HANDSHAKE_TIMEOUT)
784 .await?
785 {
786 RelayResponsePayload::ReviewerStarted {
787 native_session_id,
788 config_options,
789 reused,
790 state,
791 } => Ok(StartedReviewer {
792 native_session_id,
793 config_options,
794 reused,
795 state: *state,
796 }),
797 _ => bail!("relay returned an unexpected reviewer start response"),
798 }
799 }
800
801 pub async fn attach_reviewer(
804 &mut self,
805 role: Option<&str>,
806 after_ordinal: u64,
807 after_digest: impl Into<String>,
808 ) -> Result<RelayAttachment> {
809 let after_digest = after_digest.into();
810 let request = self.reviewer_request(
811 role,
812 ReviewerRequest::Attach {
813 after_ordinal,
814 after_digest: after_digest.clone(),
815 },
816 )?;
817 let payload = self
818 .call_with_timeout(request, RELAY_HISTORY_TIMEOUT)
819 .await?;
820 let RelayResponsePayload::Attached {
821 state,
822 events,
823 through_ordinal,
824 through_digest,
825 } = payload
826 else {
827 bail!("relay returned an unexpected reviewer attach response");
828 };
829 let mut cursor = RelayCursor {
832 ordinal: after_ordinal,
833 digest: after_digest,
834 };
835 for event in &events {
836 validate_relay_event(cursor.ordinal, &cursor.digest, event)
837 .context("verify reviewer attachment event chain")?;
838 cursor.ordinal = event.ordinal;
839 cursor.digest.clone_from(&event.digest);
840 }
841 if cursor.ordinal != through_ordinal || cursor.digest != through_digest {
842 bail!("reviewer attachment frontier does not match its event chain");
843 }
844 Ok(RelayAttachment {
845 state,
846 events,
847 through_ordinal,
848 through_digest,
849 })
850 }
851
852 pub async fn acknowledge_reviewer(
855 &mut self,
856 role: Option<&str>,
857 through_ordinal: u64,
858 through_digest: impl Into<String>,
859 ) -> Result<RelayCursor> {
860 let request = self.reviewer_request(
861 role,
862 ReviewerRequest::Acknowledge {
863 through_ordinal,
864 through_digest: through_digest.into(),
865 },
866 )?;
867 match self
868 .call_with_timeout(request, RELAY_ACKNOWLEDGE_TIMEOUT)
869 .await?
870 {
871 RelayResponsePayload::Acknowledged {
872 through_ordinal,
873 through_digest,
874 } => Ok(RelayCursor {
875 ordinal: through_ordinal,
876 digest: through_digest,
877 }),
878 _ => bail!("relay returned an unexpected reviewer acknowledgement response"),
879 }
880 }
881
882 pub async fn submit_to_reviewer(
884 &mut self,
885 role: Option<&str>,
886 command_id: impl Into<String>,
887 command: RelayCommand,
888 ) -> Result<u64> {
889 let command_id = command_id.into();
890 let request = self.reviewer_request(
891 role,
892 ReviewerRequest::Submit {
893 command_id: command_id.clone(),
894 command,
895 },
896 )?;
897 match self.call(request).await? {
898 RelayResponsePayload::Accepted {
899 command_id: accepted_id,
900 ordinal,
901 } if accepted_id == command_id => Ok(ordinal),
902 RelayResponsePayload::Accepted {
903 command_id: accepted_id,
904 ..
905 } => bail!("reviewer accepted command under ID {accepted_id}, expected {command_id}"),
906 _ => bail!("relay returned an unexpected reviewer command response"),
907 }
908 }
909
910 pub async fn reviewer_status(&mut self, role: Option<&str>) -> Result<RelayOperationalState> {
911 let request = self.reviewer_request(role, ReviewerRequest::Status)?;
912 match self.call(request).await? {
913 RelayResponsePayload::Status(status) => Ok(status),
914 _ => bail!("relay returned an unexpected reviewer status response"),
915 }
916 }
917
918 pub async fn respond_to_reviewer(
920 &mut self,
921 role: Option<&str>,
922 elicitation_id: String,
923 response: ElicitationResponse,
924 ) -> Result<()> {
925 let request = self.reviewer_request(
926 role,
927 ReviewerRequest::RespondElicitation {
928 elicitation_id: elicitation_id.clone(),
929 response,
930 },
931 )?;
932 match self.call(request).await? {
933 RelayResponsePayload::ElicitationResolved {
934 elicitation_id: resolved,
935 } if resolved == elicitation_id => Ok(()),
936 RelayResponsePayload::ElicitationResolved {
937 elicitation_id: resolved,
938 } => bail!("reviewer resolved elicitation {resolved:?}, expected {elicitation_id:?}"),
939 _ => bail!("relay returned an unexpected reviewer elicitation response"),
940 }
941 }
942
943 pub async fn pause_reviewer(&mut self, role: Option<&str>) -> Result<()> {
946 let request = self.reviewer_request(role, ReviewerRequest::Pause)?;
947 match self
948 .call_with_timeout(request, RELAY_ACKNOWLEDGE_TIMEOUT)
949 .await?
950 {
951 RelayResponsePayload::ReviewerPaused => Ok(()),
952 _ => bail!("relay returned an unexpected reviewer pause response"),
953 }
954 }
955
956 pub async fn capture_review_delta(
959 &mut self,
960 role: Option<&str>,
961 baselines: std::collections::BTreeMap<std::path::PathBuf, String>,
962 ) -> Result<Vec<mj_core::relay::RepoDelta>> {
963 let request = self.reviewer_request(role, ReviewerRequest::CaptureDelta { baselines })?;
964 match self
965 .call_with_timeout(request, REVIEW_CAPTURE_TIMEOUT)
966 .await?
967 {
968 RelayResponsePayload::ReviewDelta { repositories } => Ok(repositories),
969 _ => bail!("relay returned an unexpected review capture response"),
970 }
971 }
972
973 pub async fn advance_review_baseline(
976 &mut self,
977 role: Option<&str>,
978 trees: std::collections::BTreeMap<std::path::PathBuf, String>,
979 ) -> Result<()> {
980 let request = self.reviewer_request(role, ReviewerRequest::AdvanceBaseline { trees })?;
981 match self
982 .call_with_timeout(request, REVIEW_CAPTURE_TIMEOUT)
983 .await?
984 {
985 RelayResponsePayload::ReviewBaselineAdvanced => Ok(()),
986 _ => bail!("relay returned an unexpected review baseline response"),
987 }
988 }
989
990 pub async fn analyze_review_delta(
993 &mut self,
994 role: Option<&str>,
995 repositories: Vec<mj_core::relay::AnalyzeDeltaRepository>,
996 ) -> Result<String> {
997 let request =
998 self.reviewer_request(role, ReviewerRequest::AnalyzeDelta { repositories })?;
999 match self
1000 .call_with_timeout(request, REVIEW_ANALYSIS_TIMEOUT)
1001 .await?
1002 {
1003 RelayResponsePayload::ReviewChangedFunctions { packet } => Ok(packet),
1004 _ => bail!("relay returned an unexpected review analysis response"),
1005 }
1006 }
1007
1008 pub async fn take_lane_dispatches(
1011 &mut self,
1012 ) -> Result<Vec<mj_core::review::lanes::ReviewSubagentRequest>> {
1013 let request = self.reviewer_request(None, ReviewerRequest::TakeLaneDispatches)?;
1014 match self.call(request).await? {
1015 RelayResponsePayload::LaneDispatches { requests } => Ok(requests),
1016 _ => bail!("relay returned an unexpected lane dispatch response"),
1017 }
1018 }
1019
1020 fn reviewer_request(
1023 &self,
1024 role: Option<&str>,
1025 request: ReviewerRequest,
1026 ) -> Result<RelayRequest> {
1027 let request = RelayRequest::Reviewer {
1028 role: role.map(str::to_owned),
1029 request,
1030 };
1031 if !request.supported_at(self.protocol_version) {
1032 bail!(
1033 "a second opinion requires relay protocol {}; this session negotiated {}",
1034 request.minimum_protocol(),
1035 self.protocol_version
1036 );
1037 }
1038 Ok(request)
1039 }
1040
1041 pub async fn respond_elicitation(
1044 &mut self,
1045 elicitation_id: String,
1046 response: ElicitationResponse,
1047 ) -> Result<()> {
1048 let request = RelayRequest::RespondElicitation {
1049 elicitation_id: elicitation_id.clone(),
1050 response,
1051 };
1052 if !request.supported_at(self.protocol_version) {
1053 bail!(
1054 "elicitation responses require relay protocol {}; this session negotiated {}",
1055 request.minimum_protocol(),
1056 self.protocol_version
1057 );
1058 }
1059 match self.call(request).await? {
1060 RelayResponsePayload::ElicitationResolved {
1061 elicitation_id: resolved,
1062 } if resolved == elicitation_id => Ok(()),
1063 RelayResponsePayload::ElicitationResolved {
1064 elicitation_id: resolved,
1065 } => bail!("relay resolved elicitation {resolved:?}, expected {elicitation_id:?}"),
1066 _ => bail!("relay returned an unexpected elicitation response"),
1067 }
1068 }
1069
1070 pub async fn stop_background_task(&mut self, background_task_id: String) -> Result<()> {
1072 let request = RelayRequest::StopBackgroundTask {
1073 background_task_id: background_task_id.clone(),
1074 };
1075 if !request.supported_at(self.protocol_version) {
1076 bail!(
1077 "background task controls require relay protocol {}; this session negotiated {}",
1078 request.minimum_protocol(),
1079 self.protocol_version
1080 );
1081 }
1082 match self.call(request).await? {
1083 RelayResponsePayload::BackgroundTaskStopRequested {
1084 background_task_id: stopped,
1085 } if stopped == background_task_id => Ok(()),
1086 RelayResponsePayload::BackgroundTaskStopRequested {
1087 background_task_id: stopped,
1088 } => {
1089 bail!("relay stopped background task {stopped:?}, expected {background_task_id:?}")
1090 }
1091 _ => bail!("relay returned an unexpected background task stop response"),
1092 }
1093 }
1094
1095 pub async fn subagent_requests(
1096 &mut self,
1097 ) -> Result<(
1098 Vec<mj_core::subagent::SubagentToolRequest>,
1099 Vec<mj_core::subagent::SubagentToolResult>,
1100 )> {
1101 let request = RelayRequest::SubagentRequests;
1102 if !request.supported_at(self.protocol_version) {
1103 return Ok((Vec::new(), Vec::new()));
1104 }
1105 match self.call(request).await? {
1106 RelayResponsePayload::SubagentRequests { requests, results } => Ok((requests, results)),
1107 _ => bail!("relay returned an unexpected sub-agent request response"),
1108 }
1109 }
1110
1111 pub async fn complete_subagent_request(
1112 &mut self,
1113 result: mj_core::subagent::SubagentToolResult,
1114 ) -> Result<()> {
1115 let request = RelayRequest::CompleteSubagentRequest { result };
1116 if !request.supported_at(self.protocol_version) {
1117 bail!(
1118 "sub-agent tools require relay protocol {}",
1119 request.minimum_protocol()
1120 );
1121 }
1122 match self.call(request).await? {
1123 RelayResponsePayload::SubagentRequestCompleted => Ok(()),
1124 _ => bail!("relay returned an unexpected sub-agent completion response"),
1125 }
1126 }
1127
1128 pub async fn detach(mut self) -> Result<()> {
1129 self.input
1130 .take()
1131 .expect("connected relay owns proxy stdin")
1132 .shutdown()
1133 .await
1134 .context("close relay proxy stdin")?;
1135 let mut child = self.child.take().expect("connected relay owns proxy child");
1136 match tokio::time::timeout(RELAY_PROXY_DETACH_GRACE, child.wait()).await {
1137 Ok(status) => {
1138 status.context("wait for relay proxy")?;
1139 }
1140 Err(_) => {
1141 if let Err(error) = child.start_kill().context("stop relay proxy") {
1142 tracing::warn!(
1143 session_id = %self.session_id,
1144 operation = "detach",
1145 %error,
1146 "could not stop relay proxy after detach timeout"
1147 );
1148 return Err(error);
1149 }
1150 if let Err(error) = child.wait().await {
1151 tracing::warn!(
1152 session_id = %self.session_id,
1153 operation = "detach",
1154 %error,
1155 "could not reap relay proxy after stopping it"
1156 );
1157 }
1158 }
1159 }
1160 Ok(())
1161 }
1162
1163 async fn call(&mut self, request: RelayRequest) -> Result<RelayResponsePayload> {
1164 self.call_with_timeout(request, self.request_timeout).await
1165 }
1166
1167 async fn call_with_timeout(
1168 &mut self,
1169 request: RelayRequest,
1170 timeout: Duration,
1171 ) -> Result<RelayResponsePayload> {
1172 let operation = request.method_name();
1173 let request_id = self.request_id();
1174 let envelope = RelayRequestEnvelope {
1175 request_id: request_id.clone(),
1176 protocol_version: self.protocol_version,
1177 request,
1178 };
1179 let line = match self
1180 .exchange(&envelope, operation, timeout, ExchangeKind::Call)
1181 .await
1182 {
1183 Ok(line) => line,
1184 Err(error) => {
1185 log_relay_client_failure(self, operation, &request_id, &error);
1186 return Err(error);
1187 }
1188 };
1189 let result = decode_relay_response(&line, &request_id, self.protocol_version)
1190 .with_context(|| format!("relay {} could not perform {operation}", self.relay_version));
1191 if let Err(error) = &result {
1192 log_relay_client_failure(self, operation, &request_id, error);
1193 }
1194 result
1195 }
1196
1197 async fn call_hello(
1198 &mut self,
1199 request: RelayRequest,
1200 timeout: Duration,
1201 ) -> Result<RelayResponsePayload> {
1202 let operation = request.method_name();
1203 let request_id = self.request_id();
1204 let envelope = RelayRequestEnvelope {
1205 request_id: request_id.clone(),
1206 protocol_version: RELAY_PROTOCOL_VERSION,
1207 request,
1208 };
1209 let line = match self
1210 .exchange(&envelope, operation, timeout, ExchangeKind::Handshake)
1211 .await
1212 {
1213 Ok(line) => line,
1214 Err(error) => {
1215 log_relay_client_failure(self, operation, &request_id, &error);
1216 return Err(error);
1217 }
1218 };
1219 let result = decode_relay_hello_response(&line, &request_id);
1220 if let Err(error) = &result {
1221 log_relay_client_failure(self, operation, &request_id, error);
1222 }
1223 result
1224 }
1225
1226 async fn exchange(
1236 &mut self,
1237 envelope: &RelayRequestEnvelope,
1238 operation: &str,
1239 timeout: Duration,
1240 kind: ExchangeKind,
1241 ) -> Result<String> {
1242 if let Some(reason) = &self.abandoned {
1243 bail!("{reason}");
1244 }
1245 let mut frame = serde_json::to_vec(envelope)?;
1246 if frame.len() > MAX_FRAME_BYTES {
1247 bail!("relay {operation} request frame is too large");
1248 }
1249 frame.push(b'\n');
1250 let session_id = self.session_id.clone();
1251 let exchanged = tokio::time::timeout(timeout, async {
1252 self.input
1253 .as_mut()
1254 .expect("connected relay owns proxy stdin")
1255 .write_all(&frame)
1256 .await
1257 .map_err(|error| RelayTransportDead::from_io(error, kind))
1258 .with_context(|| format!("write relay {operation} request"))?;
1259 self.input
1260 .as_mut()
1261 .expect("connected relay owns proxy stdin")
1262 .flush()
1263 .await
1264 .map_err(|error| RelayTransportDead::from_io(error, kind))
1265 .with_context(|| format!("flush relay {operation} request"))?;
1266 let response = read_bounded_frame(&mut self.output, kind);
1267 tokio::pin!(response);
1268 let response = tokio::select! {
1269 response = &mut response => response,
1270 () = tokio::time::sleep(RELAY_SLOW_OPERATION_WARNING) => {
1271 tracing::warn!(
1272 %session_id,
1273 %operation,
1274 warning_after_seconds = RELAY_SLOW_OPERATION_WARNING.as_secs_f64(),
1275 timeout_seconds = timeout.as_secs_f64(),
1276 "relay operation is still waiting for its response"
1277 );
1278 response.await
1279 }
1280 };
1281 response
1282 .with_context(|| format!("read relay {operation} response"))?
1283 .ok_or_else(|| {
1284 anyhow::Error::new(RelayTransportDead::during_exchange(
1285 format!("relay proxy disconnected during {operation}"),
1286 kind,
1287 ))
1288 })
1289 })
1290 .await;
1291 match exchanged {
1292 Ok(line) => line,
1293 Err(_elapsed) => {
1294 let seconds = timeout.as_secs_f64();
1295 tracing::warn!(
1296 %session_id,
1297 %operation,
1298 timeout_seconds = seconds,
1299 "relay operation timed out; abandoning its sequential connection"
1300 );
1301 self.abandoned = Some(format!(
1302 "relay connection abandoned after {operation} timed out after {seconds} seconds"
1303 ));
1304 let timed_out = format!("relay {operation} timed out after {seconds} seconds");
1305 Err(anyhow!(timed_out))
1306 }
1307 }
1308 }
1309
1310 fn request_id(&mut self) -> String {
1311 let id = format!("relay-{:016x}-{}", self.connection_nonce, self.next_request);
1312 self.next_request = self.next_request.wrapping_add(1);
1313 id
1314 }
1315}
1316
1317fn log_relay_client_failure(
1322 client: &RelayClient,
1323 operation: &str,
1324 request_id: &str,
1325 error: &anyhow::Error,
1326) {
1327 let rejection = error.chain().find_map(|cause| {
1328 cause
1329 .downcast_ref::<RelayRejected>()
1330 .map(|rejected| &rejected.0)
1331 });
1332 let transport_dead = RelayTransportDead::marks(error);
1333 match rejection {
1334 Some(rejection) => tracing::warn!(
1335 session_id = %client.session_id,
1336 relay_version = %client.relay_version,
1337 %operation,
1338 %request_id,
1339 relay_error_code = ?rejection.code,
1340 relay_retryable = rejection.retryable,
1341 transport_dead,
1342 error = %error,
1343 "relay request rejected"
1344 ),
1345 None => tracing::warn!(
1346 session_id = %client.session_id,
1347 relay_version = %client.relay_version,
1348 %operation,
1349 %request_id,
1350 transport_dead,
1351 error = %error,
1352 "relay request failed"
1353 ),
1354 }
1355}
1356
1357impl Drop for RelayClient {
1358 fn drop(&mut self) {
1359 drop(self.input.take());
1364 let Some(child) = self.child.take() else {
1365 return;
1366 };
1367 let session_id = self.session_id.clone();
1368 if let Err(error) = std::thread::Builder::new()
1369 .name("hel-relay-reaper".into())
1370 .spawn(move || reap_dropped_relay_proxy(child, session_id))
1371 {
1372 tracing::warn!(
1373 session_id = %self.session_id,
1374 %error,
1375 "could not start dropped relay proxy reaper"
1376 );
1377 }
1378 }
1379}
1380
1381fn reap_dropped_relay_proxy(mut child: Child, session_id: String) {
1384 let deadline = Instant::now() + RELAY_PROXY_DETACH_GRACE;
1385 loop {
1386 match child.try_wait() {
1387 Ok(Some(status)) => {
1388 if !status.success() {
1389 tracing::warn!(
1390 %session_id,
1391 %status,
1392 "dropped relay proxy exited unsuccessfully"
1393 );
1394 }
1395 return;
1396 }
1397 Ok(None) if Instant::now() < deadline => {
1398 std::thread::sleep(RELAY_PROXY_REAP_POLL);
1399 }
1400 Ok(None) => break,
1401 Err(error) => {
1402 tracing::warn!(%session_id, %error, "could not reap dropped relay proxy");
1403 return;
1404 }
1405 }
1406 }
1407
1408 if let Err(error) = child.start_kill()
1409 && error.kind() != std::io::ErrorKind::NotFound
1410 {
1411 tracing::warn!(%session_id, %error, "could not stop dropped relay proxy");
1412 return;
1413 }
1414 let deadline = Instant::now() + RELAY_PROXY_DETACH_GRACE;
1415 loop {
1416 match child.try_wait() {
1417 Ok(Some(_)) => return,
1418 Ok(None) if Instant::now() < deadline => {
1419 std::thread::sleep(RELAY_PROXY_REAP_POLL);
1420 }
1421 Ok(None) => {
1422 tracing::warn!(%session_id, "stopped relay proxy could not be reaped in time");
1423 return;
1424 }
1425 Err(error) => {
1426 tracing::warn!(%session_id, %error, "could not reap stopped relay proxy");
1427 return;
1428 }
1429 }
1430 }
1431}
1432
1433fn credential_snapshot(payload: RelayResponsePayload) -> Result<CredentialSnapshot> {
1434 match payload {
1435 RelayResponsePayload::CredentialState {
1436 present,
1437 fingerprint,
1438 freshness_epoch_ms,
1439 } => Ok(CredentialSnapshot {
1440 present,
1441 fingerprint,
1442 freshness_epoch_ms,
1443 }),
1444 _ => bail!("relay returned an unexpected credential state response"),
1445 }
1446}
1447
1448fn skills_sync_state(payload: RelayResponsePayload) -> Result<mj_core::skills::SkillsSyncState> {
1449 match payload {
1450 RelayResponsePayload::SkillsState {
1451 present,
1452 fingerprint,
1453 } => Ok(mj_core::skills::SkillsSyncState {
1454 present,
1455 fingerprint,
1456 }),
1457 _ => bail!("relay returned an unexpected skills state response"),
1458 }
1459}
1460
1461fn github_token_snapshot(
1462 payload: RelayResponsePayload,
1463) -> Result<mj_core::credentials::GithubTokenSnapshot> {
1464 match payload {
1465 RelayResponsePayload::GithubTokenState {
1466 present,
1467 fingerprint,
1468 } => Ok(mj_core::credentials::GithubTokenSnapshot {
1469 present,
1470 fingerprint,
1471 }),
1472 _ => bail!("relay returned an unexpected GitHub token state response"),
1473 }
1474}
1475
1476async fn read_bounded_frame(
1477 reader: &mut (impl AsyncBufRead + Unpin),
1478 kind: ExchangeKind,
1479) -> Result<Option<String>> {
1480 read_bounded_frame_with_limit(reader, MAX_FRAME_BYTES, kind).await
1481}
1482
1483async fn read_bounded_frame_with_limit(
1484 reader: &mut (impl AsyncBufRead + Unpin),
1485 maximum_bytes: usize,
1486 kind: ExchangeKind,
1487) -> Result<Option<String>> {
1488 let mut frame = Vec::new();
1489 loop {
1490 let available = reader
1494 .fill_buf()
1495 .await
1496 .map_err(|error| RelayTransportDead::from_io(error, kind))?;
1497 if available.is_empty() {
1498 if frame.is_empty() {
1499 return Ok(None);
1500 }
1501 return Err(anyhow::Error::new(RelayTransportDead::during_exchange(
1502 "relay proxy disconnected in the middle of a response frame",
1503 kind,
1504 )));
1505 }
1506 let newline = available.iter().position(|byte| *byte == b'\n');
1507 let consumed = newline.map_or(available.len(), |position| position + 1);
1508 let payload = newline.map_or(available, |position| &available[..position]);
1509 if frame.len().saturating_add(payload.len()) > maximum_bytes {
1510 bail!("relay response frame is too large");
1511 }
1512 frame.extend_from_slice(payload);
1513 reader.consume(consumed);
1514 if newline.is_some() {
1515 if frame.last() == Some(&b'\r') {
1516 frame.pop();
1517 }
1518 return String::from_utf8(frame)
1519 .context("relay response is not UTF-8")
1520 .map(Some);
1521 }
1522 }
1523}
1524
1525fn clip_catch_up_page(
1526 page: RelayAttachment,
1527 previous: &RelayCursor,
1528 frontier: &RelayCursor,
1529) -> Result<RelayEventPage> {
1530 if previous.ordinal > frontier.ordinal {
1531 bail!("relay catch-up starts beyond its fixed frontier");
1532 }
1533 if previous.ordinal == frontier.ordinal {
1534 if previous != frontier {
1535 bail!("relay catch-up cursor digest differs from its fixed frontier");
1536 }
1537 if !page.events.is_empty() || page.through_ordinal != previous.ordinal {
1538 bail!("relay attachment advanced beyond its advertised frontier");
1539 }
1540 return Ok(RelayEventPage {
1541 events: Vec::new(),
1542 through_ordinal: previous.ordinal,
1543 through_digest: previous.digest.clone(),
1544 });
1545 }
1546 if page.through_ordinal <= previous.ordinal || page.events.is_empty() {
1547 bail!("relay catch-up page did not advance");
1548 }
1549 if page.through_ordinal <= frontier.ordinal {
1550 let through = RelayCursor {
1551 ordinal: page.through_ordinal,
1552 digest: page.through_digest.clone(),
1553 };
1554 if through.ordinal == frontier.ordinal && through != *frontier {
1555 bail!("relay catch-up page digest differs from its fixed frontier");
1556 }
1557 return Ok(RelayEventPage {
1558 events: page.events,
1559 through_ordinal: through.ordinal,
1560 through_digest: through.digest,
1561 });
1562 }
1563
1564 let events = page
1565 .events
1566 .into_iter()
1567 .take_while(|event| event.ordinal <= frontier.ordinal)
1568 .collect::<Vec<_>>();
1569 let reached = events
1570 .last()
1571 .map(|event| RelayCursor {
1572 ordinal: event.ordinal,
1573 digest: event.digest.clone(),
1574 })
1575 .ok_or_else(|| anyhow!("relay catch-up page skipped its fixed frontier"))?;
1576 if reached != *frontier {
1577 bail!("relay catch-up page does not contain its fixed frontier");
1578 }
1579 Ok(RelayEventPage {
1580 events,
1581 through_ordinal: reached.ordinal,
1582 through_digest: reached.digest,
1583 })
1584}
1585
1586fn decode_relay_response(
1587 line: &str,
1588 request_id: &str,
1589 protocol: u32,
1590) -> Result<RelayResponsePayload> {
1591 let response: RelayResponseEnvelope =
1592 serde_json::from_str(line).context("decode relay response")?;
1593 if response.request_id != request_id {
1594 bail!(
1595 "relay response ID mismatch: expected {request_id}, got {}",
1596 response.request_id
1597 );
1598 }
1599 if response.protocol_version != protocol {
1600 bail!(
1601 "relay response protocol mismatch: expected {protocol}, got {}",
1602 response.protocol_version
1603 );
1604 }
1605 match response.body {
1606 RelayResponseBody::Ok { payload } => Ok(payload),
1607 RelayResponseBody::Error { error } => Err(RelayRejected(error).into()),
1608 }
1609}
1610
1611fn decode_relay_hello_response(line: &str, request_id: &str) -> Result<RelayResponsePayload> {
1612 let response: RelayResponseEnvelope =
1613 serde_json::from_str(line).context("decode relay hello response")?;
1614 if response.request_id != request_id {
1615 bail!(
1616 "relay response ID mismatch: expected {request_id}, got {}",
1617 response.request_id
1618 );
1619 }
1620 match response.body {
1621 RelayResponseBody::Ok {
1622 payload: payload @ RelayResponsePayload::Hello { negotiated, .. },
1623 } => {
1624 if response.protocol_version != negotiated {
1625 bail!(
1626 "relay hello envelope uses protocol {}, negotiated {negotiated}",
1627 response.protocol_version
1628 );
1629 }
1630 Ok(payload)
1631 }
1632 RelayResponseBody::Ok { .. } => bail!("relay returned an unexpected hello response"),
1633 RelayResponseBody::Error { error } => Err(RelayRejected(error).into()),
1634 }
1635}
1636
1637pub struct CredentialSyncCoordinator {
1638 handle: CredentialSyncHandle,
1639 results: mpsc::UnboundedReceiver<CredentialSyncResult>,
1640}
1641
1642impl CredentialSyncCoordinator {
1643 pub fn spawn() -> Self {
1644 let (targets_tx, mut targets_rx) = watch::channel(Vec::new());
1645 let (triggers_tx, mut triggers_rx) = mpsc::unbounded_channel::<SyncTrigger>();
1646 let (completed_tx, mut completed_rx) = mpsc::unbounded_channel::<CredentialSyncResult>();
1647 let (results_tx, results_rx) = mpsc::unbounded_channel();
1648 tokio::spawn(async move {
1649 let mut tick = tokio::time::interval_at(
1650 tokio::time::Instant::now() + SYNC_INTERVAL,
1651 SYNC_INTERVAL,
1652 );
1653 tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
1654 let mut busy = BTreeSet::<String>::new();
1657 let mut queue = VecDeque::<SyncTrigger>::new();
1658 loop {
1659 tokio::select! {
1660 _ = tick.tick() => {
1661 for profile_id in profiles_with_targets(&targets_rx.borrow()) {
1662 enqueue(&mut queue, SyncTrigger { profile_id, cause: None });
1663 }
1664 }
1665 changed = targets_rx.changed() => {
1666 if changed.is_err() { break; }
1667 for profile_id in profiles_with_targets(&targets_rx.borrow()) {
1668 enqueue(&mut queue, SyncTrigger { profile_id, cause: None });
1669 }
1670 }
1671 trigger = triggers_rx.recv() => {
1672 let Some(trigger) = trigger else { break };
1673 enqueue(&mut queue, trigger);
1674 }
1675 completed = completed_rx.recv() => {
1676 let Some(result) = completed else { break };
1677 busy.remove(&result.profile_id);
1678 if result.trigger.is_some()
1679 || result.failure.is_some()
1680 || !result.outcomes.is_empty()
1681 {
1682 let profile_id = result.profile_id.clone();
1683 if results_tx.send(result).is_err() {
1684 tracing::debug!(
1685 %profile_id,
1686 operation = "credential_sync_result",
1687 "credential sync result receiver was already closed"
1688 );
1689 }
1690 }
1691 }
1692 }
1693
1694 let mut deferred = VecDeque::new();
1695 while let Some(trigger) = queue.pop_front() {
1696 if busy.contains(&trigger.profile_id) {
1697 deferred.push_back(trigger);
1698 continue;
1699 }
1700 let targets: Vec<_> = targets_rx
1701 .borrow()
1702 .iter()
1703 .filter(|target| target.profile_id == trigger.profile_id)
1704 .cloned()
1705 .collect();
1706 if targets.is_empty() {
1707 if trigger.cause.is_some() {
1708 let profile_id = trigger.profile_id.clone();
1709 if results_tx
1710 .send(CredentialSyncResult {
1711 profile_id: trigger.profile_id,
1712 trigger: trigger.cause,
1713 failure: None,
1714 outcomes: Vec::new(),
1715 })
1716 .is_err()
1717 {
1718 tracing::debug!(
1719 %profile_id,
1720 operation = "credential_sync_result",
1721 "credential sync result receiver was already closed"
1722 );
1723 }
1724 }
1725 continue;
1726 }
1727 busy.insert(trigger.profile_id.clone());
1728 let completed_tx = completed_tx.clone();
1729 let handle = tokio::runtime::Handle::current();
1730 tokio::spawn(async move {
1733 let joined = tokio::task::spawn_blocking(move || {
1734 handle.block_on(reconcile_profile(&targets))
1735 })
1736 .await;
1737 let (failure, outcomes) = match joined {
1738 Ok(outcomes) => (None, outcomes),
1739 Err(error) => (Some(format!("sync task stopped: {error}")), Vec::new()),
1740 };
1741 let profile_id = trigger.profile_id.clone();
1742 if completed_tx
1743 .send(CredentialSyncResult {
1744 profile_id: trigger.profile_id,
1745 trigger: trigger.cause,
1746 failure,
1747 outcomes,
1748 })
1749 .is_err()
1750 {
1751 tracing::debug!(
1752 %profile_id,
1753 operation = "credential_sync_completion",
1754 "credential sync coordinator stopped before receiving completion"
1755 );
1756 }
1757 });
1758 }
1759 queue = deferred;
1760 }
1761 });
1762 Self {
1763 handle: CredentialSyncHandle {
1764 targets: Arc::new(targets_tx),
1765 triggers: triggers_tx,
1766 },
1767 results: results_rx,
1768 }
1769 }
1770
1771 pub fn handle(&self) -> CredentialSyncHandle {
1772 self.handle.clone()
1773 }
1774
1775 pub fn try_result(&mut self) -> Option<CredentialSyncResult> {
1776 self.results.try_recv().ok()
1777 }
1778
1779 pub async fn result(&mut self) -> Option<CredentialSyncResult> {
1785 self.results.recv().await
1786 }
1787}
1788
1789async fn reconcile_profile(targets: &[CredentialSyncTarget]) -> Vec<CredentialSyncOutcome> {
1796 let github_token = targets
1797 .iter()
1798 .any(|target| target.sync_github_token)
1799 .then(crate::controller::controller_github_token)
1800 .flatten();
1801 let mut outcomes = BTreeMap::<String, CredentialSyncOutcome>::new();
1802 for pass in 0..2 {
1803 let mut pulled = false;
1804 for target in targets {
1805 match reconcile_session(target, github_token.as_deref()).await {
1806 Ok(actions) if actions.is_empty() => {}
1807 Ok(actions) => {
1808 pulled |= actions.contains(&CredentialSyncAction::Pulled);
1809 outcomes.insert(
1810 target.session_id.clone(),
1811 CredentialSyncOutcome {
1812 session_id: target.session_id.clone(),
1813 outcome: Ok(actions),
1814 },
1815 );
1816 }
1817 Err(error) => {
1818 tracing::warn!(
1819 session_id = %target.session_id,
1820 profile_id = %target.profile_id,
1821 pass = pass + 1,
1822 error = %error,
1823 "credential synchronization failed for relay session"
1824 );
1825 outcomes.insert(
1826 target.session_id.clone(),
1827 CredentialSyncOutcome {
1828 session_id: target.session_id.clone(),
1829 outcome: Err(format!("{error:#}")),
1830 },
1831 );
1832 }
1833 }
1834 }
1835 if !pulled || pass == 1 {
1836 break;
1837 }
1838 }
1839 outcomes.into_values().collect()
1840}
1841
1842async fn reconcile_session(
1844 target: &CredentialSyncTarget,
1845 github_token: Option<&str>,
1846) -> Result<Vec<CredentialSyncAction>> {
1847 let canonical_path = harness_authentication_marker(target.harness, &target.profile_home);
1848 let (canonical, canonical_bytes) = read_credential_file(target.harness, &canonical_path)?;
1849 let canonical_skills = mj_core::skills::collect_skills(target.harness, &target.profile_home)
1850 .with_context(|| {
1851 format!(
1852 "collect canonical skills for profile {} from {}",
1853 target.profile_id,
1854 target.profile_home.display()
1855 )
1856 })?;
1857 let mut client = RelayClient::connect(&target.spec, &target.session_id).await?;
1858 let result = reconcile_connected(
1859 &mut client,
1860 target,
1861 &canonical_path,
1862 &canonical,
1863 &canonical_bytes,
1864 &canonical_skills,
1865 github_token,
1866 )
1867 .await;
1868 if let Err(error) = client.detach().await {
1872 tracing::warn!(
1873 session_id = %target.session_id,
1874 "could not close the credential sync connection: {error:#}"
1875 );
1876 }
1877 result
1878}
1879
1880async fn reconcile_connected(
1881 client: &mut RelayClient,
1882 target: &CredentialSyncTarget,
1883 canonical_path: &Path,
1884 canonical: &CredentialSnapshot,
1885 canonical_bytes: &[u8],
1886 canonical_skills: &mj_core::skills::SkillsArchive,
1887 github_token: Option<&str>,
1888) -> Result<Vec<CredentialSyncAction>> {
1889 let mut actions = Vec::new();
1890 let session = client.credential_state().await?;
1891 match reconcile(canonical, &session) {
1892 SyncAction::None => {
1893 if canonical.present
1894 && session.present
1895 && canonical.fingerprint != session.fingerprint
1896 && canonical.freshness_epoch_ms.is_none()
1897 && session.freshness_epoch_ms.is_none()
1898 {
1899 tracing::warn!(
1900 session_id = %target.session_id,
1901 profile_id = %target.profile_id,
1902 "credential copies differ but neither reports a refresh time; leaving both alone"
1903 );
1904 }
1905 }
1906 SyncAction::Push => {
1907 client.install_credentials(canonical_bytes).await?;
1908 actions.push(CredentialSyncAction::Pushed);
1909 }
1910 SyncAction::Pull => {
1911 let bytes = client.read_credentials().await?;
1912 validate_credential_payload(target.harness, &bytes).with_context(|| {
1913 format!(
1914 "session {} returned an unusable credential file",
1915 target.session_id
1916 )
1917 })?;
1918 write_credential_file(target.harness, canonical_path, &bytes).with_context(|| {
1919 format!(
1920 "install fresher credentials from session {} for profile {}",
1921 target.session_id, target.profile_id
1922 )
1923 })?;
1924 actions.push(CredentialSyncAction::Pulled);
1925 }
1926 }
1927 if reconcile_skills(client, target, canonical_skills).await? {
1928 actions.push(CredentialSyncAction::SkillsPushed);
1929 }
1930 if target.sync_github_token
1931 && let Some(action) = reconcile_github_token(client, target, github_token).await?
1932 {
1933 actions.push(action);
1934 }
1935 Ok(actions)
1936}
1937
1938async fn reconcile_github_token(
1939 client: &mut RelayClient,
1940 target: &CredentialSyncTarget,
1941 canonical: Option<&str>,
1942) -> Result<Option<CredentialSyncAction>> {
1943 let session = match client.github_token_state().await {
1944 Ok(state) => state,
1945 Err(error) if sync_method_unsupported(&error) => {
1946 tracing::debug!(
1947 session_id = %target.session_id,
1948 profile_id = %target.profile_id,
1949 "worker predates GitHub token sync; skipping until the target is re-provisioned"
1950 );
1951 return Ok(None);
1952 }
1953 Err(error) => return Err(error),
1954 };
1955 match canonical {
1956 Some(token) => {
1957 let canonical = mj_core::credentials::GithubTokenSnapshot::of(token);
1958 if session == canonical {
1959 return Ok(None);
1960 }
1961 let installed = client.install_github_token(token).await?;
1962 if installed != canonical {
1963 bail!(
1964 "session {} GitHub token fingerprint does not match the controller after install",
1965 target.session_id
1966 );
1967 }
1968 Ok(Some(CredentialSyncAction::GithubTokenPushed))
1969 }
1970 None if session.present => {
1971 let removed = client.remove_github_token().await?;
1972 if removed.present {
1973 bail!(
1974 "session {} retained its GitHub token after removal",
1975 target.session_id
1976 );
1977 }
1978 Ok(Some(CredentialSyncAction::GithubTokenRemoved))
1979 }
1980 None => Ok(None),
1981 }
1982}
1983
1984async fn reconcile_skills(
1989 client: &mut RelayClient,
1990 target: &CredentialSyncTarget,
1991 canonical: &mj_core::skills::SkillsArchive,
1992) -> Result<bool> {
1993 let canonical_state = canonical.state();
1994 let session = match client.skills_state().await {
1995 Ok(state) => state,
1996 Err(error) if sync_method_unsupported(&error) => {
1997 tracing::debug!(
1998 session_id = %target.session_id,
1999 profile_id = %target.profile_id,
2000 "worker predates skills sync; skipping until the target is re-provisioned"
2001 );
2002 return Ok(false);
2003 }
2004 Err(error) => return Err(error),
2005 };
2006 if session == canonical_state {
2007 return Ok(false);
2008 }
2009 let installed = client.install_skills(&canonical.encode()).await?;
2010 if installed != canonical_state {
2011 bail!(
2012 "session {} skills fingerprint {} does not match the canonical {} after install",
2013 target.session_id,
2014 installed.fingerprint,
2015 canonical_state.fingerprint
2016 );
2017 }
2018 Ok(true)
2019}
2020
2021fn sync_method_unsupported(error: &anyhow::Error) -> bool {
2022 error
2023 .downcast_ref::<RelayRejected>()
2024 .is_some_and(|rejected| rejected.0.code == RelayErrorCode::InvalidRequest)
2025}
2026
2027#[cfg(test)]
2028mod tests {
2029 use super::*;
2030 use mj_core::relay::RelayObservation;
2031 use mj_worker::relay::DurableRelay;
2032 const SESSION_ID: &str = "018f9dd2-a3b4-7c8d-9000-123456789abc";
2033
2034 #[test]
2035 fn relay_decoder_preserves_explicit_desynchronization() {
2036 let response = RelayResponseEnvelope {
2037 request_id: "relay-1".into(),
2038 protocol_version: RELAY_PROTOCOL_VERSION,
2039 body: RelayResponseBody::Error {
2040 error: RelayProtocolError {
2041 code: RelayErrorCode::Desynchronized,
2042 message: "journal gap".into(),
2043 retryable: false,
2044 detail: None,
2045 },
2046 },
2047 };
2048 let encoded = serde_json::to_string(&response).unwrap();
2049 let error = decode_relay_response(&encoded, "relay-1", RELAY_PROTOCOL_VERSION).unwrap_err();
2050 assert!(
2051 error
2052 .downcast_ref::<RelayRejected>()
2053 .is_some_and(RelayRejected::is_desynchronized)
2054 );
2055 }
2056
2057 #[test]
2058 fn relay_decoder_rejects_crossed_request_ids() {
2059 let response = RelayResponseEnvelope {
2060 request_id: "other".into(),
2061 protocol_version: RELAY_PROTOCOL_VERSION,
2062 body: RelayResponseBody::Ok {
2063 payload: RelayResponsePayload::Acknowledged {
2064 through_ordinal: 4,
2065 through_digest: "a".repeat(64),
2066 },
2067 },
2068 };
2069 let encoded = serde_json::to_string(&response).unwrap();
2070 assert!(
2071 decode_relay_response(&encoded, "wanted", RELAY_PROTOCOL_VERSION)
2072 .unwrap_err()
2073 .to_string()
2074 .contains("ID mismatch")
2075 );
2076 }
2077
2078 #[test]
2079 fn command_spec_preserves_argv_boundaries() {
2080 let spec = CommandSpec::new("ssh", ["host", "hel worker proxy --root '/odd path'"]);
2081 assert_eq!(spec.program, "ssh");
2082 assert_eq!(spec.args.len(), 2);
2083 assert_eq!(spec.args[1], "hel worker proxy --root '/odd path'");
2084 }
2085
2086 #[test]
2087 fn relay_protocol_version_range_contains_current_version() {
2088 assert_eq!(
2089 RelayVersionRange::CURRENT.negotiate(RelayVersionRange::CURRENT),
2090 Some(RELAY_PROTOCOL_VERSION)
2091 );
2092 assert_eq!(
2093 RelayVersionRange::CURRENT.negotiate(RelayVersionRange { min: 1, max: 1 }),
2094 Some(1)
2095 );
2096 }
2097
2098 #[cfg(unix)]
2099 #[tokio::test]
2100 async fn controller_accepts_negotiated_protocol_v1() {
2101 let script = format!(
2102 r#"python3 -c '
2103import json, sys
2104session = {session:?}
2105req = json.loads(sys.stdin.readline())
2106assert req["request"]["method"] == "hello"
2107supported = req["request"]["params"]["supported"]
2108assert supported["min"] <= 1 <= supported["max"]
2109print(json.dumps({{
2110 "request_id": req["request_id"],
2111 "protocol_version": 1,
2112 "result": "ok",
2113 "payload": {{
2114 "type": "hello",
2115 "data": {{
2116 "negotiated": 1,
2117 "relay_version": "v1-fixture",
2118 "session_id": session,
2119 }},
2120 }},
2121}}), flush=True)
2122sys.stdin.read()
2123'"#,
2124 session = SESSION_ID
2125 );
2126 let spec = CommandSpec::new("sh", ["-c", &script]).purpose("v1 relay fixture");
2127 let client = RelayClient::connect_with_timeout(&spec, SESSION_ID, Duration::from_secs(5))
2128 .await
2129 .expect("protocol v1 hello must be accepted");
2130 assert_eq!(client.protocol_version(), 1);
2131 assert_eq!(client.relay_version(), "v1-fixture");
2132 }
2133
2134 #[cfg(unix)]
2138 #[tokio::test]
2139 async fn a_hello_reports_the_worker_build_or_none_from_an_older_worker() {
2140 let hello = |build: Option<&str>| {
2141 let data = match build {
2142 Some(build) => format!(
2143 r#"{{"negotiated":1,"relay_version":"build-fixture","session_id":"%s","worker_build":"{build}"}}"#
2144 ),
2145 None => r#"{"negotiated":1,"relay_version":"build-fixture","session_id":"%s"}"#
2146 .to_owned(),
2147 };
2148 format!(
2149 r#"
2150IFS= read -r hello
2151id=$(printf '%s' "$hello" | sed -n 's/.*"request_id":"\([^"]*\)".*/\1/p')
2152printf '{{"request_id":"%s","protocol_version":1,"result":"ok","payload":{{"type":"hello","data":{data}}}}}
2153' "$id" "$1"
2154sh -c 'while :; do sleep 30; done'
2155"#
2156 )
2157 };
2158 for reported in [None, Some("a".repeat(64).as_str())] {
2159 let spec = CommandSpec::new(
2160 "sh",
2161 [
2162 "-c".to_owned(),
2163 hello(reported),
2164 "hel-relay-build-fixture".to_owned(),
2165 SESSION_ID.to_owned(),
2166 ],
2167 )
2168 .purpose("relay worker build fixture");
2169 let client =
2170 RelayClient::connect_with_timeout(&spec, SESSION_ID, Duration::from_secs(5))
2171 .await
2172 .expect("hello must be accepted with and without a worker build");
2173 assert_eq!(client.worker_build(), reported);
2174 }
2175 }
2176
2177 #[cfg(unix)]
2178 #[tokio::test]
2179 async fn dropping_a_client_delivers_eof_before_stopping_its_proxy_launcher() {
2180 let directory = tempfile::tempdir().unwrap();
2181 let eof = directory.path().join("proxy-saw-eof");
2182 let script = r#"
2183IFS= read -r hello
2184id=$(printf '%s' "$hello" | sed -n 's/.*"request_id":"\([^"]*\)".*/\1/p')
2185printf '{"request_id":"%s","protocol_version":1,"result":"ok","payload":{"type":"hello","data":{"negotiated":1,"relay_version":"eof-fixture","session_id":"%s"}}}\n' "$id" "$1"
2186if IFS= read -r _; then exit 9; fi
2187: > "$2"
2188"#;
2189 let spec = CommandSpec::new(
2190 "sh",
2191 [
2192 "-c".to_owned(),
2193 script.to_owned(),
2194 "hel-relay-eof-fixture".to_owned(),
2195 SESSION_ID.to_owned(),
2196 eof.to_string_lossy().into_owned(),
2197 ],
2198 )
2199 .purpose("relay proxy EOF fixture");
2200 let client = RelayClient::connect_with_timeout(&spec, SESSION_ID, Duration::from_secs(5))
2201 .await
2202 .unwrap();
2203
2204 drop(client);
2205 tokio::time::timeout(Duration::from_secs(2), async {
2206 while !eof.exists() {
2207 tokio::time::sleep(Duration::from_millis(10)).await;
2208 }
2209 })
2210 .await
2211 .expect("proxy launcher was killed before it observed stdin EOF");
2212 }
2213
2214 #[cfg(unix)]
2215 #[tokio::test]
2216 async fn controller_rejects_negotiated_protocol_outside_supported_range() {
2217 let future_protocol = RELAY_PROTOCOL_VERSION + 1;
2218 let script = format!(
2219 r#"python3 -c '
2220import json, sys
2221session = {session:?}
2222req = json.loads(sys.stdin.readline())
2223print(json.dumps({{
2224 "request_id": req["request_id"],
2225 "protocol_version": {future_protocol},
2226 "result": "ok",
2227 "payload": {{
2228 "type": "hello",
2229 "data": {{
2230 "negotiated": {future_protocol},
2231 "relay_version": "future",
2232 "session_id": session,
2233 }},
2234 }},
2235}}), flush=True)
2236sys.stdin.read()
2237'"#,
2238 session = SESSION_ID,
2239 future_protocol = future_protocol,
2240 );
2241 let spec = CommandSpec::new("sh", ["-c", &script]).purpose("future relay fixture");
2242 let error = RelayClient::connect_with_timeout(&spec, SESSION_ID, Duration::from_secs(5))
2243 .await
2244 .err()
2245 .expect("a future protocol hello must be rejected");
2246 assert!(
2247 error.to_string().contains(&format!(
2248 "negotiated unsupported protocol {future_protocol}"
2249 )),
2250 "{error:#}"
2251 );
2252 assert!(!RelayTransportDead::marks(&error), "{error:#}");
2255 }
2256
2257 #[cfg(unix)]
2260 #[tokio::test]
2261 async fn a_proxy_that_exits_before_hello_reports_a_dead_transport() {
2262 let spec = CommandSpec::new("sh", ["-c", "exit 1"]).purpose("exiting relay proxy");
2263
2264 let error = RelayClient::connect_with_timeout(&spec, SESSION_ID, Duration::from_secs(5))
2265 .await
2266 .err()
2267 .expect("a proxy that exits cannot complete hello");
2268
2269 assert!(RelayTransportDead::marks(&error), "{error:#}");
2270 assert!(RelayTransportDead::marks_failed_handshake(&error));
2271 }
2272
2273 #[cfg(unix)]
2274 #[tokio::test]
2275 async fn silent_proxy_handshake_has_a_bounded_deadline() {
2276 let spec = CommandSpec::new("sh", ["-c", "sleep 30"]).purpose("test silent relay proxy");
2277 let started = std::time::Instant::now();
2278
2279 let error = RelayClient::connect_with_timeout(&spec, SESSION_ID, Duration::from_millis(50))
2280 .await
2281 .err()
2282 .expect("silent relay must time out");
2283
2284 assert!(error.to_string().contains("relay hello timed out"));
2285 assert!(!RelayTransportDead::marks(&error), "{error:#}");
2289 assert!(!RelayTransportDead::marks_failed_handshake(&error));
2290 assert!(started.elapsed() < Duration::from_secs(2));
2291 }
2292
2293 #[cfg(unix)]
2296 const STALLING_RELAY: &str = r#"
2297IFS= read -r hello
2298id=$(printf '%s' "$hello" | sed -n 's/.*"request_id":"\([^"]*\)".*/\1/p')
2299printf '{"request_id":"%s","protocol_version":1,"result":"ok","payload":{"type":"hello","data":{"negotiated":1,"relay_version":"stalling-fixture","session_id":"%s"}}}\n' "$id" "$1"
2300IFS= read -r stalled
2301id=$(printf '%s' "$stalled" | sed -n 's/.*"request_id":"\([^"]*\)".*/\1/p')
2302sleep 5
2303printf '{"request_id":"%s","protocol_version":1,"result":"error","error":{"code":"internal","message":"late reply","retryable":false}}\n' "$id"
2304cat > /dev/null
2305"#;
2306
2307 #[cfg(unix)]
2308 #[tokio::test]
2309 async fn a_timed_out_call_abandons_the_connection_instead_of_desynchronizing_it() {
2310 let spec = CommandSpec::new(
2311 "sh",
2312 ["-c", STALLING_RELAY, "hel-relay-fixture", SESSION_ID],
2313 )
2314 .purpose("stalling relay fixture");
2315 let mut client =
2316 RelayClient::connect_with_timeout(&spec, SESSION_ID, Duration::from_millis(500))
2317 .await
2318 .expect("the fixture answers hello immediately");
2319
2320 let timed_out = client
2321 .status()
2322 .await
2323 .expect_err("the stalled status call must time out");
2324 assert!(
2325 format!("{timed_out:#}").contains("relay status timed out"),
2326 "{timed_out:#}"
2327 );
2328 assert!(!RelayTransportDead::marks(&timed_out), "{timed_out:#}");
2331
2332 let started = std::time::Instant::now();
2337 let subsequent = client
2338 .status()
2339 .await
2340 .expect_err("a call on an abandoned connection must fail");
2341 let elapsed = started.elapsed();
2342 assert!(
2343 format!("{subsequent:#}").contains("relay connection abandoned after status timed out"),
2344 "{subsequent:#}"
2345 );
2346 assert!(
2347 elapsed < Duration::from_millis(250),
2348 "an abandoned connection must fail fast, took {elapsed:?}"
2349 );
2350
2351 let repeated = client
2352 .status()
2353 .await
2354 .expect_err("the connection stays abandoned");
2355 assert!(
2356 format!("{repeated:#}").contains("relay connection abandoned after status timed out"),
2357 "{repeated:#}"
2358 );
2359 }
2360
2361 #[test]
2362 fn an_unsupported_method_answer_still_reads_as_missing_skills_sync() {
2363 let response = mj_core::relay::unsupported_relay_method_response(
2367 "relay-1".into(),
2368 RELAY_PROTOCOL_VERSION,
2369 "skills_state".into(),
2370 );
2371 let encoded = serde_json::to_string(&response).unwrap();
2372 let error = decode_relay_response(&encoded, "relay-1", RELAY_PROTOCOL_VERSION).unwrap_err();
2373 assert!(sync_method_unsupported(&error), "{error:#}");
2374 }
2375
2376 #[tokio::test]
2377 async fn publishing_new_targets_starts_reconciliation_without_waiting_for_the_tick() {
2378 let profile = tempfile::tempdir().unwrap();
2379 let mut coordinator = CredentialSyncCoordinator::spawn();
2380 coordinator.handle().set_targets(vec![CredentialSyncTarget {
2381 session_id: SESSION_ID.into(),
2382 profile_id: "work".into(),
2383 harness: mj_core::config::HarnessKind::Codex,
2384 profile_home: profile.path().to_path_buf(),
2385 sync_github_token: false,
2386 spec: CommandSpec::new("sh", ["-c", "exit 1"]),
2387 }]);
2388
2389 let result = tokio::time::timeout(Duration::from_secs(5), coordinator.result())
2390 .await
2391 .expect("target publication must not wait for the 60-second periodic tick")
2392 .expect("credential coordinator stopped");
2393 assert_eq!(result.profile_id, "work");
2394 assert_eq!(result.outcomes.len(), 1);
2395 assert!(result.outcomes[0].outcome.is_err());
2396 }
2397
2398 #[tokio::test]
2399 async fn response_frame_limit_is_enforced_before_newline() {
2400 let (mut writer, reader) = tokio::io::duplex(32);
2401 let write = tokio::spawn(async move {
2402 writer.write_all(b"123456789\n").await.unwrap();
2403 });
2404 let mut reader = BufReader::new(reader);
2405
2406 let error = read_bounded_frame_with_limit(&mut reader, 8, ExchangeKind::Call)
2407 .await
2408 .unwrap_err();
2409
2410 write.await.unwrap();
2411 assert!(error.to_string().contains("frame is too large"));
2412 assert!(!RelayTransportDead::marks(&error), "{error:#}");
2415 }
2416
2417 #[tokio::test]
2418 async fn a_half_written_response_frame_reports_a_dead_transport() {
2419 let (mut writer, reader) = tokio::io::duplex(32);
2420 writer.write_all(b"{\"partial\":").await.unwrap();
2421 drop(writer);
2422 let mut reader = BufReader::new(reader);
2423
2424 let error = read_bounded_frame(&mut reader, ExchangeKind::Call)
2425 .await
2426 .unwrap_err();
2427
2428 assert!(RelayTransportDead::marks(&error), "{error:#}");
2429 assert!(!RelayTransportDead::marks_failed_handshake(&error));
2430 }
2431
2432 #[test]
2433 fn catch_up_page_stops_at_the_frontier_captured_before_stream_growth() {
2434 let temp = tempfile::tempdir().unwrap();
2435 let mut relay = DurableRelay::open(temp.path(), SESSION_ID, "1.0.0").unwrap();
2436 for message in ["one", "two", "arrived concurrently"] {
2437 relay
2438 .record_observation(RelayObservation::Warning {
2439 message: message.into(),
2440 })
2441 .unwrap();
2442 }
2443 let all = relay.events_after(0, RELAY_EVENT_GENESIS_DIGEST).unwrap();
2444 let previous = RelayCursor {
2445 ordinal: all[0].ordinal,
2446 digest: all[0].digest.clone(),
2447 };
2448 let frontier = RelayCursor {
2449 ordinal: all[1].ordinal,
2450 digest: all[1].digest.clone(),
2451 };
2452 let page = RelayAttachment {
2453 state: relay.operational_state(),
2454 events: all[1..].to_vec(),
2455 through_ordinal: all[2].ordinal,
2456 through_digest: all[2].digest.clone(),
2457 };
2458 let clipped = clip_catch_up_page(page, &previous, &frontier).unwrap();
2459 assert_eq!(clipped.through_ordinal, frontier.ordinal);
2460 assert_eq!(clipped.through_digest, frontier.digest);
2461 assert_eq!(clipped.events.len(), 1);
2462 assert_eq!(clipped.events.last().unwrap().ordinal, 2);
2463 }
2464}