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