Skip to main content

inferlab_runtime/
server.rs

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    /// The container this process launched, when the command is a
57    /// containerized substitution: the daemon-owned cleanup handle the
58    /// process-group kill cannot reach ([[RFC-0003:C-RUNTIME-WORKFLOWS]]).
59    #[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    /// The container this process launched, when the command is a
101    /// containerized substitution: the daemon-owned cleanup handle the
102    /// process-group kill cannot reach ([[RFC-0003:C-RUNTIME-WORKFLOWS]]).
103    #[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    /// Confirmed removal of the process's container on its launch machine
169    /// ([[RFC-0003:C-RUNTIME-WORKFLOWS]]); present when the cleanup path
170    /// attempted to remove a known container — from a running server's
171    /// handle, or from a launch failure whose command already named one
172    /// before any handle existed.
173    #[serde(default, skip_serializing_if = "Option::is_none")]
174    pub container_removal: Option<ContainerRemovalEvidence>,
175    /// Per-device residual compute-memory probes on the process's launch
176    /// machine, run after its process cleanup verified
177    /// ([[RFC-0005:C-EVIDENCE]]); present when the cleanup path probed the
178    /// assigned devices — a device-less process or an unverified cleanup
179    /// carries no probe claims.
180    #[serde(default, skip_serializing_if = "Option::is_none")]
181    pub device_residuals: Option<Vec<DeviceResidualEvidence>>,
182}
183
184/// The post-cleanup probe outcome for one assigned device
185/// ([[RFC-0005:C-EVIDENCE]]): a dead framework can leave compute memory held
186/// after its process group is reaped, so a verified cleanup probes each
187/// assigned device and records what it found — freed, still held (with the
188/// observed bytes), or undeterminable on a machine without the probe tool
189/// (which must not fail the cleanup).
190#[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    /// The resolver-assigned container name when the command is a
301    /// containerized substitution.
302    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    /// The structured outcome of removing the container this launch may
374    /// have created, when the failure attempted one; the record's cleanup
375    /// evidence carries the actual container and reason rather than a
376    /// generic note ([[RFC-0003:C-RUNTIME-WORKFLOWS]]).
377    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    /// A process group whose leader forks a member and is then SIGKILLed,
1340    /// leaving the member orphaned: the managed engine died on its own and
1341    /// the workers outlived it ([[RFC-0003:C-RUNTIME-WORKFLOWS]] cleanup).
1342    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        // A handle whose recorded leader start postdates the actual member:
1407        // the pgid reads as recycled, so ownership stays unverifiable.
1408        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        // The fixture ssh shim execs the remote script through bash; running
1443        // it directly exercises the same bytes against local /proc and ps.
1444        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        // A schema-10-era entry predates the device-residual probe: the
1508        // optional member decodes as absent and stays absent on re-encode
1509        // ([[RFC-0005:C-EVIDENCE]]).
1510        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}