返回 CodeWhale
host.rs
根目录 / crates / tui / src / fleet / host.rs
1 //! Fleet worker host adapters.
2 //!
3 //! Adapters own process boundaries for worker hosts. The manager can lease and
4 //! observe work through this trait without knowing whether the worker is a
5 //! local child process or an SSH-backed remote command.
6
7 #![allow(dead_code)]
8
9 use std::collections::{BTreeMap, BTreeSet};
10 use std::fs::File;
11 use std::io::{Read, Seek, SeekFrom};
12 use std::path::{Path, PathBuf};
13 use std::process::{Child, Command, ExitStatus, Stdio};
14 use std::thread;
15 use std::time::{Duration, Instant};
16
17 use codewhale_protocol::fleet::FleetHostSpec;
18 use thiserror::Error;
19
20 #[cfg(unix)]
21 use std::os::unix::process::CommandExt;
22 #[cfg(windows)]
23 use std::os::windows::io::AsRawHandle;
24 #[cfg(unix)]
25 use std::sync::OnceLock;
26 #[cfg(windows)]
27 use windows::Win32::Foundation::{CloseHandle, HANDLE};
28 #[cfg(windows)]
29 use windows::Win32::System::JobObjects::{
30 AssignProcessToJobObject, CreateJobObjectW, JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE,
31 JOBOBJECT_BASIC_ACCOUNTING_INFORMATION, JOBOBJECT_EXTENDED_LIMIT_INFORMATION,
32 JobObjectBasicAccountingInformation, JobObjectExtendedLimitInformation,
33 QueryInformationJobObject, SetInformationJobObject, TerminateJobObject,
34 };
35 #[cfg(windows)]
36 use windows::core::PCWSTR;
37
38 const DEFAULT_LOG_LIMIT_BYTES: usize = 64 * 1024;
39 const DEFAULT_CONNECT_TIMEOUT_SECONDS: u64 = 10;
40 const WORKER_STOP_GRACE: Duration = Duration::from_millis(750);
41
42 pub type FleetHostResult<T> = Result<T, FleetHostError>;
43
44 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
45 pub enum FleetHostErrorKind {
46 Retryable,
47 Terminal,
48 Configuration,
49 }
50
51 #[derive(Debug, Error)]
52 #[error("{kind:?}: {message}")]
53 pub struct FleetHostError {
54 pub kind: FleetHostErrorKind,
55 pub message: String,
56 }
57
58 impl FleetHostError {
59 fn retryable(message: impl Into<String>) -> Self {
60 Self {
61 kind: FleetHostErrorKind::Retryable,
62 message: message.into(),
63 }
64 }
65
66 fn terminal(message: impl Into<String>) -> Self {
67 Self {
68 kind: FleetHostErrorKind::Terminal,
69 message: message.into(),
70 }
71 }
72
73 fn configuration(message: impl Into<String>) -> Self {
74 Self {
75 kind: FleetHostErrorKind::Configuration,
76 message: message.into(),
77 }
78 }
79 }
80
81 #[derive(Debug, Clone, PartialEq, Eq)]
82 pub struct FleetWorkerCommand {
83 pub program: String,
84 pub args: Vec<String>,
85 }
86
87 impl FleetWorkerCommand {
88 pub fn new<S, I, A>(program: S, args: I) -> Self
89 where
90 S: Into<String>,
91 I: IntoIterator<Item = A>,
92 A: Into<String>,
93 {
94 Self {
95 program: program.into(),
96 args: args.into_iter().map(Into::into).collect(),
97 }
98 }
99 }
100
101 #[derive(Debug, Clone)]
102 pub struct FleetWorkerStartRequest {
103 pub worker_id: String,
104 pub command: FleetWorkerCommand,
105 pub cwd: Option<PathBuf>,
106 pub env: BTreeMap<String, String>,
107 pub env_allowlist: BTreeSet<String>,
108 pub log_limit_bytes: usize,
109 }
110
111 impl FleetWorkerStartRequest {
112 pub fn new(worker_id: impl Into<String>, command: FleetWorkerCommand) -> Self {
113 Self {
114 worker_id: worker_id.into(),
115 command,
116 cwd: None,
117 env: BTreeMap::new(),
118 env_allowlist: BTreeSet::new(),
119 log_limit_bytes: DEFAULT_LOG_LIMIT_BYTES,
120 }
121 }
122 }
123
124 #[derive(Debug, Clone, PartialEq, Eq)]
125 pub struct FleetWorkerHandle {
126 pub worker_id: String,
127 pub host_kind: FleetHostKind,
128 pub pid: Option<u32>,
129 pub log_path: PathBuf,
130 }
131
132 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
133 pub enum FleetHostKind {
134 LocalProcess,
135 Ssh,
136 }
137
138 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
139 pub enum FleetHostWorkerState {
140 Running,
141 /// The dispatcher stopped but the owned process session/job is not yet empty.
142 Draining,
143 Exited,
144 Failed,
145 Stopped,
146 Unknown,
147 }
148
149 #[derive(Debug, Clone, PartialEq, Eq)]
150 pub struct FleetHostWorkerStatus {
151 pub worker_id: String,
152 pub state: FleetHostWorkerState,
153 pub pid: Option<u32>,
154 pub exit_code: Option<i32>,
155 pub memory_mb: Option<u64>,
156 pub retryable: bool,
157 }
158
159 pub trait FleetHostAdapter {
160 fn start_worker(
161 &mut self,
162 request: FleetWorkerStartRequest,
163 ) -> FleetHostResult<FleetWorkerHandle>;
164 fn read_status(&mut self, worker_id: &str) -> FleetHostResult<FleetHostWorkerStatus>;
165 fn read_logs(&self, worker_id: &str, max_bytes: usize) -> FleetHostResult<String>;
166 fn interrupt_worker(&mut self, worker_id: &str) -> FleetHostResult<FleetHostWorkerStatus>;
167 fn restart_worker(&mut self, worker_id: &str) -> FleetHostResult<FleetWorkerHandle>;
168 fn stop_worker(&mut self, worker_id: &str) -> FleetHostResult<FleetHostWorkerStatus>;
169 fn cleanup_worker(&mut self, worker_id: &str) -> FleetHostResult<()>;
170 }
171
172 #[derive(Debug)]
173 pub struct LocalProcessFleetHostAdapter {
174 workspace: PathBuf,
175 processes: BTreeMap<String, LocalWorkerProcess>,
176 }
177
178 #[derive(Debug)]
179 struct LocalWorkerProcess {
180 request: FleetWorkerStartRequest,
181 child: Child,
182 #[cfg(unix)]
183 session_id: libc::pid_t,
184 #[cfg(unix)]
185 parent_death_writer: Option<std::io::PipeWriter>,
186 #[cfg(windows)]
187 windows_job: FleetWindowsJob,
188 host_kind: FleetHostKind,
189 log_path: PathBuf,
190 stopped: bool,
191 last_exit: Option<ExitStatus>,
192 last_memory_mb: Option<u64>,
193 /// When `ps` last sampled this worker. Status polls can run every few
194 /// milliseconds; memory is display data, so one sample a second is ample.
195 last_memory_sample: Option<std::time::Instant>,
196 }
197
198 /// Minimum spacing between `ps` memory samples for one worker.
199 const MEMORY_SAMPLE_INTERVAL: Duration = Duration::from_secs(1);
200
201 impl LocalProcessFleetHostAdapter {
202 pub fn new(workspace: impl AsRef<Path>) -> Self {
203 Self {
204 workspace: workspace.as_ref().to_path_buf(),
205 processes: BTreeMap::new(),
206 }
207 }
208
209 fn start_with_kind(
210 &mut self,
211 request: FleetWorkerStartRequest,
212 host_kind: FleetHostKind,
213 ) -> FleetHostResult<FleetWorkerHandle> {
214 validate_worker_id(&request.worker_id)?;
215 if self.processes.contains_key(&request.worker_id) {
216 let status = self.read_status(&request.worker_id)?;
217 if matches!(status.state, FleetHostWorkerState::Running) {
218 return Err(FleetHostError::terminal(format!(
219 "worker {} is already running",
220 request.worker_id
221 )));
222 }
223 // A draining worker's dispatcher exited but its tree still runs;
224 // forgetting it here would overlap a replacement with it.
225 if matches!(status.state, FleetHostWorkerState::Draining) {
226 return Err(FleetHostError::retryable(format!(
227 "worker {} is still draining; stop it before starting a replacement",
228 request.worker_id
229 )));
230 }
231 self.processes.remove(&request.worker_id);
232 }
233
234 let env = worker_env(&request.env, &request.env_allowlist)?;
235 let log_path = self.log_path_for(&request.worker_id, host_kind);
236 let log = open_worker_log(&self.workspace, &log_path)?;
237 let stderr = log
238 .try_clone()
239 .map_err(|err| FleetHostError::retryable(format!("cloning worker log: {err}")))?;
240
241 let mut command = Command::new(&request.command.program);
242 // Parent-death watch (R7): the worker's stdin is the read end of a pipe
243 // whose write end lives only in this adapter process. If the manager
244 // dies (crash, kill, power loss), the kernel closes the write end, the
245 // worker sees stdin EOF, and `--parent-death-watch` shuts the worker
246 // tree down instead of letting it spend forever. Windows workers are
247 // contained in a Job Object that the OS terminates on parent death.
248 #[cfg(unix)]
249 let (parent_death_writer, stdin) = {
250 let (reader, writer) = std::io::pipe().map_err(|err| {
251 FleetHostError::retryable(format!("creating parent-death pipe: {err}"))
252 })?;
253 (Some(writer), std::process::Stdio::from(reader))
254 };
255 #[cfg(not(unix))]
256 let stdin = Stdio::null();
257 command
258 .args(&request.command.args)
259 .stdin(stdin)
260 .stdout(Stdio::from(log))
261 .stderr(Stdio::from(stderr))
262 .env_clear()
263 .envs(env);
264 if let Some(cwd) = &request.cwd {
265 command.current_dir(cwd);
266 }
267
268 // Fleet owns the complete worker tree, not only the dispatcher PID.
269 // `codewhale` spawns `codewhale-tui`, which can in turn spawn tool
270 // processes; isolating the root prevents a stop from signalling the
271 // operator's own process group.
272 #[cfg(unix)]
273 // SAFETY: `setsid` is async-signal-safe and the closure does not touch
274 // allocator or parent-held state between fork and exec.
275 unsafe {
276 command.pre_exec(|| {
277 if libc::setsid() == -1 {
278 Err(std::io::Error::last_os_error())
279 } else {
280 Ok(())
281 }
282 });
283 }
284
285 let child = command.spawn().map_err(|err| {
286 classify_spawn_error(err, format!("starting worker {}", request.worker_id))
287 })?;
288 #[cfg(windows)]
289 let (child, windows_job) = attach_fleet_windows_job(child).map_err(|err| {
290 FleetHostError::retryable(format!(
291 "containing worker {} in a Windows Job Object: {err}",
292 request.worker_id
293 ))
294 })?;
295 let pid = child.id();
296 let handle = FleetWorkerHandle {
297 worker_id: request.worker_id.clone(),
298 host_kind,
299 pid: Some(pid),
300 log_path: log_path.clone(),
301 };
302 self.processes.insert(
303 request.worker_id.clone(),
304 LocalWorkerProcess {
305 request,
306 child,
307 // SAFETY: `setsid` is async-signal-safe and the closure does not
308 // touch allocator or parent-held state between fork and exec.
309 #[cfg(unix)]
310 session_id: pid as libc::pid_t,
311 #[cfg(unix)]
312 parent_death_writer,
313 #[cfg(windows)]
314 windows_job,
315 host_kind,
316 log_path,
317 stopped: false,
318 last_exit: None,
319 last_memory_mb: None,
320 last_memory_sample: None,
321 },
322 );
323 Ok(handle)
324 }
325
326 fn log_path_for(&self, worker_id: &str, host_kind: FleetHostKind) -> PathBuf {
327 let host_dir = match host_kind {
328 FleetHostKind::LocalProcess => "local",
329 FleetHostKind::Ssh => "ssh",
330 };
331 self.workspace
332 .join(".codewhale")
333 .join("fleet-host")
334 .join(host_dir)
335 .join(format!("{}.log", safe_path_segment(worker_id)))
336 }
337 }
338
339 impl FleetHostAdapter for LocalProcessFleetHostAdapter {
340 fn start_worker(
341 &mut self,
342 request: FleetWorkerStartRequest,
343 ) -> FleetHostResult<FleetWorkerHandle> {
344 self.start_with_kind(request, FleetHostKind::LocalProcess)
345 }
346
347 fn read_status(&mut self, worker_id: &str) -> FleetHostResult<FleetHostWorkerStatus> {
348 let process = self
349 .processes
350 .get_mut(worker_id)
351 .ok_or_else(|| FleetHostError::terminal(format!("unknown worker {worker_id}")))?;
352 if let Some(status) = process.last_exit {
353 if local_worker_tree_alive(process)? {
354 return Ok(FleetHostWorkerStatus {
355 worker_id: worker_id.to_string(),
356 state: FleetHostWorkerState::Draining,
357 pid: Some(process.child.id()),
358 exit_code: status.code(),
359 memory_mb: process.last_memory_mb,
360 retryable: true,
361 });
362 }
363 return Ok(status_from_exit(
364 worker_id,
365 Some(process.child.id()),
366 status,
367 process.stopped,
368 process.last_memory_mb,
369 ));
370 }
371 match process.child.try_wait() {
372 Ok(None) => {
373 let pid = process.child.id();
374 let due = process
375 .last_memory_sample
376 .is_none_or(|sampled| sampled.elapsed() >= MEMORY_SAMPLE_INTERVAL);
377 let memory_mb = if process.host_kind == FleetHostKind::LocalProcess && due {
378 process.last_memory_sample = Some(std::time::Instant::now());
379 sample_process_memory_mb(pid)
380 } else {
381 None
382 };
383 process.last_memory_mb = memory_mb.or(process.last_memory_mb);
384 Ok(FleetHostWorkerStatus {
385 worker_id: worker_id.to_string(),
386 state: FleetHostWorkerState::Running,
387 pid: Some(pid),
388 exit_code: None,
389 // Report the retained value, not the raw sample: a
390 // transient ps failure must not flicker a live worker's
391 // memory to None (the Exited arm already does this).
392 memory_mb: process.last_memory_mb,
393 retryable: false,
394 })
395 }
396 Ok(Some(status)) => {
397 process.last_exit = Some(status);
398 if local_worker_tree_alive(process)? {
399 return Ok(FleetHostWorkerStatus {
400 worker_id: worker_id.to_string(),
401 state: FleetHostWorkerState::Draining,
402 pid: Some(process.child.id()),
403 exit_code: status.code(),
404 memory_mb: process.last_memory_mb,
405 retryable: true,
406 });
407 }
408 Ok(status_from_exit(
409 worker_id,
410 Some(process.child.id()),
411 status,
412 process.stopped,
413 process.last_memory_mb,
414 ))
415 }
416 Err(err) => Err(FleetHostError::retryable(format!(
417 "reading worker {worker_id} status: {err}"
418 ))),
419 }
420 }
421
422 fn read_logs(&self, worker_id: &str, max_bytes: usize) -> FleetHostResult<String> {
423 let process = self
424 .processes
425 .get(worker_id)
426 .ok_or_else(|| FleetHostError::terminal(format!("unknown worker {worker_id}")))?;
427 let max_bytes = max_bytes.min(process.request.log_limit_bytes.max(1));
428 read_bounded_log(&self.workspace, &process.log_path, max_bytes)
429 }
430
431 fn interrupt_worker(&mut self, worker_id: &str) -> FleetHostResult<FleetHostWorkerStatus> {
432 {
433 let process = self
434 .processes
435 .get_mut(worker_id)
436 .ok_or_else(|| FleetHostError::terminal(format!("unknown worker {worker_id}")))?;
437 // The direct dispatcher may already be reaped while delegated
438 // session/job descendants remain. Interrupt the containment
439 // boundary unconditionally.
440 interrupt_worker_tree(process)?;
441 }
442 wait_for_exit(self, worker_id, WORKER_STOP_GRACE)
443 }
444
445 fn restart_worker(&mut self, worker_id: &str) -> FleetHostResult<FleetWorkerHandle> {
446 let request = self
447 .processes
448 .get(worker_id)
449 .map(|process| process.request.clone())
450 .ok_or_else(|| FleetHostError::terminal(format!("unknown worker {worker_id}")))?;
451 restart_after_confirmed_stop(self, worker_id, |adapter| {
452 adapter.processes.remove(worker_id);
453 adapter.start_worker(request)
454 })
455 }
456
457 fn stop_worker(&mut self, worker_id: &str) -> FleetHostResult<FleetHostWorkerStatus> {
458 {
459 let process = self
460 .processes
461 .get_mut(worker_id)
462 .ok_or_else(|| FleetHostError::terminal(format!("unknown worker {worker_id}")))?;
463 process.stopped = true;
464 if process.last_exit.is_none() {
465 match process.child.try_wait() {
466 Ok(Some(status)) => {
467 process.last_exit = Some(status);
468 }
469 Ok(None) => {}
470 Err(err) => {
471 return Err(FleetHostError::retryable(format!(
472 "reading worker {worker_id} status before stop: {err}"
473 )));
474 }
475 }
476 }
477 // Always tear down the containment boundary. A dispatcher can
478 // exit before a delegated TUI/tool child, so direct-child status
479 // is not proof that the complete worker tree is gone.
480 stop_worker_tree(process).map_err(|err| FleetHostError {
481 kind: err.kind,
482 message: format!("stopping worker {worker_id}: {}", err.message),
483 })?;
484 }
485 self.read_status(worker_id)
486 }
487
488 fn cleanup_worker(&mut self, worker_id: &str) -> FleetHostResult<()> {
489 if self.processes.contains_key(worker_id) {
490 // Cleanup is the final containment boundary. Even when the direct
491 // dispatcher already exited, delegated children may still occupy
492 // its Unix session or Windows Job Object.
493 let _ = self.stop_worker(worker_id)?;
494 }
495 self.processes.remove(worker_id);
496 Ok(())
497 }
498 }
499
500 #[derive(Debug, Clone)]
501 pub struct SshFleetHostConfig {
502 pub host: String,
503 pub user: Option<String>,
504 pub port: Option<u16>,
505 pub identity: Option<PathBuf>,
506 pub known_hosts: Option<PathBuf>,
507 pub host_key_fingerprint: Option<String>,
508 pub working_directory: PathBuf,
509 pub env_allowlist: BTreeSet<String>,
510 pub codewhale_binary: String,
511 pub ssh_binary: String,
512 pub connect_timeout_seconds: u64,
513 }
514
515 impl SshFleetHostConfig {
516 pub fn new(host: impl Into<String>, working_directory: impl Into<PathBuf>) -> Self {
517 Self {
518 host: host.into(),
519 user: None,
520 port: None,
521 identity: None,
522 known_hosts: None,
523 host_key_fingerprint: None,
524 working_directory: working_directory.into(),
525 env_allowlist: BTreeSet::new(),
526 codewhale_binary: "codewhale".to_string(),
527 ssh_binary: "ssh".to_string(),
528 connect_timeout_seconds: DEFAULT_CONNECT_TIMEOUT_SECONDS,
529 }
530 }
531
532 pub fn from_host_spec(spec: &FleetHostSpec) -> FleetHostResult<Self> {
533 let FleetHostSpec::Ssh {
534 host,
535 port,
536 user,
537 identity,
538 known_hosts,
539 host_key_fingerprint,
540 working_directory,
541 env_allowlist,
542 codewhale_binary,
543 } = spec
544 else {
545 return Err(FleetHostError::configuration(
546 "expected SSH Fleet host spec",
547 ));
548 };
549 let working_directory = working_directory.clone().ok_or_else(|| {
550 FleetHostError::configuration("SSH Fleet host spec requires working_directory")
551 })?;
552 let codewhale_binary = codewhale_binary.clone().ok_or_else(|| {
553 FleetHostError::configuration("SSH Fleet host spec requires codewhale_binary")
554 })?;
555 let mut config = Self::new(host.clone(), working_directory);
556 config.port = *port;
557 config.user = user.clone();
558 config.identity = identity.clone();
559 config.known_hosts = known_hosts.clone();
560 config.host_key_fingerprint = host_key_fingerprint.clone();
561 config.env_allowlist = env_allowlist.iter().cloned().collect();
562 config.codewhale_binary = codewhale_binary;
563 config.validate()?;
564 Ok(config)
565 }
566
567 fn validate(&self) -> FleetHostResult<()> {
568 if self.host.trim().is_empty() {
569 return Err(FleetHostError::configuration(
570 "SSH Fleet host requires an explicit host",
571 ));
572 }
573 // The destination is a single ssh argument. A leading '-' would be read
574 // as an ssh option, and whitespace/control characters would change what
575 // ssh or the remote shell receives.
576 validate_ssh_destination_part("host", &self.host)?;
577 if let Some(user) = self.user.as_ref().filter(|user| !user.trim().is_empty()) {
578 validate_ssh_destination_part("user", user)?;
579 }
580 if self.codewhale_binary.trim().is_empty() {
581 return Err(FleetHostError::configuration(
582 "SSH Fleet host requires an explicit codewhale binary path",
583 ));
584 }
585 if self.working_directory.as_os_str().is_empty() {
586 return Err(FleetHostError::configuration(
587 "SSH Fleet host requires an explicit working directory",
588 ));
589 }
590 // Fingerprint-only verification is not implemented by the OpenSSH adapter.
591 // Refuse a configured pin rather than silently substituting another trust source.
592 if self.host_key_fingerprint.is_some() {
593 return Err(FleetHostError::configuration(
594 "SSH Fleet host_key_fingerprint is unsupported; configure known_hosts instead",
595 ));
596 }
597 if self
598 .known_hosts
599 .as_ref()
600 .is_some_and(|path| path.as_os_str().is_empty())
601 {
602 return Err(FleetHostError::configuration(
603 "SSH Fleet known_hosts must not be empty",
604 ));
605 }
606 if self
607 .known_hosts
608 .as_ref()
609 .is_some_and(|path| path.to_string_lossy().contains(['"', '\n', '\r']))
610 {
611 return Err(FleetHostError::configuration(
612 "SSH Fleet known_hosts contains unsupported characters",
613 ));
614 }
615 validate_env_allowlist(&self.env_allowlist)
616 }
617
618 fn target(&self) -> String {
619 self.user
620 .as_ref()
621 .filter(|user| !user.trim().is_empty())
622 .map(|user| format!("{user}@{}", self.host))
623 .unwrap_or_else(|| self.host.clone())
624 }
625 }
626
627 #[derive(Debug)]
628 pub struct SshFleetHostAdapter {
629 config: SshFleetHostConfig,
630 local: LocalProcessFleetHostAdapter,
631 }
632
633 impl SshFleetHostAdapter {
634 pub fn new(workspace: impl AsRef<Path>, config: SshFleetHostConfig) -> FleetHostResult<Self> {
635 config.validate()?;
636 Ok(Self {
637 config,
638 local: LocalProcessFleetHostAdapter::new(workspace),
639 })
640 }
641
642 pub fn build_ssh_command(
643 &self,
644 request: &FleetWorkerStartRequest,
645 ) -> FleetHostResult<FleetWorkerCommand> {
646 self.config.validate()?;
647 let env = filtered_env(&request.env, &self.config.env_allowlist)?;
648 let mut args = vec![
649 "-o".to_string(),
650 "BatchMode=yes".to_string(),
651 "-o".to_string(),
652 "StrictHostKeyChecking=yes".to_string(),
653 "-o".to_string(),
654 format!("ConnectTimeout={}", self.config.connect_timeout_seconds),
655 ];
656 if let Some(known_hosts) = &self.config.known_hosts {
657 args.push("-o".to_string());
658 args.push(format!("UserKnownHostsFile=\"{}\"", known_hosts.display()));
659 args.push("-o".to_string());
660 args.push("GlobalKnownHostsFile=none".to_string());
661 args.push("-o".to_string());
662 // IgnoreUnknown predates OpenSSH 7.x. Older clients have no
663 // KnownHostsCommand; newer ones must disable that extra trust source.
664 args.push("IgnoreUnknown=KnownHostsCommand".to_string());
665 args.push("-o".to_string());
666 args.push("KnownHostsCommand=none".to_string());
667 args.push("-o".to_string());
668 args.push("VerifyHostKeyDNS=no".to_string());
669 }
670 for key in env.keys() {
671 args.push("-o".to_string());
672 args.push(format!("SendEnv={key}"));
673 }
674 if let Some(port) = self.config.port {
675 args.push("-p".to_string());
676 args.push(port.to_string());
677 }
678 if let Some(identity) = &self.config.identity {
679 args.push("-i".to_string());
680 args.push(identity.display().to_string());
681 }
682 // End option parsing so the destination can never be read as an option.
683 args.push("--".to_string());
684 args.push(self.config.target());
685 args.push(self.remote_command(request));
686 Ok(FleetWorkerCommand::new(
687 self.config.ssh_binary.clone(),
688 args,
689 ))
690 }
691
692 fn ssh_start_request(
693 &self,
694 request: FleetWorkerStartRequest,
695 ) -> FleetHostResult<FleetWorkerStartRequest> {
696 let command = self.build_ssh_command(&request)?;
697 let mut env = ssh_client_env();
698 env.extend(filtered_env(&request.env, &self.config.env_allowlist)?);
699 let env_allowlist = env.keys().cloned().collect();
700 Ok(FleetWorkerStartRequest {
701 worker_id: request.worker_id,
702 command,
703 cwd: None,
704 env,
705 env_allowlist,
706 log_limit_bytes: request.log_limit_bytes,
707 })
708 }
709
710 fn remote_command(&self, request: &FleetWorkerStartRequest) -> String {
711 let mut parts = vec![
712 "cd".to_string(),
713 shell_quote(&self.config.working_directory.display().to_string()),
714 "&&".to_string(),
715 "exec".to_string(),
716 shell_quote(&self.config.codewhale_binary),
717 ];
718 parts.extend(request.command.args.iter().map(|arg| shell_quote(arg)));
719 parts.join(" ")
720 }
721 }
722
723 impl FleetHostAdapter for SshFleetHostAdapter {
724 fn start_worker(
725 &mut self,
726 request: FleetWorkerStartRequest,
727 ) -> FleetHostResult<FleetWorkerHandle> {
728 let request = self.ssh_start_request(request)?;
729 self.local.start_with_kind(request, FleetHostKind::Ssh)
730 }
731
732 fn read_status(&mut self, worker_id: &str) -> FleetHostResult<FleetHostWorkerStatus> {
733 self.local.read_status(worker_id)
734 }
735
736 fn read_logs(&self, worker_id: &str, max_bytes: usize) -> FleetHostResult<String> {
737 self.local.read_logs(worker_id, max_bytes)
738 }
739
740 fn interrupt_worker(&mut self, worker_id: &str) -> FleetHostResult<FleetHostWorkerStatus> {
741 self.local.interrupt_worker(worker_id)
742 }
743
744 fn restart_worker(&mut self, worker_id: &str) -> FleetHostResult<FleetWorkerHandle> {
745 let request = self
746 .local
747 .processes
748 .get(worker_id)
749 .map(|process| process.request.clone())
750 .ok_or_else(|| FleetHostError::terminal(format!("unknown worker {worker_id}")))?;
751 restart_after_confirmed_stop(self, worker_id, |adapter| {
752 adapter.local.processes.remove(worker_id);
753 adapter.local.start_with_kind(request, FleetHostKind::Ssh)
754 })
755 }
756
757 fn stop_worker(&mut self, worker_id: &str) -> FleetHostResult<FleetHostWorkerStatus> {
758 self.local.stop_worker(worker_id)
759 }
760
761 fn cleanup_worker(&mut self, worker_id: &str) -> FleetHostResult<()> {
762 self.local.cleanup_worker(worker_id)
763 }
764 }
765
766 /// Restart policy shared by the process-backed adapters. The previous worker
767 /// is released — and `replace` spawns its successor — only once a stop has
768 /// confirmed the whole worker tree is gone. A failed or unconfirmed stop keeps
769 /// the old handle so status, logs, and cleanup still reach a worker that may
770 /// be running, and no replacement starts beside it (duplicate execution). The
771 /// process-tree lifecycle in `stop_worker` stays the authority on "gone".
772 ///
773 /// Known limitations:
774 /// - For [`SshFleetHostAdapter`], "confirmed stopped" proves only that the
775 /// local `ssh` client's process tree is gone. The remote `codewhale`
776 /// process is not observed: if the remote side does not tear the command
777 /// down when the connection drops, a restart can start a second remote
778 /// worker beside it.
779 /// - `start_worker` on an id the adapter still holds refuses only `Running`
780 /// and `Draining` (see `LocalProcessFleetHostAdapter::start_with_kind`),
781 /// while this restart also refuses `Unknown`. The process-backed adapters
782 /// never report `Unknown`, so the gap is latent; an adapter that does must
783 /// refuse `Unknown` in its start path too, or a start can overlap a worker
784 /// in an unknown state.
785 fn restart_after_confirmed_stop<A: FleetHostAdapter>(
786 adapter: &mut A,
787 worker_id: &str,
788 replace: impl FnOnce(&mut A) -> FleetHostResult<FleetWorkerHandle>,
789 ) -> FleetHostResult<FleetWorkerHandle> {
790 let stopped = adapter.stop_worker(worker_id).map_err(|err| FleetHostError {
791 kind: err.kind,
792 message: format!(
793 "restart of worker {worker_id} refused: the previous worker was not confirmed stopped ({})",
794 err.message
795 ),
796 })?;
797 if matches!(
798 stopped.state,
799 FleetHostWorkerState::Running
800 | FleetHostWorkerState::Draining
801 | FleetHostWorkerState::Unknown
802 ) {
803 return Err(FleetHostError::retryable(format!(
804 "restart of worker {worker_id} refused: the previous worker is still {:?} after stop",
805 stopped.state
806 )));
807 }
808 replace(adapter)
809 }
810
811 fn open_worker_log(workspace: &Path, path: &Path) -> FleetHostResult<File> {
812 // Pinned, no-follow open under the workspace: a link in `.codewhale` or at
813 // the log's name is refused instead of redirecting the worker's output.
814 let relative = path.strip_prefix(workspace).map_err(|_| {
815 FleetHostError::retryable(format!(
816 "worker log {} is outside the workspace",
817 path.display()
818 ))
819 })?;
820 let file = super::files::WorkspaceFile::open(workspace, relative, true)
821 .and_then(|target| target.open_write(false))
822 .map_err(|err| FleetHostError::retryable(format!("opening worker log: {err}")))?;
823 // Validated handle first, truncation second.
824 file.set_len(0)
825 .map_err(|err| FleetHostError::retryable(format!("truncating worker log: {err}")))?;
826 Ok(file)
827 }
828
829 fn read_bounded_log(workspace: &Path, path: &Path, max_bytes: usize) -> FleetHostResult<String> {
830 // The worker may still hold its log open for writing.
831 let mut file = crate::fs_confined::open_read_shared(workspace, path).map_err(|err| {
832 FleetHostError::retryable(format!("opening worker log {}: {err}", path.display()))
833 })?;
834 let len = file
835 .metadata()
836 .map_err(|err| FleetHostError::retryable(format!("reading worker log metadata: {err}")))?
837 .len();
838 let max_bytes = max_bytes.max(1) as u64;
839 if len > max_bytes {
840 file.seek(SeekFrom::Start(len - max_bytes))
841 .map_err(|err| FleetHostError::retryable(format!("seeking worker log: {err}")))?;
842 }
843 let mut bytes = Vec::new();
844 file.read_to_end(&mut bytes)
845 .map_err(|err| FleetHostError::retryable(format!("reading worker log: {err}")))?;
846 Ok(String::from_utf8_lossy(&bytes).into_owned())
847 }
848
849 fn status_from_exit(
850 worker_id: &str,
851 pid: Option<u32>,
852 status: ExitStatus,
853 stopped: bool,
854 memory_mb: Option<u64>,
855 ) -> FleetHostWorkerStatus {
856 let success = status.success();
857 FleetHostWorkerStatus {
858 worker_id: worker_id.to_string(),
859 state: if stopped {
860 FleetHostWorkerState::Stopped
861 } else if success {
862 FleetHostWorkerState::Exited
863 } else {
864 FleetHostWorkerState::Failed
865 },
866 pid,
867 exit_code: status.code(),
868 memory_mb,
869 retryable: !success && !stopped,
870 }
871 }
872
873 #[cfg(unix)]
874 fn sample_process_memory_mb(pid: u32) -> Option<u64> {
875 // Resolve `ps` via PATH like every other external command in the
876 // codebase: /bin/ps does not exist on NixOS and some minimal containers,
877 // which would silently report permanent None for live workers. Restricted
878 // sandboxes may also deny process-table inspection with EPERM; treat that
879 // as unavailable rather than panicking or inventing a sample.
880 if !process_table_inspection_available() {
881 return None;
882 }
883 let output = match Command::new("ps")
884 .args(["-o", "rss=", "-p", &pid.to_string()])
885 .output()
886 {
887 Ok(output) => output,
888 Err(err) if is_permission_denied(&err) => {
889 mark_process_table_unavailable();
890 return None;
891 }
892 Err(_) => return None,
893 };
894 if !output.status.success() {
895 return None;
896 }
897 let rss_kb = String::from_utf8_lossy(&output.stdout)
898 .split_whitespace()
899 .next()?
900 .parse::<u64>()
901 .ok()?;
902 (rss_kb > 0).then_some(rss_kb.div_ceil(1024))
903 }
904
905 #[cfg(not(unix))]
906 fn sample_process_memory_mb(_pid: u32) -> Option<u64> {
907 None
908 }
909
910 fn classify_spawn_error(err: std::io::Error, context: String) -> FleetHostError {
911 match err.kind() {
912 std::io::ErrorKind::NotFound => FleetHostError::configuration(format!("{context}: {err}")),
913 std::io::ErrorKind::PermissionDenied => {
914 FleetHostError::terminal(format!("{context}: {err}"))
915 }
916 _ => FleetHostError::retryable(format!("{context}: {err}")),
917 }
918 }
919
920 fn wait_for_exit(
921 adapter: &mut LocalProcessFleetHostAdapter,
922 worker_id: &str,
923 timeout: Duration,
924 ) -> FleetHostResult<FleetHostWorkerStatus> {
925 let deadline = Instant::now() + timeout;
926 loop {
927 let status = adapter.read_status(worker_id)?;
928 if !matches!(
929 status.state,
930 FleetHostWorkerState::Running | FleetHostWorkerState::Draining
931 ) {
932 return Ok(status);
933 }
934 if Instant::now() >= deadline {
935 return Ok(status);
936 }
937 thread::sleep(Duration::from_millis(25));
938 }
939 }
940
941 #[cfg(unix)]
942 fn local_worker_tree_alive(process: &LocalWorkerProcess) -> FleetHostResult<bool> {
943 Ok(!unix_session_members(process.session_id, Some(process.session_id))?.is_empty())
944 }
945
946 #[cfg(windows)]
947 fn local_worker_tree_alive(process: &LocalWorkerProcess) -> FleetHostResult<bool> {
948 process.windows_job.has_active_processes().map_err(|err| {
949 FleetHostError::retryable(format!("querying Windows worker job activity: {err}"))
950 })
951 }
952
953 #[cfg(not(any(unix, windows)))]
954 fn local_worker_tree_alive(_process: &LocalWorkerProcess) -> FleetHostResult<bool> {
955 Ok(false)
956 }
957
958 #[cfg(unix)]
959 fn interrupt_worker_tree(process: &mut LocalWorkerProcess) -> FleetHostResult<()> {
960 shutdown_unix_worker_session(process, &[libc::SIGINT, libc::SIGTERM])
961 }
962
963 #[cfg(windows)]
964 fn interrupt_worker_tree(process: &mut LocalWorkerProcess) -> FleetHostResult<()> {
965 process.windows_job.terminate().map_err(|err| {
966 FleetHostError::retryable(format!("interrupting Windows worker tree: {err}"))
967 })
968 }
969
970 #[cfg(not(any(unix, windows)))]
971 fn interrupt_worker_tree(process: &mut LocalWorkerProcess) -> FleetHostResult<()> {
972 process
973 .child
974 .kill()
975 .map_err(|err| FleetHostError::retryable(format!("interrupting worker: {err}")))
976 }
977
978 #[cfg(unix)]
979 fn stop_worker_tree(process: &mut LocalWorkerProcess) -> FleetHostResult<()> {
980 shutdown_unix_worker_session(process, &[libc::SIGTERM])
981 }
982
983 #[cfg(windows)]
984 fn stop_worker_tree(process: &mut LocalWorkerProcess) -> FleetHostResult<()> {
985 process.windows_job.terminate().map_err(|err| {
986 FleetHostError::retryable(format!("terminating Windows worker job: {err}"))
987 })?;
988 if process.last_exit.is_none() {
989 process.last_exit =
990 Some(process.child.wait().map_err(|err| {
991 FleetHostError::retryable(format!("reaping Windows worker: {err}"))
992 })?);
993 }
994 Ok(())
995 }
996
997 #[cfg(not(any(unix, windows)))]
998 fn stop_worker_tree(process: &mut LocalWorkerProcess) -> FleetHostResult<()> {
999 process
1000 .child
1001 .kill()
1002 .map_err(|err| FleetHostError::retryable(format!("killing worker: {err}")))?;
1003 process.last_exit = Some(
1004 process
1005 .child
1006 .wait()
1007 .map_err(|err| FleetHostError::retryable(format!("reaping worker: {err}")))?,
1008 );
1009 Ok(())
1010 }
1011
1012 #[cfg(unix)]
1013 fn shutdown_unix_worker_session(
1014 process: &mut LocalWorkerProcess,
1015 graceful_signals: &[libc::c_int],
1016 ) -> FleetHostResult<()> {
1017 let mut signal_errors = Vec::new();
1018 let known_leader = process.session_id;
1019 for signal in graceful_signals {
1020 signal_errors.extend(signal_unix_session(
1021 process.session_id,
1022 *signal,
1023 Some(known_leader),
1024 )?);
1025 if wait_for_unix_session_exit(process, WORKER_STOP_GRACE)? {
1026 return Ok(());
1027 }
1028 }
1029
1030 signal_errors.extend(signal_unix_session(
1031 process.session_id,
1032 libc::SIGKILL,
1033 Some(known_leader),
1034 )?);
1035 if wait_for_unix_session_exit(process, WORKER_STOP_GRACE)? {
1036 return Ok(());
1037 }
1038
1039 // Without a process table we can only reason about the tracked session
1040 // leader/dispatcher. Prefer an honest degraded success once that known
1041 // pid is gone instead of looping forever on ps EPERM.
1042 if !process_table_inspection_available() {
1043 if process.last_exit.is_some() && !unix_pid_exists(process.session_id) {
1044 return Ok(());
1045 }
1046 return Err(FleetHostError::retryable(format!(
1047 "Fleet session {} still has a live tracked leader after SIGKILL and process-table inspection is unavailable{}",
1048 process.session_id,
1049 if signal_errors.is_empty() {
1050 String::new()
1051 } else {
1052 format!("; signal errors: {}", signal_errors.join("; "))
1053 }
1054 )));
1055 }
1056
1057 let alive = unix_session_members(process.session_id, Some(known_leader))?;
1058 Err(FleetHostError::retryable(format!(
1059 "Fleet session {} still has live processes after SIGKILL: {alive:?}{}",
1060 process.session_id,
1061 if signal_errors.is_empty() {
1062 String::new()
1063 } else {
1064 format!("; signal errors: {}", signal_errors.join("; "))
1065 }
1066 )))
1067 }
1068
1069 #[cfg(unix)]
1070 fn wait_for_unix_session_exit(
1071 process: &mut LocalWorkerProcess,
1072 timeout: Duration,
1073 ) -> FleetHostResult<bool> {
1074 let deadline = Instant::now() + timeout;
1075 let known_leader = process.session_id;
1076 loop {
1077 if process.last_exit.is_none() {
1078 process.last_exit = process.child.try_wait().map_err(|err| {
1079 FleetHostError::retryable(format!("checking Fleet dispatcher exit: {err}"))
1080 })?;
1081 }
1082 if process.last_exit.is_some() {
1083 let members = unix_session_members(process.session_id, Some(known_leader))?;
1084 if members.is_empty() {
1085 return Ok(true);
1086 }
1087 // When process-table inspection is denied we can only track the
1088 // known session leader. Treat an empty known-pid set as success.
1089 if !process_table_inspection_available()
1090 && members.iter().all(|pid| !unix_pid_exists(*pid))
1091 {
1092 return Ok(true);
1093 }
1094 }
1095 if Instant::now() >= deadline {
1096 return Ok(false);
1097 }
1098 thread::sleep(Duration::from_millis(25));
1099 }
1100 }
1101
1102 #[cfg(unix)]
1103 fn unix_session_members(
1104 session_id: libc::pid_t,
1105 known_pids: Option<libc::pid_t>,
1106 ) -> FleetHostResult<Vec<libc::pid_t>> {
1107 match unix_process_ids() {
1108 Ok(pids) => {
1109 let mut members = Vec::new();
1110 for pid in pids {
1111 if pid > 0 {
1112 // Revalidate against the kernel after parsing the snapshot. A PID
1113 // reused by an unrelated process must never receive our signal.
1114 // SAFETY: getsid(2) dereferences no pointers.
1115 if unsafe { libc::getsid(pid) } == session_id {
1116 members.push(pid);
1117 }
1118 }
1119 }
1120 Ok(members)
1121 }
1122 Err(err) if process_table_error_is_unavailable(&err) => {
1123 // Restricted sandboxes may deny full process-table walks. Fall back
1124 // to the known session leader so stop/interrupt still reaches the
1125 // tracked dispatcher without inventing a process census.
1126 Ok(known_pids
1127 .into_iter()
1128 .filter(|pid| *pid > 0 && unix_pid_in_session(*pid, session_id))
1129 .collect())
1130 }
1131 Err(err) => Err(err),
1132 }
1133 }
1134
1135 #[cfg(unix)]
1136 fn unix_pid_in_session(pid: libc::pid_t, session_id: libc::pid_t) -> bool {
1137 // SAFETY: getsid(2) dereferences no pointers.
1138 unsafe { libc::getsid(pid) == session_id }
1139 }
1140
1141 #[cfg(unix)]
1142 fn unix_pid_exists(pid: libc::pid_t) -> bool {
1143 if pid <= 0 {
1144 return false;
1145 }
1146 // SAFETY: kill(2) dereferences no pointers; signal 0 sends nothing.
1147 if unsafe { libc::kill(pid, 0) } == 0 {
1148 return true;
1149 }
1150 std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
1151 }
1152
1153 #[cfg(unix)]
1154 fn is_permission_denied(err: &std::io::Error) -> bool {
1155 err.kind() == std::io::ErrorKind::PermissionDenied || err.raw_os_error() == Some(libc::EPERM)
1156 }
1157
1158 #[cfg(unix)]
1159 fn process_table_error_is_unavailable(err: &FleetHostError) -> bool {
1160 err.message.contains("process-table inspection unavailable")
1161 || err.message.contains("Operation not permitted")
1162 || err.message.contains("Permission denied")
1163 || err.message.contains("EPERM")
1164 }
1165
1166 #[cfg(unix)]
1167 fn mark_process_table_unavailable() {
1168 // Record the denial when nothing has been cached yet. OnceLock cannot flip
1169 // a prior true; live call sites still degrade on the immediate EPERM path.
1170 let _ = process_table_probe_cell().get_or_init(|| false);
1171 }
1172
1173 #[cfg(unix)]
1174 fn process_table_probe_cell() -> &'static OnceLock<bool> {
1175 static PROCESS_TABLE_AVAILABLE: OnceLock<bool> = OnceLock::new();
1176 &PROCESS_TABLE_AVAILABLE
1177 }
1178
1179 /// Returns whether full process-table inspection (`ps` / `/proc`) works here.
1180 /// Cached after the first probe so tests and production share one answer.
1181 #[cfg(unix)]
1182 pub(crate) fn process_table_inspection_available() -> bool {
1183 *process_table_probe_cell().get_or_init(|| match unix_process_ids_uncached() {
1184 Ok(_) => true,
1185 Err(err) if process_table_error_is_unavailable(&err) => false,
1186 // Missing `ps` binary is also unavailable inspection, not a transient
1187 // retryable blip for memory sampling / session census.
1188 Err(err)
1189 if err.message.contains("os error 2")
1190 || err.message.contains("No such file")
1191 || err.message.contains("not found") =>
1192 {
1193 false
1194 }
1195 Err(_) => false,
1196 })
1197 }
1198
1199 #[cfg(all(unix, target_os = "linux"))]
1200 fn unix_process_ids() -> FleetHostResult<Vec<libc::pid_t>> {
1201 unix_process_ids_uncached()
1202 }
1203
1204 #[cfg(all(unix, target_os = "linux"))]
1205 fn unix_process_ids_uncached() -> FleetHostResult<Vec<libc::pid_t>> {
1206 let entries = std::fs::read_dir("/proc").map_err(|err| {
1207 if is_permission_denied(&err) {
1208 FleetHostError::retryable(format!(
1209 "listing Fleet session through /proc: process-table inspection unavailable: {err}"
1210 ))
1211 } else {
1212 FleetHostError::retryable(format!("listing Fleet session through /proc: {err}"))
1213 }
1214 })?;
1215 Ok(entries
1216 .filter_map(Result::ok)
1217 .filter_map(|entry| entry.file_name().to_string_lossy().parse().ok())
1218 .collect())
1219 }
1220
1221 #[cfg(all(unix, not(target_os = "linux")))]
1222 fn unix_process_ids() -> FleetHostResult<Vec<libc::pid_t>> {
1223 if let Some(available) = process_table_probe_cell().get()
1224 && !*available
1225 {
1226 return Err(FleetHostError::retryable(
1227 "listing Fleet session with ps: process-table inspection unavailable",
1228 ));
1229 }
1230 match unix_process_ids_uncached() {
1231 Ok(pids) => {
1232 let _ = process_table_probe_cell().get_or_init(|| true);
1233 Ok(pids)
1234 }
1235 Err(err) => {
1236 if process_table_error_is_unavailable(&err) {
1237 let _ = process_table_probe_cell().get_or_init(|| false);
1238 }
1239 Err(err)
1240 }
1241 }
1242 }
1243
1244 #[cfg(all(unix, not(target_os = "linux")))]
1245 fn unix_process_ids_uncached() -> FleetHostResult<Vec<libc::pid_t>> {
1246 let output = Command::new("ps")
1247 .args(["-A", "-o", "pid="])
1248 .output()
1249 .map_err(|err| {
1250 if is_permission_denied(&err) {
1251 FleetHostError::retryable(format!(
1252 "listing Fleet session with ps: process-table inspection unavailable: {err}"
1253 ))
1254 } else {
1255 FleetHostError::retryable(format!("listing Fleet session with ps: {err}"))
1256 }
1257 })?;
1258 if !output.status.success() {
1259 let stderr = String::from_utf8_lossy(&output.stderr);
1260 let denied = stderr.contains("Operation not permitted")
1261 || stderr.contains("Permission denied")
1262 || output.status.code() == Some(1)
1263 && stderr.to_ascii_lowercase().contains("not permitted");
1264 if denied {
1265 return Err(FleetHostError::retryable(format!(
1266 "listing Fleet session with ps: process-table inspection unavailable: {stderr}"
1267 )));
1268 }
1269 return Err(FleetHostError::retryable(format!(
1270 "listing Fleet session with ps exited {:?}",
1271 output.status.code()
1272 )));
1273 }
1274
1275 Ok(String::from_utf8_lossy(&output.stdout)
1276 .lines()
1277 .filter_map(|line| line.trim().parse().ok())
1278 .collect())
1279 }
1280
1281 #[cfg(unix)]
1282 fn signal_unix_session(
1283 session_id: libc::pid_t,
1284 signal: libc::c_int,
1285 known_leader: Option<libc::pid_t>,
1286 ) -> FleetHostResult<Vec<String>> {
1287 // SAFETY: getsid(2) dereferences no pointers.
1288 let own_session = unsafe { libc::getsid(0) };
1289 if session_id <= 0 || session_id == own_session {
1290 return Err(FleetHostError::terminal(format!(
1291 "refusing to signal unsafe Fleet session {session_id}"
1292 )));
1293 }
1294
1295 let mut errors = Vec::new();
1296 // Prefer known leader first so stop/interrupt still works when the full
1297 // process table cannot be enumerated under a restricted sandbox.
1298 let mut candidates = unix_session_members(session_id, known_leader)?;
1299 if candidates.is_empty()
1300 && let Some(leader) = known_leader.filter(|pid| *pid > 0)
1301 {
1302 candidates.push(leader);
1303 }
1304 for pid in candidates {
1305 // Verify identity again immediately before signalling. Session IDs
1306 // remain stable across reparenting and separate process groups.
1307 // SAFETY: getsid(2) dereferences no pointers.
1308 if unsafe { libc::getsid(pid) } != session_id {
1309 // Leader may already be gone; still try kill on known leader when
1310 // getsid fails only with ESRCH-equivalent absence.
1311 if Some(pid) != known_leader || !unix_pid_exists(pid) {
1312 continue;
1313 }
1314 }
1315 // SAFETY: kill(2) dereferences no pointers.
1316 if unsafe { libc::kill(pid, signal) } != 0 {
1317 let err = std::io::Error::last_os_error();
1318 if err.raw_os_error() != Some(libc::ESRCH) {
1319 errors.push(format!("pid {pid}: {err}"));
1320 }
1321 }
1322 }
1323 Ok(errors)
1324 }
1325
1326 #[cfg(unix)]
1327 fn unix_pid_is_running(pid: libc::pid_t) -> bool {
1328 if !unix_pid_exists(pid) {
1329 return false;
1330 }
1331
1332 // `kill(pid, 0)` also succeeds for zombies. The Fleet containment code
1333 // has already finished its job once a descendant is dead; on macOS an
1334 // orphan can remain visible as a zombie briefly while launchd reaps it.
1335 // Ask `ps` for the process state so test assertions do not mistake that
1336 // transient kernel bookkeeping for a live leaked worker. If `ps` itself
1337 // is denied, stay conservative and treat the PID as running.
1338 if !process_table_inspection_available() {
1339 return true;
1340 }
1341 match Command::new("ps")
1342 .args(["-o", "stat=", "-p", &pid.to_string()])
1343 .output()
1344 {
1345 Ok(output) if output.status.success() => String::from_utf8_lossy(&output.stdout)
1346 .split_whitespace()
1347 .next()
1348 .is_some_and(|state| !state.starts_with('Z')),
1349 // A failed status does not prove exit: preserve the positive kernel
1350 // visibility result and let the bounded waiter retry.
1351 Ok(_) => true,
1352 Err(err) if is_permission_denied(&err) => {
1353 let _ = process_table_probe_cell().get_or_init(|| false);
1354 true
1355 }
1356 Err(_) => true,
1357 }
1358 }
1359
1360 #[cfg(all(unix, test))]
1361 fn wait_for_unix_pid_exit(pid: libc::pid_t, timeout: Duration) -> bool {
1362 let deadline = Instant::now() + timeout;
1363 loop {
1364 if !unix_pid_is_running(pid) {
1365 return true;
1366 }
1367 if Instant::now() >= deadline {
1368 return false;
1369 }
1370 thread::sleep(Duration::from_millis(25));
1371 }
1372 }
1373
1374 #[cfg(windows)]
1375 #[derive(Debug)]
1376 struct FleetWindowsJob {
1377 handle: HANDLE,
1378 }
1379
1380 #[cfg(windows)]
1381 // SAFETY: Job handles are process-wide kernel handles. The adapter owns this
1382 // wrapper exclusively and mutates workers through `&mut self`.
1383 unsafe impl Send for FleetWindowsJob {}
1384
1385 #[cfg(windows)]
1386 // SAFETY: The wrapper exposes only kernel job operations; shared access does
1387 // not mutate Rust-owned memory.
1388 unsafe impl Sync for FleetWindowsJob {}
1389
1390 #[cfg(windows)]
1391 impl FleetWindowsJob {
1392 fn attach_to_child(child: &Child) -> std::io::Result<Self> {
1393 // SAFETY: returned handle is owned by the new wrapper.
1394 let handle = unsafe { CreateJobObjectW(None, PCWSTR::null()).map_err(windows_io_error)? };
1395 let job = Self { handle };
1396 let mut limits = JOBOBJECT_EXTENDED_LIMIT_INFORMATION::default();
1397 limits.BasicLimitInformation.LimitFlags = JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE;
1398 // SAFETY: `limits` is live with matching size; both handles are live.
1399 unsafe {
1400 SetInformationJobObject(
1401 job.handle,
1402 JobObjectExtendedLimitInformation,
1403 &limits as *const _ as *const core::ffi::c_void,
1404 std::mem::size_of::<JOBOBJECT_EXTENDED_LIMIT_INFORMATION>() as u32,
1405 )
1406 .map_err(windows_io_error)?;
1407 AssignProcessToJobObject(job.handle, HANDLE(child.as_raw_handle()))
1408 .map_err(windows_io_error)?;
1409 }
1410 Ok(job)
1411 }
1412
1413 fn terminate(&self) -> std::io::Result<()> {
1414 // SAFETY: `self.handle` is a live owned job handle.
1415 unsafe { TerminateJobObject(self.handle, 1).map_err(windows_io_error) }
1416 }
1417
1418 fn has_active_processes(&self) -> std::io::Result<bool> {
1419 let mut accounting = JOBOBJECT_BASIC_ACCOUNTING_INFORMATION::default();
1420 // SAFETY: `accounting` is live with matching size.
1421 unsafe {
1422 QueryInformationJobObject(
1423 Some(self.handle),
1424 JobObjectBasicAccountingInformation,
1425 &mut accounting as *mut _ as *mut core::ffi::c_void,
1426 std::mem::size_of::<JOBOBJECT_BASIC_ACCOUNTING_INFORMATION>() as u32,
1427 None,
1428 )
1429 .map_err(windows_io_error)?;
1430 }
1431 Ok(accounting.ActiveProcesses > 0)
1432 }
1433 }
1434
1435 #[cfg(windows)]
1436 impl Drop for FleetWindowsJob {
1437 fn drop(&mut self) {
1438 // SAFETY: `self.handle` is owned here; Drop runs once.
1439 unsafe {
1440 let _ = CloseHandle(self.handle);
1441 }
1442 }
1443 }
1444
1445 #[cfg(windows)]
1446 fn attach_fleet_windows_job(mut child: Child) -> std::io::Result<(Child, FleetWindowsJob)> {
1447 match FleetWindowsJob::attach_to_child(&child) {
1448 Ok(job) => Ok((child, job)),
1449 Err(err) => {
1450 let _ = child.kill();
1451 let _ = child.wait();
1452 Err(err)
1453 }
1454 }
1455 }
1456
1457 #[cfg(windows)]
1458 fn windows_io_error(error: windows::core::Error) -> std::io::Error {
1459 std::io::Error::other(error)
1460 }
1461
1462 fn filtered_env(
1463 env: &BTreeMap<String, String>,
1464 allowlist: &BTreeSet<String>,
1465 ) -> FleetHostResult<BTreeMap<String, String>> {
1466 validate_env_allowlist(allowlist)?;
1467 Ok(env
1468 .iter()
1469 .filter(|(key, _)| allowlist.contains(*key))
1470 .map(|(key, value)| (key.clone(), value.clone()))
1471 .collect())
1472 }
1473
1474 /// Characters OpenSSH itself refuses in a command-line user or host name
1475 /// because an `ssh_config` `%h`/`%r` expansion can hand them to a shell.
1476 const SSH_DESTINATION_METACHARACTERS: &str = "'`\"$\\;&<>|(){}";
1477
1478 fn validate_ssh_destination_part(label: &str, value: &str) -> FleetHostResult<()> {
1479 let allowed = |ch: char| {
1480 if label == "host" {
1481 // Host names, IPv4/IPv6 literals (with a zone id) and ssh_config
1482 // aliases. `@` would move the user/host split.
1483 ch.is_ascii_alphanumeric() || matches!(ch, '.' | '-' | '_' | ':' | '%')
1484 } else {
1485 ch.is_ascii_graphic() && !SSH_DESTINATION_METACHARACTERS.contains(ch)
1486 }
1487 };
1488 if value.starts_with('-') || !value.chars().all(allowed) {
1489 return Err(FleetHostError::configuration(format!(
1490 "SSH Fleet {label} must not start with '-' or contain whitespace, control, non-ASCII or shell characters"
1491 )));
1492 }
1493 Ok(())
1494 }
1495
1496 fn validate_env_allowlist(allowlist: &BTreeSet<String>) -> FleetHostResult<()> {
1497 for key in allowlist {
1498 if !is_safe_env_key(key) {
1499 return Err(FleetHostError::configuration(format!(
1500 "Fleet host env allowlist key {key} looks secret-bearing; pass secrets through config providers, not worker argv/env"
1501 )));
1502 }
1503 }
1504 Ok(())
1505 }
1506
1507 fn is_safe_env_key(key: &str) -> bool {
1508 let upper = key.to_ascii_uppercase();
1509 ![
1510 "SECRET",
1511 "TOKEN",
1512 "PASSWORD",
1513 "PASSWD",
1514 "API_KEY",
1515 "CREDENTIAL",
1516 "PRIVATE_KEY",
1517 ]
1518 .iter()
1519 .any(|needle| upper.contains(needle))
1520 }
1521
1522 fn ssh_client_env() -> BTreeMap<String, String> {
1523 ["HOME", "PATH", "SSH_AUTH_SOCK"]
1524 .into_iter()
1525 .filter_map(|key| {
1526 std::env::var(key)
1527 .ok()
1528 .map(|value| (key.to_string(), value))
1529 })
1530 .collect()
1531 }
1532
1533 fn process_base_env() -> BTreeMap<String, String> {
1534 base_env_from(std::env::vars_os())
1535 }
1536
1537 /// The worker's inherited environment: the repository's scrubbed child-env
1538 /// allowlist (`child_env`), so proxy routing, CA bundles, temp directories,
1539 /// Windows system/profile roots, locale and toolchain paths survive
1540 /// `env_clear()` while provider keys and other secret-shaped names do not.
1541 /// Proxy URLs keep their credentials for the worker process itself; see
1542 /// [`crate::child_env::sanitized_runtime_env_from`] for that policy.
1543 fn base_env_from<K, V>(parent: impl IntoIterator<Item = (K, V)>) -> BTreeMap<String, String>
1544 where
1545 K: AsRef<std::ffi::OsStr>,
1546 V: AsRef<std::ffi::OsStr>,
1547 {
1548 let mut env: BTreeMap<String, String> = crate::child_env::sanitized_runtime_env_from(parent)
1549 .into_iter()
1550 .filter_map(|(key, value)| Some((key.into_string().ok()?, value.into_string().ok()?)))
1551 .collect();
1552 force_worker_telemetry_off(&mut env);
1553 env
1554 }
1555
1556 /// Fleet workers are an implementation detail of the parent session, not
1557 /// sessions of their own — the parent already accounts for the dispatch.
1558 ///
1559 /// The spawn path `env_clear()`s and rebuilds from [`process_base_env`], so an
1560 /// operator's opt-out would otherwise never reach a worker at all. Hard-off
1561 /// here means a worker can never emit, can never inherit an ambient "on", and
1562 /// can never write telemetry state into the operator's home.
1563 fn force_worker_telemetry_off(env: &mut BTreeMap<String, String>) {
1564 env.insert("CODEWHALE_TELEMETRY".to_string(), "false".to_string());
1565 env.insert("DEEPSEEK_TELEMETRY".to_string(), "false".to_string());
1566 }
1567
1568 /// Build the complete environment a worker is spawned with.
1569 ///
1570 /// The caller's allowlisted entries are merged over the process base, then
1571 /// telemetry is forced off again: an allowlist that happens to name
1572 /// `CODEWHALE_TELEMETRY` must not be able to switch a worker back on.
1573 fn worker_env(
1574 request_env: &BTreeMap<String, String>,
1575 allowlist: &BTreeSet<String>,
1576 ) -> FleetHostResult<BTreeMap<String, String>> {
1577 let mut env = process_base_env();
1578 env.extend(filtered_env(request_env, allowlist)?);
1579 force_worker_telemetry_off(&mut env);
1580 Ok(env)
1581 }
1582
1583 fn shell_quote(value: &str) -> String {
1584 if value.is_empty() {
1585 return "''".to_string();
1586 }
1587 format!("'{}'", value.replace('\'', "'\\''"))
1588 }
1589
1590 fn validate_worker_id(worker_id: &str) -> FleetHostResult<()> {
1591 if worker_id.trim().is_empty() {
1592 return Err(FleetHostError::configuration("worker id cannot be empty"));
1593 }
1594 Ok(())
1595 }
1596
1597 fn safe_path_segment(value: &str) -> String {
1598 value
1599 .chars()
1600 .map(|ch| {
1601 if ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.') {
1602 ch
1603 } else {
1604 '_'
1605 }
1606 })
1607 .collect()
1608 }
1609
1610 #[cfg(test)]
1611 mod tests {
1612 use super::*;
1613
1614 /// Worker output is opened through the pinned no-follow writer: a link in
1615 /// `.codewhale` or at the log's name is refused, and a plain log is
1616 /// truncated and readable back through the same confinement.
1617 #[cfg(unix)]
1618 #[test]
1619 fn worker_logs_are_never_opened_through_links() {
1620 use std::io::Write as _;
1621 use std::os::unix::fs::symlink;
1622 let workspace = tempfile::tempdir().unwrap();
1623 let outside = tempfile::tempdir().unwrap();
1624 let log = workspace
1625 .path()
1626 .join(".codewhale")
1627 .join("fleet-host")
1628 .join("local")
1629 .join("w1.log");
1630
1631 let mut file = open_worker_log(workspace.path(), &log).expect("plain log opens");
1632 file.write_all(b"first run output").unwrap();
1633 drop(file);
1634 assert_eq!(
1635 read_bounded_log(workspace.path(), &log, 1024).unwrap(),
1636 "first run output"
1637 );
1638 drop(open_worker_log(workspace.path(), &log).expect("reopen truncates"));
1639 assert_eq!(read_bounded_log(workspace.path(), &log, 1024).unwrap(), "");
1640
1641 // A linked directory and a linked file are both refused.
1642 let linked_root = tempfile::tempdir().unwrap();
1643 std::fs::create_dir_all(linked_root.path().join(".codewhale")).unwrap();
1644 symlink(
1645 outside.path(),
1646 linked_root.path().join(".codewhale").join("fleet-host"),
1647 )
1648 .unwrap();
1649 let linked_dir_log = linked_root
1650 .path()
1651 .join(".codewhale")
1652 .join("fleet-host")
1653 .join("local")
1654 .join("w1.log");
1655 assert!(open_worker_log(linked_root.path(), &linked_dir_log).is_err());
1656
1657 let secret = outside.path().join("secret.txt");
1658 std::fs::write(&secret, "do not read").unwrap();
1659 let parent = log.parent().unwrap();
1660 symlink(&secret, parent.join("w2.log")).unwrap();
1661 assert!(open_worker_log(workspace.path(), &parent.join("w2.log")).is_err());
1662 assert!(read_bounded_log(workspace.path(), &parent.join("w2.log"), 1024).is_err());
1663 assert_eq!(std::fs::read_to_string(&secret).unwrap(), "do not read");
1664 }
1665 use tempfile::TempDir;
1666
1667 #[cfg(unix)]
1668 fn skip_if_process_table_unavailable() -> bool {
1669 if process_table_inspection_available() {
1670 return false;
1671 }
1672 eprintln!("skipping: process-table inspection unavailable (ps/proc denied or missing)");
1673 true
1674 }
1675
1676 #[cfg(unix)]
1677 #[test]
1678 fn sample_process_memory_reports_nonzero_for_self() {
1679 if skip_if_process_table_unavailable() {
1680 return;
1681 }
1682 // The current test process is alive, so its RSS must sample to Some(>0).
1683 let mb = sample_process_memory_mb(std::process::id());
1684 assert!(
1685 matches!(mb, Some(v) if v > 0),
1686 "expected Some(>0) MB for the live self process, got {mb:?}"
1687 );
1688 }
1689
1690 #[cfg(unix)]
1691 #[test]
1692 fn sample_process_memory_is_none_for_dead_pid() {
1693 // Use a PID beyond every mainstream kernel's default pid ceiling
1694 // (Linux pid_max default 4M/32k, macOS ~99998, BSDs 99999): PID 0 is
1695 // kernel_task on macOS and semantically special to `ps -p`, so it is
1696 // not a portable "no such process" probe. When process-table inspection
1697 // is denied the sampler also returns None — same observable result.
1698 assert_eq!(sample_process_memory_mb(999_999_999), None);
1699 }
1700
1701 fn shell_command(script: &str) -> FleetWorkerCommand {
1702 if cfg!(windows) {
1703 FleetWorkerCommand::new("cmd", ["/C", script])
1704 } else {
1705 FleetWorkerCommand::new("sh", ["-c", script])
1706 }
1707 }
1708
1709 #[cfg(unix)]
1710 const DESCENDANT_HELPER_TEST: &str =
1711 "fleet::host::tests::fleet_host_stop_reaps_dispatcher_descendants";
1712
1713 #[cfg(unix)]
1714 fn run_descendant_helper_if_requested() -> bool {
1715 let Ok(mode) = std::env::var("FLEET_DESCENDANT_HELPER") else {
1716 return false;
1717 };
1718 let test_binary = std::env::current_exe().expect("current test binary");
1719 let pid_file = std::env::var("FLEET_DESCENDANT_PID_FILE").expect("helper pid file");
1720 match mode.as_str() {
1721 "dispatcher" | "detached-dispatcher" => {
1722 let mut command = Command::new(&test_binary);
1723 command
1724 .args(["--exact", DESCENDANT_HELPER_TEST, "--nocapture"])
1725 .env("FLEET_DESCENDANT_HELPER", "worker")
1726 .env("FLEET_DESCENDANT_PID_FILE", &pid_file);
1727 if mode == "detached-dispatcher" {
1728 command.spawn().expect("spawn detached dispatcher child");
1729 std::process::exit(0);
1730 }
1731 let status = command.status().expect("spawn dispatcher child");
1732 std::process::exit(status.code().unwrap_or(1));
1733 }
1734 "worker" => {
1735 let mut command = Command::new(&test_binary);
1736 command
1737 .args(["--exact", DESCENDANT_HELPER_TEST, "--nocapture"])
1738 .env("FLEET_DESCENDANT_HELPER", "tool")
1739 .env("FLEET_DESCENDANT_PID_FILE", &pid_file);
1740 // Real shell tools deliberately own a separate process group.
1741 // This makes a root-group-only Fleet stop leak the helper.
1742 command.process_group(0);
1743 let status = command.status().expect("spawn worker tool");
1744 std::process::exit(status.code().unwrap_or(1));
1745 }
1746 "tool" => {
1747 // A shell tool can ignore graceful signals and live in its own
1748 // process group. Fleet's session boundary must still reap it.
1749 unsafe {
1750 libc::signal(libc::SIGINT, libc::SIG_IGN);
1751 libc::signal(libc::SIGTERM, libc::SIG_IGN);
1752 }
1753 std::fs::write(&pid_file, std::process::id().to_string()).expect("write tool pid");
1754 thread::sleep(Duration::from_secs(30));
1755 true
1756 }
1757 other => panic!("unknown descendant helper mode {other}"),
1758 }
1759 }
1760
1761 #[cfg(unix)]
1762 fn start_dispatcher_tree(
1763 adapter: &mut LocalProcessFleetHostAdapter,
1764 tmp: &TempDir,
1765 worker_id: &str,
1766 helper_mode: &str,
1767 ) -> (libc::pid_t, libc::pid_t) {
1768 let pid_file = tmp.path().join(format!("{worker_id}-tool.pid"));
1769 let test_binary = std::env::current_exe().expect("current test binary");
1770 let mut request = FleetWorkerStartRequest::new(
1771 worker_id,
1772 FleetWorkerCommand::new(
1773 test_binary.display().to_string(),
1774 ["--exact", DESCENDANT_HELPER_TEST, "--nocapture"],
1775 ),
1776 );
1777 request.env.insert(
1778 "FLEET_DESCENDANT_HELPER".to_string(),
1779 helper_mode.to_string(),
1780 );
1781 request.env.insert(
1782 "FLEET_DESCENDANT_PID_FILE".to_string(),
1783 pid_file.display().to_string(),
1784 );
1785 request.env_allowlist = BTreeSet::from([
1786 "FLEET_DESCENDANT_HELPER".to_string(),
1787 "FLEET_DESCENDANT_PID_FILE".to_string(),
1788 ]);
1789
1790 let handle = adapter.start_worker(request).expect("start dispatcher");
1791 let root_pid = handle.pid.expect("dispatcher pid") as libc::pid_t;
1792 let tool_pid = wait_for_valid_pid_file(&pid_file, Duration::from_secs(5));
1793 if helper_mode != "detached-dispatcher" {
1794 assert!(unix_pid_is_running(root_pid));
1795 }
1796 assert!(unix_pid_is_running(tool_pid));
1797 (root_pid, tool_pid)
1798 }
1799
1800 #[cfg(unix)]
1801 fn wait_for_host_state(
1802 adapter: &mut LocalProcessFleetHostAdapter,
1803 worker_id: &str,
1804 expected: FleetHostWorkerState,
1805 timeout: Duration,
1806 ) -> FleetHostWorkerStatus {
1807 let deadline = Instant::now() + timeout;
1808 loop {
1809 let status = adapter.read_status(worker_id).expect("worker status");
1810 if status.state == expected || Instant::now() >= deadline {
1811 return status;
1812 }
1813 thread::sleep(Duration::from_millis(25));
1814 }
1815 }
1816
1817 #[cfg(unix)]
1818 fn wait_for_valid_pid_file(pid_file: &Path, timeout: Duration) -> libc::pid_t {
1819 let deadline = Instant::now() + timeout;
1820 let mut last_observation = "file not created".to_string();
1821 loop {
1822 match std::fs::read_to_string(pid_file) {
1823 Ok(contents) => {
1824 let trimmed = contents.trim();
1825 match trimmed.parse::<libc::pid_t>() {
1826 Ok(pid) if pid > 0 => return pid,
1827 _ => last_observation = format!("invalid contents {trimmed:?}"),
1828 }
1829 }
1830 Err(err) if err.kind() == std::io::ErrorKind::NotFound => {}
1831 Err(err) => last_observation = format!("read failed: {err}"),
1832 }
1833 if Instant::now() >= deadline {
1834 panic!(
1835 "separate-group tool never published a valid PID to {} ({last_observation})",
1836 pid_file.display()
1837 );
1838 }
1839 thread::sleep(Duration::from_millis(25));
1840 }
1841 }
1842
1843 #[cfg(unix)]
1844 #[test]
1845 fn pid_file_wait_ignores_created_but_incomplete_file() {
1846 let tmp = TempDir::new().unwrap();
1847 let pid_file = tmp.path().join("worker.pid");
1848 std::fs::write(&pid_file, "pid=").unwrap();
1849 let expected_pid = std::process::id() as libc::pid_t;
1850 let writer_path = pid_file.clone();
1851 let writer = thread::spawn(move || {
1852 thread::sleep(Duration::from_millis(50));
1853 std::fs::write(writer_path, expected_pid.to_string()).unwrap();
1854 });
1855
1856 assert_eq!(
1857 wait_for_valid_pid_file(&pid_file, Duration::from_secs(1)),
1858 expected_pid
1859 );
1860 writer.join().unwrap();
1861 }
1862
1863 #[cfg(unix)]
1864 #[test]
1865 fn unix_pid_running_treats_zombie_as_exited() {
1866 if skip_if_process_table_unavailable() {
1867 return;
1868 }
1869 let mut child = Command::new("sh")
1870 .args(["-c", "exit 0"])
1871 .spawn()
1872 .expect("spawn short-lived child");
1873 let pid = child.id() as libc::pid_t;
1874 let deadline = Instant::now() + Duration::from_secs(2);
1875 let saw_zombie = loop {
1876 let state = match Command::new("ps")
1877 .args(["-o", "stat=", "-p", &pid.to_string()])
1878 .output()
1879 {
1880 Ok(output) => output,
1881 Err(err) if is_permission_denied(&err) => {
1882 mark_process_table_unavailable();
1883 child.wait().ok();
1884 eprintln!("skipping: ps denied while inspecting zombie state");
1885 return;
1886 }
1887 Err(err) => panic!("inspect child state: {err}"),
1888 };
1889 let is_zombie = String::from_utf8_lossy(&state.stdout)
1890 .split_whitespace()
1891 .next()
1892 .is_some_and(|state| state.starts_with('Z'));
1893 if is_zombie {
1894 break true;
1895 }
1896 if Instant::now() >= deadline {
1897 break false;
1898 }
1899 thread::sleep(Duration::from_millis(10));
1900 };
1901
1902 let reported_running = unix_pid_is_running(pid);
1903 child.wait().expect("reap zombie child");
1904 assert!(saw_zombie, "child never became a zombie");
1905 assert!(!reported_running, "zombie was reported as running");
1906 }
1907
1908 fn wait_for_log(
1909 adapter: &LocalProcessFleetHostAdapter,
1910 worker_id: &str,
1911 needle: &str,
1912 ) -> String {
1913 let deadline = Instant::now() + Duration::from_secs(3);
1914 loop {
1915 let logs = adapter.read_logs(worker_id, 4096).unwrap();
1916 if logs.contains(needle) || Instant::now() > deadline {
1917 return logs;
1918 }
1919 thread::sleep(Duration::from_millis(25));
1920 }
1921 }
1922
1923 #[test]
1924 fn fleet_host_local_adapter_starts_reads_bounded_logs_and_stops() {
1925 let tmp = TempDir::new().unwrap();
1926 let mut adapter = LocalProcessFleetHostAdapter::new(tmp.path());
1927 let script = if cfg!(windows) {
1928 "echo 0123456789abcdef& ping -n 30 127.0.0.1 >NUL"
1929 } else {
1930 "printf 0123456789abcdef; sleep 30"
1931 };
1932 let mut request = FleetWorkerStartRequest::new("local-1", shell_command(script));
1933 let line_ending_bytes = if cfg!(windows) { 2 } else { 0 };
1934 request.log_limit_bytes = 16 + line_ending_bytes;
1935
1936 let handle = adapter.start_worker(request).unwrap();
1937 #[cfg(unix)]
1938 let direct_pid = handle.pid.expect("local worker pid");
1939 assert_eq!(handle.host_kind, FleetHostKind::LocalProcess);
1940 assert!(handle.pid.is_some());
1941 let status = adapter.read_status("local-1").unwrap();
1942 assert_eq!(status.state, FleetHostWorkerState::Running);
1943
1944 let logs = wait_for_log(&adapter, "local-1", "abcdef");
1945 let logs = logs.trim_end_matches(&['\r', '\n'][..]);
1946 assert!(logs.ends_with("0123456789abcdef"), "{logs:?}");
1947 let bounded = adapter.read_logs("local-1", 6 + line_ending_bytes).unwrap();
1948 let bounded = bounded.trim_end_matches(&['\r', '\n'][..]);
1949 assert!(bounded.ends_with("abcdef"), "{bounded:?}");
1950
1951 let status = adapter.stop_worker("local-1").unwrap();
1952 assert_eq!(status.state, FleetHostWorkerState::Stopped);
1953 #[cfg(unix)]
1954 assert!(
1955 wait_for_unix_pid_exit(direct_pid as libc::pid_t, Duration::from_secs(1)),
1956 "stopped direct worker was not reaped"
1957 );
1958 adapter.cleanup_worker("local-1").unwrap();
1959 assert_eq!(
1960 adapter.read_status("local-1").unwrap_err().kind,
1961 FleetHostErrorKind::Terminal
1962 );
1963 }
1964
1965 #[cfg(unix)]
1966 #[test]
1967 fn fleet_host_stop_reaps_dispatcher_descendants() {
1968 if run_descendant_helper_if_requested() {
1969 return;
1970 }
1971 // Full-session reaping of separate process-group tools requires a
1972 // process-table walk; without it production still signals the known
1973 // session leader and these assertions cannot be proven.
1974 if skip_if_process_table_unavailable() {
1975 return;
1976 }
1977
1978 let tmp = TempDir::new().unwrap();
1979 let mut adapter = LocalProcessFleetHostAdapter::new(tmp.path());
1980 let (root_pid, tool_pid) =
1981 start_dispatcher_tree(&mut adapter, &tmp, "dispatcher-tree", "dispatcher");
1982
1983 let status = adapter
1984 .stop_worker("dispatcher-tree")
1985 .expect("stop complete worker tree");
1986
1987 assert_eq!(status.state, FleetHostWorkerState::Stopped);
1988 assert!(
1989 wait_for_unix_pid_exit(root_pid, Duration::from_secs(1)),
1990 "dispatcher survived stop"
1991 );
1992 assert!(
1993 wait_for_unix_pid_exit(tool_pid, Duration::from_secs(1)),
1994 "separate-process-group tool survived stop"
1995 );
1996 }
1997
1998 #[cfg(unix)]
1999 #[test]
2000 fn fleet_host_interrupt_reaps_dispatcher_descendants() {
2001 if skip_if_process_table_unavailable() {
2002 return;
2003 }
2004 let tmp = TempDir::new().unwrap();
2005 let mut adapter = LocalProcessFleetHostAdapter::new(tmp.path());
2006 let (root_pid, tool_pid) =
2007 start_dispatcher_tree(&mut adapter, &tmp, "interrupt-tree", "dispatcher");
2008
2009 let status = adapter
2010 .interrupt_worker("interrupt-tree")
2011 .expect("interrupt complete worker session");
2012
2013 assert_ne!(status.state, FleetHostWorkerState::Running);
2014 assert!(
2015 wait_for_unix_pid_exit(root_pid, Duration::from_secs(1)),
2016 "dispatcher survived interrupt"
2017 );
2018 assert!(
2019 wait_for_unix_pid_exit(tool_pid, Duration::from_secs(1)),
2020 "separate-process-group tool survived interrupt"
2021 );
2022 }
2023
2024 #[cfg(unix)]
2025 #[test]
2026 fn fleet_host_reports_draining_after_dispatcher_exits_with_live_descendant() {
2027 if skip_if_process_table_unavailable() {
2028 return;
2029 }
2030 let tmp = TempDir::new().unwrap();
2031 let mut adapter = LocalProcessFleetHostAdapter::new(tmp.path());
2032 let (root_pid, tool_pid) = start_dispatcher_tree(
2033 &mut adapter,
2034 &tmp,
2035 "draining-dispatcher-tree",
2036 "detached-dispatcher",
2037 );
2038
2039 assert!(
2040 wait_for_unix_pid_exit(root_pid, Duration::from_secs(3)),
2041 "dispatcher did not exit"
2042 );
2043 let status = wait_for_host_state(
2044 &mut adapter,
2045 "draining-dispatcher-tree",
2046 FleetHostWorkerState::Draining,
2047 Duration::from_secs(3),
2048 );
2049 assert_eq!(status.state, FleetHostWorkerState::Draining);
2050 assert!(
2051 unix_pid_is_running(tool_pid),
2052 "descendant exited before draining check"
2053 );
2054
2055 let stopped = adapter
2056 .stop_worker("draining-dispatcher-tree")
2057 .expect("stop draining worker tree");
2058 assert_eq!(stopped.state, FleetHostWorkerState::Stopped);
2059 assert!(
2060 wait_for_unix_pid_exit(tool_pid, Duration::from_secs(1)),
2061 "draining descendant survived bounded stop"
2062 );
2063 }
2064
2065 #[cfg(unix)]
2066 #[test]
2067 fn fleet_host_cleanup_reaps_session_after_dispatcher_exits() {
2068 if skip_if_process_table_unavailable() {
2069 return;
2070 }
2071 let tmp = TempDir::new().unwrap();
2072 let mut adapter = LocalProcessFleetHostAdapter::new(tmp.path());
2073 let (root_pid, tool_pid) = start_dispatcher_tree(
2074 &mut adapter,
2075 &tmp,
2076 "exited-dispatcher-tree",
2077 "detached-dispatcher",
2078 );
2079 let deadline = Instant::now() + Duration::from_secs(3);
2080 loop {
2081 let status = adapter.read_status("exited-dispatcher-tree").unwrap();
2082 if status.state != FleetHostWorkerState::Running || Instant::now() >= deadline {
2083 assert_ne!(status.state, FleetHostWorkerState::Running);
2084 break;
2085 }
2086 thread::sleep(Duration::from_millis(25));
2087 }
2088 assert!(
2089 wait_for_unix_pid_exit(root_pid, Duration::from_secs(1)),
2090 "dispatcher should have exited"
2091 );
2092 assert!(
2093 unix_pid_is_running(tool_pid),
2094 "delegated tool exited too early"
2095 );
2096
2097 adapter
2098 .cleanup_worker("exited-dispatcher-tree")
2099 .expect("clean up surviving dispatcher session");
2100
2101 assert!(
2102 wait_for_unix_pid_exit(tool_pid, Duration::from_secs(1)),
2103 "tool survived after its dispatcher exited"
2104 );
2105 assert_eq!(
2106 adapter
2107 .read_status("exited-dispatcher-tree")
2108 .unwrap_err()
2109 .kind,
2110 FleetHostErrorKind::Terminal
2111 );
2112 }
2113
2114 #[test]
2115 fn fleet_host_local_adapter_restarts_worker_with_same_request() {
2116 let tmp = TempDir::new().unwrap();
2117 let mut adapter = LocalProcessFleetHostAdapter::new(tmp.path());
2118 let script = if cfg!(windows) {
2119 "echo restart-ready & ping -n 30 127.0.0.1 >NUL"
2120 } else {
2121 "printf restart-ready; sleep 30"
2122 };
2123 let request = FleetWorkerStartRequest::new("local-restart", shell_command(script));
2124 let first = adapter.start_worker(request).unwrap();
2125 let restarted = adapter.restart_worker("local-restart").unwrap();
2126
2127 assert_eq!(restarted.worker_id, first.worker_id);
2128 assert_eq!(restarted.host_kind, FleetHostKind::LocalProcess);
2129 assert_ne!(restarted.pid, first.pid);
2130 let logs = wait_for_log(&adapter, "local-restart", "restart-ready");
2131 assert!(logs.contains("restart-ready"));
2132 adapter.stop_worker("local-restart").unwrap();
2133 }
2134
2135 /// A host whose stop outcome is scripted: `None` fails the stop, `Some`
2136 /// reports that state. Only the restart path is exercised.
2137 struct ScriptedStopHost {
2138 stop_state: Option<FleetHostWorkerState>,
2139 owned: bool,
2140 spawned: usize,
2141 }
2142
2143 impl FleetHostAdapter for ScriptedStopHost {
2144 fn start_worker(
2145 &mut self,
2146 request: FleetWorkerStartRequest,
2147 ) -> FleetHostResult<FleetWorkerHandle> {
2148 self.spawned += 1;
2149 self.owned = true;
2150 Ok(FleetWorkerHandle {
2151 worker_id: request.worker_id,
2152 host_kind: FleetHostKind::LocalProcess,
2153 pid: None,
2154 log_path: PathBuf::new(),
2155 })
2156 }
2157 fn read_status(&mut self, _: &str) -> FleetHostResult<FleetHostWorkerStatus> {
2158 unreachable!("restart must not consult status outside stop")
2159 }
2160 fn read_logs(&self, _: &str, _: usize) -> FleetHostResult<String> {
2161 unreachable!()
2162 }
2163 fn interrupt_worker(&mut self, _: &str) -> FleetHostResult<FleetHostWorkerStatus> {
2164 unreachable!()
2165 }
2166 fn restart_worker(&mut self, worker_id: &str) -> FleetHostResult<FleetWorkerHandle> {
2167 restart_after_confirmed_stop(self, worker_id, |host| {
2168 host.owned = false;
2169 host.start_worker(FleetWorkerStartRequest::new(
2170 worker_id,
2171 shell_command("true"),
2172 ))
2173 })
2174 }
2175 fn stop_worker(&mut self, worker_id: &str) -> FleetHostResult<FleetHostWorkerStatus> {
2176 let Some(state) = self.stop_state else {
2177 return Err(FleetHostError::retryable(
2178 "Fleet session still has a live tracked leader after SIGKILL",
2179 ));
2180 };
2181 Ok(FleetHostWorkerStatus {
2182 worker_id: worker_id.to_string(),
2183 state,
2184 pid: Some(4242),
2185 exit_code: None,
2186 memory_mb: None,
2187 retryable: false,
2188 })
2189 }
2190 fn cleanup_worker(&mut self, _: &str) -> FleetHostResult<()> {
2191 unreachable!()
2192 }
2193 }
2194
2195 #[test]
2196 fn fleet_host_restart_keeps_the_old_worker_when_stop_is_unconfirmed() {
2197 for stop_state in [None, Some(FleetHostWorkerState::Draining)] {
2198 let mut host = ScriptedStopHost {
2199 stop_state,
2200 owned: true,
2201 spawned: 0,
2202 };
2203 let err = host
2204 .restart_worker("overlap")
2205 .expect_err("an unconfirmed stop must refuse the restart");
2206 assert!(err.message.contains("refused"), "{}", err.message);
2207 assert_eq!(
2208 host.spawned, 0,
2209 "no replacement may start beside a surviving worker ({stop_state:?})"
2210 );
2211 assert!(host.owned, "the old worker's handle is kept");
2212 }
2213
2214 let mut host = ScriptedStopHost {
2215 stop_state: Some(FleetHostWorkerState::Stopped),
2216 owned: true,
2217 spawned: 0,
2218 };
2219 host.restart_worker("overlap")
2220 .expect("a confirmed stop restarts");
2221 assert_eq!(host.spawned, 1);
2222 }
2223
2224 #[cfg(unix)]
2225 #[test]
2226 fn fleet_host_local_adapter_reports_running_worker_memory_usage() {
2227 if skip_if_process_table_unavailable() {
2228 return;
2229 }
2230 let tmp = TempDir::new().unwrap();
2231 let mut adapter = LocalProcessFleetHostAdapter::new(tmp.path());
2232 let request =
2233 FleetWorkerStartRequest::new("local-memory", shell_command("printf ready; sleep 30"));
2234
2235 adapter.start_worker(request).unwrap();
2236 let _ = wait_for_log(&adapter, "local-memory", "ready");
2237
2238 let status = adapter.read_status("local-memory").unwrap();
2239
2240 assert_eq!(status.state, FleetHostWorkerState::Running);
2241 assert!(
2242 status.memory_mb.is_some_and(|memory_mb| memory_mb > 0),
2243 "running local worker status should include RSS memory_mb, got {status:?}"
2244 );
2245
2246 adapter.stop_worker("local-memory").unwrap();
2247 }
2248
2249 #[cfg(unix)]
2250 #[test]
2251 fn fleet_host_stop_signals_known_leader_without_process_table() {
2252 // Even when full session census is unavailable, stop must still reach
2253 // the tracked session leader/dispatcher pid.
2254 let tmp = TempDir::new().unwrap();
2255 let mut adapter = LocalProcessFleetHostAdapter::new(tmp.path());
2256 let request = FleetWorkerStartRequest::new(
2257 "known-leader-stop",
2258 shell_command("printf ready; sleep 30"),
2259 );
2260 let handle = adapter.start_worker(request).unwrap();
2261 let pid = handle.pid.expect("pid") as libc::pid_t;
2262 let _ = wait_for_log(&adapter, "known-leader-stop", "ready");
2263
2264 let status = adapter
2265 .stop_worker("known-leader-stop")
2266 .expect("stop known leader without process table");
2267 assert_eq!(status.state, FleetHostWorkerState::Stopped);
2268 assert!(
2269 wait_for_unix_pid_exit(pid, Duration::from_secs(2)),
2270 "known session leader survived stop without process-table census"
2271 );
2272 adapter.cleanup_worker("known-leader-stop").unwrap();
2273 }
2274
2275 #[test]
2276 fn fleet_host_ssh_kind_does_not_report_local_process_memory() {
2277 let tmp = TempDir::new().unwrap();
2278 let mut adapter = LocalProcessFleetHostAdapter::new(tmp.path());
2279 let script = if cfg!(windows) {
2280 "echo ready & ping -n 30 127.0.0.1 >NUL"
2281 } else {
2282 "printf ready; sleep 30"
2283 };
2284 let request = FleetWorkerStartRequest::new("ssh-memory", shell_command(script));
2285
2286 adapter
2287 .start_with_kind(request, FleetHostKind::Ssh)
2288 .unwrap();
2289 let _ = wait_for_log(&adapter, "ssh-memory", "ready");
2290
2291 let status = adapter.read_status("ssh-memory").unwrap();
2292
2293 assert_eq!(status.state, FleetHostWorkerState::Running);
2294 assert_eq!(status.memory_mb, None);
2295
2296 adapter.stop_worker("ssh-memory").unwrap();
2297 }
2298
2299 #[test]
2300 fn fleet_host_rejects_secret_like_env_allowlist_keys() {
2301 let mut env = BTreeMap::new();
2302 env.insert("DEEPSEEK_API_KEY".to_string(), "secret".to_string());
2303 let allowlist = BTreeSet::from(["DEEPSEEK_API_KEY".to_string()]);
2304
2305 let err = filtered_env(&env, &allowlist).unwrap_err();
2306
2307 assert_eq!(err.kind, FleetHostErrorKind::Configuration);
2308 assert!(err.message.contains("looks secret-bearing"));
2309 }
2310
2311 #[test]
2312 fn fleet_host_ssh_command_uses_sendenv_without_argv_secret_values() {
2313 let tmp = TempDir::new().unwrap();
2314 let mut config = SshFleetHostConfig::new("builder.example.test", "/srv/codewhale");
2315 config.user = Some("fleet".to_string());
2316 config.port = Some(2222);
2317 config.identity = Some(PathBuf::from("/tmp/fleet_id"));
2318 config.codewhale_binary = "/usr/local/bin/codewhale".to_string();
2319 config.env_allowlist = BTreeSet::from(["FLEET_PROFILE".to_string()]);
2320 let adapter = SshFleetHostAdapter::new(tmp.path(), config).unwrap();
2321 let mut request = FleetWorkerStartRequest::new(
2322 "ssh-1",
2323 FleetWorkerCommand::new("codewhale", ["fleet-worker", "noop"]),
2324 );
2325 request.env.insert(
2326 "FLEET_PROFILE".to_string(),
2327 "super-secret-profile-value".to_string(),
2328 );
2329
2330 let command = adapter.build_ssh_command(&request).unwrap();
2331 let argv = command.args.join(" ");
2332
2333 assert_eq!(command.program, "ssh");
2334 assert!(argv.contains("BatchMode=yes"));
2335 assert!(argv.contains("SendEnv=FLEET_PROFILE"));
2336 assert!(argv.contains("fleet@builder.example.test"));
2337 assert!(argv.contains("/usr/local/bin/codewhale"));
2338 assert!(argv.contains("fleet-worker"));
2339 assert!(!argv.contains("super-secret-profile-value"));
2340 // Option parsing ends right before the destination.
2341 let target = command
2342 .args
2343 .iter()
2344 .position(|arg| arg == "fleet@builder.example.test")
2345 .expect("destination argument");
2346 assert_eq!(command.args[target - 1], "--");
2347 }
2348
2349 #[test]
2350 fn fleet_host_ssh_refuses_option_like_or_malformed_destination() {
2351 let tmp = TempDir::new().unwrap();
2352 for host in [
2353 "-oProxyCommand=true",
2354 "builder example.test",
2355 "builder\nexample.test",
2356 "builder\0example.test",
2357 "builder;true",
2358 "builder$(true)",
2359 "builder`true`",
2360 "fleet@builder.example.test",
2361 "builder\u{202e}.example.test",
2362 ] {
2363 let config = SshFleetHostConfig::new(host, "/srv/codewhale");
2364 let err = SshFleetHostAdapter::new(tmp.path(), config)
2365 .expect_err("malformed SSH host must be refused");
2366 assert_eq!(err.kind, FleetHostErrorKind::Configuration, "{host:?}");
2367 assert!(err.message.contains("SSH Fleet host"), "{}", err.message);
2368 }
2369 for user in [
2370 "-oProxyCommand=true",
2371 "fleet user",
2372 "fleet\tuser",
2373 "fleet\0user",
2374 "fleet;true",
2375 "fleet$(true)",
2376 "fleet|true",
2377 "fl\u{e9}et",
2378 ] {
2379 let mut config = SshFleetHostConfig::new("builder.example.test", "/srv/codewhale");
2380 config.user = Some(user.to_string());
2381 let err = SshFleetHostAdapter::new(tmp.path(), config)
2382 .expect_err("malformed SSH user must be refused");
2383 assert_eq!(err.kind, FleetHostErrorKind::Configuration, "{user:?}");
2384 assert!(err.message.contains("SSH Fleet user"), "{}", err.message);
2385 }
2386 let spec = FleetHostSpec::Ssh {
2387 host: "-oProxyCommand=true".to_string(),
2388 port: None,
2389 user: None,
2390 identity: None,
2391 known_hosts: None,
2392 host_key_fingerprint: None,
2393 working_directory: Some(PathBuf::from("/srv/codewhale")),
2394 env_allowlist: Vec::new(),
2395 codewhale_binary: Some("codewhale".to_string()),
2396 };
2397 assert!(SshFleetHostConfig::from_host_spec(&spec).is_err());
2398
2399 // Ordinary destinations stay accepted, including a directory-style
2400 // user name and IPv6 literals with a zone id.
2401 for (user, host) in [
2402 (Some("fleet"), "builder.example.test"),
2403 (Some("alice@corp.example"), "10.0.0.7"),
2404 (None, "fe80::1%en0"),
2405 (Some("ci_bot-2"), "build_box-01"),
2406 ] {
2407 let mut config = SshFleetHostConfig::new(host, "/srv/codewhale");
2408 config.user = user.map(str::to_string);
2409 SshFleetHostAdapter::new(tmp.path(), config)
2410 .unwrap_or_else(|err| panic!("{user:?}@{host} must be accepted: {err:?}"));
2411 }
2412 }
2413
2414 #[test]
2415 fn worker_base_env_uses_the_child_env_allowlist_and_keeps_proxy_route() {
2416 let parent = [
2417 ("HTTPS_PROXY", "http://fleet:pass@proxy.example.test:8080"),
2418 ("no_proxy", "localhost"),
2419 ("SSL_CERT_FILE", "/etc/ssl/corp.pem"),
2420 ("CURL_CA_BUNDLE", "/etc/ssl/corp.pem"),
2421 ("REQUESTS_CA_BUNDLE", "/etc/ssl/corp.pem"),
2422 ("NODE_EXTRA_CA_CERTS", "/etc/ssl/corp.pem"),
2423 ("PATHEXT", ".COM;.EXE;.BAT"),
2424 ("WINDIR", "C:\\Windows"),
2425 ("ProgramFiles", "C:\\Program Files"),
2426 ("USERPROFILE", "C:\\Users\\fleet"),
2427 ("TEMP", "C:\\Temp"),
2428 ("USER", "fleet"),
2429 ("TERM", "xterm-256color"),
2430 ("CARGO_HOME", "/home/fleet/.cargo"),
2431 ("CARGO_TARGET_DIR", "/tmp/target"),
2432 ("DEEPSEEK_API_KEY", "secret"),
2433 ("GITHUB_TOKEN", "secret"),
2434 ("AWS_SECRET_ACCESS_KEY", "secret"),
2435 ("CARGO_REGISTRY_TOKEN", "secret"),
2436 ("DATABASE_URL", "postgres://u:secret@db/app"),
2437 ("CODEWHALE_TELEMETRY", "true"),
2438 ];
2439 let env = base_env_from(parent);
2440 for (key, value) in &parent[..15] {
2441 assert_eq!(
2442 env.get(*key).map(String::as_str),
2443 Some(*value),
2444 "{key} must reach the worker unchanged"
2445 );
2446 }
2447 for key in [
2448 "DEEPSEEK_API_KEY",
2449 "GITHUB_TOKEN",
2450 "AWS_SECRET_ACCESS_KEY",
2451 "CARGO_REGISTRY_TOKEN",
2452 "DATABASE_URL",
2453 ] {
2454 assert!(!env.contains_key(key), "{key} must not reach the worker");
2455 }
2456 assert_eq!(
2457 env.get("CODEWHALE_TELEMETRY").map(String::as_str),
2458 Some("false")
2459 );
2460 }
2461
2462 #[test]
2463 fn runtime_surface_hardening_ssh_requires_known_host_verification() {
2464 let tmp = TempDir::new().unwrap();
2465 let request = FleetWorkerStartRequest::new(
2466 "ssh-1",
2467 FleetWorkerCommand::new("codewhale", ["fleet-worker"]),
2468 );
2469 for known_hosts in [None, Some(PathBuf::from("/tmp/fleet keys/known_hosts"))] {
2470 let mut config = SshFleetHostConfig::new("builder.example.test", "/srv/codewhale");
2471 config.known_hosts = known_hosts.clone();
2472 let adapter = SshFleetHostAdapter::new(tmp.path(), config).unwrap();
2473 let command = adapter.build_ssh_command(&request).unwrap();
2474 assert!(
2475 command
2476 .args
2477 .contains(&"StrictHostKeyChecking=yes".to_string())
2478 );
2479 if let Some(path) = known_hosts {
2480 assert!(
2481 command
2482 .args
2483 .contains(&format!("UserKnownHostsFile=\"{}\"", path.display()))
2484 );
2485 assert!(
2486 command
2487 .args
2488 .contains(&"GlobalKnownHostsFile=none".to_string())
2489 );
2490 let ignore = command
2491 .args
2492 .iter()
2493 .position(|arg| arg == "IgnoreUnknown=KnownHostsCommand")
2494 .expect("OpenSSH 7.x must ignore the newer KnownHostsCommand option");
2495 let command_option = command
2496 .args
2497 .iter()
2498 .position(|arg| arg == "KnownHostsCommand=none")
2499 .expect("disable additional known-host sources on newer clients");
2500 assert!(
2501 ignore < command_option,
2502 "IgnoreUnknown applies only to later options"
2503 );
2504 assert!(command.args.contains(&"VerifyHostKeyDNS=no".to_string()));
2505 }
2506 }
2507 let mut config = SshFleetHostConfig::new("builder.example.test", "/srv/codewhale");
2508 config.host_key_fingerprint = Some("SHA256:configured-pin".to_string());
2509 let err = SshFleetHostAdapter::new(tmp.path(), config).unwrap_err();
2510 assert_eq!(err.kind, FleetHostErrorKind::Configuration);
2511 assert!(err.message.contains("configure known_hosts"));
2512 }
2513
2514 #[test]
2515 fn runtime_surface_review_documented_ssh_host_loads() {
2516 for docs in [
2517 include_str!("../../../../docs/FLEET.md"),
2518 include_str!("../../../../docs/zh_hans/FLEET.md"),
2519 ] {
2520 // Match the actual example, not a translated section heading.
2521 // Windows checkouts may carry CRLF line endings.
2522 let docs = docs.replace("\r\n", "\n");
2523 let example = docs
2524 .split("```json\n")
2525 .skip(1)
2526 .filter_map(|block| {
2527 serde_json::from_str::<serde_json::Value>(block.split("```").next()?).ok()
2528 })
2529 .find(|value| value["id"] == "builder-1")
2530 .expect("documented SSH worker example");
2531 let host: FleetHostSpec = serde_json::from_value(example["host"].clone()).unwrap();
2532 let config = SshFleetHostConfig::from_host_spec(&host)
2533 .expect("documented host must load without migration errors");
2534 assert!(config.known_hosts.is_some());
2535 }
2536 }
2537
2538 #[test]
2539 fn fleet_host_ssh_config_requires_explicit_safe_fields() {
2540 let tmp = TempDir::new().unwrap();
2541 let mut config = SshFleetHostConfig::new("", "/srv/codewhale");
2542 config.env_allowlist = BTreeSet::from(["SAFE_FLAG".to_string()]);
2543
2544 let err = SshFleetHostAdapter::new(tmp.path(), config).unwrap_err();
2545
2546 assert_eq!(err.kind, FleetHostErrorKind::Configuration);
2547 assert!(err.message.contains("explicit host"));
2548 }
2549
2550 #[test]
2551 fn fleet_host_ssh_config_maps_from_protocol_host_spec() {
2552 let spec = FleetHostSpec::Ssh {
2553 host: "builder.example.test".to_string(),
2554 port: Some(2222),
2555 user: Some("fleet".to_string()),
2556 identity: Some(PathBuf::from("/tmp/fleet_id")),
2557 known_hosts: None,
2558 host_key_fingerprint: None,
2559 working_directory: Some(PathBuf::from("/srv/codewhale")),
2560 env_allowlist: vec!["FLEET_PROFILE".to_string()],
2561 codewhale_binary: Some("/usr/local/bin/codewhale".to_string()),
2562 };
2563
2564 let config = SshFleetHostConfig::from_host_spec(&spec).unwrap();
2565
2566 assert_eq!(config.host, "builder.example.test");
2567 assert_eq!(config.port, Some(2222));
2568 assert_eq!(config.user.as_deref(), Some("fleet"));
2569 assert_eq!(config.working_directory, PathBuf::from("/srv/codewhale"));
2570 assert!(config.env_allowlist.contains("FLEET_PROFILE"));
2571 assert_eq!(config.codewhale_binary, "/usr/local/bin/codewhale");
2572 }
2573
2574 #[test]
2575 fn worker_env_forces_telemetry_off_even_when_the_allowlist_says_otherwise() {
2576 // The spawn path env_clear()s and rebuilds from this map, so anything
2577 // absent here simply does not exist inside the worker.
2578 let base = process_base_env();
2579 assert_eq!(
2580 base.get("CODEWHALE_TELEMETRY").map(String::as_str),
2581 Some("false")
2582 );
2583 assert_eq!(
2584 base.get("DEEPSEEK_TELEMETRY").map(String::as_str),
2585 Some("false")
2586 );
2587
2588 // A caller that allowlists the switch cannot switch it back on.
2589 let request_env = BTreeMap::from([
2590 ("CODEWHALE_TELEMETRY".to_string(), "true".to_string()),
2591 ("DEEPSEEK_TELEMETRY".to_string(), "1".to_string()),
2592 ("FLEET_PROFILE".to_string(), "builder".to_string()),
2593 ]);
2594 let allowlist = BTreeSet::from([
2595 "CODEWHALE_TELEMETRY".to_string(),
2596 "DEEPSEEK_TELEMETRY".to_string(),
2597 "FLEET_PROFILE".to_string(),
2598 ]);
2599
2600 let env = worker_env(&request_env, &allowlist).expect("worker env");
2601 assert_eq!(
2602 env.get("CODEWHALE_TELEMETRY").map(String::as_str),
2603 Some("false")
2604 );
2605 assert_eq!(
2606 env.get("DEEPSEEK_TELEMETRY").map(String::as_str),
2607 Some("false")
2608 );
2609 // Unrelated allowlisted entries still come through.
2610 assert_eq!(
2611 env.get("FLEET_PROFILE").map(String::as_str),
2612 Some("builder")
2613 );
2614 }
2615
2616 /// The env map above is only a claim about a function. This dispatches a
2617 /// real worker through the real spawn path and reads what the worker
2618 /// process actually received, because that is the environment a Codewhale
2619 /// worker would resolve telemetry from.
2620 ///
2621 /// Also asserts the operator's own home stays clean: a worker is an
2622 /// implementation detail of the parent session, and the parent already
2623 /// accounts for the dispatch, so a worker that wrote telemetry state into
2624 /// `$CODEWHALE_HOME` would double-count every fleet run.
2625 #[cfg(unix)]
2626 #[test]
2627 fn fleet_worker_env_carries_telemetry_off() {
2628 let fixture = TempDir::new().expect("fixture root");
2629 let workspace = fixture.path().join("workspace");
2630 let operator_home = fixture.path().join("operator-codewhale-home");
2631 std::fs::create_dir_all(&workspace).expect("workspace");
2632 std::fs::create_dir_all(&operator_home).expect("operator home");
2633 let receipt = fixture.path().join("worker-env.txt");
2634
2635 let mut adapter = LocalProcessFleetHostAdapter::new(&workspace);
2636 let mut request = FleetWorkerStartRequest::new(
2637 "telemetry-env-probe",
2638 FleetWorkerCommand::new(
2639 "/bin/sh",
2640 ["-c".to_string(), format!("env > {}", receipt.display())],
2641 ),
2642 );
2643 // A caller that both sets and allowlists the switch still cannot turn
2644 // a worker on.
2645 request
2646 .env
2647 .insert("CODEWHALE_TELEMETRY".to_string(), "true".to_string());
2648 request.env.insert(
2649 "CODEWHALE_HOME".to_string(),
2650 operator_home.display().to_string(),
2651 );
2652 request
2653 .env_allowlist
2654 .insert("CODEWHALE_TELEMETRY".to_string());
2655 request.env_allowlist.insert("CODEWHALE_HOME".to_string());
2656
2657 adapter.start_worker(request).expect("start worker");
2658 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
2659 let mut dumped = None;
2660 while std::time::Instant::now() < deadline {
2661 if let Ok(contents) = std::fs::read_to_string(&receipt)
2662 && contents.contains("CODEWHALE_TELEMETRY=")
2663 && contents.contains("DEEPSEEK_TELEMETRY=")
2664 {
2665 dumped = Some(contents);
2666 break;
2667 }
2668 std::thread::sleep(std::time::Duration::from_millis(20));
2669 }
2670 let dumped = dumped.unwrap_or_else(|| {
2671 std::fs::read_to_string(&receipt).expect("worker must dump its environment")
2672 });
2673
2674 let value = |key: &str| {
2675 dumped
2676 .lines()
2677 .find_map(|line| line.strip_prefix(&format!("{key}=")))
2678 .map(str::to_string)
2679 };
2680 assert_eq!(
2681 value("CODEWHALE_TELEMETRY").as_deref(),
2682 Some("false"),
2683 "worker environment:\n{dumped}"
2684 );
2685 assert_eq!(
2686 value("DEEPSEEK_TELEMETRY").as_deref(),
2687 Some("false"),
2688 "worker environment:\n{dumped}"
2689 );
2690 assert!(
2691 !operator_home.join("telemetry").exists(),
2692 "a fleet worker must not write telemetry state into the operator's home"
2693 );
2694 }
2695 }
2696
2696 lines RUST