Coverage Report

Created: 2026-10-01 05:28

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/build/source/nativelink-worker/src/reaper.rs
Line
Count
Source
1
// Copyright 2026 The NativeLink Authors. All rights reserved.
2
//
3
// Licensed under the Functional Source License, Version 1.1, Apache 2.0 Future License (the "License");
4
// you may not use this file except in compliance with the License.
5
// You may obtain a copy of the License at
6
//
7
//    See LICENSE file for details
8
//
9
// Unless required by applicable law or agreed to in writing, software
10
// distributed under the License is distributed on an "AS IS" BASIS,
11
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12
// See the License for the specific language governing permissions and
13
// limitations under the License.
14
15
//! What an action leaves behind is the worker's to reap.
16
//!
17
//! The worker waits on the one process it spawned per action. Anything that
18
//! process forked and did not wait for (a shell's background job, a helper
19
//! that daemonized) is reparented when it exits: to the worker when the
20
//! worker is the container's PID 1, which it usually is, and otherwise to
21
//! whatever init the container has. The worker never waited on those, so
22
//! each became a zombie the moment it exited and stayed one for the life of
23
//! the pod; a worker seen with hundreds of them had run one build whose
24
//! actions forked and did not wait. This module makes the worker their
25
//! parent wherever it runs and reaps them on a timer.
26
//!
27
//! Inside `use_namespaces` the stub is the init of the action's PID
28
//! namespace and reaps there; this covers the actions that run without it.
29
30
use std::collections::{BTreeSet, HashSet};
31
use std::sync::Mutex;
32
33
/// How often the worker reaps the zombies its actions left behind.
34
pub const REAP_INTERVAL: core::time::Duration = core::time::Duration::from_secs(5);
35
36
/// The pids of every child this process spawned and still holds a handle
37
/// for. The reaper never touches these: their exit status belongs to the
38
/// code that waits on them, however long it takes to get round to it (a
39
/// persistent worker is only `try_wait`ed on demand, an action's wait can
40
/// be starved). Ownership is decided here, not by timing.
41
static OWNED: Mutex<Option<HashSet<u32>>> = Mutex::new(None);
42
43
/// A child this process owns, for as long as the guard lives.
44
#[derive(Debug)]
45
pub struct OwnedChild(Option<u32>);
46
47
impl OwnedChild {
48
    /// Registers `pid` as ours; `None` (a child that already exited at
49
    /// spawn) registers nothing. Drop the guard once the child has been
50
    /// waited on, or with the handle that will wait on it.
51
55
    pub fn new(pid: Option<u32>) -> Self {
52
55
        if let Some(pid) = pid {
53
55
            OWNED
54
55
                .lock()
55
55
                .unwrap_or_else(std::sync::PoisonError::into_inner)
56
55
                .get_or_insert_with(HashSet::new)
57
55
                .insert(pid);
58
55
        
}0
59
55
        Self(pid)
60
55
    }
61
}
62
63
impl Drop for OwnedChild {
64
55
    fn drop(&mut self) {
65
55
        if let Some(pid) = self.0
66
55
            && let Some(owned) = OWNED
67
55
                .lock()
68
55
                .unwrap_or_else(std::sync::PoisonError::into_inner)
69
55
                .as_mut()
70
55
        {
71
55
            owned.remove(&pid);
72
55
        
}0
73
55
    }
74
}
75
76
#[cfg(target_os = "linux")]
77
2
fn is_owned(pid: u32) -> bool {
78
2
    OWNED
79
2
        .lock()
80
2
        .unwrap_or_else(std::sync::PoisonError::into_inner)
81
2
        .as_ref()
82
2
        .is_some_and(|owned| owned.contains(&pid))
83
2
}
84
85
/// Make this process the parent of every descendant an action leaves
86
/// behind, even when it is not PID 1, so the zombies land here where they
87
/// can be reaped rather than with an init that never will.
88
#[cfg(target_os = "linux")]
89
3
pub fn become_subreaper() -> Result<(), std::io::Error> {
90
    // SAFETY: prctl with PR_SET_CHILD_SUBREAPER takes integers only and
91
    // changes a flag on this process.
92
3
    if unsafe { libc::prctl(libc::PR_SET_CHILD_SUBREAPER, 1, 0, 0, 0) } == 0 {
93
3
        Ok(())
94
    } else {
95
0
        Err(std::io::Error::last_os_error())
96
    }
97
3
}
98
99
#[cfg(not(target_os = "linux"))]
100
pub const fn become_subreaper() -> Result<(), std::io::Error> {
101
    Ok(())
102
}
103
104
/// The parent pid out of a `/proc/<pid>/stat` line, when the process is a
105
/// zombie; `None` for a live process or unreadable text. `comm` can hold
106
/// spaces and parentheses, so the fields are read after the final `)`:
107
/// state first, then the parent pid.
108
23
pub fn parse_zombie_ppid(stat: &str) -> Option<u32> {
109
23
    let 
after_comm22
= stat.rsplit_once(')')
?1
.1;
110
22
    let mut fields = after_comm.split_whitespace();
111
22
    let state = fields.next()
?0
;
112
22
    let 
ppid21
= fields.next()
?1
.
parse21
().
ok21
()
?0
;
113
21
    (state == "Z").then_some(ppid)
114
23
}
115
116
/// The zombies whose parent is this process, by pid. One short read per
117
/// process; call it off the runtime, since `/proc` can be large.
118
#[cfg(target_os = "linux")]
119
3
pub fn zombie_children() -> BTreeSet<u32> {
120
    // SAFETY: getpid has no memory safety considerations.
121
3
    let me = u32::try_from(unsafe { libc::getpid() }).unwrap_or(0);
122
3
    let Ok(entries) = std::fs::read_dir("/proc") else {
123
0
        return BTreeSet::new();
124
    };
125
3
    entries
126
3
        .flatten()
127
204
        .
filter_map3
(|entry| entry.file_name().to_string_lossy().parse::<u32>().ok())
128
18
        .
filter3
(|pid| {
129
18
            std::fs::read_to_string(format!("/proc/{pid}/stat"))
130
18
                .ok()
131
18
                .and_then(|stat| parse_zombie_ppid(&stat))
132
18
                == Some(me)
133
18
        })
134
3
        .collect()
135
3
}
136
137
/// Reaps the zombies that are not ours to wait on and were already zombies
138
/// at the previous sweep; returns how many, and this sweep's zombies for
139
/// the next call. A child this process spawned is never touched, whatever
140
/// its state: its exit status belongs to whoever holds its handle. The
141
/// second sweep is a guard on top of that, for a pid we never registered
142
/// (a fork of ours outside these paths); a pid reused within an interval
143
/// starts over, since a reaped pid is not carried forward. Blocking; run it
144
/// on the blocking pool.
145
#[cfg(target_os = "linux")]
146
2
pub fn reap_orphaned_zombies(seen_last: &BTreeSet<u32>) -> (usize, BTreeSet<u32>) {
147
2
    let mut now = zombie_children();
148
2
    let mut reaped = 0;
149
2
    let candidates: Vec<u32> = now
150
2
        .iter()
151
2
        .copied()
152
4
        .
filter2
(|pid| seen_last.contains(pid) &&
!is_owned(*pid)2
)
153
2
        .collect();
154
2
    for 
pid1
in candidates {
155
1
        let Ok(pid_t) = libc::pid_t::try_from(pid) else {
156
0
            continue;
157
        };
158
1
        let mut status = 0;
159
        // SAFETY: waitpid takes a pid and a pointer to an int on the stack;
160
        // a pid that is not our child is reported as ECHILD, not acted on.
161
1
        if unsafe { libc::waitpid(pid_t, &raw mut status, libc::WNOHANG) } == pid_t {
162
1
            reaped += 1;
163
1
            now.remove(&pid);
164
1
        
}0
165
    }
166
2
    (reaped, now)
167
2
}
168
169
#[cfg(not(target_os = "linux"))]
170
pub const fn reap_orphaned_zombies(_seen_last: &BTreeSet<u32>) -> (usize, BTreeSet<u32>) {
171
    (0, BTreeSet::new())
172
}