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/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(&current))
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
}