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