1use super::observation::{remote_group_alive_script, run_cleanup_command};
2use super::{
3 CleanupEvidence, CleanupTrigger, ContainerRemovalEvidence, HostProcessHandle, ProcessCleanup,
4 ProcessHandle, SshProcessHandle, SystemProcessRuntime,
5};
6use crate::operation_bound::{OperationBound, duration_millis};
7use crate::process_group::{LocalProcessGroup, SignalEvidence, TerminationSignal, VerifiedStatus};
8use crate::ssh::{SSH_ENV_REMOVE, ssh_argv};
9use std::time::{Duration, Instant};
10use wait_timeout::ChildExt;
11
12const POLL_INTERVAL: Duration = Duration::from_millis(100);
13const TERM_GRACE: Duration = Duration::from_secs(2);
14const KILL_GRACE: Duration = Duration::from_secs(10);
15const SERVER_CLEANUP_STATUS_DEADLINE: Duration = Duration::from_secs(2);
16pub(super) const REMOTE_SERVER_CLEANUP_DEADLINE: Duration = Duration::from_secs(30);
17const LOCAL_LAUNCH_FAILURE_REAP_GRACE: Duration = Duration::from_secs(5);
18pub(super) const TERM_POLL_LIMIT: u128 = TERM_GRACE.as_millis() / POLL_INTERVAL.as_millis();
19pub(super) const KILL_POLL_LIMIT: u128 = KILL_GRACE.as_millis() / POLL_INTERVAL.as_millis();
20const CLEANUP_MARKER: &str = "INFERLAB_CLEANUP\t";
21
22impl CleanupEvidence {
23 pub fn unavailable(trigger: CleanupTrigger, message: String) -> Self {
24 Self {
25 trigger,
26 elapsed_ms: 0,
27 status_deadline_ms: duration_millis(SERVER_CLEANUP_STATUS_DEADLINE),
28 term_grace_ms: duration_millis(TERM_GRACE),
29 kill_grace_ms: duration_millis(KILL_GRACE),
30 reap_grace_ms: None,
31 remote_deadline_ms: None,
32 verified: false,
33 already_exited: false,
34 forced: false,
35 signals: Vec::new(),
36 error: Some(message),
37 container_removal: None,
38 device_residuals: None,
39 }
40 }
41
42 pub fn from_launch_removal(
49 trigger: CleanupTrigger,
50 verified: bool,
51 removal: ContainerRemovalEvidence,
52 error: Option<String>,
53 ) -> Self {
54 Self {
55 trigger,
56 elapsed_ms: removal.elapsed_ms,
57 status_deadline_ms: duration_millis(SERVER_CLEANUP_STATUS_DEADLINE),
58 term_grace_ms: duration_millis(TERM_GRACE),
59 kill_grace_ms: duration_millis(KILL_GRACE),
60 reap_grace_ms: None,
61 remote_deadline_ms: None,
62 verified,
63 already_exited: false,
64 forced: false,
65 signals: Vec::new(),
66 error,
67 container_removal: Some(removal),
68 device_residuals: None,
69 }
70 }
71}
72
73pub(super) fn cleanup_failed_local_launch(child: &mut std::process::Child) -> CleanupEvidence {
74 let started = Instant::now();
75 let initial_status_error = match child.try_wait() {
76 Ok(Some(_)) => {
77 let mut evidence =
78 completed_cleanup(CleanupTrigger::StartupRollback, true, false, Vec::new());
79 evidence.elapsed_ms = duration_millis(started.elapsed());
80 evidence.status_deadline_ms = 0;
81 evidence.term_grace_ms = 0;
82 evidence.reap_grace_ms = Some(duration_millis(LOCAL_LAUNCH_FAILURE_REAP_GRACE));
83 return evidence;
84 }
85 Ok(None) => None,
86 Err(error) => Some(format!("failed to inspect failed launch child: {error}")),
90 };
91 let group = match LocalProcessGroup::capture_child(child) {
92 Ok(group) => group,
93 Err(error) => {
94 let mut evidence = CleanupEvidence::unavailable(
95 CleanupTrigger::StartupRollback,
96 format!("failed to capture failed launch process-group identity: {error}"),
97 );
98 evidence.elapsed_ms = duration_millis(started.elapsed());
99 evidence.status_deadline_ms = 0;
100 evidence.term_grace_ms = 0;
101 evidence.reap_grace_ms = Some(duration_millis(LOCAL_LAUNCH_FAILURE_REAP_GRACE));
102 return evidence;
103 }
104 };
105 let bound = OperationBound::finite(KILL_GRACE);
106 let signal = group.send_signal(TerminationSignal::Kill, &bound);
107 let reaped = match child.wait_timeout(LOCAL_LAUNCH_FAILURE_REAP_GRACE) {
108 Ok(Some(_)) => Ok(()),
109 Ok(None) => Err(format!(
110 "child did not reap within {} seconds",
111 LOCAL_LAUNCH_FAILURE_REAP_GRACE.as_secs()
112 )),
113 Err(error) => Err(format!("failed to reap failed launch child: {error}")),
114 };
115 let mut evidence = match reaped {
116 Ok(()) => completed_cleanup(CleanupTrigger::StartupRollback, false, true, vec![signal]),
117 Err(error) => {
118 let error = initial_status_error
119 .map(|status_error| format!("{status_error}; {error}"))
120 .unwrap_or(error);
121 cleanup_error(CleanupTrigger::StartupRollback, true, vec![signal], error)
122 }
123 };
124 evidence.elapsed_ms = duration_millis(started.elapsed());
125 evidence.status_deadline_ms = 0;
126 evidence.term_grace_ms = 0;
127 evidence.reap_grace_ms = Some(duration_millis(LOCAL_LAUNCH_FAILURE_REAP_GRACE));
128 evidence
129}
130
131pub(super) fn removal_summary(removal: &ContainerRemovalEvidence) -> String {
132 match (removal.confirmed, removal.already_absent, &removal.error) {
133 (true, true, _) => format!("container {} was already absent", removal.container),
134 (true, _, _) => format!("container {} was removed", removal.container),
135 (false, _, Some(error)) => {
136 format!(
137 "container {} removal was not confirmed: {error}",
138 removal.container
139 )
140 }
141 (false, _, None) => format!("container {} removal was not confirmed", removal.container),
142 }
143}
144
145pub(super) fn terminate_local(
146 handle: &HostProcessHandle,
147 trigger: CleanupTrigger,
148) -> CleanupEvidence {
149 let started = Instant::now();
150 if let Err(error) = handle.validate() {
151 let mut evidence = CleanupEvidence::unavailable(trigger, error);
152 evidence.elapsed_ms = duration_millis(started.elapsed());
153 return evidence;
154 }
155 let group = match LocalProcessGroup::new(
156 handle.leader_pid,
157 handle.process_group,
158 handle.leader_start_time_ticks,
159 ) {
160 Ok(group) => group,
161 Err(error) => {
162 let mut evidence = CleanupEvidence::unavailable(trigger, error.to_string());
163 evidence.elapsed_ms = duration_millis(started.elapsed());
164 return evidence;
165 }
166 };
167 let status_bound = OperationBound::finite(SERVER_CLEANUP_STATUS_DEADLINE);
168 match group.verified_status(&status_bound) {
169 Ok(VerifiedStatus::Alive) => {}
170 Ok(VerifiedStatus::Exited | VerifiedStatus::Reused) => {
171 let mut evidence = completed_cleanup(trigger, true, false, Vec::new());
172 evidence.elapsed_ms = duration_millis(started.elapsed());
173 return evidence;
174 }
175 Ok(VerifiedStatus::LeaderMissingWithMembers) => {
176 let violations = match group.cohort_violations(&status_bound) {
180 Ok(violations) => violations,
181 Err(error) => {
182 let mut evidence = CleanupEvidence::unavailable(trigger, error.to_string());
183 evidence.elapsed_ms = duration_millis(started.elapsed());
184 return evidence;
185 }
186 };
187 if !violations.is_empty() {
188 let members = violations
189 .iter()
190 .map(|(pid, ticks)| format!("member {pid} started at {ticks}"))
191 .collect::<Vec<_>>()
192 .join(", ");
193 let mut evidence = CleanupEvidence::unavailable(
194 trigger,
195 format!(
196 "process-group {} still has members but recorded leader {} no longer exists and the cohort check failed: {} before the recorded leader start {}",
197 handle.process_group,
198 handle.leader_pid,
199 members,
200 handle.leader_start_time_ticks
201 ),
202 );
203 evidence.elapsed_ms = duration_millis(started.elapsed());
204 return evidence;
205 }
206 }
207 Err(error) => {
208 let mut evidence = CleanupEvidence::unavailable(trigger, error.to_string());
209 evidence.elapsed_ms = duration_millis(started.elapsed());
210 return evidence;
211 }
212 }
213 let term_bound = OperationBound::finite(TERM_GRACE);
214 let mut signals = vec![group.send_signal(TerminationSignal::Term, &term_bound)];
215 let mut evidence = match group.wait_until_stopped(None, &term_bound, POLL_INTERVAL) {
216 Ok(true) => completed_cleanup(trigger, false, false, signals),
217 Ok(false) => {
218 let kill_bound = OperationBound::finite(KILL_GRACE);
219 signals.push(group.send_signal(TerminationSignal::Kill, &kill_bound));
220 match group.wait_until_stopped(None, &kill_bound, POLL_INTERVAL) {
221 Ok(true) => completed_cleanup(trigger, false, true, signals),
222 Ok(false) => CleanupEvidence {
223 trigger,
224 elapsed_ms: 0,
225 status_deadline_ms: duration_millis(SERVER_CLEANUP_STATUS_DEADLINE),
226 term_grace_ms: duration_millis(TERM_GRACE),
227 kill_grace_ms: duration_millis(KILL_GRACE),
228 reap_grace_ms: None,
229 remote_deadline_ms: None,
230 verified: false,
231 already_exited: false,
232 forced: true,
233 signals,
234 error: Some(format!(
235 "server process group {} did not exit after SIGKILL",
236 handle.process_group
237 )),
238 container_removal: None,
239 device_residuals: None,
240 },
241 Err(error) => cleanup_error(trigger, true, signals, error.to_string()),
242 }
243 }
244 Err(error) => cleanup_error(trigger, false, signals, error.to_string()),
245 };
246 evidence.elapsed_ms = duration_millis(started.elapsed());
247 evidence
248}
249
250pub(super) fn terminate_ssh(handle: &SshProcessHandle, trigger: CleanupTrigger) -> CleanupEvidence {
251 let started = Instant::now();
252 let bound = OperationBound::finite(REMOTE_SERVER_CLEANUP_DEADLINE);
253 let mut evidence = terminate_ssh_under(handle, trigger, &bound);
254 evidence.elapsed_ms = duration_millis(started.elapsed());
255 evidence.remote_deadline_ms = Some(duration_millis(REMOTE_SERVER_CLEANUP_DEADLINE));
256 evidence
257}
258
259pub(super) fn remote_cleanup_script(handle: &SshProcessHandle) -> String {
264 format!(
265 "set +e; pgid={}; pid={}; expected={}; if [ -r /proc/$pid/stat ]; then actual=$(awk '{{print $22}}' /proc/$pid/stat); if [ $? -ne 0 ]; then printf '{marker}unknown\\t-\\t0\\t-\\t1\\tstat-unreadable\\n'; exit 0; fi; if [ \"$actual\" != \"$expected\" ]; then printf '{marker}stale\\t-\\t0\\t-\\t0\\t%s\\n' \"$actual\"; exit 0; fi; elif {}; then bad=\"\"; for mpid in $(ps -eo pid=,pgid=,stat= | awk -v pgid=\"$pgid\" '$2 == pgid && $3 !~ /^Z/ {{print $1}}'); do mticks=$(awk '{{print $22}}' /proc/$mpid/stat 2>/dev/null) || mticks=\"\"; if [ -n \"$mticks\" ] && [ \"$mticks\" -lt \"$expected\" ]; then bad=\"$bad $mpid:$mticks\"; fi; done; if [ -n \"$bad\" ]; then printf '{marker}unknown\\t-\\t0\\t-\\t1\\tleader-missing; cohort members%s predate recorded leader start %s\\n' \"$bad\" \"$expected\"; exit 0; fi; else printf '{marker}already\\t-\\t0\\t-\\t0\\t-\\n'; exit 0; fi; if ! {}; then printf '{marker}already\\t-\\t0\\t-\\t0\\t-\\n'; exit 0; fi; kill -TERM -- -$pgid; term_code=$?; i=0; while {} && [ $i -lt {term_limit} ]; do sleep 0.1; i=$((i+1)); done; forced=0; kill_code=-; if {}; then forced=1; kill -KILL -- -$pgid; kill_code=$?; i=0; while {} && [ $i -lt {kill_limit} ]; do sleep 0.1; i=$((i+1)); done; fi; alive=0; if {}; then alive=1; fi; printf '{marker}cleanup\\t%s\\t%s\\t%s\\t%s\\t-\\n' \"$term_code\" \"$forced\" \"$kill_code\" \"$alive\"",
266 handle.process_group,
267 handle.leader_pid,
268 handle.leader_start_time_ticks,
269 remote_group_alive_script("$pgid"),
270 remote_group_alive_script("$pgid"),
271 remote_group_alive_script("$pgid"),
272 remote_group_alive_script("$pgid"),
273 remote_group_alive_script("$pgid"),
274 remote_group_alive_script("$pgid"),
275 term_limit = TERM_POLL_LIMIT,
276 kill_limit = KILL_POLL_LIMIT,
277 marker = CLEANUP_MARKER,
278 )
279}
280
281pub(super) fn terminate_ssh_under(
282 handle: &SshProcessHandle,
283 trigger: CleanupTrigger,
284 bound: &OperationBound,
285) -> CleanupEvidence {
286 let script = remote_cleanup_script(handle);
287 match run_cleanup_command(
288 &ssh_argv(&handle.target, &script),
289 SSH_ENV_REMOVE,
290 bound,
291 "SSH process cleanup",
292 ) {
293 Ok(output) if output.status.success() => {
294 let stdout = String::from_utf8_lossy(&output.stdout);
295 let Some(result) = parse_cleanup_output(&stdout) else {
296 return cleanup_error(
297 trigger,
298 false,
299 Vec::new(),
300 "SSH cleanup returned no cleanup result".to_owned(),
301 );
302 };
303 match result.state {
304 RemoteCleanupState::Already => {
305 return completed_cleanup(trigger, true, false, Vec::new());
306 }
307 RemoteCleanupState::Stale => {
308 return CleanupEvidence::unavailable(
309 trigger,
310 format!(
311 "managed SSH process {} exited and its pid was reused: observed start time {}",
312 handle.leader_pid, result.detail
313 ),
314 );
315 }
316 RemoteCleanupState::Unknown => {
317 return CleanupEvidence::unavailable(
318 trigger,
319 format!(
320 "SSH process-group {} ownership could not be verified: {}",
321 handle.process_group, result.detail
322 ),
323 );
324 }
325 RemoteCleanupState::Cleanup => {}
326 }
327 let Some(term_code) = result.term_code else {
328 return cleanup_error(
329 trigger,
330 false,
331 Vec::new(),
332 "SSH cleanup returned no SIGTERM status".to_owned(),
333 );
334 };
335 let stderr = String::from_utf8_lossy(&output.stderr).trim().to_owned();
336 let mut signals = vec![remote_signal_evidence(
337 TerminationSignal::Term,
338 handle.process_group,
339 term_code,
340 &stderr,
341 )];
342 if let Some(kill_code) = result.kill_code {
343 signals.push(remote_signal_evidence(
344 TerminationSignal::Kill,
345 handle.process_group,
346 kill_code,
347 &stderr,
348 ));
349 }
350 if result.alive {
351 cleanup_error(
352 trigger,
353 result.forced,
354 signals,
355 format!(
356 "SSH process group {} did not exit after cleanup",
357 handle.process_group
358 ),
359 )
360 } else {
361 completed_cleanup(trigger, false, result.forced, signals)
362 }
363 }
364 Ok(output) => cleanup_error(
365 trigger,
366 false,
367 Vec::new(),
368 format!(
369 "SSH cleanup exited with {}: {}",
370 output.status,
371 String::from_utf8_lossy(&output.stderr).trim()
372 ),
373 ),
374 Err(error) => cleanup_error(trigger, false, Vec::new(), error.to_string()),
375 }
376}
377
378pub(super) fn remove_server_container(
382 target: Option<&str>,
383 container: &str,
384) -> ContainerRemovalEvidence {
385 use crate::container::{Removal, RemovalFailure, remove_container};
386 let started = Instant::now();
387 let evidence =
388 |confirmed: bool,
389 already_absent: bool,
390 error: Option<String>,
391 operation_elapsed_ms: u64,
392 client_cleanup: Option<crate::container::CommandCleanupEvidence>| {
393 ContainerRemovalEvidence {
394 container: container.to_owned(),
395 elapsed_ms: duration_millis(started.elapsed()),
396 operation_elapsed_ms,
397 deadline_ms: duration_millis(crate::container::REMOVAL_TIMEOUT),
398 client_cleanup,
399 confirmed,
400 already_absent,
401 error,
402 }
403 };
404 match remove_container(target, container) {
405 Removal::Confirmed { already_absent } => evidence(
406 true,
407 already_absent,
408 None,
409 duration_millis(started.elapsed()),
410 None,
411 ),
412 Removal::Unconfirmed(RemovalFailure::Exit { status, stderr }) => evidence(
413 false,
414 false,
415 Some(format!(
416 "docker rm -f exited with {status}: {}",
417 stderr.trim()
418 )),
419 duration_millis(started.elapsed()),
420 None,
421 ),
422 Removal::Unconfirmed(RemovalFailure::Deadline {
423 operation_elapsed_ms,
424 client_cleanup,
425 }) => evidence(
426 false,
427 false,
428 Some(format!(
429 "docker rm -f {container} exceeded its {}s deadline",
430 crate::container::REMOVAL_TIMEOUT.as_secs()
431 )),
432 operation_elapsed_ms,
433 client_cleanup,
434 ),
435 Removal::Unconfirmed(RemovalFailure::Launch(error)) => evidence(
436 false,
437 false,
438 Some(format!("docker rm failed to launch: {error}")),
439 duration_millis(started.elapsed()),
440 None,
441 ),
442 Removal::Unconfirmed(RemovalFailure::Wait(error)) => evidence(
443 false,
444 false,
445 Some(format!("docker rm wait failed: {error}")),
446 duration_millis(started.elapsed()),
447 None,
448 ),
449 Removal::Unconfirmed(RemovalFailure::WaitCleanup {
450 source,
451 operation_elapsed_ms,
452 client_cleanup,
453 }) => evidence(
454 false,
455 false,
456 Some(format!("docker rm wait failed: {source}")),
457 operation_elapsed_ms,
458 Some(client_cleanup),
459 ),
460 Removal::Unconfirmed(RemovalFailure::Ssh(error)) => evidence(
461 false,
462 false,
463 Some(error),
464 duration_millis(started.elapsed()),
465 None,
466 ),
467 }
468}
469
470pub(super) fn completed_cleanup(
471 trigger: CleanupTrigger,
472 already_exited: bool,
473 forced: bool,
474 signals: Vec<SignalEvidence>,
475) -> CleanupEvidence {
476 CleanupEvidence {
477 trigger,
478 elapsed_ms: 0,
479 status_deadline_ms: duration_millis(SERVER_CLEANUP_STATUS_DEADLINE),
480 term_grace_ms: duration_millis(TERM_GRACE),
481 kill_grace_ms: duration_millis(KILL_GRACE),
482 reap_grace_ms: None,
483 remote_deadline_ms: None,
484 verified: true,
485 already_exited,
486 forced,
487 signals,
488 error: None,
489 container_removal: None,
490 device_residuals: None,
491 }
492}
493
494pub(super) fn cleanup_error(
495 trigger: CleanupTrigger,
496 forced: bool,
497 signals: Vec<SignalEvidence>,
498 error: String,
499) -> CleanupEvidence {
500 CleanupEvidence {
501 trigger,
502 elapsed_ms: 0,
503 status_deadline_ms: duration_millis(SERVER_CLEANUP_STATUS_DEADLINE),
504 term_grace_ms: duration_millis(TERM_GRACE),
505 kill_grace_ms: duration_millis(KILL_GRACE),
506 reap_grace_ms: None,
507 remote_deadline_ms: None,
508 verified: false,
509 already_exited: false,
510 forced,
511 signals,
512 error: Some(error),
513 container_removal: None,
514 device_residuals: None,
515 }
516}
517
518enum RemoteCleanupState {
519 Cleanup,
520 Already,
521 Stale,
522 Unknown,
523}
524
525struct RemoteCleanupOutput {
526 state: RemoteCleanupState,
527 term_code: Option<i32>,
528 forced: bool,
529 kill_code: Option<i32>,
530 alive: bool,
531 detail: String,
532}
533
534fn parse_cleanup_output(output: &str) -> Option<RemoteCleanupOutput> {
535 let result = output
536 .lines()
537 .rev()
538 .find_map(|line| line.strip_prefix(CLEANUP_MARKER))?;
539 let mut fields = result.split('\t');
540 let state = match fields.next()? {
541 "cleanup" => RemoteCleanupState::Cleanup,
542 "already" => RemoteCleanupState::Already,
543 "stale" => RemoteCleanupState::Stale,
544 "unknown" => RemoteCleanupState::Unknown,
545 _ => return None,
546 };
547 let term_code = match fields.next()? {
548 "-" => None,
549 value => Some(value.parse().ok()?),
550 };
551 let forced = fields.next()? == "1";
552 let kill_code = match fields.next()? {
553 "-" => None,
554 value => Some(value.parse().ok()?),
555 };
556 let alive = fields.next()? == "1";
557 let detail = fields.next()?.to_owned();
558 Some(RemoteCleanupOutput {
559 state,
560 term_code,
561 forced,
562 kill_code,
563 alive,
564 detail,
565 })
566}
567
568fn remote_signal_evidence(
569 signal: TerminationSignal,
570 process_group: u32,
571 exit_code: i32,
572 stderr: &str,
573) -> SignalEvidence {
574 SignalEvidence {
575 signal,
576 process_group,
577 exit_code: Some(exit_code),
578 stderr: (!stderr.is_empty()).then(|| stderr.to_owned()),
579 error: None,
580 }
581}
582
583impl ProcessCleanup for SystemProcessRuntime {
584 fn terminate(
585 &self,
586 handle: &ProcessHandle,
587 trigger: CleanupTrigger,
588 on_container_removal: &mut dyn FnMut(&str),
589 ) -> CleanupEvidence {
590 let mut evidence = match handle {
591 ProcessHandle::Local(handle) => terminate_local(handle, trigger),
592 ProcessHandle::Ssh(handle) => terminate_ssh(handle, trigger),
593 };
594 let (container, target) = match handle {
600 ProcessHandle::Local(handle) => (handle.container.as_deref(), None),
601 ProcessHandle::Ssh(handle) => (handle.container.as_deref(), Some(&*handle.target)),
602 };
603 if let Some(container) = container {
604 on_container_removal(container);
605 let removal = remove_server_container(target, container);
606 if !removal.confirmed {
607 evidence.verified = false;
608 if evidence.error.is_none() {
609 evidence.error = Some(format!(
610 "container {container} removal was not confirmed: {}",
611 removal.error.as_deref().unwrap_or("unknown outcome")
612 ));
613 }
614 }
615 evidence.container_removal = Some(removal);
616 }
617 evidence
618 }
619}