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