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 | |
| 9 | use crate::PROC_ROOT; |
| 10 | |
| 11 | use oxedyne_fe2o3_core::prelude::*; |
| 12 | |
| 13 | use std::{ |
| 14 | fs, |
| 15 | path::Path, |
| 16 | }; |
| 17 | |
| 18 | pub const CGROUP_ROOT: &str = "/sys/fs/cgroup"; |
| 19 | pub 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)] |
| 24 | pub 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 | |
| 34 | impl 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`? |
| 132 | pub 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. |
| 139 | pub 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. |
| 162 | pub 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. |
| 170 | pub 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. |
| 177 | fn 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)] |
| 193 | mod 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 | } |