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, PoisonError};
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::{
17    BoundedProcessExecutor, CommandSpec, SSH_MASTER_OPEN_TIMEOUT, SSH_RETRY_ATTEMPTS, SshAdmission,
18    SshPermit, SshSessionLease, is_transport_rejection,
19};
20use mj_core::config::harness_authentication_marker;
21use mj_core::credentials::{
22    CredentialSnapshot, CredentialSyncAction, CredentialSyncHandle, CredentialSyncOutcome,
23    CredentialSyncResult, CredentialSyncTarget, SYNC_INTERVAL, SyncAction, SyncTrigger, enqueue,
24    profiles_with_targets, read_credential_file, reconcile, validate_credential_payload,
25    write_credential_file,
26};
27use mj_core::elicitation::ElicitationResponse;
28use mj_core::relay::{
29    MAX_FRAME_BYTES, RELAY_EVENT_GENESIS_DIGEST, RELAY_MIN_PROTOCOL_VERSION,
30    RELAY_PROTOCOL_VERSION, RelayCommand, RelayCursor, RelayErrorCode, RelayEvent,
31    RelayOperationalState, RelayProtocolError, RelayRequest, RelayRequestEnvelope,
32    RelayResponseBody, RelayResponseEnvelope, RelayResponsePayload, RelayVersionRange,
33    ReviewerRequest, validate_relay_event,
34};
35
36pub use mj_client::session::{RelayAttachment, StartedReviewer};
37use mj_core::worker_launch::ReviewerLaunchConfig;
38
39const RELAY_RPC_TIMEOUT: Duration = Duration::from_secs(15);
40const RELAY_SLOW_OPERATION_WARNING: Duration = Duration::from_secs(5);
41/// Starting a target-side proxy may page the full worker executable in and
42/// traverse a container runtime before the relay sees `hello`. That is worker
43/// startup latency, not an ordinary in-connection RPC.
44const RELAY_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(300);
45/// An attachment can decompress a transport-sized page from cold journal
46/// segments. It remains bounded by the relay frame budget, but cold or loaded
47/// storage needs a filesystem deadline rather than an in-memory RPC deadline.
48const RELAY_HISTORY_TIMEOUT: Duration = Duration::from_secs(900);
49/// Advancing an acknowledgement can durably prune a large relay journal. The
50/// worker performs that maintenance before replying, so it needs a deadline
51/// sized for filesystem work rather than ordinary relay bookkeeping.
52const RELAY_ACKNOWLEDGE_TIMEOUT: Duration = Duration::from_secs(300);
53/// Capturing a review delta runs Git over every workspace repository, which is
54/// filesystem work on a possibly large tree rather than relay bookkeeping.
55const REVIEW_CAPTURE_TIMEOUT: Duration = Duration::from_secs(300);
56/// Bifrost's semantic diff analysis has its own 600-second budget inside the
57/// worker; this leaves room for it to report a timeout as an error rather than
58/// having the call time out underneath it.
59const REVIEW_ANALYSIS_TIMEOUT: Duration = Duration::from_secs(660);
60const RELAY_PROXY_DETACH_GRACE: Duration = Duration::from_millis(500);
61const RELAY_PROXY_REAP_POLL: Duration = Duration::from_millis(10);
62
63/// How many trailing stderr lines a failed connect reports back to its caller.
64const RELAY_PROXY_STDERR_TAIL: usize = 10;
65
66/// The proxy's last [`RELAY_PROXY_STDERR_TAIL`] non-empty stderr lines, shared
67/// with whoever has to report them.
68///
69/// The drain publishes each line here as it reads it, rather than returning
70/// the whole tail when it finishes. A failed connect has to bound how long it
71/// waits for the drain, because a proxy that leaves a grandchild holding
72/// stderr never reaches EOF. Reading the tail from here means that bound costs
73/// only the lines not yet read, instead of discarding every line already
74/// collected.
75type ProxyStderrTail = Arc<std::sync::Mutex<VecDeque<String>>>;
76
77/// Forward a relay proxy's stderr to the log, one line at a time, until the
78/// child closes it, keeping the tail in `tail`. Reporting rather than dropping
79/// keeps connect failures diagnosable now that the controller no longer shares
80/// its terminal, and lets a failed connect put the proxy's own complaint in
81/// the error the caller sees rather than only in the log.
82async fn drain_proxy_stderr(
83    errors: tokio::process::ChildStderr,
84    purpose: String,
85    session_id: String,
86    tail: ProxyStderrTail,
87) {
88    let mut lines = BufReader::new(errors).lines();
89    loop {
90        match lines.next_line().await {
91            Ok(Some(line)) if line.trim().is_empty() => continue,
92            Ok(Some(line)) => {
93                tracing::warn!(%session_id, %purpose, %line, "relay proxy stderr");
94                let mut tail = tail.lock().unwrap_or_else(PoisonError::into_inner);
95                if tail.len() == RELAY_PROXY_STDERR_TAIL {
96                    tail.pop_front();
97                }
98                tail.push_back(line);
99            }
100            Ok(None) => return,
101            Err(error) => {
102                tracing::warn!(%session_id, %purpose, %error, "read relay proxy stderr");
103                return;
104            }
105        }
106    }
107}
108
109mod errors;
110pub use errors::*;
111mod connect;
112mod exchange;
113mod relay;
114mod reviewer;
115mod transport;
116use transport::*;
117mod credential_sync;
118pub use credential_sync::*;
119
120/// Controller-side connection to the durable ACP relay protocol.
121///
122/// This type does not construct transcript state or request unbounded history.
123/// Callers persist bounded attachment pages, then acknowledge only a frontier
124/// that is already durable locally.
125pub struct RelayClient {
126    child: Option<Child>,
127    input: Option<ChildStdin>,
128    output: BufReader<ChildStdout>,
129    request_timeout: Duration,
130    /// Why this connection can no longer be used, once a call gave up on a
131    /// reply that is still in flight. See [`RelayClient::exchange`].
132    abandoned: Option<String>,
133    next_request: u64,
134    connection_nonce: u64,
135    protocol_version: u32,
136    session_id: String,
137    relay_version: String,
138    /// Content address of the executable the worker is running, as reported in
139    /// hello. `None` from a worker built before the field existed.
140    worker_build: Option<String>,
141    latest_ordinal: u64,
142    latest_digest: String,
143    /// The shared-connection session the proxy runs on. Unlike the admission
144    /// permit, which is released once hello completes, the session is in use
145    /// for as long as the proxy runs, so the lease lives as long as `child`.
146    ssh_session: Option<SshSessionLease>,
147}
148
149#[cfg(test)]
150mod tests;