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 CommandSpec, SSH_RETRY_ATTEMPTS, SshAdmission, SshPermit, is_transport_rejection,
18};
19use mj_core::config::harness_authentication_marker;
20use mj_core::credentials::{
21 CredentialSnapshot, CredentialSyncAction, CredentialSyncHandle, CredentialSyncOutcome,
22 CredentialSyncResult, CredentialSyncTarget, SYNC_INTERVAL, SyncAction, SyncTrigger, enqueue,
23 profiles_with_targets, read_credential_file, reconcile, validate_credential_payload,
24 write_credential_file,
25};
26use mj_core::elicitation::ElicitationResponse;
27use mj_core::relay::{
28 MAX_FRAME_BYTES, RELAY_EVENT_GENESIS_DIGEST, RELAY_MIN_PROTOCOL_VERSION,
29 RELAY_PROTOCOL_VERSION, RelayCommand, RelayCursor, RelayErrorCode, RelayEvent,
30 RelayOperationalState, RelayProtocolError, RelayRequest, RelayRequestEnvelope,
31 RelayResponseBody, RelayResponseEnvelope, RelayResponsePayload, RelayVersionRange,
32 ReviewerRequest, validate_relay_event,
33};
34
35pub use mj_client::session::{RelayAttachment, StartedReviewer};
36use mj_core::worker_launch::ReviewerLaunchConfig;
37
38const RELAY_RPC_TIMEOUT: Duration = Duration::from_secs(15);
39const RELAY_SLOW_OPERATION_WARNING: Duration = Duration::from_secs(5);
40const RELAY_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(300);
44const RELAY_HISTORY_TIMEOUT: Duration = Duration::from_secs(900);
48const RELAY_ACKNOWLEDGE_TIMEOUT: Duration = Duration::from_secs(300);
52const REVIEW_CAPTURE_TIMEOUT: Duration = Duration::from_secs(300);
55const REVIEW_ANALYSIS_TIMEOUT: Duration = Duration::from_secs(660);
59const RELAY_PROXY_DETACH_GRACE: Duration = Duration::from_millis(500);
60const RELAY_PROXY_REAP_POLL: Duration = Duration::from_millis(10);
61
62const RELAY_PROXY_STDERR_TAIL: usize = 10;
64
65type ProxyStderrTail = Arc<std::sync::Mutex<VecDeque<String>>>;
75
76async fn drain_proxy_stderr(
82 errors: tokio::process::ChildStderr,
83 purpose: String,
84 session_id: String,
85 tail: ProxyStderrTail,
86) {
87 let mut lines = BufReader::new(errors).lines();
88 loop {
89 match lines.next_line().await {
90 Ok(Some(line)) if line.trim().is_empty() => continue,
91 Ok(Some(line)) => {
92 tracing::warn!(%session_id, %purpose, %line, "relay proxy stderr");
93 let mut tail = tail.lock().unwrap_or_else(PoisonError::into_inner);
94 if tail.len() == RELAY_PROXY_STDERR_TAIL {
95 tail.pop_front();
96 }
97 tail.push_back(line);
98 }
99 Ok(None) => return,
100 Err(error) => {
101 tracing::warn!(%session_id, %purpose, %error, "read relay proxy stderr");
102 return;
103 }
104 }
105 }
106}
107
108mod errors;
109pub use errors::*;
110mod connect;
111mod exchange;
112mod relay;
113mod reviewer;
114mod transport;
115use transport::*;
116mod credential_sync;
117pub use credential_sync::*;
118
119pub struct RelayClient {
125 child: Option<Child>,
126 input: Option<ChildStdin>,
127 output: BufReader<ChildStdout>,
128 request_timeout: Duration,
129 abandoned: Option<String>,
132 next_request: u64,
133 connection_nonce: u64,
134 protocol_version: u32,
135 session_id: String,
136 relay_version: String,
137 worker_build: Option<String>,
140 latest_ordinal: u64,
141 latest_digest: String,
142}
143
144#[cfg(test)]
145mod tests;