Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_sys/src/resident.rs

12.6 KiB, 1 run

created by r1870400018:61448, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! The kernel cuts a command name to fifteen bytes (`TASK_COMM_LEN` less its
2//! terminator), so a resident named longer than that is matched on its first
3//! fifteen. `daimond_gateway` fits exactly; a longer name compared whole would
4//! never match anything, and would read as a service that is not running.
5//!
6//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
7//! Anthropic Claude
8
9use crate::PROC_ROOT;
10
11use oxedyne_fe2o3_core::prelude::*;
12
13use std::{
14 fs,
15 path::Path,
16};
17
18pub const CGROUP_ROOT: &str = "/sys/fs/cgroup";
19pub const COMM_MAX: usize = 15; // TASK_COMM_LEN less its terminator
20
21/// The memory held by one named service: every process with that command name,
22/// summed.
23#[derive(Clone, Debug, Default, Eq, PartialEq)]
24pub struct Resident {
25 pub name: String,
26 pub procs: u32, // processes matched
27 pub rss_kib: u64, // summed `VmRSS`, which the kernel labels kB
28 // The tighter of `memory.high` and `memory.max` on the cgroup the matched
29 // processes share. `None` when neither is set, or when the processes do not
30 // share one cgroup, since one cap over several services is not this one's.
31 pub cap_kib: Option<u64>,
32}
33
34impl Resident {
35 /// Per cent of the cap in use, rounded down.
36 pub fn cap_pct(&self) -> Option<u64> {
37 match self.cap_kib {
38 Some(cap) if cap > 0 => Some(self.rss_kib.saturating_mul(100) / cap),
39 _ => None,
40 }
41 }
42
43 /// Reads every process named in `names` in one walk of `/proc`, and each
44 /// one's cgroup cap from `/sys/fs/cgroup`. A name nothing matches comes back
45 /// with no processes rather than as an error: a service that is not running
46 /// is the reading, not a failure to take one.
47 pub fn sample(names: &[String]) -> Outcome<Vec<Self>> {
48 Self::sample_under(names, Path::new(PROC_ROOT), Path::new(CGROUP_ROOT))
49 }
50
51 /// [`Self::sample`] against another process and cgroup root, so the walk can
52 /// be exercised on a tree built for the purpose.
53 pub fn sample_under(
54 names: &[String],
55 proc_root: &Path,
56 cgroup_root: &Path,
57 )
58 -> Outcome<Vec<Self>>
59 {
60 if names.is_empty() {
61 return Ok(Vec::new());
62 }
63 // Each resident with the cgroup path of every process it matched.
64 let mut found: Vec<(Self, Vec<String>)> = names.iter()
65 .map(|n| (Self { name: n.clone(), ..Self::default() }, Vec::new()))
66 .collect();
67 let dir = match fs::read_dir(proc_root) {
68 Ok(d) => d,
69 Err(e) => return Err(err!(e,
70 "Listing {:?} to find resident processes.", proc_root;
71 IO, File, Read)),
72 };
73 for entry in dir {
74 let entry = match entry {
75 Ok(e) => e,
76 Err(_) => continue,
77 };
78 let file_name = entry.file_name();
79 let pid = match file_name.to_str() {
80 Some(s) if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) => s,
81 _ => continue,
82 };
83 // A process can exit between the listing and the read. That is a
84 // race with the scheduler, not a fault in the reading.
85 let status = match fs::read_to_string(proc_root.join(pid).join("status")) {
86 Ok(s) => s,
87 Err(_) => continue,
88 };
89 let (comm, rss) = match status_name_rss(&status) {
90 (Some(c), Some(r)) => (c, r),
91 // No `VmRSS` is a kernel thread, which holds no user memory.
92 _ => continue,
93 };
94 let mut cgroup: Option<String> = None;
95 for (res, paths) in found.iter_mut() {
96 if !comm_matches(comm, &res.name) {
97 continue;
98 }
99 res.procs = res.procs.saturating_add(1);
100 res.rss_kib = res.rss_kib.saturating_add(rss);
101 if cgroup.is_none() {
102 // Unreadable is recorded as its own path, so it can never be
103 // mistaken for agreement with the other processes.
104 cgroup = Some(
105 match fs::read_to_string(proc_root.join(pid).join("cgroup")) {
106 Ok(c) => match cgroup_v2_path(&c) {
107 Some(p) => p.to_string(),
108 None => fmt!("?{}", pid),
109 },
110 Err(_) => fmt!("?{}", pid),
111 });
112 }
113 if let Some(p) = &cgroup {
114 paths.push(p.clone());
115 }
116 }
117 }
118 let mut out = Vec::with_capacity(found.len());
119 for (mut res, mut paths) in found {
120 paths.sort();
121 paths.dedup();
122 if paths.len() == 1 && !paths[0].starts_with('?') {
123 res.cap_kib = cgroup_cap_kib(cgroup_root, &paths[0]);
124 }
125 out.push(res);
126 }
127 Ok(out)
128 }
129}
130
131/// Does a kernel command name belong to the resident called `want`?
132pub fn comm_matches(comm: &str, want: &str) -> bool {
133 let want = want.as_bytes();
134 let want = if want.len() > COMM_MAX { &want[..COMM_MAX] } else { want };
135 comm.as_bytes() == want
136}
137
138/// The `Name` and the `VmRSS` in kibibytes from a `/proc/<pid>/status` body.
139pub fn status_name_rss(content: &str) -> (Option<&str>, Option<u64>) {
140 let mut name = None;
141 let mut rss = None;
142 for line in content.lines() {
143 let (key, rest) = match line.split_once(':') {
144 Some(kv) => kv,
145 None => continue,
146 };
147 match key {
148 "Name" => name = Some(rest.trim()),
149 "VmRSS" => rss = rest.split_whitespace().next()
150 .and_then(|t| t.parse::<u64>().ok()),
151 _ => (),
152 }
153 if name.is_some() && rss.is_some() {
154 break;
155 }
156 }
157 (name, rss)
158}
159
160/// The unified-hierarchy path from a `/proc/<pid>/cgroup` body, the line that
161/// begins `0::`. A host on the legacy hierarchy alone has no such line.
162pub fn cgroup_v2_path(content: &str) -> Option<&str> {
163 content.lines()
164 .find_map(|l| l.strip_prefix("0::"))
165 .map(|p| p.trim())
166}
167
168/// A `memory.max` or `memory.high` body in kibibytes, or `None` for `max`,
169/// which is the kernel's word for no limit.
170pub fn cgroup_limit_kib(content: &str) -> Option<u64> {
171 content.trim().parse::<u64>().ok().map(|bytes| bytes / 1024)
172}
173
174/// The tighter of the two limits the kernel enforces on a cgroup. `memory.high`
175/// throttles and reclaims, `memory.max` kills; either is where the service
176/// stops having the memory it asked for.
177fn cgroup_cap_kib(cgroup_root: &Path, path: &str) -> Option<u64> {
178 // `Path::join` with an absolute path replaces the root outright.
179 let dir = cgroup_root.join(path.trim_start_matches('/'));
180 let read = |f: &str| -> Option<u64> {
181 fs::read_to_string(dir.join(f)).ok().and_then(|c| cgroup_limit_kib(&c))
182 };
183 match (read("memory.high"), read("memory.max")) {
184 (Some(h), Some(m)) => Some(h.min(m)),
185 (Some(h), None) => Some(h),
186 (None, Some(m)) => Some(m),
187 (None, None) => None,
188 }
189}
190
191
192#[cfg(test)]
193mod tests {
194 use super::*;
195
196 use std::path::PathBuf;
197
198 /// A fresh, empty directory under the system temp root.
199 fn scratch(tag: &str) -> Outcome<PathBuf> {
200 let nanos = match std::time::SystemTime::now()
201 .duration_since(std::time::UNIX_EPOCH)
202 {
203 Ok(d) => d.as_nanos(),
204 Err(_) => 0,
205 };
206 let dir = std::env::temp_dir().join(fmt!(
207 "fe2o3_sys_resident_{}_{}_{}", tag, std::process::id(), nanos));
208 res!(fs::create_dir_all(&dir));
209 Ok(dir)
210 }
211
212 fn put(root: &Path, rel: &str, body: &str) -> Outcome<()> {
213 let p = root.join(rel);
214 if let Some(parent) = p.parent() {
215 res!(fs::create_dir_all(parent));
216 }
217 res!(fs::write(&p, body));
218 Ok(())
219 }
220
221 fn status(name: &str, rss_kb: Option<u64>) -> String {
222 match rss_kb {
223 Some(k) => fmt!("Name:\t{}\nUmask:\t0022\nState:\tS (sleeping)\n\
224 VmPeak:\t 900000 kB\nVmRSS:\t {} kB\nThreads:\t4\n", name, k),
225 None => fmt!("Name:\t{}\nState:\tI (idle)\nThreads:\t1\n", name),
226 }
227 }
228
229 #[test]
230 fn status_yields_the_name_and_the_resident_set() {
231 let s = status("steel", Some(412_000));
232 assert_eq!(status_name_rss(&s), (Some("steel"), Some(412_000)));
233 let k = status("kworker/0:1", None);
234 assert_eq!(status_name_rss(&k), (Some("kworker/0:1"), None),
235 "a kernel thread has a name and no resident set");
236 }
237
238 #[test]
239 fn the_unified_hierarchy_line_is_the_one_read() {
240 let hybrid = "12:memory:/system.slice/x.service\n0::/system.slice/steel.service\n";
241 assert_eq!(cgroup_v2_path(hybrid), Some("/system.slice/steel.service"));
242 assert_eq!(cgroup_v2_path("4:cpu:/\n"), None);
243 assert_eq!(cgroup_limit_kib("max\n"), None, "max is no limit, not a number");
244 assert_eq!(cgroup_limit_kib("1073741824\n"), Some(1_048_576));
245 }
246
247 /// A configured name longer than the kernel keeps is matched on the part the
248 /// kernel keeps, and on nothing shorter.
249 #[test]
250 fn a_long_name_is_matched_on_the_fifteen_bytes_the_kernel_keeps() {
251 assert!(comm_matches("daimond_gateway", "daimond_gateway"));
252 assert!(comm_matches("a_rather_long_s", "a_rather_long_service"));
253 assert!(!comm_matches("a_rather_long", "a_rather_long_service"));
254 assert!(!comm_matches("steel", "steel2"));
255 }
256
257 /// Two processes of one service in one cgroup are summed and judged against
258 /// that cgroup's tighter limit; a name nothing matches is reported as not
259 /// running rather than dropped.
260 #[test]
261 fn a_walk_sums_the_processes_and_reads_their_shared_cap() -> Outcome<()> {
262 let root = res!(scratch("walk"));
263 let procs = root.join("proc");
264 let cgroups = root.join("cgroup");
265 res!(put(&procs, "101/status", &status("steel", Some(300_000))));
266 res!(put(&procs, "101/cgroup", "0::/system.slice/steel.service\n"));
267 res!(put(&procs, "102/status", &status("steel", Some(100_000))));
268 res!(put(&procs, "102/cgroup", "0::/system.slice/steel.service\n"));
269 res!(put(&procs, "200/status", &status("bash", Some(5_000))));
270 res!(put(&procs, "200/cgroup", "0::/user.slice\n"));
271 res!(put(&procs, "3/status", &status("steel", None)));
272 res!(put(&procs, "self/status", &status("steel", Some(1))));
273 // 800 MiB max, 600 MiB high: the tighter one is the cap.
274 res!(put(&cgroups, "system.slice/steel.service/memory.max", "838860800\n"));
275 res!(put(&cgroups, "system.slice/steel.service/memory.high", "629145600\n"));
276
277 let names = vec![fmt!("steel"), fmt!("absent")];
278 let got = res!(Resident::sample_under(&names, &procs, &cgroups));
279 assert_eq!(got.len(), 2);
280 assert_eq!(got[0].name, "steel");
281 assert_eq!(got[0].procs, 2,
282 "two user processes match; the kernel thread and the non-numeric entry do not");
283 assert_eq!(got[0].rss_kib, 400_000);
284 assert_eq!(got[0].cap_kib, Some(614_400), "memory.high is the tighter limit");
285 assert_eq!(got[0].cap_pct(), Some(65));
286 assert_eq!(got[1], Resident { name: fmt!("absent"), ..Resident::default() },
287 "a service that is not running is a reading of nothing, not an error");
288
289 let _ = fs::remove_dir_all(&root);
290 Ok(())
291 }
292
293 /// One cap over processes in different cgroups would describe neither.
294 #[test]
295 fn processes_in_different_cgroups_report_no_cap() -> Outcome<()> {
296 let root = res!(scratch("split"));
297 let procs = root.join("proc");
298 let cgroups = root.join("cgroup");
299 res!(put(&procs, "11/status", &status("worker", Some(1_000))));
300 res!(put(&procs, "11/cgroup", "0::/system.slice/a.service\n"));
301 res!(put(&procs, "12/status", &status("worker", Some(2_000))));
302 res!(put(&procs, "12/cgroup", "0::/system.slice/b.service\n"));
303 res!(put(&cgroups, "system.slice/a.service/memory.max", "1048576\n"));
304 res!(put(&cgroups, "system.slice/b.service/memory.max", "1048576\n"));
305
306 let got = res!(Resident::sample_under(&[fmt!("worker")], &procs, &cgroups));
307 assert_eq!(got[0].procs, 2);
308 assert_eq!(got[0].rss_kib, 3_000);
309 assert_eq!(got[0].cap_kib, None);
310 assert_eq!(got[0].cap_pct(), None);
311
312 let _ = fs::remove_dir_all(&root);
313 Ok(())
314 }
315}