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 CommandSpec, SSH_MASTER_OPEN_TIMEOUT, SSH_RETRY_ATTEMPTS, SshAdmission, SshPermit, SshRefusal,
19 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 RELAY_PROXY_DETACH_GRACE: Duration = Duration::from_millis(500);
58const RELAY_PROXY_REAP_POLL: Duration = Duration::from_millis(10);
59
60const RELAY_PROXY_STDERR_TAIL: usize = 10;
62
63const WORKER_SOCKET_RETRY_DELAYS: [Duration; 5] = [
69 Duration::from_millis(50),
70 Duration::from_millis(100),
71 Duration::from_millis(200),
72 Duration::from_millis(400),
73 Duration::from_millis(800),
74];
75
76type ProxyStderrTail = Arc<std::sync::Mutex<VecDeque<String>>>;
86
87async fn drain_proxy_stderr(
98 errors: tokio::process::ChildStderr,
99 purpose: String,
100 session_id: String,
101 tail: ProxyStderrTail,
102 handshake_done: Arc<AtomicBool>,
103) {
104 let mut lines = BufReader::new(errors).lines();
105 loop {
106 match lines.next_line().await {
107 Ok(Some(line)) if line.trim().is_empty() => continue,
108 Ok(Some(line)) => {
109 if handshake_done.load(Ordering::Acquire) {
110 tracing::warn!(%session_id, %purpose, %line, "relay proxy stderr");
111 } else {
112 tracing::debug!(%session_id, %purpose, %line, "relay proxy stderr");
113 }
114 let mut tail = tail.lock().unwrap_or_else(PoisonError::into_inner);
115 if tail.len() == RELAY_PROXY_STDERR_TAIL {
116 tail.pop_front();
117 }
118 tail.push_back(line);
119 }
120 Ok(None) => return,
121 Err(error) => {
122 tracing::warn!(%session_id, %purpose, %error, "read relay proxy stderr");
123 return;
124 }
125 }
126 }
127}
128
129mod errors;
130pub use errors::*;
131mod connect;
132mod exchange;
133mod relay;
134mod reviewer;
135mod transport;
136use transport::*;
137mod credential_sync;
138pub use credential_sync::*;
139
140pub struct RelayClient {
146 child: Option<Child>,
147 input: Option<ChildStdin>,
148 output: BufReader<ChildStdout>,
149 request_timeout: Duration,
150 abandoned: Option<String>,
153 next_request: u64,
154 connection_nonce: u64,
155 protocol_version: u32,
156 session_id: String,
157 relay_version: String,
158 worker_build: Option<String>,
161 latest_ordinal: u64,
162 latest_digest: String,
163 ssh_session: Option<SshSessionLease>,
167}
168
169#[cfg(test)]
170mod tests;