Skip to main content

inferlab_runtime/server/
cleanup.rs

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    /// Cleanup evidence for a launch failure that removed (or tried to
42    /// remove) the container it created. `verified` is the caller's
43    /// conjunction of process cleanup AND container removal — a confirmed
44    /// removal alone is not verified cleanup if the launcher stop was not
45    /// confirmed — and the structured outcome names the actual container
46    /// ([[RFC-0003:C-RUNTIME-WORKFLOWS]]).
47    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        // The subsequent kill and bounded reap are authoritative cleanup
85        // verification. Preserve this diagnostic only if that verification
86        // also fails; a successful reap resolves the transient status error.
87        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
345/// Confirm a server container is gone from its launch machine
346/// ([[RFC-0003:C-RUNTIME-WORKFLOWS]]), mapping the shared removal outcome
347/// onto this record's evidence shape.
348pub(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        // The container is a daemon-owned object: the group kill reaches
560        // only the docker client, so a known container must be confirmed
561        // removed on its launch machine — unconditionally, because it can
562        // survive every group state observed above
563        // ([[RFC-0003:C-RUNTIME-WORKFLOWS]]).
564        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}