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