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 }
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 SSH_ENV_REMOVE,
257 bound,
258 "SSH process cleanup",
259 ) {
260 Ok(output) if output.status.success() => {
261 let stdout = String::from_utf8_lossy(&output.stdout);
262 let Some(result) = parse_cleanup_output(&stdout) else {
263 return cleanup_error(
264 trigger,
265 false,
266 Vec::new(),
267 "SSH cleanup returned no cleanup result".to_owned(),
268 );
269 };
270 match result.state {
271 RemoteCleanupState::Already => {
272 return completed_cleanup(trigger, true, false, Vec::new());
273 }
274 RemoteCleanupState::Stale => {
275 return CleanupEvidence::unavailable(
276 trigger,
277 format!(
278 "managed SSH process {} exited and its pid was reused: observed start time {}",
279 handle.leader_pid, result.detail
280 ),
281 );
282 }
283 RemoteCleanupState::Unknown => {
284 return CleanupEvidence::unavailable(
285 trigger,
286 format!(
287 "SSH process-group {} ownership could not be verified: {}",
288 handle.process_group, result.detail
289 ),
290 );
291 }
292 RemoteCleanupState::Cleanup => {}
293 }
294 let Some(term_code) = result.term_code else {
295 return cleanup_error(
296 trigger,
297 false,
298 Vec::new(),
299 "SSH cleanup returned no SIGTERM status".to_owned(),
300 );
301 };
302 let stderr = String::from_utf8_lossy(&output.stderr).trim().to_owned();
303 let mut signals = vec![remote_signal_evidence(
304 TerminationSignal::Term,
305 handle.process_group,
306 term_code,
307 &stderr,
308 )];
309 if let Some(kill_code) = result.kill_code {
310 signals.push(remote_signal_evidence(
311 TerminationSignal::Kill,
312 handle.process_group,
313 kill_code,
314 &stderr,
315 ));
316 }
317 if result.alive {
318 cleanup_error(
319 trigger,
320 result.forced,
321 signals,
322 format!(
323 "SSH process group {} did not exit after cleanup",
324 handle.process_group
325 ),
326 )
327 } else {
328 completed_cleanup(trigger, false, result.forced, signals)
329 }
330 }
331 Ok(output) => cleanup_error(
332 trigger,
333 false,
334 Vec::new(),
335 format!(
336 "SSH cleanup exited with {}: {}",
337 output.status,
338 String::from_utf8_lossy(&output.stderr).trim()
339 ),
340 ),
341 Err(error) => cleanup_error(trigger, false, Vec::new(), error.to_string()),
342 }
343}
344
345pub(super) fn remove_server_container(
349 target: Option<&str>,
350 container: &str,
351) -> ContainerRemovalEvidence {
352 use crate::container::{Removal, RemovalFailure, remove_container};
353 let started = Instant::now();
354 let evidence =
355 |confirmed: bool,
356 already_absent: bool,
357 error: Option<String>,
358 operation_elapsed_ms: u64,
359 client_cleanup: Option<crate::container::CommandCleanupEvidence>| {
360 ContainerRemovalEvidence {
361 container: container.to_owned(),
362 elapsed_ms: duration_millis(started.elapsed()),
363 operation_elapsed_ms,
364 deadline_ms: duration_millis(crate::container::REMOVAL_TIMEOUT),
365 client_cleanup,
366 confirmed,
367 already_absent,
368 error,
369 }
370 };
371 match remove_container(target, container) {
372 Removal::Confirmed { already_absent } => evidence(
373 true,
374 already_absent,
375 None,
376 duration_millis(started.elapsed()),
377 None,
378 ),
379 Removal::Unconfirmed(RemovalFailure::Exit { status, stderr }) => evidence(
380 false,
381 false,
382 Some(format!(
383 "docker rm -f exited with {status}: {}",
384 stderr.trim()
385 )),
386 duration_millis(started.elapsed()),
387 None,
388 ),
389 Removal::Unconfirmed(RemovalFailure::Deadline {
390 operation_elapsed_ms,
391 client_cleanup,
392 }) => evidence(
393 false,
394 false,
395 Some(format!(
396 "docker rm -f {container} exceeded its {}s deadline",
397 crate::container::REMOVAL_TIMEOUT.as_secs()
398 )),
399 operation_elapsed_ms,
400 client_cleanup,
401 ),
402 Removal::Unconfirmed(RemovalFailure::Launch(error)) => evidence(
403 false,
404 false,
405 Some(format!("docker rm failed to launch: {error}")),
406 duration_millis(started.elapsed()),
407 None,
408 ),
409 Removal::Unconfirmed(RemovalFailure::Wait(error)) => evidence(
410 false,
411 false,
412 Some(format!("docker rm wait failed: {error}")),
413 duration_millis(started.elapsed()),
414 None,
415 ),
416 Removal::Unconfirmed(RemovalFailure::WaitCleanup {
417 source,
418 operation_elapsed_ms,
419 client_cleanup,
420 }) => evidence(
421 false,
422 false,
423 Some(format!("docker rm wait failed: {source}")),
424 operation_elapsed_ms,
425 Some(client_cleanup),
426 ),
427 Removal::Unconfirmed(RemovalFailure::Ssh(error)) => evidence(
428 false,
429 false,
430 Some(error),
431 duration_millis(started.elapsed()),
432 None,
433 ),
434 }
435}
436
437pub(super) fn completed_cleanup(
438 trigger: CleanupTrigger,
439 already_exited: bool,
440 forced: bool,
441 signals: Vec<SignalEvidence>,
442) -> CleanupEvidence {
443 CleanupEvidence {
444 trigger,
445 elapsed_ms: 0,
446 status_deadline_ms: duration_millis(SERVER_CLEANUP_STATUS_DEADLINE),
447 term_grace_ms: duration_millis(TERM_GRACE),
448 kill_grace_ms: duration_millis(KILL_GRACE),
449 reap_grace_ms: None,
450 remote_deadline_ms: None,
451 verified: true,
452 already_exited,
453 forced,
454 signals,
455 error: None,
456 container_removal: None,
457 }
458}
459
460pub(super) fn cleanup_error(
461 trigger: CleanupTrigger,
462 forced: bool,
463 signals: Vec<SignalEvidence>,
464 error: String,
465) -> CleanupEvidence {
466 CleanupEvidence {
467 trigger,
468 elapsed_ms: 0,
469 status_deadline_ms: duration_millis(SERVER_CLEANUP_STATUS_DEADLINE),
470 term_grace_ms: duration_millis(TERM_GRACE),
471 kill_grace_ms: duration_millis(KILL_GRACE),
472 reap_grace_ms: None,
473 remote_deadline_ms: None,
474 verified: false,
475 already_exited: false,
476 forced,
477 signals,
478 error: Some(error),
479 container_removal: None,
480 }
481}
482
483enum RemoteCleanupState {
484 Cleanup,
485 Already,
486 Stale,
487 Unknown,
488}
489
490struct RemoteCleanupOutput {
491 state: RemoteCleanupState,
492 term_code: Option<i32>,
493 forced: bool,
494 kill_code: Option<i32>,
495 alive: bool,
496 detail: String,
497}
498
499fn parse_cleanup_output(output: &str) -> Option<RemoteCleanupOutput> {
500 let result = output
501 .lines()
502 .rev()
503 .find_map(|line| line.strip_prefix(CLEANUP_MARKER))?;
504 let mut fields = result.split('\t');
505 let state = match fields.next()? {
506 "cleanup" => RemoteCleanupState::Cleanup,
507 "already" => RemoteCleanupState::Already,
508 "stale" => RemoteCleanupState::Stale,
509 "unknown" => RemoteCleanupState::Unknown,
510 _ => return None,
511 };
512 let term_code = match fields.next()? {
513 "-" => None,
514 value => Some(value.parse().ok()?),
515 };
516 let forced = fields.next()? == "1";
517 let kill_code = match fields.next()? {
518 "-" => None,
519 value => Some(value.parse().ok()?),
520 };
521 let alive = fields.next()? == "1";
522 let detail = fields.next()?.to_owned();
523 Some(RemoteCleanupOutput {
524 state,
525 term_code,
526 forced,
527 kill_code,
528 alive,
529 detail,
530 })
531}
532
533fn remote_signal_evidence(
534 signal: TerminationSignal,
535 process_group: u32,
536 exit_code: i32,
537 stderr: &str,
538) -> SignalEvidence {
539 SignalEvidence {
540 signal,
541 process_group,
542 exit_code: Some(exit_code),
543 stderr: (!stderr.is_empty()).then(|| stderr.to_owned()),
544 error: None,
545 }
546}
547
548impl ProcessCleanup for SystemProcessRuntime {
549 fn terminate(
550 &self,
551 handle: &ProcessHandle,
552 trigger: CleanupTrigger,
553 on_container_removal: &mut dyn FnMut(&str),
554 ) -> CleanupEvidence {
555 let mut evidence = match handle {
556 ProcessHandle::Local(handle) => terminate_local(handle, trigger),
557 ProcessHandle::Ssh(handle) => terminate_ssh(handle, trigger),
558 };
559 let (container, target) = match handle {
565 ProcessHandle::Local(handle) => (handle.container.as_deref(), None),
566 ProcessHandle::Ssh(handle) => (handle.container.as_deref(), Some(&*handle.target)),
567 };
568 if let Some(container) = container {
569 on_container_removal(container);
570 let removal = remove_server_container(target, container);
571 if !removal.confirmed {
572 evidence.verified = false;
573 if evidence.error.is_none() {
574 evidence.error = Some(format!(
575 "container {container} removal was not confirmed: {}",
576 removal.error.as_deref().unwrap_or("unknown outcome")
577 ));
578 }
579 }
580 evidence.container_removal = Some(removal);
581 }
582 evidence
583 }
584}