mj_controller/
worker_client.rs1use 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);
41const RELAY_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(300);
45const RELAY_HISTORY_TIMEOUT: Duration = Duration::from_secs(900);
49const RELAY_ACKNOWLEDGE_TIMEOUT: Duration = Duration::from_secs(300);
53const REVIEW_CAPTURE_TIMEOUT: Duration = Duration::from_secs(300);
56const 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
63const RELAY_PROXY_STDERR_TAIL: usize = 10;
65
66type ProxyStderrTail = Arc<std::sync::Mutex<VecDeque<String>>>;
76
77async 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
120pub struct RelayClient {
126 child: Option<Child>,
127 input: Option<ChildStdin>,
128 output: BufReader<ChildStdout>,
129 request_timeout: Duration,
130 abandoned: Option<String>,
133 next_request: u64,
134 connection_nonce: u64,
135 protocol_version: u32,
136 session_id: String,
137 relay_version: String,
138 worker_build: Option<String>,
141 latest_ordinal: u64,
142 latest_digest: String,
143 ssh_session: Option<SshSessionLease>,
147}
148
149#[cfg(test)]
150mod tests;