/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 | | } |