1use 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);
42const RELAY_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(300);
46const RELAY_HISTORY_TIMEOUT: Duration = Duration::from_secs(900);
50const RELAY_ACKNOWLEDGE_TIMEOUT: Duration = Duration::from_secs(300);
54const REVIEW_CAPTURE_TIMEOUT: Duration = Duration::from_secs(300);
57const 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
64const RELAY_PROXY_STDERR_TAIL: usize = 10;
66
67const 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
80type ProxyStderrTail = Arc<std::sync::Mutex<VecDeque<String>>>;
90
91async 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
144pub struct RelayClient {
150 child: Option<Child>,
151 input: Option<ChildStdin>,
152 output: BufReader<ChildStdout>,
153 request_timeout: Duration,
154 abandoned: Option<String>,
157 next_request: u64,
158 connection_nonce: u64,
159 protocol_version: u32,
160 session_id: String,
161 relay_version: String,
162 worker_build: Option<String>,
165 latest_ordinal: u64,
166 latest_digest: String,
167 ssh_session: Option<SshSessionLease>,
171}
172
173#[cfg(test)]
174mod tests;