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