/build/source/nativelink-worker/src/capacity.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 this worker can see of its own CPU and memory, and what it should |
16 | | //! advertise from that. |
17 | | //! |
18 | | //! Every deployment so far typed the advertisement by hand next to the pod's |
19 | | //! limits, and the two drifted: chinchilla advertised 16 cores and 60 GiB |
20 | | //! inside 15-core, 56 GiB limits. The limit is the number the kernel |
21 | | //! enforces, so the advertisement is derived from it here. |
22 | | |
23 | | use core::hash::BuildHasher; |
24 | | use std::collections::HashMap; |
25 | | use std::path::{Path, PathBuf}; |
26 | | |
27 | | use nativelink_config::cas_server::{ |
28 | | CapacityConfig, CpuUnit, MemoryEnforcement, ResourceEnforcementConfig, WorkerProperty, |
29 | | }; |
30 | | use nativelink_error::{Code, Error, make_err}; |
31 | | use tracing::{info, warn}; |
32 | | |
33 | | /// CPU and memory as the cgroup (or, failing a limit, the host) reports them. |
34 | | #[derive(Debug, Clone, Copy, PartialEq, Eq)] |
35 | | pub struct ObservedCapacity { |
36 | | pub cpu_millicores: u64, |
37 | | pub memory_kb: u64, |
38 | | } |
39 | | |
40 | | /// `cpu.max` is `"<quota> <period>"` in microseconds, or `"max <period>"` when |
41 | | /// unlimited; unlimited means the host's cores. |
42 | 4 | pub fn parse_cpu_max(cpu_max: &str, host_cpus: u64) -> Option<u64> { |
43 | 4 | let mut fields = cpu_max.split_whitespace(); |
44 | 4 | let quota = fields.next()?0 ; |
45 | 4 | if quota == "max" { |
46 | 1 | return Some(host_cpus.saturating_mul(1000)); |
47 | 3 | } |
48 | 3 | let quota2 : u642 = quota.parse().ok()?1 ; |
49 | 2 | let period: u64 = fields.next()?0 .parse().ok()?0 ; |
50 | 2 | if period == 0 { |
51 | 1 | return None; |
52 | 1 | } |
53 | 1 | Some(quota.saturating_mul(1000) / period) |
54 | 4 | } |
55 | | |
56 | | /// `memory.max` is bytes, or `"max"` when unlimited; unlimited means the |
57 | | /// host's memory. |
58 | 3 | pub fn parse_memory_max(memory_max: &str, host_memory_kb: u64) -> Option<u64> { |
59 | 3 | let value = memory_max.trim(); |
60 | 3 | if value == "max" { |
61 | 1 | return Some(host_memory_kb); |
62 | 2 | } |
63 | 2 | value.parse::<u64>().ok().map(|bytes| bytes1 / 1024) |
64 | 3 | } |
65 | | |
66 | | /// `MemTotal` from `/proc/meminfo`, in KiB. |
67 | 64 | pub fn parse_meminfo_total_kb(meminfo: &str) -> Option<u64> { |
68 | 64 | meminfo.lines().find_map(|line| { |
69 | 64 | let rest = line.strip_prefix("MemTotal:")?0 ; |
70 | 64 | rest.split_whitespace().next()?0 .parse().ok() |
71 | 64 | }) |
72 | 64 | } |
73 | | |
74 | | /// The headroom the advertisement is divided by: the block's own figure |
75 | | /// when set, else the enforcement headroom while memory enforcement is on, |
76 | | /// so the two numbers cannot drift apart, else nothing. |
77 | 4 | pub fn memory_headroom_percent( |
78 | 4 | config: &CapacityConfig, |
79 | 4 | enforcement: Option<&ResourceEnforcementConfig>, |
80 | 4 | ) -> u64 { |
81 | 4 | config.memory_headroom_percent.unwrap_or_else(|| {3 |
82 | 3 | enforcement |
83 | 3 | .filter(|enforcement| enforcement.memory2 == MemoryEnforcement::Soft2 ) |
84 | 3 | .map_or(0, |enforcement| enforcement.memory_headroom_percent) |
85 | 3 | }) |
86 | 4 | } |
87 | | |
88 | | /// The observed capacity less the worker's own share, memory divided by |
89 | | /// `memory_headroom_percent`. Returns `(cpu, memory_kb)` with the CPU on |
90 | | /// the scale `cpu_unit` names: whole cores rounded down, so a 14-core pod |
91 | | /// keeping one core back advertises `13`, not `13000`, to a scheduler |
92 | | /// whose actions ask for `cpu_count=1`. |
93 | 0 | pub const fn advertised( |
94 | 0 | observed: ObservedCapacity, |
95 | 0 | config: &CapacityConfig, |
96 | 0 | memory_headroom_percent: u64, |
97 | 0 | ) -> (u64, u64) { |
98 | 0 | let cpu_millicores = observed |
99 | 0 | .cpu_millicores |
100 | 0 | .saturating_sub(config.overhead_cpu_millicores); |
101 | 0 | let cpu = match config.cpu_unit { |
102 | 0 | CpuUnit::Cores => cpu_millicores / 1000, |
103 | 0 | CpuUnit::Millicores => cpu_millicores, |
104 | | }; |
105 | 0 | let memory = observed |
106 | 0 | .memory_kb |
107 | 0 | .saturating_sub(config.overhead_memory_kb) |
108 | 0 | .saturating_mul(100) |
109 | 0 | / (100 + memory_headroom_percent); |
110 | 0 | (cpu, memory) |
111 | 0 | } |
112 | | |
113 | | /// What to do when the cgroup cannot be read: the configured properties |
114 | | /// stand if they carry both numbers, since an operator who kept them has a |
115 | | /// worker that still takes work; a worker with neither would register and |
116 | | /// then satisfy no action that asks for CPU or memory, idling for good on |
117 | | /// one warning line, so that is refused. |
118 | 3 | pub fn without_cgroup<S: BuildHasher>( |
119 | 3 | config: &CapacityConfig, |
120 | 3 | properties: &HashMap<String, WorkerProperty, S>, |
121 | 3 | ) -> Result<(), Error> { |
122 | 3 | let missing: Vec<&str> = [&config.cpu_property_name, &config.memory_property_name] |
123 | 3 | .into_iter() |
124 | 6 | .filter3 (|name| !properties.contains_key(name.as_str())) |
125 | 3 | .map(String::as_str) |
126 | 3 | .collect(); |
127 | 3 | if missing.is_empty() { |
128 | 1 | warn!( |
129 | | "capacity is configured but no cgroup v2 limit applies to the worker (see the line above for the cgroup looked at); advertising the configured platform_properties instead" |
130 | | ); |
131 | 1 | return Ok(()); |
132 | 2 | } |
133 | 2 | Err(make_err!( |
134 | 2 | Code::FailedPrecondition, |
135 | 2 | "capacity is configured but no cgroup v2 limit applies to the worker (not Linux, cgroup v1, no permission, or no limit at any level; see the line above), and platform_properties carries no {}: set them, or drop the capacity block", |
136 | 2 | missing.join(" or ") |
137 | 2 | )) |
138 | 3 | } |
139 | | |
140 | | /// The cgroup v2 path of this process from `/proc/self/cgroup`: the `0::` |
141 | | /// line, which is `/` inside a container with its own cgroup namespace and |
142 | | /// the full `/kubepods.slice/.../cri-containerd-<id>.scope` path in one that |
143 | | /// shares the host's, as a privileged container does. |
144 | 77 | pub fn parse_self_cgroup(self_cgroup: &str) -> Option<&str> { |
145 | 77 | self_cgroup |
146 | 77 | .lines() |
147 | 77 | .find_map(|line| line.strip_prefix("0::")) |
148 | 77 | .map(str::trim) |
149 | 77 | .filter(|path| path75 .starts_with75 ('/')) |
150 | 77 | } |
151 | | |
152 | | /// Where this process's own `cpu.max`, `memory.max` and `memory.current` |
153 | | /// live: the mount root joined with the path from `/proc/self/cgroup` when |
154 | | /// that directory exists under the mount, the mount root otherwise. A |
155 | | /// container with a private cgroup namespace sees itself at the root, so its |
156 | | /// path is `/` and both agree. A privileged container sees the host's tree, |
157 | | /// where the root's files are the host's (or absent, as on a systemd host), |
158 | | /// and its own numbers sit further down. |
159 | 73 | pub fn resolve_cgroup_dir( |
160 | 73 | mount_root: &Path, |
161 | 73 | self_cgroup: Option<&str>, |
162 | 73 | is_dir: impl Fn(&Path) -> bool, |
163 | 73 | ) -> PathBuf { |
164 | 73 | if let Some(path72 ) = self_cgroup.and_then(parse_self_cgroup) { |
165 | 72 | let own = mount_root.join(path.trim_start_matches('/')); |
166 | 72 | if is_dir(&own) { |
167 | 2 | return own; |
168 | 70 | } |
169 | 1 | } |
170 | 71 | mount_root.to_path_buf() |
171 | 73 | } |
172 | | |
173 | | /// This process's cgroup directory on the live system. |
174 | 69 | fn cgroup_dir() -> PathBuf { |
175 | 69 | let self_cgroup = std::fs::read_to_string("/proc/self/cgroup").ok(); |
176 | 69 | resolve_cgroup_dir(Path::new("/sys/fs/cgroup"), self_cgroup.as_deref(), |dir| { |
177 | 69 | dir.join("memory.max").is_file() |
178 | 69 | }) |
179 | 69 | } |
180 | | |
181 | | /// Whether a `cpu.max` or `memory.max` value is "no limit": `max`, or for |
182 | | /// `cpu.max` `max <period>`. |
183 | 8 | fn is_unlimited(value: &str) -> bool { |
184 | 8 | value.split_whitespace().next() == Some("max") |
185 | 8 | } |
186 | | |
187 | | /// The first limit for `file` from `dir` up to and including `mount_root`, |
188 | | /// with the directory it was found in. The process's own cgroup is often |
189 | | /// an unlimited child of a limited parent (a systemd scope in a pod, a |
190 | | /// split cgroup, a privileged container under `kubepods.slice`), and the |
191 | | /// parent's limit is the one the kernel enforces on it. `None` when every |
192 | | /// level says `max` or nothing is readable: no limit applies. |
193 | 138 | pub fn find_limit( |
194 | 138 | dir: &Path, |
195 | 138 | mount_root: &Path, |
196 | 138 | file: &str, |
197 | 138 | read: impl Fn(&Path) -> Option<String>, |
198 | 138 | ) -> Option<(PathBuf, String)> { |
199 | 138 | let mut at = dir; |
200 | | loop { |
201 | 147 | if let Some(value8 ) = read(&at.join(file)) |
202 | 8 | && !is_unlimited(value.trim()) |
203 | | { |
204 | 3 | return Some((at.to_path_buf(), value)); |
205 | 144 | } |
206 | 144 | if at == mount_root || !9 at9 .starts_with(mount_root) { |
207 | 135 | return None; |
208 | 9 | } |
209 | 9 | at = at.parent()?0 ; |
210 | | } |
211 | 138 | } |
212 | | |
213 | 132 | fn read_file(path: &Path) -> Option<String> { |
214 | 132 | std::fs::read_to_string(path).ok() |
215 | 132 | } |
216 | | |
217 | | /// Reads the limits that apply to the worker: for CPU and for memory, the |
218 | | /// nearest limit from its own cgroup up to the mount root. `None` where no |
219 | | /// limit applies at any level (not Linux, cgroup v1, no permission, or a |
220 | | /// host where nothing is limited), so the configured properties stand. |
221 | 63 | pub fn observe_cgroup() -> Option<ObservedCapacity> { |
222 | 63 | let host_cpus = std::thread::available_parallelism().map_or(1, |n| n.get() as u64); |
223 | 63 | let host_memory_kb = std::fs::read_to_string("/proc/meminfo") |
224 | 63 | .ok() |
225 | 63 | .and_then(|meminfo| parse_meminfo_total_kb(&meminfo)) |
226 | 63 | .unwrap_or(0); |
227 | 63 | let mount_root = Path::new("/sys/fs/cgroup"); |
228 | 63 | let dir = cgroup_dir(); |
229 | 63 | let cpu = find_limit(&dir, mount_root, "cpu.max", read_file); |
230 | 63 | let memory = find_limit(&dir, mount_root, "memory.max", read_file); |
231 | 63 | if cpu.is_none() && memory.is_none() { |
232 | 63 | warn!( |
233 | 63 | cgroup = %dir.display(), |
234 | | "no cgroup v2 cpu.max or memory.max limit from this cgroup up to the mount root; keeping the configured platform_properties" |
235 | | ); |
236 | 63 | return None; |
237 | 0 | } |
238 | 0 | let cpu_max = cpu.map_or_else(|| "max".to_string(), |(_, value)| value); |
239 | 0 | let memory_max = memory.map_or_else(|| "max".to_string(), |(_, value)| value); |
240 | | Some(ObservedCapacity { |
241 | 0 | cpu_millicores: parse_cpu_max(&cpu_max, host_cpus)?, |
242 | 0 | memory_kb: parse_memory_max(&memory_max, host_memory_kb)?, |
243 | | }) |
244 | 63 | } |
245 | | |
246 | | /// Sets the CPU and memory properties to what the cgroup allows. When |
247 | | /// nothing can be read, the configured properties stand if they carry both, |
248 | | /// and the worker fails to start if they do not. Returns what was |
249 | | /// advertised. |
250 | 0 | pub fn apply<S: BuildHasher>( |
251 | 0 | config: &CapacityConfig, |
252 | 0 | memory_headroom_percent: u64, |
253 | 0 | properties: &mut HashMap<String, WorkerProperty, S>, |
254 | 0 | ) -> Result<Option<(u64, u64)>, Error> { |
255 | 0 | let Some(observed) = observe_cgroup() else { |
256 | 0 | without_cgroup(config, properties)?; |
257 | 0 | return Ok(None); |
258 | | }; |
259 | 0 | let (cpu, memory_kb) = advertised(observed, config, memory_headroom_percent); |
260 | 0 | properties.insert( |
261 | 0 | config.cpu_property_name.clone(), |
262 | 0 | WorkerProperty::Values(vec![cpu.to_string()]), |
263 | | ); |
264 | 0 | properties.insert( |
265 | 0 | config.memory_property_name.clone(), |
266 | 0 | WorkerProperty::Values(vec![memory_kb.to_string()]), |
267 | | ); |
268 | 0 | info!( |
269 | | observed_cpu_millicores = observed.cpu_millicores, |
270 | | observed_memory_kb = observed.memory_kb, |
271 | | cpu, |
272 | | cpu_unit = ?config.cpu_unit, |
273 | | memory_kb, |
274 | | memory_headroom_percent, |
275 | | "Advertising capacity from the cgroup" |
276 | | ); |
277 | 0 | Ok(Some((cpu, memory_kb))) |
278 | 0 | } |
279 | | |
280 | 2 | pub fn parse_memory_current_kb(memory_current: &str) -> Option<u64> { |
281 | 2 | memory_current |
282 | 2 | .trim() |
283 | 2 | .parse::<u64>() |
284 | 2 | .ok() |
285 | 2 | .map(|bytes| bytes1 / 1024) |
286 | 2 | } |
287 | | |
288 | | /// `MemAvailable` from `/proc/meminfo`, in KiB. |
289 | 7 | pub fn parse_meminfo_available_kb(meminfo: &str) -> Option<u64> { |
290 | 20 | meminfo.lines()7 .find_map7 (|line| { |
291 | 20 | let rest7 = line.strip_prefix("MemAvailable:")?13 ; |
292 | 7 | rest.split_whitespace().next()?0 .parse().ok() |
293 | 20 | }) |
294 | 7 | } |
295 | | |
296 | | /// Memory the worker could still give an action: the cgroup limit less its |
297 | | /// current usage when there is a limit, the host's `MemAvailable` when |
298 | | /// there is not, nothing when neither can be read. |
299 | 10 | pub const fn free_memory_kb_from( |
300 | 10 | limit_kb: Option<u64>, |
301 | 10 | current_kb: Option<u64>, |
302 | 10 | available_kb: Option<u64>, |
303 | 10 | ) -> Option<u64> { |
304 | 10 | match (limit_kb, current_kb) { |
305 | 2 | (Some(limit), Some(current)) => Some(limit.saturating_sub(current)), |
306 | 8 | _ => available_kb, |
307 | | } |
308 | 10 | } |
309 | | |
310 | | /// What the worker reports on each keepalive. Linux only; elsewhere the |
311 | | /// worker reports nothing and is never vetoed. |
312 | | #[cfg(target_os = "linux")] |
313 | 6 | pub fn free_memory_kb() -> Option<u64> { |
314 | | // The limit and the usage counted against it come from the same level: |
315 | | // the nearest limited ancestor, which charges everything below it. |
316 | 6 | let limited = find_limit( |
317 | 6 | &cgroup_dir(), |
318 | 6 | Path::new("/sys/fs/cgroup"), |
319 | 6 | "memory.max", |
320 | | read_file, |
321 | | ); |
322 | 6 | let limit_kb = limited |
323 | 6 | .as_ref() |
324 | 6 | .and_then(|(_, max)| parse_memory_current_kb0 (max.trim()0 )); |
325 | 6 | let current_kb = limited.as_ref().and_then(|(dir, _)| {0 |
326 | 0 | std::fs::read_to_string(dir.join("memory.current")) |
327 | 0 | .ok() |
328 | 0 | .and_then(|current| parse_memory_current_kb(¤t)) |
329 | 0 | }); |
330 | 6 | let available_kb = std::fs::read_to_string("/proc/meminfo") |
331 | 6 | .ok() |
332 | 6 | .and_then(|meminfo| parse_meminfo_available_kb(&meminfo)); |
333 | 6 | free_memory_kb_from(limit_kb, current_kb, available_kb) |
334 | 6 | } |
335 | | |
336 | | #[cfg(not(target_os = "linux"))] |
337 | | pub const fn free_memory_kb() -> Option<u64> { |
338 | | None |
339 | | } |