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}
176
177#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
178#[serde(deny_unknown_fields)]
179pub struct ContainerRemovalEvidence {
180 pub container: String,
181 pub elapsed_ms: u64,
182 pub operation_elapsed_ms: u64,
183 pub deadline_ms: u64,
184 pub client_cleanup: Option<crate::container::CommandCleanupEvidence>,
185 pub confirmed: bool,
186 pub already_absent: bool,
187 pub error: Option<String>,
188}
189
190#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
191#[serde(deny_unknown_fields)]
192pub struct ProcessStatus {
193 pub queried: bool,
194 pub alive: bool,
195 pub error: Option<String>,
196}
197
198#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
199#[serde(deny_unknown_fields)]
200pub struct TargetRegistryMatchEvidence {
201 pub url: String,
202 pub role: String,
203 pub healthy: bool,
204 #[serde(default, skip_serializing_if = "Option::is_none")]
205 pub bootstrap_port: Option<u16>,
206}
207
208#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
209#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
210pub enum ReadinessEvidence {
211 Http {
212 url: String,
213 attempts: u32,
214 ready_unix_ms: u64,
215 timing: OperationTimingEvidence,
216 diagnostic_attempts: Vec<ReadinessAttemptEvidence>,
217 },
218 HttpTargetRegistry {
219 readiness_url: String,
220 registry_url: String,
221 attempts: u32,
222 ready_unix_ms: u64,
223 matched_targets: Vec<TargetRegistryMatchEvidence>,
224 timing: OperationTimingEvidence,
225 diagnostic_attempts: Vec<ReadinessAttemptEvidence>,
226 },
227 ProcessAlive {
228 ready_unix_ms: u64,
229 timing: OperationTimingEvidence,
230 diagnostic_attempts: Vec<ReadinessAttemptEvidence>,
231 },
232}
233
234#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
235#[serde(deny_unknown_fields)]
236pub struct ReadinessAttemptEvidence {
237 pub operation: String,
238 pub effective_bound_ms: u64,
239 pub succeeded: bool,
240 pub error: Option<String>,
241}
242
243#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
244#[serde(rename_all = "snake_case")]
245pub enum ReadinessFailureKind {
246 Exited,
247 Interrupted,
248 Timeout,
249}
250
251#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
252#[serde(deny_unknown_fields)]
253pub struct ReadinessFailure {
254 pub kind: ReadinessFailureKind,
255 pub message: String,
256 pub timing: Option<OperationTimingEvidence>,
257 pub diagnostic_attempts: Vec<ReadinessAttemptEvidence>,
258}
259
260pub struct ProcessSpec<'a> {
261 pub launch: &'a LaunchPlan,
262 pub command: &'a CommandPlan,
263 pub launch_files: &'a [LaunchFilePlan],
264 pub cache_root: &'a Path,
265 pub stdout: &'a Path,
266 pub stderr: &'a Path,
267 pub remote_dir: &'a Path,
268 pub container: Option<&'a str>,
271}
272
273#[derive(Debug, thiserror::Error)]
274pub enum ServerLaunchError {
275 #[error("{message}")]
276 Preparation { message: String },
277 #[error("failed to {operation} {path}: {source}")]
278 FileIo {
279 operation: &'static str,
280 path: PathBuf,
281 #[source]
282 source: std::io::Error,
283 },
284 #[error("failed to launch {program:?}: {source}")]
285 Process {
286 program: String,
287 #[source]
288 source: std::io::Error,
289 },
290 #[error(transparent)]
291 Ssh(#[from] crate::ssh::SshError),
292 #[error("{operation} exited with {status}: {diagnostics}")]
293 Exit {
294 operation: String,
295 status: std::process::ExitStatus,
296 diagnostics: String,
297 },
298 #[error("SSH launch on {target:?} returned non-UTF-8 identity: {source}")]
299 NonUtf8Identity {
300 target: String,
301 #[source]
302 source: std::string::FromUtf8Error,
303 },
304 #[error("SSH launch on {target:?} returned no process id")]
305 MissingProcessId { target: String },
306 #[error("SSH launch on {target:?} returned no process start time")]
307 MissingStartTime { target: String },
308 #[error("SSH launch on {target:?} returned invalid process id {value:?}: {source}")]
309 InvalidProcessId {
310 target: String,
311 value: String,
312 #[source]
313 source: std::num::ParseIntError,
314 },
315 #[error("SSH launch on {target:?} returned invalid process start time {value:?}: {source}")]
316 InvalidStartTime {
317 target: String,
318 value: String,
319 #[source]
320 source: std::num::ParseIntError,
321 },
322 #[error("SSH launch on {target:?} returned an invalid process identity: {details}")]
323 InvalidIdentity { target: String, details: String },
324 #[error("existing launch file target {path} is not a regular file", path = path.display())]
325 NotRegularFile { path: PathBuf },
326 #[error(
327 "existing launch file {path} does not match declared digest {expected}; found {actual}",
328 path = path.display()
329 )]
330 FileDigestMismatch {
331 path: PathBuf,
332 expected: String,
333 actual: String,
334 },
335}
336
337#[derive(Debug)]
338pub struct LaunchFailure {
339 pub error: ServerLaunchError,
340 pub ownership_unknown: bool,
341 pub container_removal: Option<Box<ContainerRemovalEvidence>>,
346 pub cleanup: Option<Box<CleanupEvidence>>,
347 cleanup_note: Option<String>,
348}
349
350impl LaunchFailure {
351 pub fn before_launch(message: String) -> Self {
352 Self {
353 error: ServerLaunchError::Preparation { message },
354 ownership_unknown: false,
355 container_removal: None,
356 cleanup: None,
357 cleanup_note: None,
358 }
359 }
360
361 fn from_error(error: ServerLaunchError) -> Self {
362 Self {
363 error,
364 ownership_unknown: false,
365 container_removal: None,
366 cleanup: None,
367 cleanup_note: None,
368 }
369 }
370
371 pub fn message(&self) -> String {
372 let mut message = self.error.to_string();
373 if let Some(note) = &self.cleanup_note {
374 message = format!("{message}; {note}");
375 } else if let Some(cleanup) = &self.cleanup
376 && let Some(error) = &cleanup.error
377 {
378 message = format!("{message}; local launch cleanup was not verified: {error}");
379 }
380 if let Some(removal) = &self.container_removal {
381 message = format!("{message}; {}", removal_summary(removal));
382 }
383 message
384 }
385
386 pub fn unresolved_ownership(message: String) -> Self {
387 Self {
388 error: ServerLaunchError::Preparation { message },
389 ownership_unknown: true,
390 container_removal: None,
391 cleanup: None,
392 cleanup_note: None,
393 }
394 }
395}
396
397pub trait ProcessLauncher {
398 fn spawn(&self, spec: ProcessSpec<'_>) -> Result<ProcessHandle, LaunchFailure>;
399}
400
401pub trait ProcessObserver {
402 fn status(&self, handle: &ProcessHandle) -> ProcessStatus;
403 fn status_with_bound(&self, handle: &ProcessHandle, bound: &OperationBound) -> ProcessStatus;
404 fn sync_logs(
405 &self,
406 handle: &ProcessHandle,
407 stdout: &Path,
408 stderr: &Path,
409 cleanup: bool,
410 ) -> Result<(), LogSyncError>;
411}
412
413pub trait ReadinessObserver {
414 fn wait_ready(
415 &self,
416 handle: &ProcessHandle,
417 endpoint: &ProcessEndpointPlan,
418 readiness: &ReadinessPlan,
419 bound: &OperationBound,
420 on_probe_failure: &mut dyn FnMut(&str),
421 ) -> Result<ReadinessEvidence, ReadinessFailure>;
422}
423
424pub trait ProcessCleanup {
425 fn terminate(
426 &self,
427 handle: &ProcessHandle,
428 trigger: CleanupTrigger,
429 on_container_removal: &mut dyn FnMut(&str),
430 ) -> CleanupEvidence;
431}
432
433pub trait ServerRuntime:
434 ProcessLauncher + ProcessObserver + ReadinessObserver + ProcessCleanup
435{
436}
437
438impl<T> ServerRuntime for T where
439 T: ProcessLauncher + ProcessObserver + ReadinessObserver + ProcessCleanup
440{
441}
442
443#[derive(Debug, thiserror::Error)]
444pub enum ProcessCommandError {
445 #[error("{operation} deadline expired")]
446 Deadline { operation: String },
447 #[error("{operation} was interrupted")]
448 Interrupted { operation: String },
449 #[error("failed to launch {operation}: {source}")]
450 Launch {
451 operation: String,
452 #[source]
453 source: std::io::Error,
454 },
455 #[error("failed to launch {operation}: {source}")]
456 Ssh {
457 operation: String,
458 #[source]
459 source: crate::ssh::SshError,
460 },
461 #[error("{operation} failed: {source}")]
462 Io {
463 operation: String,
464 #[source]
465 source: std::io::Error,
466 },
467 #[error("{operation} wait failed: {source}; child cleanup: {cleanup}")]
468 WaitCleanup {
469 operation: String,
470 #[source]
471 source: std::io::Error,
472 cleanup: String,
473 },
474 #[error("{operation} exited with {status}: {stderr}")]
475 Exit {
476 operation: String,
477 status: std::process::ExitStatus,
478 stderr: String,
479 },
480}
481
482#[derive(Debug, thiserror::Error)]
483pub enum LogSyncError {
484 #[error("failed to read remote log {path}: {source}")]
485 ReadRemote {
486 path: PathBuf,
487 #[source]
488 source: ProcessCommandError,
489 },
490 #[error("failed to read remote log {path}: command exited with {status}: {stderr}")]
491 RemoteExit {
492 path: PathBuf,
493 status: std::process::ExitStatus,
494 stderr: String,
495 },
496 #[error("failed to write local log {path}: {source}")]
497 WriteLocal {
498 path: PathBuf,
499 #[source]
500 source: std::io::Error,
501 },
502}
503
504#[derive(Clone, Copy, Debug, Default)]
505pub struct SystemProcessRuntime;
506
507#[cfg(test)]
508mod tests {
509 use super::*;
510 use crate::plan::LaunchFilePlan;
511 use std::cell::Cell;
512 use std::io::{BufRead, BufReader};
513 use std::net::TcpListener;
514 use std::os::unix::fs::{MetadataExt, PermissionsExt};
515
516 #[test]
517 fn expired_readiness_owner_prevents_a_fresh_network_attempt() {
518 let bound = OperationBound::finite(Duration::ZERO);
519 let error = probe_http("127.0.0.1", 9, "/ready", &bound, 30)
520 .err()
521 .map(|error| error.to_string())
522 .unwrap_or_default();
523
524 assert_eq!(error, "readiness operation deadline expired");
525 }
526
527 #[test]
528 fn readiness_attempt_deadline_bounds_a_trickled_status_line() -> Result<(), String> {
529 let listener = TcpListener::bind(("127.0.0.1", 0)).map_err(|error| error.to_string())?;
530 let port = listener
531 .local_addr()
532 .map_err(|error| error.to_string())?
533 .port();
534 let server = thread::spawn(move || -> Result<(), String> {
535 let (mut stream, _) = listener.accept().map_err(|error| error.to_string())?;
536 let mut request = [0_u8; 1024];
537 let _ = stream.read(&mut request);
538 let response = format!("HTTP/1.1 200 {}\r\n", " ".repeat(96));
539 for byte in response.bytes() {
540 if stream.write_all(&[byte]).is_err() {
541 break;
542 }
543 thread::sleep(Duration::from_millis(20));
544 }
545 Ok(())
546 });
547
548 let started = Instant::now();
549 let error = probe_http("127.0.0.1", port, "/ready", &OperationBound::unbounded(), 1)
550 .err()
551 .map(|error| error.to_string())
552 .unwrap_or_default();
553 let elapsed = started.elapsed();
554 server
555 .join()
556 .map_err(|_| "trickle fixture panicked".to_owned())??;
557
558 assert!(error.contains("deadline expired"), "{error}");
559 assert!(
560 elapsed < Duration::from_secs(2),
561 "a one-second attempt lasted {elapsed:?}"
562 );
563 Ok(())
564 }
565
566 #[test]
567 fn finite_readiness_accepts_a_response_after_250_milliseconds() -> Result<(), String> {
568 let listener = TcpListener::bind(("127.0.0.1", 0)).map_err(|error| error.to_string())?;
569 let port = listener
570 .local_addr()
571 .map_err(|error| error.to_string())?
572 .port();
573 let server = thread::spawn(move || -> Result<(), String> {
574 let (mut stream, _) = listener.accept().map_err(|error| error.to_string())?;
575 let mut request = [0_u8; 1024];
576 let _ = stream.read(&mut request);
577 thread::sleep(Duration::from_millis(350));
578 stream
579 .write_all(b"HTTP/1.1 204 No Content\r\nContent-Length: 0\r\n\r\n")
580 .map_err(|error| error.to_string())
581 });
582
583 let started = Instant::now();
584 probe_http(
585 "127.0.0.1",
586 port,
587 "/ready",
588 &OperationBound::finite(Duration::from_secs(2)),
589 1,
590 )
591 .map_err(|error| error.to_string())?;
592 let elapsed = started.elapsed();
593 server
594 .join()
595 .map_err(|_| "delayed readiness fixture panicked".to_owned())??;
596
597 assert!(elapsed >= Duration::from_millis(250), "elapsed {elapsed:?}");
598 assert!(elapsed < Duration::from_secs(2), "elapsed {elapsed:?}");
599 Ok(())
600 }
601
602 #[test]
603 fn readiness_attempt_deadline_includes_the_complete_response_body() -> Result<(), String> {
604 let listener = TcpListener::bind(("127.0.0.1", 0)).map_err(|error| error.to_string())?;
605 let port = listener
606 .local_addr()
607 .map_err(|error| error.to_string())?
608 .port();
609 let server = thread::spawn(move || -> Result<(), String> {
610 let (mut stream, _) = listener.accept().map_err(|error| error.to_string())?;
611 let mut request = [0_u8; 1024];
612 let _ = stream.read(&mut request);
613 stream
614 .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 10\r\n\r\nx")
615 .map_err(|error| error.to_string())?;
616 thread::sleep(Duration::from_millis(1_500));
617 Ok(())
618 });
619
620 let started = Instant::now();
621 let error = probe_http("127.0.0.1", port, "/ready", &OperationBound::unbounded(), 1)
622 .err()
623 .map(|error| error.to_string());
624 let elapsed = started.elapsed();
625 server
626 .join()
627 .map_err(|_| "readiness body fixture panicked".to_owned())??;
628
629 assert!(
630 error.is_some_and(|error| error == "readiness operation deadline expired"),
631 "readiness accepted an incomplete response body"
632 );
633 assert!(
634 elapsed < Duration::from_millis(1_500),
635 "elapsed {elapsed:?}"
636 );
637 Ok(())
638 }
639
640 #[test]
641 fn registry_attempt_deadline_bounds_a_trickled_body() -> Result<(), String> {
642 let listener = TcpListener::bind(("127.0.0.1", 0)).map_err(|error| error.to_string())?;
643 let port = listener
644 .local_addr()
645 .map_err(|error| error.to_string())?
646 .port();
647 let server = thread::spawn(move || -> Result<(), String> {
648 let (mut stream, _) = listener.accept().map_err(|error| error.to_string())?;
649 let mut request = [0_u8; 1024];
650 let _ = stream.read(&mut request);
651 stream
652 .write_all(b"HTTP/1.1 200 OK\r\nConnection: close\r\n\r\n")
653 .map_err(|error| error.to_string())?;
654 let body = format!("{}{{}}", " ".repeat(96));
655 for byte in body.bytes() {
656 if stream.write_all(&[byte]).is_err() {
657 break;
658 }
659 thread::sleep(Duration::from_millis(20));
660 }
661 Ok(())
662 });
663
664 let started = Instant::now();
665 let error = probe_http_json(
666 "127.0.0.1",
667 port,
668 "/workers",
669 "target registry",
670 &OperationBound::unbounded(),
671 1,
672 )
673 .err()
674 .map(|error| error.to_string())
675 .unwrap_or_default();
676 let elapsed = started.elapsed();
677 server
678 .join()
679 .map_err(|_| "registry trickle fixture panicked".to_owned())??;
680
681 assert!(error.contains("deadline expired"), "{error}");
682 assert!(
683 elapsed < Duration::from_secs(2),
684 "a one-second registry attempt lasted {elapsed:?}"
685 );
686 Ok(())
687 }
688
689 #[test]
690 fn expired_readiness_owner_rejects_a_registry_match() {
691 let expected = vec![TargetRegistryExpectedTarget {
692 url: "http://decode:30001".to_owned(),
693 role: "decode".to_owned(),
694 bootstrap_port: None,
695 }];
696 let response = serde_json::json!({
697 "workers": [{
698 "url": "http://decode:30001",
699 "worker_type": "decode",
700 "is_healthy": true
701 }]
702 });
703
704 let error = match_target_registry(
705 &response,
706 &target_registry_probe(&expected),
707 &OperationBound::finite(Duration::ZERO),
708 )
709 .err()
710 .map(|error| error.to_string())
711 .unwrap_or_default();
712
713 assert_eq!(error, "readiness operation deadline expired");
714 }
715
716 #[test]
717 fn process_status_command_cannot_outlive_the_readiness_owner() {
718 let started = Instant::now();
719 let error = run_status_command(
720 &["sh", "-c", "sleep 5"],
721 &OperationBound::finite(Duration::from_millis(50)),
722 )
723 .err()
724 .map(|error| error.to_string())
725 .unwrap_or_default();
726
727 assert_eq!(error, "process status attempt deadline expired");
728 assert!(
729 started.elapsed() < Duration::from_secs(1),
730 "bounded process status did not stop promptly"
731 );
732 }
733
734 #[test]
735 fn finite_readiness_retries_an_expired_process_status_attempt() -> Result<(), String> {
736 let calls = Cell::new(0_u32);
737 let mut failures = Vec::new();
738 let bound = OperationBound::finite(Duration::from_secs(3));
739 let evidence = wait_process_alive_ready(
740 |_| {
741 calls.set(calls.get() + 1);
742 if calls.get() == 1 {
743 thread::sleep(Duration::from_millis(1_050));
744 ProcessStatus {
745 queried: false,
746 alive: false,
747 error: Some("process status attempt deadline expired".to_owned()),
748 }
749 } else {
750 alive_status()
751 }
752 },
753 1,
754 &bound,
755 &mut |failure| failures.push(failure.to_owned()),
756 )
757 .map_err(|failure| failure.message)?;
758
759 assert_eq!(calls.get(), 2);
760 assert_eq!(failures, ["process status attempt deadline expired"]);
761 let ReadinessEvidence::ProcessAlive {
762 timing,
763 diagnostic_attempts,
764 ..
765 } = evidence
766 else {
767 return Err("process-alive readiness returned the wrong evidence kind".to_owned());
768 };
769 assert_eq!(
770 timing.start_boundary,
771 "after_process_spawn_before_readiness_attempt"
772 );
773 assert_eq!(diagnostic_attempts.len(), 1);
774 assert!(diagnostic_attempts[0].succeeded);
775 assert!((1..=1_000).contains(&diagnostic_attempts[0].effective_bound_ms));
776 Ok(())
777 }
778
779 #[test]
780 fn unbounded_process_status_command_does_not_acquire_a_timeout() -> Result<(), String> {
781 let output = run_status_command(
782 &["sh", "-c", "sleep 0.1; printf alive"],
783 &OperationBound::unbounded(),
784 )
785 .map_err(|error| error.to_string())?;
786
787 assert!(output.status.success());
788 assert_eq!(output.stdout, b"alive");
789 Ok(())
790 }
791
792 fn launch_file(root: &Path, text: &str, name: &str) -> LaunchFilePlan {
793 let sha256 = format!("{:x}", Sha256::digest(text.as_bytes()));
794 let relative_path = format!("launch-files/{sha256}/{name}");
795 LaunchFilePlan {
796 resolved_path: root.join(&relative_path),
797 relative_path,
798 text: text.to_owned(),
799 sha256,
800 }
801 }
802
803 fn run_script_with_input(script: &str, input: &[u8]) -> Result<Output, String> {
804 match crate::container::run_with_bound(
805 &["bash", "-c", script],
806 None,
807 Some(input),
808 &OperationBound::unbounded(),
809 None,
810 ) {
811 Ok(crate::container::BoundedWait::Exited {
812 status,
813 stdout,
814 stderr,
815 }) => Ok(Output {
816 status,
817 stdout,
818 stderr,
819 }),
820 Ok(crate::container::BoundedWait::Expired { .. }) => {
821 Err("unbounded launch-file fixture expired".to_owned())
822 }
823 Ok(crate::container::BoundedWait::Interrupted { .. }) => {
824 Err("launch-file fixture was interrupted".to_owned())
825 }
826 Err(_) => Err("launch-file fixture failed".to_owned()),
827 }
828 }
829
830 fn target_registry_endpoint(
831 registry_body: String,
832 ) -> Result<(ProcessEndpointPlan, thread::JoinHandle<Result<(), String>>), String> {
833 let listener = TcpListener::bind(("127.0.0.1", 0)).map_err(|error| error.to_string())?;
834 let port = listener
835 .local_addr()
836 .map_err(|error| error.to_string())?
837 .port();
838 let server = thread::spawn(move || {
839 for _ in 0..2 {
840 let (mut stream, _) = listener.accept().map_err(|error| error.to_string())?;
841 let mut request_line = String::new();
842 let mut reader =
843 BufReader::new(stream.try_clone().map_err(|error| error.to_string())?);
844 reader
845 .read_line(&mut request_line)
846 .map_err(|error| error.to_string())?;
847 loop {
848 let mut header = String::new();
849 reader
850 .read_line(&mut header)
851 .map_err(|error| error.to_string())?;
852 if header == "\r\n" || header.is_empty() {
853 break;
854 }
855 }
856 let body = if request_line.starts_with("GET /workers ") {
857 registry_body.as_bytes()
858 } else {
859 b""
860 };
861 let mut response = format!(
862 "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
863 body.len()
864 )
865 .into_bytes();
866 response.extend_from_slice(body);
867 stream
868 .write_all(&response)
869 .map_err(|error| error.to_string())?;
870 }
871 Ok(())
872 });
873 Ok((
874 ProcessEndpointPlan {
875 host: "127.0.0.1".to_owned(),
876 port,
877 },
878 server,
879 ))
880 }
881
882 fn target_registry_probe<'a>(
883 expected_targets: &'a [TargetRegistryExpectedTarget],
884 ) -> HttpTargetRegistryProbe<'a> {
885 HttpTargetRegistryProbe {
886 readiness_path: "/readiness",
887 registry_path: "/workers",
888 targets_field: "workers",
889 target_url_field: "url",
890 target_role_field: "worker_type",
891 target_healthy_field: "is_healthy",
892 target_bootstrap_port_field: "bootstrap_port",
893 expected_targets,
894 }
895 }
896
897 fn alive_status() -> ProcessStatus {
898 ProcessStatus {
899 queried: true,
900 alive: true,
901 error: None,
902 }
903 }
904
905 #[test]
906 fn local_launch_file_publication_reuses_the_immutable_target() -> Result<(), String> {
907 let root = tempfile::tempdir().map_err(|error| error.to_string())?;
908 let launch_file = launch_file(
909 root.path(),
910 "worker: \u{2603}\nmode: context\n",
911 "worker.yaml",
912 );
913
914 materialize_local_launch_files(std::slice::from_ref(&launch_file))
915 .map_err(|error| error.to_string())?;
916 let first_metadata =
917 fs::metadata(&launch_file.resolved_path).map_err(|error| error.to_string())?;
918 materialize_local_launch_files(std::slice::from_ref(&launch_file))
919 .map_err(|error| error.to_string())?;
920 let second_metadata =
921 fs::metadata(&launch_file.resolved_path).map_err(|error| error.to_string())?;
922
923 assert_eq!(
924 fs::read_to_string(&launch_file.resolved_path).map_err(|error| error.to_string())?,
925 launch_file.text
926 );
927 assert_eq!(first_metadata.ino(), second_metadata.ino());
928 assert_eq!(second_metadata.permissions().mode() & 0o222, 0);
929 Ok(())
930 }
931
932 #[test]
933 fn local_launch_file_mismatch_fails_before_spawn_without_replacing_it() -> Result<(), String> {
934 let root = tempfile::tempdir().map_err(|error| error.to_string())?;
935 let cache = root.path().join("cache");
936 let launch_file = launch_file(&cache, "expected\n", "worker.yaml");
937 let parent = launch_file
938 .resolved_path
939 .parent()
940 .ok_or_else(|| "launch file has no parent".to_owned())?;
941 fs::create_dir_all(parent).map_err(|error| error.to_string())?;
942 fs::write(&launch_file.resolved_path, "stale\n").map_err(|error| error.to_string())?;
943 let marker = root.path().join("spawned");
944 let command = CommandPlan {
945 argv: vec![
946 "sh".to_owned(),
947 "-c".to_owned(),
948 format!("printf launched > {}", shell_quote_path(&marker)),
949 ],
950 env: BTreeMap::new(),
951 explicit_env: Vec::new(),
952 pass_env: Vec::new(),
953 cwd: root.path().to_path_buf(),
954 };
955 let launch_files = vec![launch_file.clone()];
956
957 let result = spawn_local(ProcessSpec {
958 launch: &LaunchPlan::Local,
959 command: &command,
960 launch_files: &launch_files,
961 cache_root: &cache,
962 stdout: &root.path().join("stdout.log"),
963 stderr: &root.path().join("stderr.log"),
964 remote_dir: &root.path().join("remote"),
965 container: None,
966 });
967
968 let failure = match result {
969 Err(failure) => failure,
970 Ok(handle) => {
971 let _ = terminate_local(&handle, CleanupTrigger::StartupRollback);
972 return Err("mismatched launch file unexpectedly spawned a process".to_owned());
973 }
974 };
975 assert!(!failure.ownership_unknown, "{failure:?}");
976 assert!(failure.message().contains("does not match"), "{failure:?}");
977 assert!(!marker.exists());
978 assert_eq!(
979 fs::read_to_string(&launch_file.resolved_path).map_err(|error| error.to_string())?,
980 "stale\n"
981 );
982 Ok(())
983 }
984
985 #[test]
986 fn remote_launch_file_script_publishes_stdin_without_replacing_targets() -> Result<(), String> {
987 let root = tempfile::tempdir().map_err(|error| error.to_string())?;
988 let published = launch_file(
989 root.path(),
990 "worker: \u{96ea}\nmode: context\n",
991 "worker.yaml",
992 );
993 let script = remote_launch_file_script(&published).map_err(|error| error.to_string())?;
994
995 let first = run_script_with_input(&script, published.text.as_bytes())?;
996 assert!(
997 first.status.success(),
998 "{}",
999 String::from_utf8_lossy(&first.stderr)
1000 );
1001 let first_metadata =
1002 fs::metadata(&published.resolved_path).map_err(|error| error.to_string())?;
1003 let first_inode = first_metadata.ino();
1004 assert_eq!(first_metadata.permissions().mode() & 0o222, 0);
1005 let reused = run_script_with_input(&script, published.text.as_bytes())?;
1006 assert!(
1007 reused.status.success(),
1008 "{}",
1009 String::from_utf8_lossy(&reused.stderr)
1010 );
1011 assert_eq!(
1012 fs::read_to_string(&published.resolved_path).map_err(|error| error.to_string())?,
1013 published.text
1014 );
1015 assert_eq!(
1016 fs::metadata(&published.resolved_path)
1017 .map_err(|error| error.to_string())?
1018 .ino(),
1019 first_inode
1020 );
1021
1022 let corrupt = launch_file(root.path(), "expected\n", "corrupt.yaml");
1023 let corrupt_parent = corrupt
1024 .resolved_path
1025 .parent()
1026 .ok_or_else(|| "launch file has no parent".to_owned())?;
1027 fs::create_dir_all(corrupt_parent).map_err(|error| error.to_string())?;
1028 fs::write(&corrupt.resolved_path, "stale\n").map_err(|error| error.to_string())?;
1029 let rejected = run_script_with_input(
1030 &remote_launch_file_script(&corrupt).map_err(|error| error.to_string())?,
1031 corrupt.text.as_bytes(),
1032 )?;
1033 assert!(!rejected.status.success());
1034 assert_eq!(
1035 fs::read_to_string(&corrupt.resolved_path).map_err(|error| error.to_string())?,
1036 "stale\n"
1037 );
1038 Ok(())
1039 }
1040
1041 #[test]
1042 fn target_registry_readiness_records_all_expected_targets() -> Result<(), String> {
1043 let (endpoint, server) = target_registry_endpoint(
1044 serde_json::json!({
1045 "workers": [
1046 {
1047 "url": "http://prefill:30000",
1048 "worker_type": "prefill",
1049 "is_healthy": true,
1050 "bootstrap_port": 8998
1051 },
1052 {
1053 "url": "http://decode:30001",
1054 "worker_type": "decode",
1055 "is_healthy": true
1056 }
1057 ]
1058 })
1059 .to_string(),
1060 )?;
1061 let expected = vec![
1062 TargetRegistryExpectedTarget {
1063 url: "http://prefill:30000".to_owned(),
1064 role: "prefill".to_owned(),
1065 bootstrap_port: Some(8998),
1066 },
1067 TargetRegistryExpectedTarget {
1068 url: "http://decode:30001".to_owned(),
1069 role: "decode".to_owned(),
1070 bootstrap_port: None,
1071 },
1072 ];
1073
1074 let bound = OperationBound::unbounded();
1075 let evidence = wait_http_target_registry_ready(
1076 |_| alive_status(),
1077 &endpoint,
1078 target_registry_probe(&expected),
1079 1,
1080 &bound,
1081 &mut |_| {},
1082 )
1083 .map_err(|failure| failure.message)?;
1084 server
1085 .join()
1086 .map_err(|_| "target registry fixture panicked".to_owned())??;
1087
1088 let record_value = serde_json::to_value(&evidence).map_err(|error| error.to_string())?;
1089 assert_eq!(record_value["kind"], "http_target_registry");
1090 assert_eq!(
1091 record_value["matched_targets"].as_array().map(Vec::len),
1092 Some(2)
1093 );
1094 let ReadinessEvidence::HttpTargetRegistry {
1095 readiness_url,
1096 registry_url,
1097 attempts,
1098 matched_targets,
1099 timing,
1100 diagnostic_attempts,
1101 ready_unix_ms: _,
1102 } = evidence
1103 else {
1104 return Err("target registry readiness returned the wrong evidence kind".to_owned());
1105 };
1106 assert_eq!(
1107 readiness_url,
1108 format!("http://127.0.0.1:{}/readiness", endpoint.port)
1109 );
1110 assert_eq!(
1111 registry_url,
1112 format!("http://127.0.0.1:{}/workers", endpoint.port)
1113 );
1114 assert_eq!(attempts, 1);
1115 assert_eq!(
1116 timing.budget,
1117 crate::operation_bound::OperationBudgetEvidence::Unbounded
1118 );
1119 assert_eq!(
1120 timing.terminal_cause,
1121 crate::operation_bound::OperationTerminalCause::Succeeded
1122 );
1123 assert_eq!(diagnostic_attempts.len(), 3);
1124 assert!(diagnostic_attempts.iter().all(|attempt| {
1125 attempt.succeeded && (1..=1_000).contains(&attempt.effective_bound_ms)
1126 }));
1127 assert_eq!(
1128 matched_targets,
1129 vec![
1130 TargetRegistryMatchEvidence {
1131 url: "http://prefill:30000".to_owned(),
1132 role: "prefill".to_owned(),
1133 healthy: true,
1134 bootstrap_port: Some(8998),
1135 },
1136 TargetRegistryMatchEvidence {
1137 url: "http://decode:30001".to_owned(),
1138 role: "decode".to_owned(),
1139 healthy: true,
1140 bootstrap_port: None,
1141 },
1142 ]
1143 );
1144 Ok(())
1145 }
1146
1147 #[test]
1148 fn finite_target_registry_attempts_record_the_resolved_attempt_budget() -> Result<(), String> {
1149 let (endpoint, server) = target_registry_endpoint(
1150 serde_json::json!({
1151 "workers": [{
1152 "url": "http://decode:30001",
1153 "worker_type": "decode",
1154 "is_healthy": true
1155 }]
1156 })
1157 .to_string(),
1158 )?;
1159 let expected = vec![TargetRegistryExpectedTarget {
1160 url: "http://decode:30001".to_owned(),
1161 role: "decode".to_owned(),
1162 bootstrap_port: None,
1163 }];
1164
1165 let bound = OperationBound::finite(Duration::from_secs(2));
1166 let evidence = wait_http_target_registry_ready(
1167 |_| alive_status(),
1168 &endpoint,
1169 target_registry_probe(&expected),
1170 1,
1171 &bound,
1172 &mut |_| {},
1173 )
1174 .map_err(|failure| failure.message)?;
1175 server
1176 .join()
1177 .map_err(|_| "target registry fixture panicked".to_owned())??;
1178
1179 let ReadinessEvidence::HttpTargetRegistry {
1180 timing,
1181 diagnostic_attempts,
1182 ..
1183 } = evidence
1184 else {
1185 return Err("target registry readiness returned the wrong evidence kind".to_owned());
1186 };
1187 assert_eq!(
1188 timing.budget,
1189 crate::operation_bound::OperationBudgetEvidence::Finite {
1190 configured_ms: 2_000,
1191 }
1192 );
1193 assert_eq!(diagnostic_attempts.len(), 3);
1194 assert!(diagnostic_attempts.iter().all(|attempt| {
1195 attempt.succeeded && (1..=1_000).contains(&attempt.effective_bound_ms)
1196 }));
1197 Ok(())
1198 }
1199
1200 #[test]
1201 fn target_registry_readiness_rejects_partial_registration() -> Result<(), String> {
1202 let (endpoint, server) = target_registry_endpoint(
1203 serde_json::json!({
1204 "workers": [{
1205 "url": "http://prefill:30000",
1206 "worker_type": "prefill",
1207 "is_healthy": true,
1208 "bootstrap_port": 8998
1209 }]
1210 })
1211 .to_string(),
1212 )?;
1213 let expected = vec![
1214 TargetRegistryExpectedTarget {
1215 url: "http://prefill:30000".to_owned(),
1216 role: "prefill".to_owned(),
1217 bootstrap_port: Some(8998),
1218 },
1219 TargetRegistryExpectedTarget {
1220 url: "http://decode:30001".to_owned(),
1221 role: "decode".to_owned(),
1222 bootstrap_port: None,
1223 },
1224 ];
1225
1226 let mut probe_failures = Vec::new();
1227 let bound = OperationBound::finite(Duration::from_secs(1));
1228 let failure = match wait_http_target_registry_ready(
1229 |_| alive_status(),
1230 &endpoint,
1231 target_registry_probe(&expected),
1232 1,
1233 &bound,
1234 &mut |failure| probe_failures.push(failure.to_owned()),
1235 ) {
1236 Err(failure) => failure,
1237 Ok(evidence) => {
1238 return Err(format!(
1239 "partial target registration unexpectedly became ready: {evidence:?}"
1240 ));
1241 }
1242 };
1243 server
1244 .join()
1245 .map_err(|_| "target registry fixture panicked".to_owned())??;
1246
1247 assert_eq!(failure.kind, ReadinessFailureKind::Timeout);
1248 let timing = failure
1249 .timing
1250 .as_ref()
1251 .ok_or_else(|| "readiness timeout has no timing evidence".to_owned())?;
1252 assert_eq!(
1253 timing.budget,
1254 crate::operation_bound::OperationBudgetEvidence::Finite {
1255 configured_ms: 1_000,
1256 }
1257 );
1258 assert_eq!(timing.terminal_cause, OperationTerminalCause::TimedOut);
1259 assert!(probe_failures.iter().any(|failure| {
1260 failure.contains("target registry has no \"decode\" target at \"http://decode:30001\"")
1261 }));
1262 Ok(())
1263 }
1264
1265 #[test]
1266 fn termination_waits_for_the_group_after_the_launcher_exits() -> Result<(), String> {
1267 let mut child = Command::new("sh")
1268 .args([
1269 "-c",
1270 "trap 'exit 0' TERM; sh -c 'trap \"\" TERM; exec sleep 30' & wait",
1271 ])
1272 .stdin(Stdio::null())
1273 .stdout(Stdio::null())
1274 .stderr(Stdio::null())
1275 .process_group(0)
1276 .spawn()
1277 .map_err(|error| error.to_string())?;
1278 let handle = HostProcessHandle::new(child.id(), None)?;
1279 thread::sleep(Duration::from_millis(100));
1280 let reaper = thread::spawn(move || child.wait());
1281
1282 let cleanup = terminate_local(&handle, CleanupTrigger::Stop);
1283 if !cleanup.verified {
1284 let _ = Command::new("kill")
1285 .args(["-KILL", "--", &format!("-{}", handle.process_group)])
1286 .status();
1287 }
1288 let _ = reaper.join();
1289
1290 assert!(cleanup.verified, "{cleanup:?}");
1291 assert!(cleanup.forced);
1292 assert!(cleanup.elapsed_ms >= cleanup.term_grace_ms);
1293 assert_eq!(cleanup.status_deadline_ms, 2_000);
1294 assert_eq!(cleanup.term_grace_ms, 2_000);
1295 assert_eq!(cleanup.kill_grace_ms, 10_000);
1296 Ok(())
1297 }
1298
1299 #[test]
1300 fn rejects_a_reused_process_identity() -> Result<(), Box<dyn std::error::Error>> {
1301 let pid = std::process::id();
1302 let actual = process_start_time(pid)
1303 .map_err(std::io::Error::other)?
1304 .ok_or_else(|| std::io::Error::other("test process has no /proc identity"))?;
1305 let recorded = if actual == u64::MAX {
1306 actual - 1
1307 } else {
1308 actual + 1
1309 };
1310 let status = verified_local_status(&HostProcessHandle {
1311 leader_pid: pid,
1312 process_group: pid,
1313 leader_start_time_ticks: recorded,
1314 container: None,
1315 });
1316
1317 assert!(status.queried);
1318 assert!(!status.alive);
1319 assert!(
1320 status
1321 .error
1322 .is_some_and(|error| error.contains("pid was reused"))
1323 );
1324 Ok(())
1325 }
1326}