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 '{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
346/// Confirm a server container is gone from its launch machine
347/// ([[RFC-0003:C-RUNTIME-WORKFLOWS]]), mapping the shared removal outcome
348/// onto this record's evidence shape.
349pub(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        // The container is a daemon-owned object: the group kill reaches
561        // only the docker client, so a known container must be confirmed
562        // removed on its launch machine — unconditionally, because it can
563        // survive every group state observed above
564        // ([[RFC-0003:C-RUNTIME-WORKFLOWS]]).
565        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}