Skip to main content

mj_controller/
worker_client.rs

1//! Controller-side client for a session relay's JSON-lines proxy.
2
3use 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);
38/// Starting a target-side proxy may page the full worker executable in and
39/// traverse a container runtime before the relay sees `hello`. That is worker
40/// startup latency, not an ordinary in-connection RPC.
41const RELAY_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(300);
42/// An attachment can decompress a transport-sized page from cold journal
43/// segments. It remains bounded by the relay frame budget, but cold or loaded
44/// storage needs a filesystem deadline rather than an in-memory RPC deadline.
45const RELAY_HISTORY_TIMEOUT: Duration = Duration::from_secs(900);
46/// Advancing an acknowledgement can durably prune a large relay journal. The
47/// worker performs that maintenance before replying, so it needs a deadline
48/// sized for filesystem work rather than ordinary relay bookkeeping.
49const RELAY_ACKNOWLEDGE_TIMEOUT: Duration = Duration::from_secs(300);
50/// Capturing a review delta runs Git over every workspace repository, which is
51/// filesystem work on a possibly large tree rather than relay bookkeeping.
52const REVIEW_CAPTURE_TIMEOUT: Duration = Duration::from_secs(300);
53/// Bifrost's semantic diff analysis has its own 600-second budget inside the
54/// worker; this leaves room for it to report a timeout as an error rather than
55/// having the call time out underneath it.
56const 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
60/// Forward a relay proxy's stderr to the log, one line at a time, until the
61/// child closes it. Reporting rather than dropping keeps connect failures
62/// diagnosable now that the controller no longer shares its terminal.
63async 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/// One bounded page in a catch-up whose upper frontier was fixed before any
85/// page was applied. The relay may return newer events on later `Attach`
86/// calls; those are deliberately left for the next catch-up.
87#[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    /// Whether the relay itself said the same request could succeed later.
122    /// Validation rejections say no; transient internal failures say yes.
123    pub fn is_retryable(&self) -> bool {
124        self.0.retryable
125    }
126}
127
128/// A relay transport that can no longer carry requests: the proxy exited, one
129/// of its pipes failed, or the handshake never completed.
130///
131/// Every site that can prove this attaches the marker, and recovery decisions
132/// such as worker auto-restart downcast for it. Nothing reads the message text,
133/// so rewording a diagnostic can never silently disable recovery.
134#[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    /// Mark an I/O failure on the relay's pipes. The marker reports exactly
157    /// what the I/O error reported, so it adds a type without adding text.
158    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    /// Whether this error, or any cause behind it, is a dead relay transport.
170    pub fn marks(error: &anyhow::Error) -> bool {
171        error.downcast_ref::<Self>().is_some()
172    }
173
174    /// Whether the worker was reachable enough to run its liveness probe but
175    /// the proxy then disconnected or failed I/O during a fresh handshake.
176    /// Timeouts are deliberately not marked: a live proxy can be waiting on a
177    /// loaded container runtime or filesystem, which restarting only worsens.
178    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/// Whether an exchange is the handshake that proves the transport carries
186/// traffic at all.
187///
188/// A disconnected handshake proves the new transport never became usable. A
189/// timeout does not: the proxy launcher or worker can still be alive and slow,
190/// so timeout classification is handled separately in [`RelayClient::exchange`].
191#[derive(Clone, Copy, PartialEq, Eq)]
192enum ExchangeKind {
193    Handshake,
194    Call,
195}
196
197/// Controller-side connection to the durable ACP relay protocol.
198///
199/// This type does not construct transcript state or request unbounded history.
200/// Callers persist bounded attachment pages, then acknowledge only a frontier
201/// that is already durable locally.
202pub struct RelayClient {
203    child: Option<Child>,
204    input: Option<ChildStdin>,
205    output: BufReader<ChildStdout>,
206    request_timeout: Duration,
207    /// Why this connection can no longer be used, once a call gave up on a
208    /// reply that is still in flight. See [`RelayClient::exchange`].
209    abandoned: Option<String>,
210    next_request: u64,
211    connection_nonce: u64,
212    protocol_version: u32,
213    session_id: String,
214    relay_version: String,
215    /// Content address of the executable the worker is running, as reported in
216    /// hello. `None` from a worker built before the field existed.
217    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            // Never inherit: the controller owns a TUI alternate screen, so a
255            // child writing to the shared stderr corrupts the display outside
256            // the renderer's buffer. Drain it into the log instead.
257            .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            // Keep the expected identity from process creation onward so a
325            // handshake failure and the dropped proxy that follows it remain
326            // attributable even when Hello never returns a session ID.
327            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    /// Content address of the executable serving this connection, or `None`
387    /// from a worker too old to report one. A controller reads `None` as
388    /// outdated: it predates the field, so it predates this controller.
389    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    /// Start a bounded catch-up by capturing the relay frontier before the
454    /// caller applies anything. Callers persist `first_page`, request further
455    /// pages with [`Self::next_catch_up_page`], and may acknowledge the fixed
456    /// frontier after all of those pages are durable.
457    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    /// Fetch the next bounded page without chasing events that arrived after
482    /// `frontier` was captured. A response may contain such newer events; the
483    /// returned page is clipped at the exact ordinal-and-digest frontier.
484    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    /// Return the fingerprint and freshness of this session's harness
536    /// credentials without exposing the credential bytes.
537    pub async fn credential_state(&mut self) -> Result<CredentialSnapshot> {
538        credential_snapshot(self.call(RelayRequest::CredentialState).await?)
539    }
540
541    /// Read this session's credential file. Callers must keep these bytes out
542    /// of durable relay observations, logs, and archives.
543    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    /// Install credentials into the harness home fixed by this session's
553    /// launch config.
554    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    /// Return the fingerprint of this session's synced skills trees without
588    /// transferring the tree itself.
589    pub async fn skills_state(&mut self) -> Result<mj_core::skills::SkillsSyncState> {
590        skills_sync_state(self.call(RelayRequest::SkillsState).await?)
591    }
592
593    /// Install background text that only the target harness sees, prepended
594    /// to the next real prompt without creating a synthetic transcript turn.
595    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    /// Replace this session's synced skills trees with an encoded
651    /// `skills::SkillsArchive`. The destination directories are fixed by
652    /// the session's launch config and the harness skills whitelist.
653    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    /// Copy a verified controller blob to this session before admitting its reference.
666    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    /// Recover the local copy needed for queue editing and resubmission.
702    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    /// Start the second-opinion reviewer beside this session, or report the
766    /// running one when it already matches `config`.
767    ///
768    /// The reviewer's profile must already be staged on the target. Starting
769    /// can take as long as opening any harness session, so this uses the
770    /// handshake deadline rather than the bookkeeping one.
771    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    /// Replay the reviewer's journal from a cursor, exactly as [`Self::attach`]
802    /// does for the primary.
803    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        // The reviewer's journal is verified the same way the primary's is: a
830        // sidecar's history is not exempt from the chain check.
831        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    /// Advance the reviewer's acknowledged frontier so its journal can be
853    /// pruned once the controller has the events durably.
854    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    /// Queue one command on the reviewer's own relay.
883    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    /// Answer a form the reviewer's harness is waiting on.
919    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    /// Cancel any reviewer turn in flight and stop its process group, keeping
944    /// its staged profile, native session and journal for the next review.
945    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    /// Report what every workspace repository changed since the review
957    /// baselines the controller holds.
958    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    /// Record the trees a completed review reviewed through, so the next
974    /// review starts from them.
975    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    /// Run Bifrost's semantic diff analysis over the captured trees. It can
991    /// take minutes on a large changeset, so it carries its own budget.
992    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    /// Collect the specialist lanes the review supervisor asked for since the
1009    /// last call.
1010    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    /// Wraps a reviewer action, refusing it on a worker too old to know what a
1021    /// reviewer is rather than sending a method it would reject as unknown.
1022    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    /// Answer an ACP form over the live relay connection. User-entered content
1042    /// is intentionally excluded from the relay's durable command path.
1043    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    /// Ask the live worker to stop one process-local background task.
1071    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    /// Write one request frame and read the reply that belongs to it.
1227    ///
1228    /// The connection is strictly sequential, so giving up on a reply does not
1229    /// cancel it: the relay may still answer, and that answer would be read as
1230    /// the *next* call's response. Timeouts therefore abandon the connection
1231    /// rather than the single call. Every later call fails immediately with the
1232    /// true cause, so callers reconnect deliberately instead of chasing a
1233    /// mismatched response ID. This matters most where a short bookkeeping
1234    /// deadline and a long compaction deadline share one connection.
1235    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
1317/// Keep transport, protocol, and explicit relay rejections visible at the
1318/// point where a request fails. Callers often turn these into a user-facing
1319/// string or a retry, which otherwise loses the operation and request ID that
1320/// make concurrent session failures diagnosable.
1321fn 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        // Async owners call `detach` so EOF has a bounded chance to propagate
1360        // through Podman or SSH before the launcher is stopped. Drop is the
1361        // shutdown-safe fallback: it may run while Tokio's drivers are already
1362        // gone, so its bounded reaper cannot use runtime work or Tokio timers.
1363        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
1381/// Let EOF traverse a proxy launcher, then stop and reap it without relying on
1382/// an async runtime that may already be shutting down.
1383fn 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        // A failed read and a half-written frame are transport deaths; the
1491        // limit and encoding failures below are protocol violations that a
1492        // worker restart would not fix, so only these two carry the marker.
1493        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            // A pull rewrites the canonical file, so one profile is never
1655            // reconciled twice at once.
1656            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                    // The blocking join is awaited so a panicked reconcile is
1731                    // reported and its profile always leaves the busy set.
1732                    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    /// Waits for the next finished sync.
1780    ///
1781    /// Event-driven loops select on this instead of polling; `None` means the
1782    /// coordinator task has stopped. Cancel-safe, so a lost `select!` race
1783    /// keeps the result queued.
1784    pub async fn result(&mut self) -> Option<CredentialSyncResult> {
1785        self.results.recv().await
1786    }
1787}
1788
1789/// Reconcile one profile with every live session that runs it.
1790///
1791/// A pull makes every other session's copy stale by definition, so the pass
1792/// runs again once with the new canonical bytes. Two passes are enough: the
1793/// second cannot pull anything the first did not already see unless a harness
1794/// refreshed mid-cycle, and that lands in the next cycle.
1795async 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
1842/// Returns every action taken; an empty list means the copies already agree.
1843async 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    // Detach even when the exchange failed; the worker and harness keep
1869    // running either way. A failed detach only leaks a short-lived proxy, so it
1870    // is reported rather than turned into a sync failure.
1871    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
1984/// Converge the session's synced skills trees onto the canonical archive.
1985/// Returns true when a push happened. Workers old enough to predate skills
1986/// sync answer the unknown method with `InvalidRequest`; those sessions are
1987/// skipped quietly until their target is re-provisioned.
1988async 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    /// The build a worker reports is what decides whether it is replaced, so a
2135    /// controller has to read it from hello - and read a worker that reports
2136    /// none as exactly that, rather than failing the handshake.
2137    #[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        // The transport carried the answer perfectly well; restarting the
2253        // worker cannot make it speak a protocol it does not implement.
2254        assert!(!RelayTransportDead::marks(&error), "{error:#}");
2255    }
2256
2257    /// A proxy that exits without answering is the ordinary shape of a dead
2258    /// worker. Recovery hangs on this being typed rather than read.
2259    #[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        // The launcher is still alive. A loaded target can look exactly like
2286        // this while starting its proxy, so worker recovery must not restart
2287        // the native session merely because the deadline elapsed.
2288        assert!(!RelayTransportDead::marks(&error), "{error:#}");
2289        assert!(!RelayTransportDead::marks_failed_handshake(&error));
2290        assert!(started.elapsed() < Duration::from_secs(2));
2291    }
2292
2293    /// A relay that answers `hello` at once and then stalls, replying to the
2294    /// next request long after any controller deadline. `$1` is the session id.
2295    #[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        // A busy worker that misses one deadline is not a dead transport: it
2329        // answered the handshake, and killing it would be worse than waiting.
2330        assert!(!RelayTransportDead::marks(&timed_out), "{timed_out:#}");
2331
2332        // The abandoned reply is still in flight. A later call must not read it
2333        // as its own response, so it fails at once with the real cause. The
2334        // normal request deadline is long enough that without this the
2335        // controller would block on someone else's reply.
2336        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        // Workers that predate skills sync answer the unknown method with an
2364        // `InvalidRequest` rejection, and so does a current worker's structured
2365        // unsupported-method response. Both must skip the session quietly.
2366        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        // An oversized frame is a protocol violation, not a dead transport:
2413        // the same worker would send the same frame after a restart.
2414        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}