1mod cleanup;
2mod launch;
3mod observation;
4mod readiness;
5
6use crate::operation_bound::{OperationBound, OperationTimingEvidence};
7use crate::plan::{CommandPlan, LaunchFilePlan, LaunchPlan, ProcessEndpointPlan, ReadinessPlan};
8use crate::process_group::{SignalEvidence, process_start_time};
9use serde::{Deserialize, Serialize};
10use std::path::{Path, PathBuf};
11use std::time::Duration;
12
13#[cfg(test)]
14use crate::operation_bound::OperationTerminalCause;
15#[cfg(test)]
16use crate::plan::TargetRegistryExpectedTarget;
17#[cfg(test)]
18use crate::shell::shell_quote_path;
19use cleanup::removal_summary;
20#[cfg(test)]
21use cleanup::terminate_local;
22#[cfg(test)]
23use launch::{materialize_local_launch_files, remote_launch_file_script, spawn_local};
24#[cfg(test)]
25use observation::{run_status_command, verified_local_status};
26#[cfg(test)]
27use readiness::{
28 HttpTargetRegistryProbe, match_target_registry, probe_http, probe_http_json,
29 wait_http_target_registry_ready, wait_process_alive_ready,
30};
31#[cfg(test)]
32use sha2::{Digest, Sha256};
33#[cfg(test)]
34use std::collections::BTreeMap;
35#[cfg(test)]
36use std::fs;
37#[cfg(test)]
38use std::io::{Read, Write};
39#[cfg(test)]
40use std::os::unix::process::CommandExt;
41#[cfg(test)]
42use std::process::{Command, Output, Stdio};
43#[cfg(test)]
44use std::thread;
45#[cfg(test)]
46use std::time::Instant;
47
48pub const REMOTE_LOG_SYNC_DEADLINE: Duration = Duration::from_secs(30);
49
50#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
51#[serde(deny_unknown_fields)]
52pub struct HostProcessHandle {
53 pub leader_pid: u32,
54 pub process_group: u32,
55 pub leader_start_time_ticks: u64,
56 #[serde(default, skip_serializing_if = "Option::is_none")]
60 pub container: Option<String>,
61}
62
63impl HostProcessHandle {
64 fn new(leader_pid: u32, container: Option<String>) -> Result<Self, String> {
65 if leader_pid == 0 {
66 return Err("host process-group handle requires a non-zero leader pid".to_owned());
67 }
68 let leader_start_time_ticks = process_start_time(leader_pid)
69 .map_err(|error| error.to_string())?
70 .ok_or_else(|| {
71 format!("host process {leader_pid} exited before its identity could be recorded")
72 })?;
73 Ok(Self {
74 leader_pid,
75 process_group: leader_pid,
76 leader_start_time_ticks,
77 container,
78 })
79 }
80
81 fn validate(&self) -> Result<(), String> {
82 validate_process_identity(
83 self.leader_pid,
84 self.process_group,
85 self.leader_start_time_ticks,
86 "host",
87 )
88 }
89}
90
91#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
92#[serde(deny_unknown_fields)]
93pub struct SshProcessHandle {
94 pub target: String,
95 pub leader_pid: u32,
96 pub process_group: u32,
97 pub leader_start_time_ticks: u64,
98 pub stdout: PathBuf,
99 pub stderr: PathBuf,
100 #[serde(default, skip_serializing_if = "Option::is_none")]
104 pub container: Option<String>,
105}
106
107impl SshProcessHandle {
108 fn validate(&self) -> Result<(), String> {
109 if self.target.is_empty() {
110 return Err("SSH process handle requires a target".to_owned());
111 }
112 validate_process_identity(
113 self.leader_pid,
114 self.process_group,
115 self.leader_start_time_ticks,
116 "SSH",
117 )
118 }
119}
120
121fn validate_process_identity(
122 leader_pid: u32,
123 process_group: u32,
124 leader_start_time_ticks: u64,
125 kind: &str,
126) -> Result<(), String> {
127 if leader_pid == 0 || process_group == 0 || leader_start_time_ticks == 0 {
128 return Err(format!("{kind} process-group handle requires non-zero ids"));
129 }
130 if leader_pid != process_group {
131 return Err(format!(
132 "{kind} process-group handle requires leader_pid to equal process_group"
133 ));
134 }
135 Ok(())
136}
137
138#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
139#[serde(tag = "kind", rename_all = "kebab-case", deny_unknown_fields)]
140pub enum ProcessHandle {
141 Local(HostProcessHandle),
142 Ssh(SshProcessHandle),
143}
144
145#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
146#[serde(rename_all = "kebab-case")]
147pub enum CleanupTrigger {
148 StartupRollback,
149 Stop,
150 Recovery,
151}
152
153#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
154#[serde(deny_unknown_fields)]
155pub struct CleanupEvidence {
156 pub trigger: CleanupTrigger,
157 pub elapsed_ms: u64,
158 pub status_deadline_ms: u64,
159 pub term_grace_ms: u64,
160 pub kill_grace_ms: u64,
161 pub reap_grace_ms: Option<u64>,
162 pub remote_deadline_ms: Option<u64>,
163 pub verified: bool,
164 pub already_exited: bool,
165 pub forced: bool,
166 pub signals: Vec<SignalEvidence>,
167 pub error: Option<String>,
168 #[serde(default, skip_serializing_if = "Option::is_none")]
174 pub container_removal: Option<ContainerRemovalEvidence>,
175 #[serde(default, skip_serializing_if = "Option::is_none")]
181 pub device_residuals: Option<Vec<DeviceResidualEvidence>>,
182}
183
184#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
191#[serde(tag = "outcome", rename_all = "snake_case", deny_unknown_fields)]
192pub enum DeviceResidualEvidence {
193 Freed {
194 machine: String,
195 device: u32,
196 },
197 ResidualHeld {
198 machine: String,
199 device: u32,
200 bytes: u64,
201 },
202 ProbeUnavailable {
203 machine: String,
204 device: u32,
205 reason: String,
206 },
207}
208
209#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
210#[serde(deny_unknown_fields)]
211pub struct ContainerRemovalEvidence {
212 pub container: String,
213 pub elapsed_ms: u64,
214 pub operation_elapsed_ms: u64,
215 pub deadline_ms: u64,
216 pub client_cleanup: Option<crate::container::CommandCleanupEvidence>,
217 pub confirmed: bool,
218 pub already_absent: bool,
219 pub error: Option<String>,
220}
221
222#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
223#[serde(deny_unknown_fields)]
224pub struct ProcessStatus {
225 pub queried: bool,
226 pub alive: bool,
227 pub error: Option<String>,
228}
229
230#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
231#[serde(deny_unknown_fields)]
232pub struct TargetRegistryMatchEvidence {
233 pub url: String,
234 pub role: String,
235 pub healthy: bool,
236 #[serde(default, skip_serializing_if = "Option::is_none")]
237 pub bootstrap_port: Option<u16>,
238}
239
240#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
241#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
242pub enum ReadinessEvidence {
243 Http {
244 url: String,
245 attempts: u32,
246 ready_unix_ms: u64,
247 timing: OperationTimingEvidence,
248 diagnostic_attempts: Vec<ReadinessAttemptEvidence>,
249 },
250 HttpTargetRegistry {
251 readiness_url: String,
252 registry_url: String,
253 attempts: u32,
254 ready_unix_ms: u64,
255 matched_targets: Vec<TargetRegistryMatchEvidence>,
256 timing: OperationTimingEvidence,
257 diagnostic_attempts: Vec<ReadinessAttemptEvidence>,
258 },
259 ProcessAlive {
260 ready_unix_ms: u64,
261 timing: OperationTimingEvidence,
262 diagnostic_attempts: Vec<ReadinessAttemptEvidence>,
263 },
264}
265
266#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
267#[serde(deny_unknown_fields)]
268pub struct ReadinessAttemptEvidence {
269 pub operation: String,
270 pub effective_bound_ms: u64,
271 pub succeeded: bool,
272 pub error: Option<String>,
273}
274
275#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
276#[serde(rename_all = "snake_case")]
277pub enum ReadinessFailureKind {
278 Exited,
279 Interrupted,
280 Timeout,
281}
282
283#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
284#[serde(deny_unknown_fields)]
285pub struct ReadinessFailure {
286 pub kind: ReadinessFailureKind,
287 pub message: String,
288 pub timing: Option<OperationTimingEvidence>,
289 pub diagnostic_attempts: Vec<ReadinessAttemptEvidence>,
290}
291
292pub struct ProcessSpec<'a> {
293 pub launch: &'a LaunchPlan,
294 pub command: &'a CommandPlan,
295 pub launch_files: &'a [LaunchFilePlan],
296 pub cache_root: &'a Path,
297 pub stdout: &'a Path,
298 pub stderr: &'a Path,
299 pub remote_dir: &'a Path,
300 pub container: Option<&'a str>,
303}
304
305#[derive(Debug, thiserror::Error)]
306pub enum ServerLaunchError {
307 #[error("{message}")]
308 Preparation { message: String },
309 #[error("failed to {operation} {path}: {source}")]
310 FileIo {
311 operation: &'static str,
312 path: PathBuf,
313 #[source]
314 source: std::io::Error,
315 },
316 #[error("failed to launch {program:?}: {source}")]
317 Process {
318 program: String,
319 #[source]
320 source: std::io::Error,
321 },
322 #[error(transparent)]
323 Ssh(#[from] crate::ssh::SshError),
324 #[error("{operation} exited with {status}: {diagnostics}")]
325 Exit {
326 operation: String,
327 status: std::process::ExitStatus,
328 diagnostics: String,
329 },
330 #[error("SSH launch on {target:?} returned non-UTF-8 identity: {source}")]
331 NonUtf8Identity {
332 target: String,
333 #[source]
334 source: std::string::FromUtf8Error,
335 },
336 #[error("SSH launch on {target:?} returned no process id")]
337 MissingProcessId { target: String },
338 #[error("SSH launch on {target:?} returned no process start time")]
339 MissingStartTime { target: String },
340 #[error("SSH launch on {target:?} returned invalid process id {value:?}: {source}")]
341 InvalidProcessId {
342 target: String,
343 value: String,
344 #[source]
345 source: std::num::ParseIntError,
346 },
347 #[error("SSH launch on {target:?} returned invalid process start time {value:?}: {source}")]
348 InvalidStartTime {
349 target: String,
350 value: String,
351 #[source]
352 source: std::num::ParseIntError,
353 },
354 #[error("SSH launch on {target:?} returned an invalid process identity: {details}")]
355 InvalidIdentity { target: String, details: String },
356 #[error("existing launch file target {path} is not a regular file", path = path.display())]
357 NotRegularFile { path: PathBuf },
358 #[error(
359 "existing launch file {path} does not match declared digest {expected}; found {actual}",
360 path = path.display()
361 )]
362 FileDigestMismatch {
363 path: PathBuf,
364 expected: String,
365 actual: String,
366 },
367}
368
369#[derive(Debug)]
370pub struct LaunchFailure {
371 pub error: ServerLaunchError,
372 pub ownership_unknown: bool,
373 pub container_removal: Option<Box<ContainerRemovalEvidence>>,
378 pub cleanup: Option<Box<CleanupEvidence>>,
379 cleanup_note: Option<String>,
380}
381
382impl LaunchFailure {
383 pub fn before_launch(message: String) -> Self {
384 Self {
385 error: ServerLaunchError::Preparation { message },
386 ownership_unknown: false,
387 container_removal: None,
388 cleanup: None,
389 cleanup_note: None,
390 }
391 }
392
393 fn from_error(error: ServerLaunchError) -> Self {
394 Self {
395 error,
396 ownership_unknown: false,
397 container_removal: None,
398 cleanup: None,
399 cleanup_note: None,
400 }
401 }
402
403 pub fn message(&self) -> String {
404 let mut message = self.error.to_string();
405 if let Some(note) = &self.cleanup_note {
406 message = format!("{message}; {note}");
407 } else if let Some(cleanup) = &self.cleanup
408 && let Some(error) = &cleanup.error
409 {
410 message = format!("{message}; local launch cleanup was not verified: {error}");
411 }
412 if let Some(removal) = &self.container_removal {
413 message = format!("{message}; {}", removal_summary(removal));
414 }
415 message
416 }
417
418 pub fn unresolved_ownership(message: String) -> Self {
419 Self {
420 error: ServerLaunchError::Preparation { message },
421 ownership_unknown: true,
422 container_removal: None,
423 cleanup: None,
424 cleanup_note: None,
425 }
426 }
427}
428
429pub trait ProcessLauncher {
430 fn spawn(&self, spec: ProcessSpec<'_>) -> Result<ProcessHandle, LaunchFailure>;
431}
432
433pub trait ProcessObserver {
434 fn status(&self, handle: &ProcessHandle) -> ProcessStatus;
435 fn status_with_bound(&self, handle: &ProcessHandle, bound: &OperationBound) -> ProcessStatus;
436 fn sync_logs(
437 &self,
438 handle: &ProcessHandle,
439 stdout: &Path,
440 stderr: &Path,
441 cleanup: bool,
442 ) -> Result<(), LogSyncError>;
443}
444
445pub trait ReadinessObserver {
446 fn wait_ready(
447 &self,
448 handle: &ProcessHandle,
449 endpoint: &ProcessEndpointPlan,
450 readiness: &ReadinessPlan,
451 bound: &OperationBound,
452 on_probe_failure: &mut dyn FnMut(&str),
453 ) -> Result<ReadinessEvidence, ReadinessFailure>;
454}
455
456pub trait ProcessCleanup {
457 fn terminate(
458 &self,
459 handle: &ProcessHandle,
460 trigger: CleanupTrigger,
461 on_container_removal: &mut dyn FnMut(&str),
462 ) -> CleanupEvidence;
463}
464
465pub trait ServerRuntime:
466 ProcessLauncher + ProcessObserver + ReadinessObserver + ProcessCleanup
467{
468}
469
470impl<T> ServerRuntime for T where
471 T: ProcessLauncher + ProcessObserver + ReadinessObserver + ProcessCleanup
472{
473}
474
475#[derive(Debug, thiserror::Error)]
476pub enum ProcessCommandError {
477 #[error("{operation} deadline expired")]
478 Deadline { operation: String },
479 #[error("{operation} was interrupted")]
480 Interrupted { operation: String },
481 #[error("failed to launch {operation}: {source}")]
482 Launch {
483 operation: String,
484 #[source]
485 source: std::io::Error,
486 },
487 #[error("failed to launch {operation}: {source}")]
488 Ssh {
489 operation: String,
490 #[source]
491 source: crate::ssh::SshError,
492 },
493 #[error("{operation} failed: {source}")]
494 Io {
495 operation: String,
496 #[source]
497 source: std::io::Error,
498 },
499 #[error("{operation} wait failed: {source}; child cleanup: {cleanup}")]
500 WaitCleanup {
501 operation: String,
502 #[source]
503 source: std::io::Error,
504 cleanup: String,
505 },
506 #[error("{operation} exited with {status}: {stderr}")]
507 Exit {
508 operation: String,
509 status: std::process::ExitStatus,
510 stderr: String,
511 },
512}
513
514#[derive(Debug, thiserror::Error)]
515pub enum LogSyncError {
516 #[error("failed to read remote log {path}: {source}")]
517 ReadRemote {
518 path: PathBuf,
519 #[source]
520 source: ProcessCommandError,
521 },
522 #[error("failed to read remote log {path}: command exited with {status}: {stderr}")]
523 RemoteExit {
524 path: PathBuf,
525 status: std::process::ExitStatus,
526 stderr: String,
527 },
528 #[error("failed to write local log {path}: {source}")]
529 WriteLocal {
530 path: PathBuf,
531 #[source]
532 source: std::io::Error,
533 },
534}
535
536#[derive(Clone, Copy, Debug, Default)]
537pub struct SystemProcessRuntime;
538
539#[cfg(test)]
540mod tests {
541 use super::*;
542 use crate::plan::LaunchFilePlan;
543 use std::cell::Cell;
544 use std::io::{BufRead, BufReader};
545 use std::net::TcpListener;
546 use std::os::unix::fs::{MetadataExt, PermissionsExt};
547
548 #[test]
549 fn expired_readiness_owner_prevents_a_fresh_network_attempt() {
550 let bound = OperationBound::finite(Duration::ZERO);
551 let error = probe_http("127.0.0.1", 9, "/ready", &bound, 30)
552 .err()
553 .map(|error| error.to_string())
554 .unwrap_or_default();
555
556 assert_eq!(error, "readiness operation deadline expired");
557 }
558
559 #[test]
560 fn readiness_attempt_deadline_bounds_a_trickled_status_line() -> Result<(), String> {
561 let listener = TcpListener::bind(("127.0.0.1", 0)).map_err(|error| error.to_string())?;
562 let port = listener
563 .local_addr()
564 .map_err(|error| error.to_string())?
565 .port();
566 let server = thread::spawn(move || -> Result<(), String> {
567 let (mut stream, _) = listener.accept().map_err(|error| error.to_string())?;
568 let mut request = [0_u8; 1024];
569 let _ = stream.read(&mut request);
570 let response = format!("HTTP/1.1 200 {}\r\n", " ".repeat(96));
571 for byte in response.bytes() {
572 if stream.write_all(&[byte]).is_err() {
573 break;
574 }
575 thread::sleep(Duration::from_millis(20));
576 }
577 Ok(())
578 });
579
580 let started = Instant::now();
581 let error = probe_http("127.0.0.1", port, "/ready", &OperationBound::unbounded(), 1)
582 .err()
583 .map(|error| error.to_string())
584 .unwrap_or_default();
585 let elapsed = started.elapsed();
586 server
587 .join()
588 .map_err(|_| "trickle fixture panicked".to_owned())??;
589
590 assert!(error.contains("deadline expired"), "{error}");
591 assert!(
592 elapsed < Duration::from_secs(2),
593 "a one-second attempt lasted {elapsed:?}"
594 );
595 Ok(())
596 }
597
598 #[test]
599 fn finite_readiness_accepts_a_response_after_250_milliseconds() -> Result<(), String> {
600 let listener = TcpListener::bind(("127.0.0.1", 0)).map_err(|error| error.to_string())?;
601 let port = listener
602 .local_addr()
603 .map_err(|error| error.to_string())?
604 .port();
605 let server = thread::spawn(move || -> Result<(), String> {
606 let (mut stream, _) = listener.accept().map_err(|error| error.to_string())?;
607 let mut request = [0_u8; 1024];
608 let _ = stream.read(&mut request);
609 thread::sleep(Duration::from_millis(350));
610 stream
611 .write_all(b"HTTP/1.1 204 No Content\r\nContent-Length: 0\r\n\r\n")
612 .map_err(|error| error.to_string())
613 });
614
615 let started = Instant::now();
616 probe_http(
617 "127.0.0.1",
618 port,
619 "/ready",
620 &OperationBound::finite(Duration::from_secs(2)),
621 1,
622 )
623 .map_err(|error| error.to_string())?;
624 let elapsed = started.elapsed();
625 server
626 .join()
627 .map_err(|_| "delayed readiness fixture panicked".to_owned())??;
628
629 assert!(elapsed >= Duration::from_millis(250), "elapsed {elapsed:?}");
630 assert!(elapsed < Duration::from_secs(2), "elapsed {elapsed:?}");
631 Ok(())
632 }
633
634 #[test]
635 fn readiness_attempt_deadline_includes_the_complete_response_body() -> Result<(), String> {
636 let listener = TcpListener::bind(("127.0.0.1", 0)).map_err(|error| error.to_string())?;
637 let port = listener
638 .local_addr()
639 .map_err(|error| error.to_string())?
640 .port();
641 let server = thread::spawn(move || -> Result<(), String> {
642 let (mut stream, _) = listener.accept().map_err(|error| error.to_string())?;
643 let mut request = [0_u8; 1024];
644 let _ = stream.read(&mut request);
645 stream
646 .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 10\r\n\r\nx")
647 .map_err(|error| error.to_string())?;
648 thread::sleep(Duration::from_millis(1_500));
649 Ok(())
650 });
651
652 let started = Instant::now();
653 let error = probe_http("127.0.0.1", port, "/ready", &OperationBound::unbounded(), 1)
654 .err()
655 .map(|error| error.to_string());
656 let elapsed = started.elapsed();
657 server
658 .join()
659 .map_err(|_| "readiness body fixture panicked".to_owned())??;
660
661 assert!(
662 error.is_some_and(|error| error == "readiness operation deadline expired"),
663 "readiness accepted an incomplete response body"
664 );
665 assert!(
666 elapsed < Duration::from_millis(1_500),
667 "elapsed {elapsed:?}"
668 );
669 Ok(())
670 }
671
672 #[test]
673 fn registry_attempt_deadline_bounds_a_trickled_body() -> Result<(), String> {
674 let listener = TcpListener::bind(("127.0.0.1", 0)).map_err(|error| error.to_string())?;
675 let port = listener
676 .local_addr()
677 .map_err(|error| error.to_string())?
678 .port();
679 let server = thread::spawn(move || -> Result<(), String> {
680 let (mut stream, _) = listener.accept().map_err(|error| error.to_string())?;
681 let mut request = [0_u8; 1024];
682 let _ = stream.read(&mut request);
683 stream
684 .write_all(b"HTTP/1.1 200 OK\r\nConnection: close\r\n\r\n")
685 .map_err(|error| error.to_string())?;
686 let body = format!("{}{{}}", " ".repeat(96));
687 for byte in body.bytes() {
688 if stream.write_all(&[byte]).is_err() {
689 break;
690 }
691 thread::sleep(Duration::from_millis(20));
692 }
693 Ok(())
694 });
695
696 let started = Instant::now();
697 let error = probe_http_json(
698 "127.0.0.1",
699 port,
700 "/workers",
701 "target registry",
702 &OperationBound::unbounded(),
703 1,
704 )
705 .err()
706 .map(|error| error.to_string())
707 .unwrap_or_default();
708 let elapsed = started.elapsed();
709 server
710 .join()
711 .map_err(|_| "registry trickle fixture panicked".to_owned())??;
712
713 assert!(error.contains("deadline expired"), "{error}");
714 assert!(
715 elapsed < Duration::from_secs(2),
716 "a one-second registry attempt lasted {elapsed:?}"
717 );
718 Ok(())
719 }
720
721 #[test]
722 fn expired_readiness_owner_rejects_a_registry_match() {
723 let expected = vec![TargetRegistryExpectedTarget {
724 url: "http://decode:30001".to_owned(),
725 role: "decode".to_owned(),
726 bootstrap_port: None,
727 }];
728 let response = serde_json::json!({
729 "workers": [{
730 "url": "http://decode:30001",
731 "worker_type": "decode",
732 "is_healthy": true
733 }]
734 });
735
736 let error = match_target_registry(
737 &response,
738 &target_registry_probe(&expected),
739 &OperationBound::finite(Duration::ZERO),
740 )
741 .err()
742 .map(|error| error.to_string())
743 .unwrap_or_default();
744
745 assert_eq!(error, "readiness operation deadline expired");
746 }
747
748 #[test]
749 fn process_status_command_cannot_outlive_the_readiness_owner() {
750 let started = Instant::now();
751 let error = run_status_command(
752 &["sh", "-c", "sleep 5"],
753 &[],
754 &OperationBound::finite(Duration::from_millis(50)),
755 )
756 .err()
757 .map(|error| error.to_string())
758 .unwrap_or_default();
759
760 assert_eq!(error, "process status attempt deadline expired");
761 assert!(
762 started.elapsed() < Duration::from_secs(1),
763 "bounded process status did not stop promptly"
764 );
765 }
766
767 #[test]
768 fn finite_readiness_retries_an_expired_process_status_attempt() -> Result<(), String> {
769 let calls = Cell::new(0_u32);
770 let mut failures = Vec::new();
771 let bound = OperationBound::finite(Duration::from_secs(3));
772 let evidence = wait_process_alive_ready(
773 |_| {
774 calls.set(calls.get() + 1);
775 if calls.get() == 1 {
776 thread::sleep(Duration::from_millis(1_050));
777 ProcessStatus {
778 queried: false,
779 alive: false,
780 error: Some("process status attempt deadline expired".to_owned()),
781 }
782 } else {
783 alive_status()
784 }
785 },
786 1,
787 &bound,
788 &mut |failure| failures.push(failure.to_owned()),
789 )
790 .map_err(|failure| failure.message)?;
791
792 assert_eq!(calls.get(), 2);
793 assert_eq!(failures, ["process status attempt deadline expired"]);
794 let ReadinessEvidence::ProcessAlive {
795 timing,
796 diagnostic_attempts,
797 ..
798 } = evidence
799 else {
800 return Err("process-alive readiness returned the wrong evidence kind".to_owned());
801 };
802 assert_eq!(
803 timing.start_boundary,
804 "after_process_spawn_before_readiness_attempt"
805 );
806 assert_eq!(diagnostic_attempts.len(), 1);
807 assert!(diagnostic_attempts[0].succeeded);
808 assert!((1..=1_000).contains(&diagnostic_attempts[0].effective_bound_ms));
809 Ok(())
810 }
811
812 #[test]
813 fn unbounded_process_status_command_does_not_acquire_a_timeout() -> Result<(), String> {
814 let output = run_status_command(
815 &["sh", "-c", "sleep 0.1; printf alive"],
816 &[],
817 &OperationBound::unbounded(),
818 )
819 .map_err(|error| error.to_string())?;
820
821 assert!(output.status.success());
822 assert_eq!(output.stdout, b"alive");
823 Ok(())
824 }
825
826 fn launch_file(root: &Path, text: &str, name: &str) -> LaunchFilePlan {
827 let sha256 = format!("{:x}", Sha256::digest(text.as_bytes()));
828 let relative_path = format!("launch-files/{sha256}/{name}");
829 LaunchFilePlan {
830 resolved_path: root.join(&relative_path),
831 relative_path,
832 text: text.to_owned(),
833 sha256,
834 }
835 }
836
837 fn run_script_with_input(script: &str, input: &[u8]) -> Result<Output, String> {
838 match crate::container::run_with_bound(
839 &["bash", "-c", script],
840 &[],
841 None,
842 Some(input),
843 &OperationBound::unbounded(),
844 None,
845 ) {
846 Ok(crate::container::BoundedWait::Exited {
847 status,
848 stdout,
849 stderr,
850 }) => Ok(Output {
851 status,
852 stdout,
853 stderr,
854 }),
855 Ok(crate::container::BoundedWait::Expired { .. }) => {
856 Err("unbounded launch-file fixture expired".to_owned())
857 }
858 Ok(crate::container::BoundedWait::Interrupted { .. }) => {
859 Err("launch-file fixture was interrupted".to_owned())
860 }
861 Err(_) => Err("launch-file fixture failed".to_owned()),
862 }
863 }
864
865 fn target_registry_endpoint(
866 registry_body: String,
867 ) -> Result<(ProcessEndpointPlan, thread::JoinHandle<Result<(), String>>), String> {
868 let listener = TcpListener::bind(("127.0.0.1", 0)).map_err(|error| error.to_string())?;
869 let port = listener
870 .local_addr()
871 .map_err(|error| error.to_string())?
872 .port();
873 let server = thread::spawn(move || {
874 for _ in 0..2 {
875 let (mut stream, _) = listener.accept().map_err(|error| error.to_string())?;
876 let mut request_line = String::new();
877 let mut reader =
878 BufReader::new(stream.try_clone().map_err(|error| error.to_string())?);
879 reader
880 .read_line(&mut request_line)
881 .map_err(|error| error.to_string())?;
882 loop {
883 let mut header = String::new();
884 reader
885 .read_line(&mut header)
886 .map_err(|error| error.to_string())?;
887 if header == "\r\n" || header.is_empty() {
888 break;
889 }
890 }
891 let body = if request_line.starts_with("GET /workers ") {
892 registry_body.as_bytes()
893 } else {
894 b""
895 };
896 let mut response = format!(
897 "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
898 body.len()
899 )
900 .into_bytes();
901 response.extend_from_slice(body);
902 stream
903 .write_all(&response)
904 .map_err(|error| error.to_string())?;
905 }
906 Ok(())
907 });
908 Ok((
909 ProcessEndpointPlan {
910 host: "127.0.0.1".to_owned(),
911 port,
912 },
913 server,
914 ))
915 }
916
917 fn target_registry_probe<'a>(
918 expected_targets: &'a [TargetRegistryExpectedTarget],
919 ) -> HttpTargetRegistryProbe<'a> {
920 HttpTargetRegistryProbe {
921 readiness_path: "/readiness",
922 registry_path: "/workers",
923 targets_field: "workers",
924 target_url_field: "url",
925 target_role_field: "worker_type",
926 target_healthy_field: "is_healthy",
927 target_bootstrap_port_field: "bootstrap_port",
928 expected_targets,
929 }
930 }
931
932 fn alive_status() -> ProcessStatus {
933 ProcessStatus {
934 queried: true,
935 alive: true,
936 error: None,
937 }
938 }
939
940 #[test]
941 fn local_launch_file_publication_reuses_the_immutable_target() -> Result<(), String> {
942 let root = tempfile::tempdir().map_err(|error| error.to_string())?;
943 let launch_file = launch_file(
944 root.path(),
945 "worker: \u{2603}\nmode: context\n",
946 "worker.yaml",
947 );
948
949 materialize_local_launch_files(std::slice::from_ref(&launch_file))
950 .map_err(|error| error.to_string())?;
951 let first_metadata =
952 fs::metadata(&launch_file.resolved_path).map_err(|error| error.to_string())?;
953 materialize_local_launch_files(std::slice::from_ref(&launch_file))
954 .map_err(|error| error.to_string())?;
955 let second_metadata =
956 fs::metadata(&launch_file.resolved_path).map_err(|error| error.to_string())?;
957
958 assert_eq!(
959 fs::read_to_string(&launch_file.resolved_path).map_err(|error| error.to_string())?,
960 launch_file.text
961 );
962 assert_eq!(first_metadata.ino(), second_metadata.ino());
963 assert_eq!(second_metadata.permissions().mode() & 0o222, 0);
964 Ok(())
965 }
966
967 #[test]
968 fn local_launch_file_mismatch_fails_before_spawn_without_replacing_it() -> Result<(), String> {
969 let root = tempfile::tempdir().map_err(|error| error.to_string())?;
970 let cache = root.path().join("cache");
971 let launch_file = launch_file(&cache, "expected\n", "worker.yaml");
972 let parent = launch_file
973 .resolved_path
974 .parent()
975 .ok_or_else(|| "launch file has no parent".to_owned())?;
976 fs::create_dir_all(parent).map_err(|error| error.to_string())?;
977 fs::write(&launch_file.resolved_path, "stale\n").map_err(|error| error.to_string())?;
978 let marker = root.path().join("spawned");
979 let command = CommandPlan {
980 argv: vec![
981 "sh".to_owned(),
982 "-c".to_owned(),
983 format!("printf launched > {}", shell_quote_path(&marker)),
984 ],
985 env: BTreeMap::new(),
986 explicit_env: Vec::new(),
987 pass_env: Vec::new(),
988 cwd: root.path().to_path_buf(),
989 };
990 let launch_files = vec![launch_file.clone()];
991
992 let result = spawn_local(ProcessSpec {
993 launch: &LaunchPlan::Local,
994 command: &command,
995 launch_files: &launch_files,
996 cache_root: &cache,
997 stdout: &root.path().join("stdout.log"),
998 stderr: &root.path().join("stderr.log"),
999 remote_dir: &root.path().join("remote"),
1000 container: None,
1001 });
1002
1003 let failure = match result {
1004 Err(failure) => failure,
1005 Ok(handle) => {
1006 let _ = terminate_local(&handle, CleanupTrigger::StartupRollback);
1007 return Err("mismatched launch file unexpectedly spawned a process".to_owned());
1008 }
1009 };
1010 assert!(!failure.ownership_unknown, "{failure:?}");
1011 assert!(failure.message().contains("does not match"), "{failure:?}");
1012 assert!(!marker.exists());
1013 assert_eq!(
1014 fs::read_to_string(&launch_file.resolved_path).map_err(|error| error.to_string())?,
1015 "stale\n"
1016 );
1017 Ok(())
1018 }
1019
1020 #[test]
1021 fn remote_launch_file_script_publishes_stdin_without_replacing_targets() -> Result<(), String> {
1022 let root = tempfile::tempdir().map_err(|error| error.to_string())?;
1023 let published = launch_file(
1024 root.path(),
1025 "worker: \u{96ea}\nmode: context\n",
1026 "worker.yaml",
1027 );
1028 let script = remote_launch_file_script(&published).map_err(|error| error.to_string())?;
1029
1030 let first = run_script_with_input(&script, published.text.as_bytes())?;
1031 assert!(
1032 first.status.success(),
1033 "{}",
1034 String::from_utf8_lossy(&first.stderr)
1035 );
1036 let first_metadata =
1037 fs::metadata(&published.resolved_path).map_err(|error| error.to_string())?;
1038 let first_inode = first_metadata.ino();
1039 assert_eq!(first_metadata.permissions().mode() & 0o222, 0);
1040 let reused = run_script_with_input(&script, published.text.as_bytes())?;
1041 assert!(
1042 reused.status.success(),
1043 "{}",
1044 String::from_utf8_lossy(&reused.stderr)
1045 );
1046 assert_eq!(
1047 fs::read_to_string(&published.resolved_path).map_err(|error| error.to_string())?,
1048 published.text
1049 );
1050 assert_eq!(
1051 fs::metadata(&published.resolved_path)
1052 .map_err(|error| error.to_string())?
1053 .ino(),
1054 first_inode
1055 );
1056
1057 let corrupt = launch_file(root.path(), "expected\n", "corrupt.yaml");
1058 let corrupt_parent = corrupt
1059 .resolved_path
1060 .parent()
1061 .ok_or_else(|| "launch file has no parent".to_owned())?;
1062 fs::create_dir_all(corrupt_parent).map_err(|error| error.to_string())?;
1063 fs::write(&corrupt.resolved_path, "stale\n").map_err(|error| error.to_string())?;
1064 let rejected = run_script_with_input(
1065 &remote_launch_file_script(&corrupt).map_err(|error| error.to_string())?,
1066 corrupt.text.as_bytes(),
1067 )?;
1068 assert!(!rejected.status.success());
1069 assert_eq!(
1070 fs::read_to_string(&corrupt.resolved_path).map_err(|error| error.to_string())?,
1071 "stale\n"
1072 );
1073 Ok(())
1074 }
1075
1076 #[test]
1077 fn target_registry_readiness_records_all_expected_targets() -> Result<(), String> {
1078 let (endpoint, server) = target_registry_endpoint(
1079 serde_json::json!({
1080 "workers": [
1081 {
1082 "url": "http://prefill:30000",
1083 "worker_type": "prefill",
1084 "is_healthy": true,
1085 "bootstrap_port": 8998
1086 },
1087 {
1088 "url": "http://decode:30001",
1089 "worker_type": "decode",
1090 "is_healthy": true
1091 }
1092 ]
1093 })
1094 .to_string(),
1095 )?;
1096 let expected = vec![
1097 TargetRegistryExpectedTarget {
1098 url: "http://prefill:30000".to_owned(),
1099 role: "prefill".to_owned(),
1100 bootstrap_port: Some(8998),
1101 },
1102 TargetRegistryExpectedTarget {
1103 url: "http://decode:30001".to_owned(),
1104 role: "decode".to_owned(),
1105 bootstrap_port: None,
1106 },
1107 ];
1108
1109 let bound = OperationBound::unbounded();
1110 let evidence = wait_http_target_registry_ready(
1111 |_| alive_status(),
1112 &endpoint,
1113 target_registry_probe(&expected),
1114 1,
1115 &bound,
1116 &mut |_| {},
1117 )
1118 .map_err(|failure| failure.message)?;
1119 server
1120 .join()
1121 .map_err(|_| "target registry fixture panicked".to_owned())??;
1122
1123 let record_value = serde_json::to_value(&evidence).map_err(|error| error.to_string())?;
1124 assert_eq!(record_value["kind"], "http_target_registry");
1125 assert_eq!(
1126 record_value["matched_targets"].as_array().map(Vec::len),
1127 Some(2)
1128 );
1129 let ReadinessEvidence::HttpTargetRegistry {
1130 readiness_url,
1131 registry_url,
1132 attempts,
1133 matched_targets,
1134 timing,
1135 diagnostic_attempts,
1136 ready_unix_ms: _,
1137 } = evidence
1138 else {
1139 return Err("target registry readiness returned the wrong evidence kind".to_owned());
1140 };
1141 assert_eq!(
1142 readiness_url,
1143 format!("http://127.0.0.1:{}/readiness", endpoint.port)
1144 );
1145 assert_eq!(
1146 registry_url,
1147 format!("http://127.0.0.1:{}/workers", endpoint.port)
1148 );
1149 assert_eq!(attempts, 1);
1150 assert_eq!(
1151 timing.budget,
1152 crate::operation_bound::OperationBudgetEvidence::Unbounded
1153 );
1154 assert_eq!(
1155 timing.terminal_cause,
1156 crate::operation_bound::OperationTerminalCause::Succeeded
1157 );
1158 assert_eq!(diagnostic_attempts.len(), 3);
1159 assert!(diagnostic_attempts.iter().all(|attempt| {
1160 attempt.succeeded && (1..=1_000).contains(&attempt.effective_bound_ms)
1161 }));
1162 assert_eq!(
1163 matched_targets,
1164 vec![
1165 TargetRegistryMatchEvidence {
1166 url: "http://prefill:30000".to_owned(),
1167 role: "prefill".to_owned(),
1168 healthy: true,
1169 bootstrap_port: Some(8998),
1170 },
1171 TargetRegistryMatchEvidence {
1172 url: "http://decode:30001".to_owned(),
1173 role: "decode".to_owned(),
1174 healthy: true,
1175 bootstrap_port: None,
1176 },
1177 ]
1178 );
1179 Ok(())
1180 }
1181
1182 #[test]
1183 fn finite_target_registry_attempts_record_the_resolved_attempt_budget() -> Result<(), String> {
1184 let (endpoint, server) = target_registry_endpoint(
1185 serde_json::json!({
1186 "workers": [{
1187 "url": "http://decode:30001",
1188 "worker_type": "decode",
1189 "is_healthy": true
1190 }]
1191 })
1192 .to_string(),
1193 )?;
1194 let expected = vec![TargetRegistryExpectedTarget {
1195 url: "http://decode:30001".to_owned(),
1196 role: "decode".to_owned(),
1197 bootstrap_port: None,
1198 }];
1199
1200 let bound = OperationBound::finite(Duration::from_secs(2));
1201 let evidence = wait_http_target_registry_ready(
1202 |_| alive_status(),
1203 &endpoint,
1204 target_registry_probe(&expected),
1205 1,
1206 &bound,
1207 &mut |_| {},
1208 )
1209 .map_err(|failure| failure.message)?;
1210 server
1211 .join()
1212 .map_err(|_| "target registry fixture panicked".to_owned())??;
1213
1214 let ReadinessEvidence::HttpTargetRegistry {
1215 timing,
1216 diagnostic_attempts,
1217 ..
1218 } = evidence
1219 else {
1220 return Err("target registry readiness returned the wrong evidence kind".to_owned());
1221 };
1222 assert_eq!(
1223 timing.budget,
1224 crate::operation_bound::OperationBudgetEvidence::Finite {
1225 configured_ms: 2_000,
1226 }
1227 );
1228 assert_eq!(diagnostic_attempts.len(), 3);
1229 assert!(diagnostic_attempts.iter().all(|attempt| {
1230 attempt.succeeded && (1..=1_000).contains(&attempt.effective_bound_ms)
1231 }));
1232 Ok(())
1233 }
1234
1235 #[test]
1236 fn target_registry_readiness_rejects_partial_registration() -> Result<(), String> {
1237 let (endpoint, server) = target_registry_endpoint(
1238 serde_json::json!({
1239 "workers": [{
1240 "url": "http://prefill:30000",
1241 "worker_type": "prefill",
1242 "is_healthy": true,
1243 "bootstrap_port": 8998
1244 }]
1245 })
1246 .to_string(),
1247 )?;
1248 let expected = vec![
1249 TargetRegistryExpectedTarget {
1250 url: "http://prefill:30000".to_owned(),
1251 role: "prefill".to_owned(),
1252 bootstrap_port: Some(8998),
1253 },
1254 TargetRegistryExpectedTarget {
1255 url: "http://decode:30001".to_owned(),
1256 role: "decode".to_owned(),
1257 bootstrap_port: None,
1258 },
1259 ];
1260
1261 let mut probe_failures = Vec::new();
1262 let bound = OperationBound::finite(Duration::from_secs(1));
1263 let failure = match wait_http_target_registry_ready(
1264 |_| alive_status(),
1265 &endpoint,
1266 target_registry_probe(&expected),
1267 1,
1268 &bound,
1269 &mut |failure| probe_failures.push(failure.to_owned()),
1270 ) {
1271 Err(failure) => failure,
1272 Ok(evidence) => {
1273 return Err(format!(
1274 "partial target registration unexpectedly became ready: {evidence:?}"
1275 ));
1276 }
1277 };
1278 server
1279 .join()
1280 .map_err(|_| "target registry fixture panicked".to_owned())??;
1281
1282 assert_eq!(failure.kind, ReadinessFailureKind::Timeout);
1283 let timing = failure
1284 .timing
1285 .as_ref()
1286 .ok_or_else(|| "readiness timeout has no timing evidence".to_owned())?;
1287 assert_eq!(
1288 timing.budget,
1289 crate::operation_bound::OperationBudgetEvidence::Finite {
1290 configured_ms: 1_000,
1291 }
1292 );
1293 assert_eq!(timing.terminal_cause, OperationTerminalCause::TimedOut);
1294 assert!(probe_failures.iter().any(|failure| {
1295 failure.contains("target registry has no \"decode\" target at \"http://decode:30001\"")
1296 }));
1297 Ok(())
1298 }
1299
1300 #[test]
1301 fn termination_waits_for_the_group_after_the_launcher_exits() -> Result<(), String> {
1302 let mut child = Command::new("sh")
1303 .args([
1304 "-c",
1305 "trap 'exit 0' TERM; sh -c 'trap \"\" TERM; exec sleep 30' & wait",
1306 ])
1307 .stdin(Stdio::null())
1308 .stdout(Stdio::null())
1309 .stderr(Stdio::null())
1310 .process_group(0)
1311 .spawn()
1312 .map_err(|error| error.to_string())?;
1313 let handle = HostProcessHandle::new(child.id(), None)?;
1314 thread::sleep(Duration::from_millis(100));
1315 let reaper = thread::spawn(move || child.wait());
1316
1317 let cleanup = terminate_local(&handle, CleanupTrigger::Stop);
1318 if !cleanup.verified {
1319 let _ = Command::new("kill")
1320 .args(["-KILL", "--", &format!("-{}", handle.process_group)])
1321 .status();
1322 }
1323 let _ = reaper.join();
1324
1325 assert!(cleanup.verified, "{cleanup:?}");
1326 assert!(cleanup.forced);
1327 assert!(cleanup.elapsed_ms >= cleanup.term_grace_ms);
1328 assert_eq!(cleanup.status_deadline_ms, 2_000);
1329 assert_eq!(cleanup.term_grace_ms, 2_000);
1330 assert_eq!(cleanup.kill_grace_ms, 10_000);
1331 Ok(())
1332 }
1333
1334 struct OrphanGroup {
1335 handle: HostProcessHandle,
1336 member_pid: u32,
1337 }
1338
1339 fn spawn_orphan_group(root: &Path) -> Result<OrphanGroup, String> {
1343 let marker = root.join("member.pid");
1344 let mut child = Command::new("bash")
1345 .args([
1346 "-c",
1347 &format!("sleep 300 & echo $! > {}; exec sleep 300", marker.display()),
1348 ])
1349 .stdin(Stdio::null())
1350 .stdout(Stdio::null())
1351 .stderr(Stdio::null())
1352 .process_group(0)
1353 .spawn()
1354 .map_err(|error| error.to_string())?;
1355 let handle = HostProcessHandle::new(child.id(), None)?;
1356 let deadline = Instant::now() + Duration::from_secs(5);
1357 while !marker.exists() {
1358 if Instant::now() > deadline {
1359 let _ = Command::new("kill")
1360 .args(["-KILL", "--", &format!("-{}", handle.process_group)])
1361 .status();
1362 return Err("orphan-group member did not start".to_owned());
1363 }
1364 thread::sleep(Duration::from_millis(10));
1365 }
1366 child.kill().map_err(|error| error.to_string())?;
1367 child.wait().map_err(|error| error.to_string())?;
1368 let member_pid = fs::read_to_string(&marker)
1369 .map_err(|error| error.to_string())?
1370 .trim()
1371 .parse::<u32>()
1372 .map_err(|error| error.to_string())?;
1373 Ok(OrphanGroup { handle, member_pid })
1374 }
1375
1376 fn force_kill_group(process_group: u32) {
1377 let _ = Command::new("kill")
1378 .args(["-KILL", "--", &format!("-{process_group}")])
1379 .status();
1380 }
1381
1382 #[test]
1383 fn termination_reaps_orphaned_members_after_the_leader_exits() -> Result<(), String> {
1384 let root = tempfile::tempdir().map_err(|error| error.to_string())?;
1385 let orphan = spawn_orphan_group(root.path())?;
1386 let cleanup = terminate_local(&orphan.handle, CleanupTrigger::Stop);
1387 if !cleanup.verified {
1388 force_kill_group(orphan.handle.process_group);
1389 }
1390
1391 assert!(cleanup.verified, "{cleanup:?}");
1392 assert!(!cleanup.already_exited, "{cleanup:?}");
1393 assert!(!cleanup.signals.is_empty(), "{cleanup:?}");
1394 assert_eq!(
1395 process_start_time(orphan.member_pid).map_err(|error| error.to_string())?,
1396 None,
1397 "the orphaned member was reaped"
1398 );
1399 Ok(())
1400 }
1401
1402 #[test]
1403 fn termination_refuses_a_cohort_inconsistent_orphan_group() -> Result<(), String> {
1404 let root = tempfile::tempdir().map_err(|error| error.to_string())?;
1405 let orphan = spawn_orphan_group(root.path())?;
1406 let inflated = HostProcessHandle {
1409 leader_start_time_ticks: orphan.handle.leader_start_time_ticks + 1_000_000_000,
1410 ..orphan.handle.clone()
1411 };
1412 let cleanup = terminate_local(&inflated, CleanupTrigger::Stop);
1413 force_kill_group(orphan.handle.process_group);
1414
1415 assert!(!cleanup.verified, "{cleanup:?}");
1416 assert!(!cleanup.forced, "{cleanup:?}");
1417 let error = cleanup.error.ok_or("the refusal carries a reason")?;
1418 assert!(
1419 error.contains(&orphan.member_pid.to_string()),
1420 "the refusal names the offending member: {error}"
1421 );
1422 assert!(
1423 error.contains("before the recorded leader start"),
1424 "{error}"
1425 );
1426 Ok(())
1427 }
1428
1429 #[test]
1430 fn ssh_cleanup_script_kills_a_cohort_consistent_orphan_group() -> Result<(), String> {
1431 let root = tempfile::tempdir().map_err(|error| error.to_string())?;
1432 let orphan = spawn_orphan_group(root.path())?;
1433 let script = super::cleanup::remote_cleanup_script(&SshProcessHandle {
1434 target: "fixture".to_owned(),
1435 leader_pid: orphan.handle.leader_pid,
1436 process_group: orphan.handle.process_group,
1437 leader_start_time_ticks: orphan.handle.leader_start_time_ticks,
1438 stdout: root.path().join("stdout.log"),
1439 stderr: root.path().join("stderr.log"),
1440 container: None,
1441 });
1442 let output = Command::new("bash")
1445 .args(["-c", &script])
1446 .output()
1447 .map_err(|error| error.to_string())?;
1448 let stdout = String::from_utf8_lossy(&output.stdout);
1449 if !stdout.contains("INFERLAB_CLEANUP\tcleanup") {
1450 force_kill_group(orphan.handle.process_group);
1451 }
1452
1453 assert!(output.status.success(), "{stdout}");
1454 assert!(
1455 stdout
1456 .lines()
1457 .rev()
1458 .find_map(|line| line.strip_prefix("INFERLAB_CLEANUP\t"))
1459 .is_some_and(|line| line.starts_with("cleanup\t")),
1460 "{stdout}"
1461 );
1462 assert_eq!(
1463 process_start_time(orphan.member_pid).map_err(|error| error.to_string())?,
1464 None,
1465 "the orphaned member was reaped"
1466 );
1467 Ok(())
1468 }
1469
1470 #[test]
1471 fn ssh_cleanup_script_refuses_a_cohort_inconsistent_orphan_group() -> Result<(), String> {
1472 let root = tempfile::tempdir().map_err(|error| error.to_string())?;
1473 let orphan = spawn_orphan_group(root.path())?;
1474 let script = super::cleanup::remote_cleanup_script(&SshProcessHandle {
1475 target: "fixture".to_owned(),
1476 leader_pid: orphan.handle.leader_pid,
1477 process_group: orphan.handle.process_group,
1478 leader_start_time_ticks: orphan.handle.leader_start_time_ticks + 1_000_000_000,
1479 stdout: root.path().join("stdout.log"),
1480 stderr: root.path().join("stderr.log"),
1481 container: None,
1482 });
1483 let output = Command::new("bash")
1484 .args(["-c", &script])
1485 .output()
1486 .map_err(|error| error.to_string())?;
1487 let stdout = String::from_utf8_lossy(&output.stdout);
1488 force_kill_group(orphan.handle.process_group);
1489
1490 assert!(output.status.success(), "{stdout}");
1491 let line = stdout
1492 .lines()
1493 .rev()
1494 .find_map(|line| line.strip_prefix("INFERLAB_CLEANUP\t"))
1495 .ok_or("the script printed no cleanup result")?;
1496 assert!(line.starts_with("unknown\t"), "{stdout}");
1497 assert!(
1498 line.contains(&orphan.member_pid.to_string()),
1499 "the detail names the offending member: {line}"
1500 );
1501 Ok(())
1502 }
1503
1504 #[test]
1505 fn cleanup_evidence_without_device_residuals_still_decodes()
1506 -> Result<(), Box<dyn std::error::Error>> {
1507 let legacy = serde_json::json!({
1511 "trigger": "stop",
1512 "elapsed_ms": 3,
1513 "status_deadline_ms": 2_000,
1514 "term_grace_ms": 2_000,
1515 "kill_grace_ms": 10_000,
1516 "reap_grace_ms": null,
1517 "remote_deadline_ms": null,
1518 "verified": true,
1519 "already_exited": false,
1520 "forced": false,
1521 "signals": [],
1522 "error": null
1523 });
1524
1525 let evidence: CleanupEvidence = serde_json::from_value(legacy)?;
1526
1527 assert_eq!(evidence.device_residuals, None);
1528 assert!(
1529 serde_json::to_value(&evidence)?
1530 .get("device_residuals")
1531 .is_none()
1532 );
1533 Ok(())
1534 }
1535
1536 #[test]
1537 fn device_residual_evidence_encodes_under_its_outcome_tag()
1538 -> Result<(), Box<dyn std::error::Error>> {
1539 let freed = DeviceResidualEvidence::Freed {
1540 machine: "local".to_owned(),
1541 device: 0,
1542 };
1543 let held = DeviceResidualEvidence::ResidualHeld {
1544 machine: "local".to_owned(),
1545 device: 1,
1546 bytes: 42,
1547 };
1548 let unavailable = DeviceResidualEvidence::ProbeUnavailable {
1549 machine: "node-b".to_owned(),
1550 device: 3,
1551 reason: "nvidia-smi exited".to_owned(),
1552 };
1553
1554 assert_eq!(
1555 serde_json::to_value(&freed)?,
1556 serde_json::json!({"outcome": "freed", "machine": "local", "device": 0})
1557 );
1558 assert_eq!(
1559 serde_json::to_value(&held)?,
1560 serde_json::json!({"outcome": "residual_held", "machine": "local", "device": 1, "bytes": 42})
1561 );
1562 assert_eq!(
1563 serde_json::to_value(&unavailable)?,
1564 serde_json::json!({"outcome": "probe_unavailable", "machine": "node-b", "device": 3, "reason": "nvidia-smi exited"})
1565 );
1566 Ok(())
1567 }
1568
1569 #[test]
1570 fn rejects_a_reused_process_identity() -> Result<(), Box<dyn std::error::Error>> {
1571 let pid = std::process::id();
1572 let actual = process_start_time(pid)
1573 .map_err(std::io::Error::other)?
1574 .ok_or_else(|| std::io::Error::other("test process has no /proc identity"))?;
1575 let recorded = if actual == u64::MAX {
1576 actual - 1
1577 } else {
1578 actual + 1
1579 };
1580 let status = verified_local_status(&HostProcessHandle {
1581 leader_pid: pid,
1582 process_group: pid,
1583 leader_start_time_ticks: recorded,
1584 container: None,
1585 });
1586
1587 assert!(status.queried);
1588 assert!(!status.alive);
1589 assert!(
1590 status
1591 .error
1592 .is_some_and(|error| error.contains("pid was reused"))
1593 );
1594 Ok(())
1595 }
1596}