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