Skip to main content

scx_layered/
main.rs

1// Copyright (c) Meta Platforms, Inc. and affiliates.
2
3// This software may be used and distributed according to the terms of the
4// GNU General Public License version 2.
5mod bpf_skel;
6mod stats;
7
8use std::collections::BTreeMap;
9use std::collections::BTreeSet;
10use std::collections::HashMap;
11use std::collections::HashSet;
12use std::ffi::CString;
13use std::fs;
14use std::io::Write;
15use std::mem::MaybeUninit;
16use std::ops::Sub;
17use std::path::Path;
18use std::path::PathBuf;
19use std::sync::Arc;
20use std::sync::atomic::AtomicBool;
21use std::sync::atomic::Ordering;
22use std::thread::ThreadId;
23use std::time::Duration;
24use std::time::Instant;
25
26use inotify::{Inotify, WatchMask};
27use std::os::unix::io::AsRawFd;
28
29use anyhow::Context;
30use anyhow::Result;
31use anyhow::anyhow;
32use anyhow::bail;
33pub use bpf_skel::*;
34use clap::Parser;
35use crossbeam::channel::Receiver;
36use crossbeam::select;
37use lazy_static::lazy_static;
38use libbpf_rs::AsRawLibbpf;
39use libbpf_rs::MapCore as _;
40use libbpf_rs::OpenObject;
41use libbpf_rs::ProgramInput;
42use libbpf_rs::libbpf_sys;
43use nix::sched::CpuSet;
44use nvml_wrapper::Nvml;
45use nvml_wrapper::error::NvmlError;
46use once_cell::sync::OnceCell;
47use regex::Regex;
48use scx_layered::alloc::{LayerAlloc, LayerDemand, unified_alloc};
49use scx_layered::*;
50use scx_raw_pmu::PMUManager;
51use scx_stats::prelude::*;
52use scx_utils::CoreType;
53use scx_utils::Cpumask;
54use scx_utils::Llc;
55use scx_utils::NR_CPU_IDS;
56use scx_utils::NR_CPUS_POSSIBLE;
57use scx_utils::NetDev;
58use scx_utils::Topology;
59use scx_utils::TopologyArgs;
60use scx_utils::UserExitInfo;
61use scx_utils::build_id;
62use scx_utils::compat;
63use scx_utils::init_libbpf_logging;
64use scx_utils::libbpf_clap_opts::LibbpfOpts;
65use scx_utils::perf;
66use scx_utils::pm::{cpu_idle_resume_latency_supported, update_cpu_idle_resume_latency};
67use scx_utils::read_netdevs;
68use scx_utils::scx_enums;
69use scx_utils::scx_ops_attach;
70use scx_utils::scx_ops_load;
71use scx_utils::scx_ops_open;
72use scx_utils::uei_exited;
73use scx_utils::uei_report;
74use stats::LayerStats;
75use stats::StatsReq;
76use stats::StatsRes;
77use stats::SysStats;
78use std::collections::VecDeque;
79use sysinfo::{Pid, ProcessRefreshKind, ProcessesToUpdate, System};
80use tracing::{debug, error, info, trace, warn};
81use tracing_subscriber::filter::EnvFilter;
82use walkdir::WalkDir;
83
84const SCHEDULER_NAME: &str = "scx_layered";
85const MAX_PATH: usize = bpf_intf::consts_MAX_PATH as usize;
86const MAX_COMM: usize = bpf_intf::consts_MAX_COMM as usize;
87const MAX_LAYER_WEIGHT: u32 = bpf_intf::consts_MAX_LAYER_WEIGHT;
88const MIN_LAYER_WEIGHT: u32 = bpf_intf::consts_MIN_LAYER_WEIGHT;
89const MAX_LAYER_MATCH_ORS: usize = bpf_intf::consts_MAX_LAYER_MATCH_ORS as usize;
90const MAX_LAYER_NAME: usize = bpf_intf::consts_MAX_LAYER_NAME as usize;
91const MAX_LAYERS: usize = bpf_intf::consts_MAX_LAYERS as usize;
92const DEFAULT_LAYER_WEIGHT: u32 = bpf_intf::consts_DEFAULT_LAYER_WEIGHT;
93const USAGE_HALF_LIFE: u32 = bpf_intf::consts_USAGE_HALF_LIFE;
94const USAGE_HALF_LIFE_F64: f64 = USAGE_HALF_LIFE as f64 / 1_000_000_000.0;
95
96const LAYER_USAGE_OWNED: usize = bpf_intf::layer_usage_LAYER_USAGE_OWNED as usize;
97const LAYER_USAGE_OPEN: usize = bpf_intf::layer_usage_LAYER_USAGE_OPEN as usize;
98const LAYER_USAGE_SUM_UPTO: usize = bpf_intf::layer_usage_LAYER_USAGE_SUM_UPTO as usize;
99const LAYER_USAGE_PROTECTED: usize = bpf_intf::layer_usage_LAYER_USAGE_PROTECTED as usize;
100const LAYER_USAGE_PROTECTED_PREEMPT: usize =
101    bpf_intf::layer_usage_LAYER_USAGE_PROTECTED_PREEMPT as usize;
102const NR_LAYER_USAGES: usize = bpf_intf::layer_usage_NR_LAYER_USAGES as usize;
103
104const NR_GSTATS: usize = bpf_intf::global_stat_id_NR_GSTATS as usize;
105const NR_LSTATS: usize = bpf_intf::layer_stat_id_NR_LSTATS as usize;
106const NR_LLC_LSTATS: usize = bpf_intf::llc_layer_stat_id_NR_LLC_LSTATS as usize;
107
108const NR_LAYER_MATCH_KINDS: usize = bpf_intf::layer_match_kind_NR_LAYER_MATCH_KINDS as usize;
109
110static NVML: OnceCell<Nvml> = OnceCell::new();
111
112fn nvml() -> Result<&'static Nvml, NvmlError> {
113    NVML.get_or_try_init(Nvml::init)
114}
115
116lazy_static! {
117    static ref USAGE_DECAY: f64 = 0.5f64.powf(1.0 / USAGE_HALF_LIFE_F64);
118    static ref DFL_DISALLOW_OPEN_AFTER_US: u64 = 2 * scx_enums.SCX_SLICE_DFL / 1000;
119    static ref DFL_DISALLOW_PREEMPT_AFTER_US: u64 = 4 * scx_enums.SCX_SLICE_DFL / 1000;
120    static ref EXAMPLE_CONFIG: LayerConfig = serde_json::from_str(
121        r#"[
122          {
123            "name": "batch",
124            "comment": "tasks under system.slice or tasks with nice value > 0",
125            "matches": [[{"CgroupPrefix": "system.slice/"}], [{"NiceAbove": 0}]],
126            "kind": {"Confined": {
127              "util_range": [0.8, 0.9], "cpus_range": [0, 16],
128              "min_exec_us": 1000, "slice_us": 20000, "weight": 100,
129              "xllc_mig_min_us": 1000.0, "perf": 1024
130            }}
131          },
132          {
133            "name": "immediate",
134            "comment": "tasks under workload.slice with nice value < 0",
135            "matches": [[{"CgroupPrefix": "workload.slice/"}, {"NiceBelow": 0}]],
136            "kind": {"Open": {
137              "min_exec_us": 100, "yield_ignore": 0.25, "slice_us": 20000,
138              "preempt": true, "exclusive": true,
139              "prev_over_idle_core": true,
140              "weight": 100, "perf": 1024
141            }}
142          },
143          {
144            "name": "stress-ng",
145            "comment": "stress-ng test layer",
146            "matches": [[{"CommPrefix": "stress-ng"}], [{"PcommPrefix": "stress-ng"}]],
147            "kind": {"Confined": {
148              "util_range": [0.2, 0.8],
149              "min_exec_us": 800, "preempt": true, "slice_us": 800,
150              "weight": 100, "growth_algo": "Topo", "perf": 1024
151            }}
152          },
153          {
154            "name": "normal",
155            "comment": "the rest",
156            "matches": [[]],
157            "kind": {"Grouped": {
158              "util_range": [0.5, 0.6], "util_includes_open_cputime": true,
159              "min_exec_us": 200, "slice_us": 20000, "weight": 100,
160              "xllc_mig_min_us": 100.0, "growth_algo": "Linear", "perf": 1024
161            }}
162          }
163        ]"#,
164    )
165    .unwrap();
166}
167
168/// scx_layered: A highly configurable multi-layer sched_ext scheduler
169///
170/// scx_layered allows classifying tasks into multiple layers and applying
171/// different scheduling policies to them. The configuration is specified in
172/// json and composed of two parts - matches and policies.
173///
174/// Matches
175/// =======
176///
177/// Whenever a task is forked or its attributes are changed, the task goes
178/// through a series of matches to determine the layer it belongs to. A
179/// match set is composed of OR groups of AND blocks. An example:
180///
181///   "matches": [
182///     [
183///       {
184///         "CgroupPrefix": "system.slice/"
185///       }
186///     ],
187///     [
188///       {
189///         "CommPrefix": "fbagent"
190///       },
191///       {
192///         "NiceAbove": 0
193///       }
194///     ]
195///   ],
196///
197/// The outer array contains the OR groups and the inner AND blocks, so the
198/// above matches:
199///
200/// - Tasks which are in the cgroup sub-hierarchy under "system.slice".
201///
202/// - Or tasks whose comm starts with "fbagent" and have a nice value > 0.
203///
204/// Currently, the following matches are supported:
205///
206/// - CgroupPrefix: Matches the prefix of the cgroup that the task belongs
207///   to. As this is a string match, whether the pattern has the trailing
208///   '/' makes a difference. For example, "TOP/CHILD/" only matches tasks
209///   which are under that particular cgroup while "TOP/CHILD" also matches
210///   tasks under "TOP/CHILD0/" or "TOP/CHILD1/".
211///
212/// - CommPrefix: Matches the task's comm prefix.
213///
214/// - PcommPrefix: Matches the task's thread group leader's comm prefix.
215///
216/// - NiceAbove: Matches if the task's nice value is greater than the
217///   pattern.
218///
219/// - NiceBelow: Matches if the task's nice value is smaller than the
220///   pattern.
221///
222/// - NiceEquals: Matches if the task's nice value is exactly equal to
223///   the pattern.
224///
225/// - UIDEquals: Matches if the task's effective user id matches the value
226///
227/// - GIDEquals: Matches if the task's effective group id matches the value.
228///
229/// - PIDEquals: Matches if the task's pid matches the value.
230///
231/// - PPIDEquals: Matches if the task's ppid matches the value.
232///
233/// - TGIDEquals: Matches if the task's tgid matches the value.
234///
235/// - NSPIDEquals: Matches if the task's namespace id and pid matches the values.
236///
237/// - NSEquals: Matches if the task's namespace id matches the values.
238///
239/// - IsGroupLeader: Bool. When true, matches if the task is group leader
240///   (i.e. PID == TGID), aka the thread from which other threads are made.
241///   When false, matches if the task is *not* the group leader (i.e. the rest).
242///
243/// - CmdJoin: Matches when the task uses pthread_setname_np to send a join/leave
244///   command to the scheduler. See examples/cmdjoin.c for more details.
245///
246/// - UsedGpuTid: Bool. When true, matches if the tasks which have used
247///   gpus by tid.
248///
249/// - UsedGpuPid: Bool. When true, matches if the tasks which have used gpu
250///   by tgid/pid.
251///
252/// - [EXPERIMENTAL] AvgRuntime: (u64, u64). Match tasks whose average runtime
253///   is within the provided values [min, max).
254///
255/// - HintEquals: u64. Match tasks whose hint value equals this value.
256///   The value must be in the range [0, 1024].
257///
258/// - SystemCpuUtilBelow: f64. Match when the system CPU utilization fraction
259///   is below the specified threshold (a value in the range [0.0, 1.0]). This
260///   option can only be used in conjunction with HintEquals.
261///
262/// - DsqInsertBelow: f64. Match when the layer DSQ insertion fraction is below
263///   the specified threshold (a value in the range [0.0, 1.0]). This option can
264///   only be used in conjunction with HintEquals.
265///
266/// While there are complexity limitations as the matches are performed in
267/// BPF, it is straightforward to add more types of matches.
268///
269/// Templates
270/// ---------
271///
272/// Templates let us create a variable number of layers dynamically at initialization
273/// time out of a cgroup name suffix/prefix. Sometimes we know there are multiple
274/// applications running on a machine, each with their own cgroup but do not know the
275/// exact names of the applications or cgroups, e.g., in cloud computing contexts where
276/// workloads are placed on machines dynamically and run under cgroups whose name is
277/// autogenerated. In that case, we cannot hardcode the cgroup match rules when writing
278/// the configuration. We thus cannot easily prevent tasks from different cgroups from
279/// falling into the same layer and affecting each other's performance.
280///
281///
282/// Templates offer a solution to this problem by generating one layer for each such cgroup,
283/// provided these cgroups share a suffix, and that the suffix is unique to them. Templates
284/// have a cgroup suffix rule that we use to find the relevant cgroups in the system. For each
285/// such cgroup, we copy the layer config and add a matching rule that matches just this cgroup.
286///
287///
288/// Policies
289/// ========
290///
291/// The following is an example policy configuration for a layer.
292///
293///   "kind": {
294///     "Confined": {
295///       "cpus_range": [1, 8],
296///       "util_range": [0.8, 0.9]
297///     }
298///   }
299///
300/// It's of "Confined" kind, which tries to concentrate the layer's tasks
301/// into a limited number of CPUs. In the above case, the number of CPUs
302/// assigned to the layer is scaled between 1 and 8 so that the per-cpu
303/// utilization is kept between 80% and 90%. If the CPUs are loaded higher
304/// than 90%, more CPUs are allocated to the layer. If the utilization drops
305/// below 80%, the layer loses CPUs.
306///
307/// Currently, the following policy kinds are supported:
308///
309/// - Confined: Tasks are restricted to the allocated CPUs. The number of
310///   CPUs allocated is modulated to keep the per-CPU utilization in
311///   "util_range". The range can optionally be restricted with the
312///   "cpus_range" property.
313///
314/// - Grouped: Similar to Confined but tasks may spill outside if there are
315///   idle CPUs outside the allocated ones. The range can optionally be
316///   restricted with the "cpus_range" property. The optional
317///   "idle_confined" flag restricts idle selection to
318///   the layer's CPUs until the layer becomes saturated or the system
319///   is fully allocated, providing Confined-style cache locality under
320///   normal load with Grouped-style overflow under pressure.
321///
322/// - Open: Prefer the CPUs which are not occupied by Confined or Grouped
323///   layers. Tasks in this group will spill into occupied CPUs if there are
324///   no unoccupied idle CPUs.
325///
326/// All layers take the following options:
327///
328/// - min_exec_us: Minimum execution time in microseconds. Whenever a task
329///   is scheduled in, this is the minimum CPU time that it's charged no
330///   matter how short the actual execution time may be.
331///
332/// - yield_ignore: Yield ignore ratio. If 0.0, yield(2) forfeits a whole
333///   execution slice. 0.25 yields three quarters of an execution slice and
334///   so on. If 1.0, yield is completely ignored.
335///
336/// - slice_us: Scheduling slice duration in microseconds.
337///
338/// - fifo: Use FIFO queues within the layer instead of the default vtime.
339///
340/// - preempt: If true, tasks in the layer will preempt tasks which belong
341///   to other non-preempting layers when no idle CPUs are available.
342///
343/// - preempt_first: If true, tasks in the layer will try to preempt tasks
344///   in their previous CPUs before trying to find idle CPUs.
345///
346/// - exclusive: If true, tasks in the layer will occupy the whole core. The
347///   other logical CPUs sharing the same core will be kept idle. This isn't
348///   a hard guarantee, so don't depend on it for security purposes.
349///
350/// - allow_node_aligned: DEPRECATED. Node-aligned tasks are now always
351///   dispatched on layer DSQs. This field is ignored if specified.
352///
353/// - prev_over_idle_core: On SMT enabled systems, prefer using the same CPU
354///   when picking a CPU for tasks on this layer, even if that CPUs SMT
355///   sibling is processing a task.
356///
357/// - weight: Weight of the layer, which is a range from 1 to 10000 with a
358///   default of 100. Layer weights are used during contention to prevent
359///   starvation across layers. Weights are used in combination with
360///   utilization to determine the infeasible adjusted weight with higher
361///   weights having a larger adjustment in adjusted utilization.
362///
363/// - disallow_open_after_us: Duration to wait after machine reaches saturation
364///   before confining tasks in Open layers.
365///
366/// - cpus_range_frac: Array of 2 floats between 0 and 1.0. Lower and upper
367///   bound fractions of all CPUs to give to a layer. Mutually exclusive
368///   with cpus_range.
369///
370/// - disallow_preempt_after_us: Duration to wait after machine reaches saturation
371///   before confining tasks to preempt.
372///
373/// - xllc_mig_min_us: Skip cross-LLC migrations if they are likely to run on
374///   their existing LLC sooner than this.
375///
376/// - idle_smt: *** DEPRECATED ****
377///
378/// - growth_algo: Determines the order in which CPUs are allocated to the
379///   layer as it grows. All algorithms are NUMA-aware and produce per-node
380///   core orderings. Three placement classes exist:
381///
382///   * Locality (most algorithms): prefer the layer's home NUMA node and
383///     spill to remote nodes only when local capacity is exhausted.
384///   * Balanced (RoundRobin): distribute proportionally across all NUMA
385///     nodes via cap-weighted water_fill. A congested node reduces its
386///     share but does not cap the total — freed capacity flows to less
387///     loaded nodes.
388///   * Strict spread (NodeSpread*): enforce equal CPU counts across all
389///     NUMA nodes, capped at the least available node. Use when the
390///     workload needs strictly equal per-node CPU counts (memory-
391///     bandwidth-bound or replicated work).
392///
393///   Default: Sticky.
394///
395/// - perf: CPU performance target. 0 means no configuration. A value
396///   between 1 and 1024 indicates the performance level CPUs running tasks
397///   in this layer are configured to using scx_bpf_cpuperf_set().
398///
399/// - idle_resume_us: Sets the idle resume QoS value. CPU idle time governors are expected to
400///   regard the minimum of the global (effective) CPU latency limit and the effective resume
401///   latency constraint for the given CPU as the upper limit for the exit latency of the idle
402///   states. See the latest kernel docs for more details:
403///   https://www.kernel.org/doc/html/latest/admin-guide/pm/cpuidle.html
404///
405/// - nodes: If set the layer will use the set of NUMA nodes for scheduling
406///   decisions. If unset then all available NUMA nodes will be used. If the
407///   llcs value is set the cpuset of NUMA nodes will be or'ed with the LLC
408///   config.
409///
410/// - llcs: If set the layer will use the set of LLCs (last level caches)
411///   for scheduling decisions. If unset then all LLCs will be used. If
412///   the nodes value is set the cpuset of LLCs will be or'ed with the nodes
413///   config.
414///
415///
416/// Similar to matches, adding new policies and extending existing ones
417/// should be relatively straightforward.
418///
419/// Configuration example and running scx_layered
420/// =============================================
421///
422/// An scx_layered config is composed of layer configs. A layer config is
423/// composed of a name, a set of matches, and a policy block. Running the
424/// following will write an example configuration into example.json.
425///
426///   $ scx_layered -e example.json
427///
428/// Note that the last layer in the configuration must have an empty match set
429/// as a catch-all for tasks which haven't been matched into previous layers.
430///
431/// The configuration can be specified in multiple json files and
432/// command line arguments, which are concatenated in the specified
433/// order. Each must contain valid layer configurations.
434///
435/// By default, an argument to scx_layered is interpreted as a JSON string. If
436/// the argument is a pointer to a JSON file, it should be prefixed with file:
437/// or f: as follows:
438///
439///   $ scx_layered file:example.json
440///   ...
441///   $ scx_layered f:example.json
442///
443/// Monitoring Statistics
444/// =====================
445///
446/// Run with `--stats INTERVAL` to enable stats monitoring. There is
447/// also an scx_stat server listening on /var/run/scx/root/stat that can
448/// be monitored by running `scx_layered --monitor INTERVAL` separately.
449///
450///   ```bash
451///   $ scx_layered --monitor 1
452///   tot= 117909 local=86.20 open_idle= 0.21 affn_viol= 1.37 proc=6ms
453///   busy= 34.2 util= 1733.6 load=  21744.1 fb_cpus=[n0:1]
454///     batch    : util/frac=   11.8/  0.7 load/frac=     29.7:  0.1 tasks=  2597
455///                tot=   3478 local=67.80 open_idle= 0.00 preempt= 0.00 affn_viol= 0.00
456///                cpus=  2 [  2,  2] 04000001 00000000
457///     immediate: util/frac= 1218.8/ 70.3 load/frac=  21399.9: 98.4 tasks=  1107
458///                tot=  68997 local=90.57 open_idle= 0.26 preempt= 9.36 affn_viol= 0.00
459///                cpus= 50 [ 50, 50] fbfffffe 000fffff
460///     normal   : util/frac=  502.9/ 29.0 load/frac=    314.5:  1.4 tasks=  3512
461///                tot=  45434 local=80.97 open_idle= 0.16 preempt= 0.00 affn_viol= 3.56
462///                cpus= 50 [ 50, 50] fbfffffe 000fffff
463///   ```
464///
465/// Global statistics: see [`SysStats`]
466///
467/// Per-layer statistics: see [`LayerStats`]
468///
469#[derive(Debug, Parser)]
470#[command(verbatim_doc_comment)]
471struct Opts {
472    /// Deprecated, noop, use RUST_LOG or --log-level instead.
473    #[clap(short = 'v', long, action = clap::ArgAction::Count)]
474    verbose: u8,
475
476    /// Scheduling slice duration in microseconds.
477    #[clap(short = 's', long, default_value = "20000")]
478    slice_us: u64,
479
480    /// Maximum consecutive execution time in microseconds. A task may be
481    /// allowed to keep executing on a CPU for this long. Note that this is
482    /// the upper limit and a task may have to moved off the CPU earlier. 0
483    /// indicates default - 20 * slice_us.
484    #[clap(short = 'M', long, default_value = "0")]
485    max_exec_us: u64,
486
487    /// Scheduling interval in seconds.
488    #[clap(short = 'i', long, default_value = "0.1")]
489    interval: f64,
490
491    /// ***DEPRECATED*** Disable load-fraction based max layer CPU limit.
492    /// recommended.
493    #[clap(short = 'n', long, default_value = "false")]
494    no_load_frac_limit: bool,
495
496    /// Exit debug dump buffer length. 0 indicates default.
497    #[clap(long, default_value = "0")]
498    exit_dump_len: u32,
499
500    /// Specify the logging level. Accepts rust's envfilter syntax for modular
501    /// logging: https://docs.rs/tracing-subscriber/latest/tracing_subscriber/filter/struct.EnvFilter.html#example-syntax. Examples: ["info", "warn,tokio=info"]
502    #[clap(long, default_value = "info")]
503    log_level: String,
504
505    /// Disable topology awareness. When enabled, the "nodes" and "llcs" settings on
506    /// a layer are ignored. Defaults to false on topologies with multiple NUMA nodes
507    /// or LLCs, and true otherwise.
508    #[arg(short = 't', long, num_args = 0..=1, default_missing_value = "true", require_equals = true)]
509    disable_topology: Option<bool>,
510
511    /// Disable monitor
512    #[clap(long)]
513    monitor_disable: bool,
514
515    /// Write example layer specifications into the file and exit.
516    #[clap(short = 'e', long)]
517    example: Option<String>,
518
519    /// ***DEPRECATED*** Disables preemption if the weighted load fraction
520    /// of a layer (load_frac_adj) exceeds the threshold. The default is
521    /// disabled (0.0).
522    #[clap(long, default_value = "0.0")]
523    layer_preempt_weight_disable: f64,
524
525    /// ***DEPRECATED*** Disables layer growth if the weighted load fraction
526    /// of a layer (load_frac_adj) exceeds the threshold. The default is
527    /// disabled (0.0).
528    #[clap(long, default_value = "0.0")]
529    layer_growth_weight_disable: f64,
530
531    /// Enable stats monitoring with the specified interval.
532    #[clap(long)]
533    stats: Option<f64>,
534
535    /// Run in stats monitoring mode with the specified interval. Scheduler
536    /// is not launched.
537    #[clap(long)]
538    monitor: Option<f64>,
539
540    /// Column limit for stats monitor output.
541    #[clap(long, default_value = "95")]
542    stats_columns: usize,
543
544    /// Disable per-LLC stats in monitor output.
545    #[clap(long)]
546    stats_no_llc: bool,
547
548    /// Run with example layer specifications (useful for e.g. CI pipelines)
549    #[clap(long)]
550    run_example: bool,
551
552    /// Allocate CPUs at an SMT granularity (not core)
553    #[clap(long)]
554    allow_partial_core: bool,
555
556    /// ***DEPRECATED *** Enables iteration over local LLCs first for
557    /// dispatch.
558    #[clap(long, default_value = "false")]
559    local_llc_iteration: bool,
560
561    /// Low priority fallback DSQs are used to execute tasks with custom CPU
562    /// affinities. These DSQs are immediately executed iff a CPU is
563    /// otherwise idle. However, after the specified wait, they are
564    /// guaranteed upto --lo-fb-share fraction of each CPU.
565    #[clap(long, default_value = "10000")]
566    lo_fb_wait_us: u64,
567
568    /// The fraction of CPU time guaranteed to low priority fallback DSQs.
569    /// See --lo-fb-wait-us.
570    #[clap(long, default_value = ".05")]
571    lo_fb_share: f64,
572
573    /// Disable antistall
574    #[clap(long, default_value = "false")]
575    disable_antistall: bool,
576
577    /// Enable numa topology based gpu task affinitization.
578    #[clap(long, default_value = "false")]
579    enable_gpu_affinitize: bool,
580
581    /// Interval at which to reaffinitize gpu tasks to numa nodes.
582    /// Defaults to 900s
583    #[clap(long, default_value = "900")]
584    gpu_affinitize_secs: u64,
585
586    /// Enable match debug
587    /// This stores a mapping of task tid
588    /// to layer id such that bpftool map dump
589    /// can be used to debug layer matches.
590    #[clap(long, default_value = "false")]
591    enable_match_debug: bool,
592
593    /// Maximum task runnable_at delay (in seconds) before antistall turns on
594    #[clap(long, default_value = "3")]
595    antistall_sec: u64,
596
597    /// Enable gpu support
598    #[clap(long, default_value = "false")]
599    enable_gpu_support: bool,
600
601    /// Gpu Kprobe Level
602    /// The value set here determines how aggressive
603    /// the kprobes enabled on gpu driver functions are.
604    /// Higher values are more aggressive, incurring more system overhead
605    /// and more accurately identifying PIDs using GPUs in a more timely manner.
606    /// Lower values incur less system overhead, at the cost of less accurately
607    /// identifying GPU pids and taking longer to do so.
608    #[clap(long, default_value = "3")]
609    gpu_kprobe_level: u64,
610
611    /// Enable utilization compensation for unattributed CPU work (irq, softirq, stolen). When
612    /// enabled, each CPU's layer usage is scaled by the inverse of available capacity to account
613    /// for time lost to interrupts.
614    #[clap(long, default_value = "false")]
615    util_compensation: bool,
616
617    /// Enable netdev IRQ balancing. This is experimental and should be used with caution.
618    #[clap(long, default_value = "false")]
619    netdev_irq_balance: bool,
620
621    /// Disable queued wakeup optimization.
622    #[clap(long, default_value = "false")]
623    disable_queued_wakeup: bool,
624
625    /// Per-cpu kthreads are preempting by default. Make it not so.
626    #[clap(long, default_value = "false")]
627    disable_percpu_kthread_preempt: bool,
628
629    /// Only highpri (nice < 0) per-cpu kthreads are preempting by default.
630    /// Make every per-cpu kthread preempting. Meaningful only if
631    /// --disable-percpu-kthread-preempt is not set.
632    #[clap(long, default_value = "false")]
633    percpu_kthread_preempt_all: bool,
634
635    /// Print scheduler version and exit.
636    #[clap(short = 'V', long, action = clap::ArgAction::SetTrue)]
637    version: bool,
638
639    /// Optional run ID for tracking scheduler instances.
640    #[clap(long)]
641    run_id: Option<u64>,
642
643    /// Show descriptions for statistics.
644    #[clap(long)]
645    help_stats: bool,
646
647    /// Layer specification. See --help.
648    specs: Vec<String>,
649
650    /// Periodically force tasks in layers using the AvgRuntime match rule to reevaluate which layer they belong to. Default period of 2s.
651    /// turns this off.
652    #[clap(long, default_value = "2000")]
653    layer_refresh_ms_avgruntime: u64,
654
655    /// Set the path for pinning the task hint map.
656    #[clap(long, default_value = "")]
657    task_hint_map: String,
658
659    /// Print the config (after template expansion) and exit.
660    #[clap(long, default_value = "false")]
661    print_and_exit: bool,
662
663    /// Enable affinitized task to use hi fallback queue to get more CPU time.
664    #[clap(long, default_value = "")]
665    hi_fb_thread_name: String,
666
667    #[clap(flatten, next_help_heading = "Topology Options")]
668    topology: TopologyArgs,
669
670    #[clap(flatten, next_help_heading = "Libbpf Options")]
671    pub libbpf: LibbpfOpts,
672}
673
674// Cgroup event types for inter-thread communication
675#[derive(Debug, Clone)]
676enum CgroupEvent {
677    Created {
678        path: String,
679        cgroup_id: u64,    // inode number
680        match_bitmap: u64, // bitmap of matched regex rules
681    },
682    Removed {
683        path: String,
684        cgroup_id: u64, // inode number
685    },
686}
687
688fn read_total_cpu(reader: &fb_procfs::ProcReader) -> Result<fb_procfs::CpuStat> {
689    reader
690        .read_stat()
691        .context("Failed to read procfs")?
692        .total_cpu
693        .ok_or_else(|| anyhow!("Could not read total cpu stat in proc"))
694}
695
696fn calc_util(curr: &fb_procfs::CpuStat, prev: &fb_procfs::CpuStat) -> Result<f64> {
697    match (curr, prev) {
698        (
699            fb_procfs::CpuStat {
700                user_usec: Some(curr_user),
701                nice_usec: Some(curr_nice),
702                system_usec: Some(curr_system),
703                idle_usec: Some(curr_idle),
704                iowait_usec: Some(curr_iowait),
705                irq_usec: Some(curr_irq),
706                softirq_usec: Some(curr_softirq),
707                stolen_usec: Some(curr_stolen),
708                ..
709            },
710            fb_procfs::CpuStat {
711                user_usec: Some(prev_user),
712                nice_usec: Some(prev_nice),
713                system_usec: Some(prev_system),
714                idle_usec: Some(prev_idle),
715                iowait_usec: Some(prev_iowait),
716                irq_usec: Some(prev_irq),
717                softirq_usec: Some(prev_softirq),
718                stolen_usec: Some(prev_stolen),
719                ..
720            },
721        ) => {
722            let idle_usec = curr_idle.saturating_sub(*prev_idle);
723            let iowait_usec = curr_iowait.saturating_sub(*prev_iowait);
724            let user_usec = curr_user.saturating_sub(*prev_user);
725            let system_usec = curr_system.saturating_sub(*prev_system);
726            let nice_usec = curr_nice.saturating_sub(*prev_nice);
727            let irq_usec = curr_irq.saturating_sub(*prev_irq);
728            let softirq_usec = curr_softirq.saturating_sub(*prev_softirq);
729            let stolen_usec = curr_stolen.saturating_sub(*prev_stolen);
730
731            let busy_usec =
732                user_usec + system_usec + nice_usec + irq_usec + softirq_usec + stolen_usec;
733            let total_usec = idle_usec + busy_usec + iowait_usec;
734            if total_usec > 0 {
735                Ok(((busy_usec as f64) / (total_usec as f64)).clamp(0.0, 1.0))
736            } else {
737                Ok(1.0)
738            }
739        }
740        _ => bail!("Missing stats in cpustat"),
741    }
742}
743
744fn copy_into_cstr(dst: &mut [i8], src: &str) {
745    let cstr = CString::new(src).unwrap();
746    let bytes = unsafe { std::mem::transmute::<&[u8], &[i8]>(cstr.as_bytes_with_nul()) };
747    dst[0..bytes.len()].copy_from_slice(bytes);
748}
749
750#[allow(clippy::needless_range_loop)]
751fn read_cpu_ctxs(skel: &BpfSkel) -> Result<Vec<bpf_intf::cpu_ctx>> {
752    let mut cpu_ctxs = vec![];
753    let cpu_ctxs_vec = skel
754        .maps
755        .cpu_ctxs
756        .lookup_percpu(&0u32.to_ne_bytes(), libbpf_rs::MapFlags::ANY)
757        .context("Failed to lookup cpu_ctx")?
758        .unwrap();
759    for cpu in 0..*NR_CPUS_POSSIBLE {
760        cpu_ctxs.push(
761            *plain::from_bytes(cpu_ctxs_vec[cpu].as_slice())
762                .expect("cpu_ctx: short or misaligned buffer"),
763        );
764    }
765    Ok(cpu_ctxs)
766}
767
768#[derive(Clone, Debug)]
769struct BpfStats {
770    gstats: Vec<u64>,
771    lstats: Vec<Vec<u64>>,
772    lstats_sums: Vec<u64>,
773    llc_lstats: Vec<Vec<Vec<u64>>>, // [layer][llc][stat]
774}
775
776impl BpfStats {
777    #[allow(clippy::needless_range_loop)]
778    fn read(skel: &BpfSkel, cpu_ctxs: &[bpf_intf::cpu_ctx]) -> Self {
779        let nr_layers = skel.maps.rodata_data.as_ref().unwrap().nr_layers as usize;
780        let nr_llcs = skel.maps.rodata_data.as_ref().unwrap().nr_llcs as usize;
781        let mut gstats = vec![0u64; NR_GSTATS];
782        let mut lstats = vec![vec![0u64; NR_LSTATS]; nr_layers];
783        let mut llc_lstats = vec![vec![vec![0u64; NR_LLC_LSTATS]; nr_llcs]; nr_layers];
784
785        for cpu in 0..*NR_CPUS_POSSIBLE {
786            for (stat, value) in gstats.iter_mut().enumerate() {
787                *value += cpu_ctxs[cpu].gstats[stat];
788            }
789            for (layer, layer_stats) in lstats.iter_mut().enumerate() {
790                for (stat, value) in layer_stats.iter_mut().enumerate() {
791                    *value += cpu_ctxs[cpu].lstats[layer][stat];
792                }
793            }
794        }
795
796        let mut lstats_sums = vec![0u64; NR_LSTATS];
797        for layer_stats in lstats.iter() {
798            for (stat, value) in lstats_sums.iter_mut().enumerate() {
799                *value += layer_stats[stat];
800            }
801        }
802
803        for llc_id in 0..nr_llcs {
804            // XXX - This would be a lot easier if llc_ctx were in
805            // the bss. Unfortunately, kernel < v6.12 crashes and
806            // kernel >= v6.12 fails verification after such
807            // conversion due to seemingly verifier bugs. Convert to
808            // bss maps later.
809            let v = skel
810                .maps
811                .llc_data
812                .lookup(&(llc_id as u32).to_ne_bytes(), libbpf_rs::MapFlags::ANY)
813                .unwrap()
814                .unwrap();
815            let llcc: &bpf_intf::llc_ctx =
816                plain::from_bytes(v.as_slice()).expect("llc_ctx: short or misaligned buffer");
817
818            for (layer_id, layer_llc_stats) in llc_lstats.iter_mut().enumerate() {
819                for (stat_id, stat_value) in layer_llc_stats[llc_id].iter_mut().enumerate() {
820                    *stat_value = llcc.lstats[layer_id][stat_id];
821                }
822            }
823        }
824
825        Self {
826            gstats,
827            lstats,
828            lstats_sums,
829            llc_lstats,
830        }
831    }
832}
833
834impl<'b> Sub<&'b BpfStats> for &BpfStats {
835    type Output = BpfStats;
836
837    fn sub(self, rhs: &'b BpfStats) -> BpfStats {
838        let vec_sub = |l: &[u64], r: &[u64]| l.iter().zip(r.iter()).map(|(l, r)| *l - *r).collect();
839        BpfStats {
840            gstats: vec_sub(&self.gstats, &rhs.gstats),
841            lstats: self
842                .lstats
843                .iter()
844                .zip(rhs.lstats.iter())
845                .map(|(l, r)| vec_sub(l, r))
846                .collect(),
847            lstats_sums: vec_sub(&self.lstats_sums, &rhs.lstats_sums),
848            llc_lstats: self
849                .llc_lstats
850                .iter()
851                .zip(rhs.llc_lstats.iter())
852                .map(|(l_layer, r_layer)| {
853                    l_layer
854                        .iter()
855                        .zip(r_layer.iter())
856                        .map(|(l_llc, r_llc)| {
857                            let (l_llc, mut r_llc) = (l_llc.clone(), r_llc.clone());
858                            // Lat is not subtractable, take L side.
859                            r_llc[bpf_intf::llc_layer_stat_id_LLC_LSTAT_LAT as usize] = 0;
860                            vec_sub(&l_llc, &r_llc)
861                        })
862                        .collect()
863                })
864                .collect(),
865        }
866    }
867}
868
869#[derive(Clone, Debug)]
870struct Stats {
871    at: Instant,
872    elapsed: Duration,
873    topo: Arc<Topology>,
874    nr_layers: usize,
875    nr_layer_tasks: Vec<usize>,
876    layer_nr_node_pinned_tasks: Vec<Vec<u64>>,
877
878    total_util: f64, // Running AVG of sum of layer_utils
879    layer_utils: Vec<Vec<f64>>,
880    layer_node_pinned_utils: Vec<Vec<f64>>,
881    prev_layer_node_pinned_usages: Vec<Vec<u64>>,
882    layer_node_utils: Vec<Vec<f64>>,
883    prev_layer_node_usages: Vec<Vec<u64>>,
884    layer_node_duty_sums: Vec<Vec<f64>>, // Per-node per-layer duty in CPU units (EWMA)
885    prev_layer_node_duty_raw: Vec<Vec<u64>>, // Raw accumulated layer_duty_sum values
886
887    layer_membws: Vec<Vec<f64>>, // Estimated memory bandsidth consumption
888    prev_layer_membw_agg: Vec<Vec<u64>>, // Estimated aggregate membw consumption
889
890    cpu_busy: f64, // Read from /proc, maybe higher than total_util
891    prev_total_cpu: fb_procfs::CpuStat,
892    prev_pmu_resctrl_membw: (u64, u64), // (PMU-reported membw, resctrl-reported membw)
893
894    util_compensation: bool,
895    layer_utils_compensated: Vec<Vec<f64>>, // EWMA of per-CPU-scaled layer utils
896    prev_cpu_layer_usages: Vec<u64>,        // Per-CPU per-layer usages for computing deltas
897    prev_per_cpu_stats: BTreeMap<u32, fb_procfs::CpuStat>,
898
899    system_cpu_util_ewma: f64,       // 10s EWMA of system CPU utilization
900    layer_dsq_insert_ewma: Vec<f64>, // 10s EWMA of per-layer DSQ insertion ratio
901
902    bpf_stats: BpfStats,
903    prev_bpf_stats: BpfStats,
904
905    processing_dur: Duration,
906    prev_processing_dur: Duration,
907
908    layer_slice_us: Vec<u64>,
909
910    gpu_tasks_affinitized: u64,
911    gpu_task_affinitization_ms: u64,
912}
913
914impl Stats {
915    #[allow(clippy::needless_range_loop)]
916    fn read_layer_membw_agg(cpu_ctxs: &[bpf_intf::cpu_ctx], nr_layers: usize) -> Vec<Vec<u64>> {
917        let mut layer_membw_agg = vec![vec![0u64; NR_LAYER_USAGES]; nr_layers];
918
919        for cpu in 0..*NR_CPUS_POSSIBLE {
920            for (layer, layer_membw) in layer_membw_agg.iter_mut().enumerate() {
921                for (usage, value) in layer_membw.iter_mut().enumerate() {
922                    *value += cpu_ctxs[cpu].layer_membw_agg[layer][usage];
923                }
924            }
925        }
926
927        layer_membw_agg
928    }
929
930    #[allow(clippy::needless_range_loop)]
931    fn read_layer_node_pinned_usages(
932        cpu_ctxs: &[bpf_intf::cpu_ctx],
933        topo: &Topology,
934        nr_layers: usize,
935        nr_nodes: usize,
936    ) -> Vec<Vec<u64>> {
937        let mut usages = vec![vec![0u64; nr_nodes]; nr_layers];
938
939        for cpu in 0..*NR_CPUS_POSSIBLE {
940            let node = topo.all_cpus.get(&cpu).map_or(0, |c| c.node_id);
941            for (layer, layer_usages) in usages.iter_mut().enumerate().take(nr_layers) {
942                layer_usages[node] += cpu_ctxs[cpu].node_pinned_usage[layer];
943            }
944        }
945
946        usages
947    }
948
949    #[allow(clippy::needless_range_loop)]
950    fn read_layer_node_usages(
951        cpu_ctxs: &[bpf_intf::cpu_ctx],
952        topo: &Topology,
953        nr_layers: usize,
954        nr_nodes: usize,
955    ) -> Vec<Vec<u64>> {
956        let mut usages = vec![vec![0u64; nr_nodes]; nr_layers];
957
958        for cpu in 0..*NR_CPUS_POSSIBLE {
959            let node = topo.all_cpus.get(&cpu).map_or(0, |c| c.node_id);
960            for (layer, layer_usages) in usages.iter_mut().enumerate().take(nr_layers) {
961                for usage in 0..=LAYER_USAGE_SUM_UPTO {
962                    layer_usages[node] += cpu_ctxs[cpu].layer_usages[layer][usage];
963                }
964            }
965        }
966
967        usages
968    }
969
970    #[allow(clippy::needless_range_loop)]
971    fn read_layer_node_duty_raw(
972        cpu_ctxs: &[bpf_intf::cpu_ctx],
973        topo: &Topology,
974        nr_layers: usize,
975        nr_nodes: usize,
976    ) -> Vec<Vec<u64>> {
977        let mut sums = vec![vec![0u64; nr_nodes]; nr_layers];
978
979        for cpu in 0..*NR_CPUS_POSSIBLE {
980            let node = topo.all_cpus.get(&cpu).map_or(0, |c| c.node_id);
981            for (layer, layer_sums) in sums.iter_mut().enumerate().take(nr_layers) {
982                layer_sums[node] += cpu_ctxs[cpu].layer_duty_sum[layer];
983            }
984        }
985
986        sums
987    }
988
989    #[allow(clippy::needless_range_loop)]
990    fn read_per_cpu_layer_usages(cpu_ctxs: &[bpf_intf::cpu_ctx], nr_layers: usize) -> Vec<u64> {
991        let stride = nr_layers * NR_LAYER_USAGES;
992        let mut flat = vec![0u64; *NR_CPUS_POSSIBLE * stride];
993
994        for cpu in 0..*NR_CPUS_POSSIBLE {
995            let base = cpu * stride;
996            for layer in 0..nr_layers {
997                for usage in 0..NR_LAYER_USAGES {
998                    flat[base + layer * NR_LAYER_USAGES + usage] =
999                        cpu_ctxs[cpu].layer_usages[layer][usage];
1000                }
1001            }
1002        }
1003
1004        flat
1005    }
1006
1007    /// Use the membw reported by resctrl to normalize the values reported by hw counters.
1008    /// We have the following problem:
1009    /// 1) We want per-task memory bandwidth reporting. We cannot do this with resctrl, much
1010    ///    less transparently, since we would require different RMID for each task.
1011    /// 2) We want to directly use perf counters for tracking per-task memory bandwidth, but
1012    ///    we can't: Non-resctrl counters do not measure the right thing (e.g., they only measure
1013    ///    proxies like load operations),
1014    /// 3) Resctrl counters are not accessible directly so we cannot read them from the BPF side.
1015    ///
1016    /// Approximate per-task memory bandwidth using perf counters to measure _relative_ memory
1017    /// bandwidth usage.
1018    fn resctrl_read_total_membw() -> Result<u64> {
1019        let mut total_membw = 0u64;
1020        for entry in WalkDir::new("/sys/fs/resctrl/mon_data")
1021            .min_depth(1)
1022            .into_iter()
1023            .filter_map(Result::ok)
1024            .filter(|x| x.path().is_dir())
1025        {
1026            let mut path = entry.path().to_path_buf();
1027            path.push("mbm_total_bytes");
1028            total_membw += fs::read_to_string(path)?.trim().parse::<u64>()?;
1029        }
1030
1031        Ok(total_membw)
1032    }
1033
1034    fn new(
1035        skel: &mut BpfSkel,
1036        proc_reader: &fb_procfs::ProcReader,
1037        topo: Arc<Topology>,
1038        gpu_task_affinitizer: &GpuTaskAffinitizer,
1039        util_compensation: bool,
1040    ) -> Result<Self> {
1041        let nr_layers = skel.maps.rodata_data.as_ref().unwrap().nr_layers as usize;
1042        let nr_nodes = topo.nodes.len();
1043        let cpu_ctxs = read_cpu_ctxs(skel)?;
1044        let bpf_stats = BpfStats::read(skel, &cpu_ctxs);
1045        let pmu_membw = Self::read_layer_membw_agg(&cpu_ctxs, nr_layers);
1046
1047        Ok(Self {
1048            at: Instant::now(),
1049            elapsed: Default::default(),
1050
1051            topo: topo.clone(),
1052            nr_layers,
1053            nr_layer_tasks: vec![0; nr_layers],
1054            layer_nr_node_pinned_tasks: vec![vec![0; nr_nodes]; nr_layers],
1055
1056            total_util: 0.0,
1057            layer_utils: vec![vec![0.0; NR_LAYER_USAGES]; nr_layers],
1058            layer_node_pinned_utils: vec![vec![0.0; nr_nodes]; nr_layers],
1059            prev_layer_node_pinned_usages: Self::read_layer_node_pinned_usages(
1060                &cpu_ctxs, &topo, nr_layers, nr_nodes,
1061            ),
1062            layer_node_utils: vec![vec![0.0; nr_nodes]; nr_layers],
1063            prev_layer_node_usages: Self::read_layer_node_usages(
1064                &cpu_ctxs, &topo, nr_layers, nr_nodes,
1065            ),
1066            layer_node_duty_sums: vec![vec![0.0; nr_nodes]; nr_layers],
1067            prev_layer_node_duty_raw: Self::read_layer_node_duty_raw(
1068                &cpu_ctxs, &topo, nr_layers, nr_nodes,
1069            ),
1070            layer_membws: vec![vec![0.0; NR_LAYER_USAGES]; nr_layers],
1071            // This is not normalized because we don't have enough history to do so.
1072            // It should not matter too much, since the value is dropped on the first
1073            // iteration.
1074            prev_layer_membw_agg: pmu_membw,
1075            prev_pmu_resctrl_membw: (0, 0),
1076
1077            cpu_busy: 0.0,
1078            prev_total_cpu: read_total_cpu(proc_reader)?,
1079            util_compensation,
1080            layer_utils_compensated: vec![vec![0.0; NR_LAYER_USAGES]; nr_layers],
1081            prev_cpu_layer_usages: Self::read_per_cpu_layer_usages(&cpu_ctxs, nr_layers),
1082            prev_per_cpu_stats: BTreeMap::new(),
1083            system_cpu_util_ewma: 0.0,
1084            layer_dsq_insert_ewma: vec![0.0; nr_layers],
1085
1086            bpf_stats: bpf_stats.clone(),
1087            prev_bpf_stats: bpf_stats,
1088
1089            processing_dur: Default::default(),
1090            prev_processing_dur: Default::default(),
1091
1092            layer_slice_us: vec![0; nr_layers],
1093            gpu_tasks_affinitized: gpu_task_affinitizer.tasks_affinitized,
1094            gpu_task_affinitization_ms: gpu_task_affinitizer.last_task_affinitization_ms,
1095        })
1096    }
1097
1098    fn refresh(
1099        &mut self,
1100        skel: &mut BpfSkel,
1101        proc_reader: &fb_procfs::ProcReader,
1102        now: Instant,
1103        cur_processing_dur: Duration,
1104        gpu_task_affinitizer: &GpuTaskAffinitizer,
1105    ) -> Result<()> {
1106        let elapsed = now.duration_since(self.at);
1107        let elapsed_f64 = elapsed.as_secs_f64();
1108        let cpu_ctxs = read_cpu_ctxs(skel)?;
1109
1110        let layers = &skel.maps.bss_data.as_ref().unwrap().layers;
1111        let nr_layer_tasks: Vec<usize> = layers
1112            .iter()
1113            .take(self.nr_layers)
1114            .map(|layer| layer.nr_tasks as usize)
1115            .collect();
1116        let layer_nr_node_pinned_tasks: Vec<Vec<u64>> = layers
1117            .iter()
1118            .take(self.nr_layers)
1119            .map(|layer| {
1120                layer.node[..self.topo.nodes.len()]
1121                    .iter()
1122                    .map(|n| n.nr_pinned_tasks)
1123                    .collect()
1124            })
1125            .collect();
1126        let layer_slice_us: Vec<u64> = layers
1127            .iter()
1128            .take(self.nr_layers)
1129            .map(|layer| layer.slice_ns / 1000_u64)
1130            .collect();
1131
1132        let cur_layer_node_pinned_usages = Self::read_layer_node_pinned_usages(
1133            &cpu_ctxs,
1134            &self.topo,
1135            self.nr_layers,
1136            self.topo.nodes.len(),
1137        );
1138        let cur_layer_node_usages = Self::read_layer_node_usages(
1139            &cpu_ctxs,
1140            &self.topo,
1141            self.nr_layers,
1142            self.topo.nodes.len(),
1143        );
1144        let cur_layer_node_duty_raw = Self::read_layer_node_duty_raw(
1145            &cpu_ctxs,
1146            &self.topo,
1147            self.nr_layers,
1148            self.topo.nodes.len(),
1149        );
1150        let cur_layer_membw_agg = Self::read_layer_membw_agg(&cpu_ctxs, self.nr_layers);
1151
1152        // Memory BW normalization. It requires finding the delta according to perf, the delta
1153        // according to resctl, and finding the factor between them. This also helps in
1154        // determining whether this delta is stable.
1155        //
1156        // We scale only the raw PMC delta - the part of the counter incremented between two
1157        // periods. This is because the scaling factor is only valid for that time period.
1158        // Non-scaled samples are not comparable, while scaled ones are.
1159        let (pmu_prev, resctrl_prev) = self.prev_pmu_resctrl_membw;
1160        let pmu_cur: u64 = cur_layer_membw_agg
1161            .iter()
1162            .map(|membw_agg| membw_agg[LAYER_USAGE_OPEN] + membw_agg[LAYER_USAGE_OWNED])
1163            .sum();
1164        let resctrl_cur = Self::resctrl_read_total_membw()?;
1165        let factor = (resctrl_cur - resctrl_prev) as f64 / (pmu_cur - pmu_prev) as f64;
1166
1167        // Computes the runtime deltas and converts them from ns to s.
1168        let compute_diff = |cur_agg: &Vec<Vec<u64>>, prev_agg: &Vec<Vec<u64>>| {
1169            cur_agg
1170                .iter()
1171                .zip(prev_agg.iter())
1172                .map(|(cur, prev)| {
1173                    cur.iter()
1174                        .zip(prev.iter())
1175                        .map(|(c, p)| (c - p) as f64 / 1_000_000_000.0 / elapsed_f64)
1176                        .collect()
1177                })
1178                .collect()
1179        };
1180
1181        // Computes the total memory traffic done since last computation and normalizes to GBs.
1182        // We derive the rate of consumption elsewhere.
1183        let compute_mem_diff = |cur_agg: &Vec<Vec<u64>>, prev_agg: &Vec<Vec<u64>>| {
1184            cur_agg
1185                .iter()
1186                .zip(prev_agg.iter())
1187                .map(|(cur, prev)| {
1188                    cur.iter()
1189                        .zip(prev.iter())
1190                        .map(|(c, p)| (*c as i64 - *p as i64) as f64 / 1024_f64.powf(3.0))
1191                        .collect()
1192                })
1193                .collect()
1194        };
1195
1196        // Scale the raw value delta by the resctrl/pmc computed factor.
1197        let cur_layer_membw: Vec<Vec<f64>> =
1198            compute_mem_diff(&cur_layer_membw_agg, &self.prev_layer_membw_agg);
1199
1200        let cur_layer_membw: Vec<Vec<f64>> = cur_layer_membw
1201            .iter()
1202            .map(|x| x.iter().map(|x| *x * factor).collect())
1203            .collect();
1204
1205        let metric_decay =
1206            |cur_metric: Vec<Vec<f64>>, prev_metric: &Vec<Vec<f64>>, decay_rate: f64| {
1207                cur_metric
1208                    .iter()
1209                    .zip(prev_metric.iter())
1210                    .map(|(cur, prev)| {
1211                        cur.iter()
1212                            .zip(prev.iter())
1213                            .map(|(c, p)| {
1214                                let decay = decay_rate.powf(elapsed_f64);
1215                                p * decay + c * (1.0 - decay)
1216                            })
1217                            .collect()
1218                    })
1219                    .collect()
1220            };
1221
1222        let cur_node_pinned_utils: Vec<Vec<f64>> = compute_diff(
1223            &cur_layer_node_pinned_usages,
1224            &self.prev_layer_node_pinned_usages,
1225        );
1226        let layer_node_pinned_utils: Vec<Vec<f64>> = metric_decay(
1227            cur_node_pinned_utils,
1228            &self.layer_node_pinned_utils,
1229            *USAGE_DECAY,
1230        );
1231        let cur_node_utils: Vec<Vec<f64>> =
1232            compute_diff(&cur_layer_node_usages, &self.prev_layer_node_usages);
1233        let layer_node_utils: Vec<Vec<f64>> =
1234            metric_decay(cur_node_utils, &self.layer_node_utils, *USAGE_DECAY);
1235        let cur_node_duty: Vec<Vec<f64>> =
1236            compute_diff(&cur_layer_node_duty_raw, &self.prev_layer_node_duty_raw);
1237        let layer_node_duty_sums: Vec<Vec<f64>> =
1238            metric_decay(cur_node_duty, &self.layer_node_duty_sums, *USAGE_DECAY);
1239
1240        let layer_membws: Vec<Vec<f64>> = metric_decay(cur_layer_membw, &self.layer_membws, 0.0);
1241
1242        let proc_stat = proc_reader
1243            .read_stat()
1244            .context("Failed to read /proc/stat")?;
1245        let cur_total_cpu = proc_stat
1246            .total_cpu
1247            .ok_or_else(|| anyhow!("Could not read total cpu stat in proc"))?;
1248        let cpu_busy = calc_util(&cur_total_cpu, &self.prev_total_cpu)?;
1249
1250        // Calculate system CPU utilization EWMA (10 second window)
1251        const SYS_CPU_UTIL_EWMA_SECS: f64 = 10.0;
1252        let elapsed_f64 = elapsed.as_secs_f64();
1253        let alpha = elapsed_f64 / SYS_CPU_UTIL_EWMA_SECS.max(elapsed_f64);
1254        let system_cpu_util_ewma = alpha * cpu_busy + (1.0 - alpha) * self.system_cpu_util_ewma;
1255
1256        // Per-CPU scale factors: s[c] = Δt / (Δt - irq - softirq - stolen).
1257        // When compensation is off, all scales stay 1.0.
1258        let cur_per_cpu_stats = proc_stat.cpus_map.unwrap_or_default();
1259        let mut cpu_scales = vec![1.0f64; *NR_CPUS_POSSIBLE];
1260        if self.util_compensation {
1261            for (&cpu_id, cur_cpu_stat) in &cur_per_cpu_stats {
1262                let cpu = cpu_id as usize;
1263                if let Some(prev_cpu_stat) = self.prev_per_cpu_stats.get(&cpu_id)
1264                    && let (
1265                        fb_procfs::CpuStat {
1266                            user_usec: Some(cu),
1267                            nice_usec: Some(cn),
1268                            system_usec: Some(cs),
1269                            idle_usec: Some(ci),
1270                            iowait_usec: Some(cw),
1271                            irq_usec: Some(cq),
1272                            softirq_usec: Some(cf),
1273                            stolen_usec: Some(ct),
1274                            ..
1275                        },
1276                        fb_procfs::CpuStat {
1277                            user_usec: Some(pu),
1278                            nice_usec: Some(pn),
1279                            system_usec: Some(ps),
1280                            idle_usec: Some(pi),
1281                            iowait_usec: Some(pw),
1282                            irq_usec: Some(pq),
1283                            softirq_usec: Some(pf),
1284                            stolen_usec: Some(pt),
1285                            ..
1286                        },
1287                    ) = (cur_cpu_stat, prev_cpu_stat)
1288                {
1289                    let delta_total = cu.saturating_sub(*pu)
1290                        + cn.saturating_sub(*pn)
1291                        + cs.saturating_sub(*ps)
1292                        + ci.saturating_sub(*pi)
1293                        + cw.saturating_sub(*pw)
1294                        + cq.saturating_sub(*pq)
1295                        + cf.saturating_sub(*pf)
1296                        + ct.saturating_sub(*pt);
1297                    let overhead =
1298                        cq.saturating_sub(*pq) + cf.saturating_sub(*pf) + ct.saturating_sub(*pt);
1299                    let available = delta_total.saturating_sub(overhead);
1300                    cpu_scales[cpu] = if available > 0 {
1301                        (delta_total as f64 / available as f64).clamp(1.0, 20.0)
1302                    } else {
1303                        1.0
1304                    };
1305                }
1306            }
1307        }
1308
1309        // Single pass over all CPUs: build both raw and compensated layer
1310        // util streams. When compensation is off, all scales are 1.0 so
1311        // comp == raw naturally.
1312        let cur_cpu_layer_usages = Self::read_per_cpu_layer_usages(&cpu_ctxs, self.nr_layers);
1313        let stride = self.nr_layers * NR_LAYER_USAGES;
1314        let mut raw_sums = vec![vec![0.0f64; NR_LAYER_USAGES]; self.nr_layers];
1315        let mut scaled_sums = vec![vec![0.0f64; NR_LAYER_USAGES]; self.nr_layers];
1316        #[allow(clippy::needless_range_loop)]
1317        for cpu in 0..*NR_CPUS_POSSIBLE {
1318            let scale = cpu_scales[cpu];
1319            let base = cpu * stride;
1320            for layer in 0..self.nr_layers {
1321                for usage in 0..NR_LAYER_USAGES {
1322                    let idx = base + layer * NR_LAYER_USAGES + usage;
1323                    let delta =
1324                        cur_cpu_layer_usages[idx].saturating_sub(self.prev_cpu_layer_usages[idx]);
1325                    if delta > 0 {
1326                        let delta_f = delta as f64;
1327                        raw_sums[layer][usage] += delta_f;
1328                        scaled_sums[layer][usage] += delta_f * scale;
1329                    }
1330                }
1331            }
1332        }
1333        let normalize = |sums: Vec<Vec<f64>>| -> Vec<Vec<f64>> {
1334            sums.into_iter()
1335                .map(|layer_sums| {
1336                    layer_sums
1337                        .into_iter()
1338                        .map(|s| s / 1_000_000_000.0 / elapsed_f64)
1339                        .collect()
1340                })
1341                .collect()
1342        };
1343        let layer_utils: Vec<Vec<f64>> =
1344            metric_decay(normalize(raw_sums), &self.layer_utils, *USAGE_DECAY);
1345        let layer_utils_compensated: Vec<Vec<f64>> = metric_decay(
1346            normalize(scaled_sums),
1347            &self.layer_utils_compensated,
1348            *USAGE_DECAY,
1349        );
1350
1351        let cur_bpf_stats = BpfStats::read(skel, &cpu_ctxs);
1352        let bpf_stats = &cur_bpf_stats - &self.prev_bpf_stats;
1353
1354        // Calculate per-layer DSQ insertion EWMA (10 second window)
1355        const DSQ_INSERT_EWMA_SECS: f64 = 10.0;
1356        let dsq_alpha = elapsed_f64 / DSQ_INSERT_EWMA_SECS.max(elapsed_f64);
1357        let layer_dsq_insert_ewma: Vec<f64> = (0..self.nr_layers)
1358            .map(|layer_id| {
1359                let sel_local = bpf_stats.lstats[layer_id]
1360                    [bpf_intf::layer_stat_id_LSTAT_SEL_LOCAL as usize]
1361                    as f64;
1362                let enq_local = bpf_stats.lstats[layer_id]
1363                    [bpf_intf::layer_stat_id_LSTAT_ENQ_LOCAL as usize]
1364                    as f64;
1365                let enq_dsq = bpf_stats.lstats[layer_id]
1366                    [bpf_intf::layer_stat_id_LSTAT_ENQ_DSQ as usize]
1367                    as f64;
1368                let total_dispatches = sel_local + enq_local + enq_dsq;
1369
1370                let cur_ratio = if total_dispatches > 0.0 {
1371                    enq_dsq / total_dispatches
1372                } else {
1373                    0.0
1374                };
1375
1376                dsq_alpha * cur_ratio + (1.0 - dsq_alpha) * self.layer_dsq_insert_ewma[layer_id]
1377            })
1378            .collect();
1379
1380        let processing_dur = cur_processing_dur
1381            .checked_sub(self.prev_processing_dur)
1382            .unwrap();
1383
1384        *self = Self {
1385            at: now,
1386            elapsed,
1387            topo: self.topo.clone(),
1388            nr_layers: self.nr_layers,
1389            nr_layer_tasks,
1390            layer_nr_node_pinned_tasks,
1391
1392            total_util: layer_utils
1393                .iter()
1394                .map(|x| x.iter().take(LAYER_USAGE_SUM_UPTO + 1).sum::<f64>())
1395                .sum(),
1396            layer_utils,
1397            layer_node_pinned_utils,
1398            prev_layer_node_pinned_usages: cur_layer_node_pinned_usages,
1399            layer_node_utils,
1400            prev_layer_node_usages: cur_layer_node_usages,
1401            layer_node_duty_sums,
1402            prev_layer_node_duty_raw: cur_layer_node_duty_raw,
1403
1404            layer_membws,
1405            prev_layer_membw_agg: cur_layer_membw_agg,
1406            // Was updated during normalization.
1407            prev_pmu_resctrl_membw: (pmu_cur, resctrl_cur),
1408
1409            cpu_busy,
1410            prev_total_cpu: cur_total_cpu,
1411            util_compensation: self.util_compensation,
1412            layer_utils_compensated,
1413            prev_cpu_layer_usages: cur_cpu_layer_usages,
1414            prev_per_cpu_stats: cur_per_cpu_stats,
1415            system_cpu_util_ewma,
1416            layer_dsq_insert_ewma,
1417
1418            bpf_stats,
1419            prev_bpf_stats: cur_bpf_stats,
1420
1421            processing_dur,
1422            prev_processing_dur: cur_processing_dur,
1423
1424            layer_slice_us,
1425            gpu_tasks_affinitized: gpu_task_affinitizer.tasks_affinitized,
1426            gpu_task_affinitization_ms: gpu_task_affinitizer.last_task_affinitization_ms,
1427        };
1428        Ok(())
1429    }
1430}
1431
1432#[derive(Debug)]
1433struct Layer {
1434    name: String,
1435    kind: LayerKind,
1436    growth_algo: LayerGrowthAlgo,
1437    core_order: Vec<Vec<usize>>,
1438
1439    assigned_llcs: Vec<Vec<usize>>,
1440
1441    nr_cpus: usize,
1442    nr_llc_cpus: Vec<usize>,
1443    nr_node_cpus: Vec<usize>,
1444    cpus: Cpumask,
1445    allowed_cpus: Cpumask,
1446
1447    /// Per-node count of CPUs allocated for pinned demand.
1448    nr_pinned_cpus: Vec<usize>,
1449}
1450
1451fn get_kallsyms_addr(sym_name: &str) -> Result<u64> {
1452    fs::read_to_string("/proc/kallsyms")?
1453        .lines()
1454        .find(|line| line.contains(sym_name))
1455        .and_then(|line| line.split_whitespace().next())
1456        .and_then(|addr| u64::from_str_radix(addr, 16).ok())
1457        .ok_or_else(|| anyhow!("Symbol '{}' not found", sym_name))
1458}
1459
1460fn resolve_cpus_pct_range(
1461    cpus_range: &Option<(usize, usize)>,
1462    cpus_range_frac: &Option<(f64, f64)>,
1463    max_cpus: usize,
1464) -> Result<(usize, usize)> {
1465    match (cpus_range, cpus_range_frac) {
1466        (Some(_x), Some(_y)) => {
1467            bail!("cpus_range cannot be used with cpus_pct.");
1468        }
1469        (Some((cpus_range_min, cpus_range_max)), None) => Ok((*cpus_range_min, *cpus_range_max)),
1470        (None, Some((cpus_frac_min, cpus_frac_max))) => {
1471            if *cpus_frac_min < 0_f64
1472                || *cpus_frac_min > 1_f64
1473                || *cpus_frac_max < 0_f64
1474                || *cpus_frac_max > 1_f64
1475            {
1476                bail!("cpus_range_frac values must be between 0.0 and 1.0");
1477            }
1478            let cpus_min_count = ((max_cpus as f64) * cpus_frac_min).round_ties_even() as usize;
1479            let cpus_max_count = ((max_cpus as f64) * cpus_frac_max).round_ties_even() as usize;
1480            Ok((
1481                std::cmp::max(cpus_min_count, 1),
1482                std::cmp::min(std::cmp::max(cpus_max_count, 1), max_cpus),
1483            ))
1484        }
1485        (None, None) => Ok((0, max_cpus)),
1486    }
1487}
1488
1489fn update_peak_util(prev_peak: f64, cur_util: f64, half_life: Duration, elapsed: Duration) -> f64 {
1490    let decay = 0.5f64.powf(elapsed.as_secs_f64() / half_life.as_secs_f64());
1491    (prev_peak * decay).max(cur_util)
1492}
1493
1494fn calc_peak_aware_target_range(
1495    util: f64,
1496    peak_util: f64,
1497    util_range: (f64, f64),
1498) -> (usize, usize) {
1499    let low = (util / util_range.1).ceil() as usize;
1500    let high = ((util.max(peak_util) / util_range.0).floor() as usize).max(low);
1501    (low, high)
1502}
1503
1504impl Layer {
1505    fn new(spec: &LayerSpec, topo: &Topology, core_order: &Vec<Vec<usize>>) -> Result<Self> {
1506        let name = &spec.name;
1507        let kind = spec.kind.clone();
1508        let mut allowed_cpus = Cpumask::new();
1509        match &kind {
1510            LayerKind::Confined {
1511                cpus_range,
1512                cpus_range_frac,
1513                common: LayerCommon { nodes, llcs, .. },
1514                ..
1515            } => {
1516                let cpus_range =
1517                    resolve_cpus_pct_range(cpus_range, cpus_range_frac, topo.all_cpus.len())?;
1518                if cpus_range.0 > cpus_range.1 || cpus_range.1 == 0 {
1519                    bail!("invalid cpus_range {:?}", cpus_range);
1520                }
1521                if nodes.is_empty() && llcs.is_empty() {
1522                    allowed_cpus.set_all();
1523                } else {
1524                    // build up the cpus bitset
1525                    for (node_id, node) in &topo.nodes {
1526                        // first do the matching for nodes
1527                        if nodes.contains(node_id) {
1528                            for &id in node.all_cpus.keys() {
1529                                allowed_cpus.set_cpu(id)?;
1530                            }
1531                        }
1532                        // next match on any LLCs
1533                        for (llc_id, llc) in &node.llcs {
1534                            if llcs.contains(llc_id) {
1535                                for &id in llc.all_cpus.keys() {
1536                                    allowed_cpus.set_cpu(id)?;
1537                                }
1538                            }
1539                        }
1540                    }
1541                }
1542            }
1543            LayerKind::Grouped {
1544                common: LayerCommon { nodes, llcs, .. },
1545                ..
1546            }
1547            | LayerKind::Open {
1548                common: LayerCommon { nodes, llcs, .. },
1549                ..
1550            } => {
1551                if nodes.is_empty() && llcs.is_empty() {
1552                    allowed_cpus.set_all();
1553                } else {
1554                    // build up the cpus bitset
1555                    for (node_id, node) in &topo.nodes {
1556                        // first do the matching for nodes
1557                        if nodes.contains(node_id) {
1558                            for &id in node.all_cpus.keys() {
1559                                allowed_cpus.set_cpu(id)?;
1560                            }
1561                        }
1562                        // next match on any LLCs
1563                        for (llc_id, llc) in &node.llcs {
1564                            if llcs.contains(llc_id) {
1565                                for &id in llc.all_cpus.keys() {
1566                                    allowed_cpus.set_cpu(id)?;
1567                                }
1568                            }
1569                        }
1570                    }
1571                }
1572            }
1573        }
1574
1575        // Util can be above 1.0 for grouped layers if
1576        // util_includes_open_cputime is set.
1577        if let Some(util_range) = kind.util_range() {
1578            if util_range.0 < 0.0 || util_range.1 < 0.0 || util_range.0 >= util_range.1 {
1579                bail!("invalid util_range {:?}", util_range);
1580            }
1581        } else if kind.common().util_peak_half_life_ms != 0 {
1582            bail!("util_peak_half_life_ms requires util_range");
1583        }
1584
1585        let layer_growth_algo = kind.common().growth_algo.clone();
1586
1587        debug!(
1588            "layer: {} algo: {:?} core order: {:?}",
1589            name, &layer_growth_algo, core_order
1590        );
1591
1592        Ok(Self {
1593            name: name.into(),
1594            kind,
1595            growth_algo: layer_growth_algo,
1596            core_order: core_order.clone(),
1597
1598            assigned_llcs: vec![vec![]; topo.nodes.len()],
1599
1600            nr_cpus: 0,
1601            nr_llc_cpus: vec![0; topo.all_llcs.len()],
1602            nr_node_cpus: vec![0; topo.nodes.len()],
1603            cpus: Cpumask::new(),
1604            allowed_cpus,
1605
1606            nr_pinned_cpus: vec![0; topo.nodes.len()],
1607        })
1608    }
1609}
1610#[derive(Debug, Clone)]
1611struct NodeInfo {
1612    node_mask: nix::sched::CpuSet,
1613    _node_id: usize,
1614}
1615
1616#[derive(Debug)]
1617struct GpuTaskAffinitizer {
1618    // This struct tracks information necessary to numa affinitize
1619    // gpu tasks periodically when needed.
1620    gpu_devs_to_node_info: HashMap<u32, NodeInfo>,
1621    gpu_pids_to_devs: HashMap<Pid, u32>,
1622    last_process_time: Option<Instant>,
1623    sys: System,
1624    pid_map: HashMap<Pid, Vec<Pid>>,
1625    poll_interval: Duration,
1626    enable: bool,
1627    tasks_affinitized: u64,
1628    last_task_affinitization_ms: u64,
1629}
1630
1631impl GpuTaskAffinitizer {
1632    pub fn new(poll_interval: u64, enable: bool) -> GpuTaskAffinitizer {
1633        GpuTaskAffinitizer {
1634            gpu_devs_to_node_info: HashMap::new(),
1635            gpu_pids_to_devs: HashMap::new(),
1636            last_process_time: None,
1637            sys: System::default(),
1638            pid_map: HashMap::new(),
1639            poll_interval: Duration::from_secs(poll_interval),
1640            enable,
1641            tasks_affinitized: 0,
1642            last_task_affinitization_ms: 0,
1643        }
1644    }
1645
1646    fn find_one_cpu(&self, affinity: Vec<u64>) -> Result<u32> {
1647        for (chunk, &mask) in affinity.iter().enumerate() {
1648            let mut inner_offset: u64 = 1;
1649            for _ in 0..64 {
1650                if (mask & inner_offset) != 0 {
1651                    return Ok((64 * chunk + u64::trailing_zeros(inner_offset) as usize) as u32);
1652                }
1653                inner_offset <<= 1;
1654            }
1655        }
1656        anyhow::bail!("unable to get CPU from NVML bitmask");
1657    }
1658
1659    fn node_to_cpuset(&self, node: &scx_utils::Node) -> Result<CpuSet> {
1660        let mut cpuset = CpuSet::new();
1661        for cpu_id in node.all_cpus.keys() {
1662            cpuset.set(*cpu_id)?;
1663        }
1664        Ok(cpuset)
1665    }
1666
1667    fn init_dev_node_map(&mut self, topo: Arc<Topology>) -> Result<()> {
1668        let nvml = nvml()?;
1669        let device_count = nvml.device_count()?;
1670
1671        for idx in 0..device_count {
1672            let dev = nvml.device_by_index(idx)?;
1673            // For machines w/ up to 1024 CPUs.
1674            let cpu = dev.cpu_affinity(16)?;
1675            let ideal_cpu = self.find_one_cpu(cpu)?;
1676            if let Some(cpu) = topo.all_cpus.get(&(ideal_cpu as usize)) {
1677                self.gpu_devs_to_node_info.insert(
1678                    idx,
1679                    NodeInfo {
1680                        node_mask: self.node_to_cpuset(
1681                            topo.nodes.get(&cpu.node_id).expect("topo missing node"),
1682                        )?,
1683                        _node_id: cpu.node_id,
1684                    },
1685                );
1686            }
1687        }
1688        Ok(())
1689    }
1690
1691    fn update_gpu_pids(&mut self) -> Result<()> {
1692        let nvml = nvml()?;
1693        for i in 0..nvml.device_count()? {
1694            let device = nvml.device_by_index(i)?;
1695            for proc in device
1696                .running_compute_processes()?
1697                .into_iter()
1698                .chain(device.running_graphics_processes()?)
1699            {
1700                self.gpu_pids_to_devs.insert(Pid::from_u32(proc.pid), i);
1701            }
1702        }
1703        Ok(())
1704    }
1705
1706    fn update_process_info(&mut self) -> Result<()> {
1707        self.sys.refresh_processes_specifics(
1708            ProcessesToUpdate::All,
1709            true,
1710            ProcessRefreshKind::nothing(),
1711        );
1712        self.pid_map.clear();
1713        for (pid, proc_) in self.sys.processes() {
1714            if let Some(ppid) = proc_.parent() {
1715                self.pid_map.entry(ppid).or_default().push(*pid);
1716            }
1717        }
1718        Ok(())
1719    }
1720
1721    fn get_child_pids_and_tids(&self, root_pid: Pid) -> HashSet<Pid> {
1722        let mut work = VecDeque::from([root_pid]);
1723        let mut pids_and_tids: HashSet<Pid> = HashSet::new();
1724
1725        while let Some(pid) = work.pop_front() {
1726            if pids_and_tids.insert(pid) {
1727                if let Some(kids) = self.pid_map.get(&pid) {
1728                    work.extend(kids);
1729                }
1730                if let Some(proc_) = self.sys.process(pid)
1731                    && let Some(tasks) = proc_.tasks()
1732                {
1733                    pids_and_tids.extend(tasks.iter().copied());
1734                }
1735            }
1736        }
1737        pids_and_tids
1738    }
1739
1740    fn affinitize_gpu_pids(&mut self) -> Result<()> {
1741        if !self.enable {
1742            return Ok(());
1743        }
1744        for (pid, dev) in &self.gpu_pids_to_devs {
1745            let node_info = self
1746                .gpu_devs_to_node_info
1747                .get(dev)
1748                .expect("Unable to get gpu pid node mask");
1749            for child in self.get_child_pids_and_tids(*pid) {
1750                match nix::sched::sched_setaffinity(
1751                    nix::unistd::Pid::from_raw(child.as_u32() as i32),
1752                    &node_info.node_mask,
1753                ) {
1754                    Ok(_) => {
1755                        // Increment the global counter for successful affinitization
1756                        self.tasks_affinitized += 1;
1757                    }
1758                    Err(_) => {
1759                        debug!(
1760                            "Error affinitizing gpu pid {} to node {:#?}",
1761                            child.as_u32(),
1762                            node_info
1763                        );
1764                    }
1765                };
1766            }
1767        }
1768        Ok(())
1769    }
1770
1771    pub fn maybe_affinitize(&mut self) {
1772        if !self.enable {
1773            return;
1774        }
1775        let now = Instant::now();
1776
1777        if let Some(last_process_time) = self.last_process_time
1778            && (now - last_process_time) < self.poll_interval
1779        {
1780            return;
1781        }
1782
1783        match self.update_gpu_pids() {
1784            Ok(_) => {}
1785            Err(e) => {
1786                error!("Error updating GPU PIDs: {}", e);
1787            }
1788        };
1789        match self.update_process_info() {
1790            Ok(_) => {}
1791            Err(e) => {
1792                error!("Error updating process info to affinitize GPU PIDs: {}", e);
1793            }
1794        };
1795        match self.affinitize_gpu_pids() {
1796            Ok(_) => {}
1797            Err(e) => {
1798                error!("Error updating GPU PIDs: {}", e);
1799            }
1800        };
1801        self.last_process_time = Some(now);
1802        self.last_task_affinitization_ms = (Instant::now() - now).as_millis() as u64;
1803    }
1804
1805    pub fn init(&mut self, topo: Arc<Topology>) {
1806        if !self.enable || self.last_process_time.is_some() {
1807            return;
1808        }
1809
1810        match self.init_dev_node_map(topo) {
1811            Ok(_) => {}
1812            Err(e) => {
1813                error!("Error initializing gpu node dev map: {}", e);
1814            }
1815        };
1816        self.sys = System::new_all();
1817    }
1818}
1819
1820struct Scheduler<'a> {
1821    skel: BpfSkel<'a>,
1822    struct_ops: Option<libbpf_rs::Link>,
1823    layer_specs: Vec<LayerSpec>,
1824
1825    sched_intv: Duration,
1826    layer_refresh_intv: Duration,
1827
1828    cpu_pool: CpuPool,
1829    layers: Vec<Layer>,
1830    idle_qos_enabled: bool,
1831
1832    proc_reader: fb_procfs::ProcReader,
1833    sched_stats: Stats,
1834    layer_peak_utils: Vec<f64>,
1835
1836    cgroup_regexes: Option<HashMap<u32, Regex>>,
1837
1838    nr_layer_cpus_ranges: Vec<(usize, usize)>,
1839    xnuma_mig_src: Vec<Vec<bool>>,
1840    growth_denied: Vec<Vec<bool>>,
1841    processing_dur: Duration,
1842
1843    topo: Arc<Topology>,
1844    netdevs: BTreeMap<String, NetDev>,
1845    stats_server: StatsServer<StatsReq, StatsRes>,
1846    gpu_task_handler: GpuTaskAffinitizer,
1847}
1848
1849const DUTY_CYCLE_SCALE: f64 = (1u64 << 20) as f64;
1850const XNUMA_RATE_DAMPEN: f64 = 0.5;
1851
1852/// Result of xnuma water-fill computation for a single layer.
1853struct XnumaRates {
1854    /// rates[src][dst]: migration rate in duty-cycle-scaled units.
1855    rates: Vec<Vec<u64>>,
1856}
1857
1858/// Determine per-node migration source state with two-threshold hysteresis.
1859///
1860/// Each (layer, node) independently decides if it's a migration source.
1861/// Open (is_mig_src=true) requires all three:
1862///   1. load/alloc > threshold.1 (significant load)
1863///   2. surplus/alloc > delta.1 (significant imbalance)
1864///   3. growth_denied (allocation can't solve it)
1865///
1866/// Close (is_mig_src=false) when any one:
1867///   1. load/alloc < threshold.0 (load dropped)
1868///   2. surplus/alloc < delta.0 (imbalance resolved)
1869///   3. !growth_denied (growth succeeded)
1870fn xnuma_check_active(
1871    duty_sums: &[f64],
1872    allocs: &[usize],
1873    threshold: (f64, f64),
1874    threshold_delta: (f64, f64),
1875    growth_denied: &[bool],
1876    currently_active: &[bool],
1877) -> Vec<bool> {
1878    let nr_nodes = duty_sums.len();
1879    let total_duty: f64 = duty_sums.iter().sum();
1880    let total_alloc: f64 = allocs.iter().map(|&a| a as f64).sum();
1881    let eq_ratio = if total_alloc > 0.0 {
1882        total_duty / total_alloc
1883    } else {
1884        0.0
1885    };
1886
1887    let (thresh_lo, thresh_hi) = threshold;
1888    let (delta_lo, delta_hi) = threshold_delta;
1889
1890    let mut result = vec![false; nr_nodes];
1891    for nid in 0..nr_nodes {
1892        let alloc = allocs[nid] as f64;
1893        if alloc <= 0.0 {
1894            if duty_sums[nid] > 0.0 && growth_denied[nid] {
1895                result[nid] = true;
1896            }
1897            continue;
1898        }
1899
1900        let load_ratio = duty_sums[nid] / alloc;
1901        let surplus = duty_sums[nid] - eq_ratio * alloc;
1902        let surplus_ratio = surplus / alloc;
1903
1904        let should_activate =
1905            load_ratio > thresh_hi && surplus_ratio > delta_hi && growth_denied[nid];
1906        let should_deactivate =
1907            load_ratio < thresh_lo || surplus_ratio < delta_lo || !growth_denied[nid];
1908
1909        if should_activate {
1910            result[nid] = true;
1911        } else if should_deactivate {
1912            result[nid] = false;
1913        } else {
1914            result[nid] = currently_active[nid];
1915        }
1916    }
1917    result
1918}
1919
1920/// Compute water-fill migration rates for a single layer.
1921///
1922/// Finds the equalization ratio (water line) across all nodes, then
1923/// computes per-(src, dst) migration rates proportional to each source's
1924/// surplus and each destination's share of total deficit.
1925fn xnuma_compute_rates(duty_sums: &[f64], allocs: &[usize]) -> XnumaRates {
1926    let nr_nodes = duty_sums.len();
1927    let total_duty: f64 = duty_sums.iter().sum();
1928    let total_alloc: f64 = allocs.iter().map(|&a| a as f64).sum();
1929
1930    if total_alloc <= 0.0 {
1931        return XnumaRates {
1932            rates: vec![vec![0u64; nr_nodes]; nr_nodes],
1933        };
1934    }
1935
1936    let eq_ratio = total_duty / total_alloc;
1937
1938    let mut surpluses = vec![0.0f64; nr_nodes];
1939    let mut deficits = vec![0.0f64; nr_nodes];
1940    for nid in 0..nr_nodes {
1941        let expected = eq_ratio * allocs[nid] as f64;
1942        let delta = duty_sums[nid] - expected;
1943        if delta > 0.0 {
1944            surpluses[nid] = delta;
1945        } else {
1946            deficits[nid] = -delta;
1947        }
1948    }
1949
1950    let total_deficit: f64 = deficits.iter().sum();
1951
1952    let mut rates = vec![vec![0u64; nr_nodes]; nr_nodes];
1953    for src in 0..nr_nodes {
1954        for dst in 0..nr_nodes {
1955            if src == dst || total_deficit <= 0.0 || surpluses[src] <= 0.0 {
1956                continue;
1957            }
1958            // Dampen: transfer half the surplus per cycle so convergence
1959            // is gradual rather than a single-step overcorrection.
1960            let migration = surpluses[src] * deficits[dst] / total_deficit * XNUMA_RATE_DAMPEN;
1961            rates[src][dst] = (migration * DUTY_CYCLE_SCALE) as u64;
1962        }
1963    }
1964
1965    XnumaRates { rates }
1966}
1967
1968impl<'a> Scheduler<'a> {
1969    fn init_layers(
1970        skel: &mut OpenBpfSkel,
1971        specs: &[LayerSpec],
1972        topo: &Topology,
1973    ) -> Result<HashMap<u32, Regex>> {
1974        skel.maps.rodata_data.as_mut().unwrap().nr_layers = specs.len() as u32;
1975        let mut perf_set = false;
1976
1977        let mut layer_iteration_order = (0..specs.len()).collect::<Vec<_>>();
1978        let mut layer_weights: Vec<usize> = vec![];
1979        let mut cgroup_regex_id = 0;
1980        let mut cgroup_regexes = HashMap::new();
1981
1982        for (spec_i, spec) in specs.iter().enumerate() {
1983            let layer = &mut skel.maps.bss_data.as_mut().unwrap().layers[spec_i];
1984
1985            for (or_i, or) in spec.matches.iter().enumerate() {
1986                for (and_i, and) in or.iter().enumerate() {
1987                    let mt = &mut layer.matches[or_i].matches[and_i];
1988
1989                    // Rules are allowlist-based by default
1990                    mt.exclude.write(false);
1991
1992                    match and {
1993                        LayerMatch::CgroupPrefix(prefix) => {
1994                            mt.kind = bpf_intf::layer_match_kind_MATCH_CGROUP_PREFIX as i32;
1995                            copy_into_cstr(&mut mt.cgroup_prefix, prefix.as_str());
1996                        }
1997                        LayerMatch::CgroupSuffix(suffix) => {
1998                            mt.kind = bpf_intf::layer_match_kind_MATCH_CGROUP_SUFFIX as i32;
1999                            copy_into_cstr(&mut mt.cgroup_suffix, suffix.as_str());
2000                        }
2001                        LayerMatch::CgroupRegex(regex_str) => {
2002                            if cgroup_regex_id >= bpf_intf::consts_MAX_CGROUP_REGEXES {
2003                                bail!(
2004                                    "Too many cgroup regex rules. Maximum allowed: {}",
2005                                    bpf_intf::consts_MAX_CGROUP_REGEXES
2006                                );
2007                            }
2008
2009                            // CgroupRegex matching handled in userspace via cgroup watcher
2010                            mt.kind = bpf_intf::layer_match_kind_MATCH_CGROUP_REGEX as i32;
2011                            mt.cgroup_regex_id = cgroup_regex_id;
2012
2013                            let regex = Regex::new(regex_str).with_context(|| {
2014                                format!("Invalid regex '{}' in layer '{}'", regex_str, spec.name)
2015                            })?;
2016                            cgroup_regexes.insert(cgroup_regex_id, regex);
2017                            cgroup_regex_id += 1;
2018                        }
2019                        LayerMatch::CgroupContains(substr) => {
2020                            mt.kind = bpf_intf::layer_match_kind_MATCH_CGROUP_CONTAINS as i32;
2021                            copy_into_cstr(&mut mt.cgroup_substr, substr.as_str());
2022                        }
2023                        LayerMatch::CommPrefix(prefix) => {
2024                            mt.kind = bpf_intf::layer_match_kind_MATCH_COMM_PREFIX as i32;
2025                            copy_into_cstr(&mut mt.comm_prefix, prefix.as_str());
2026                        }
2027                        LayerMatch::CommPrefixExclude(prefix) => {
2028                            mt.kind = bpf_intf::layer_match_kind_MATCH_COMM_PREFIX as i32;
2029                            mt.exclude.write(true);
2030                            copy_into_cstr(&mut mt.comm_prefix, prefix.as_str());
2031                        }
2032                        LayerMatch::PcommPrefix(prefix) => {
2033                            mt.kind = bpf_intf::layer_match_kind_MATCH_PCOMM_PREFIX as i32;
2034                            copy_into_cstr(&mut mt.pcomm_prefix, prefix.as_str());
2035                        }
2036                        LayerMatch::PcommPrefixExclude(prefix) => {
2037                            mt.kind = bpf_intf::layer_match_kind_MATCH_PCOMM_PREFIX as i32;
2038                            mt.exclude.write(true);
2039                            copy_into_cstr(&mut mt.pcomm_prefix, prefix.as_str());
2040                        }
2041                        LayerMatch::NiceAbove(nice) => {
2042                            mt.kind = bpf_intf::layer_match_kind_MATCH_NICE_ABOVE as i32;
2043                            mt.nice = *nice;
2044                        }
2045                        LayerMatch::NiceBelow(nice) => {
2046                            mt.kind = bpf_intf::layer_match_kind_MATCH_NICE_BELOW as i32;
2047                            mt.nice = *nice;
2048                        }
2049                        LayerMatch::NiceEquals(nice) => {
2050                            mt.kind = bpf_intf::layer_match_kind_MATCH_NICE_EQUALS as i32;
2051                            mt.nice = *nice;
2052                        }
2053                        LayerMatch::UIDEquals(user_id) => {
2054                            mt.kind = bpf_intf::layer_match_kind_MATCH_USER_ID_EQUALS as i32;
2055                            mt.user_id = *user_id;
2056                        }
2057                        LayerMatch::GIDEquals(group_id) => {
2058                            mt.kind = bpf_intf::layer_match_kind_MATCH_GROUP_ID_EQUALS as i32;
2059                            mt.group_id = *group_id;
2060                        }
2061                        LayerMatch::PIDEquals(pid) => {
2062                            mt.kind = bpf_intf::layer_match_kind_MATCH_PID_EQUALS as i32;
2063                            mt.pid = *pid;
2064                        }
2065                        LayerMatch::PPIDEquals(ppid) => {
2066                            mt.kind = bpf_intf::layer_match_kind_MATCH_PPID_EQUALS as i32;
2067                            mt.ppid = *ppid;
2068                        }
2069                        LayerMatch::TGIDEquals(tgid) => {
2070                            mt.kind = bpf_intf::layer_match_kind_MATCH_TGID_EQUALS as i32;
2071                            mt.tgid = *tgid;
2072                        }
2073                        LayerMatch::NSPIDEquals(nsid, pid) => {
2074                            mt.kind = bpf_intf::layer_match_kind_MATCH_NSPID_EQUALS as i32;
2075                            mt.nsid = *nsid;
2076                            mt.pid = *pid;
2077                        }
2078                        LayerMatch::NSEquals(nsid) => {
2079                            mt.kind = bpf_intf::layer_match_kind_MATCH_NS_EQUALS as i32;
2080                            mt.nsid = *nsid as u64;
2081                        }
2082                        LayerMatch::CmdJoin(joincmd) => {
2083                            mt.kind = bpf_intf::layer_match_kind_MATCH_SCXCMD_JOIN as i32;
2084                            copy_into_cstr(&mut mt.comm_prefix, joincmd);
2085                        }
2086                        LayerMatch::IsGroupLeader(polarity) => {
2087                            mt.kind = bpf_intf::layer_match_kind_MATCH_IS_GROUP_LEADER as i32;
2088                            mt.is_group_leader.write(*polarity);
2089                        }
2090                        LayerMatch::IsKthread(polarity) => {
2091                            mt.kind = bpf_intf::layer_match_kind_MATCH_IS_KTHREAD as i32;
2092                            mt.is_kthread.write(*polarity);
2093                        }
2094                        LayerMatch::UsedGpuTid(polarity) => {
2095                            mt.kind = bpf_intf::layer_match_kind_MATCH_USED_GPU_TID as i32;
2096                            mt.used_gpu_tid.write(*polarity);
2097                        }
2098                        LayerMatch::UsedGpuPid(polarity) => {
2099                            mt.kind = bpf_intf::layer_match_kind_MATCH_USED_GPU_PID as i32;
2100                            mt.used_gpu_pid.write(*polarity);
2101                        }
2102                        LayerMatch::AvgRuntime(min, max) => {
2103                            mt.kind = bpf_intf::layer_match_kind_MATCH_AVG_RUNTIME as i32;
2104                            mt.min_avg_runtime_us = *min;
2105                            mt.max_avg_runtime_us = *max;
2106                        }
2107                        LayerMatch::HintEquals(hint) => {
2108                            mt.kind = bpf_intf::layer_match_kind_MATCH_HINT_EQUALS as i32;
2109                            mt.hint = *hint;
2110                        }
2111                        LayerMatch::SystemCpuUtilBelow(threshold) => {
2112                            mt.kind = bpf_intf::layer_match_kind_MATCH_SYSTEM_CPU_UTIL_BELOW as i32;
2113                            mt.system_cpu_util_below = (*threshold * 10000.0) as u64;
2114                        }
2115                        LayerMatch::DsqInsertBelow(threshold) => {
2116                            mt.kind = bpf_intf::layer_match_kind_MATCH_DSQ_INSERT_BELOW as i32;
2117                            mt.dsq_insert_below = (*threshold * 10000.0) as u64;
2118                        }
2119                        LayerMatch::NumaNode(node_id) => {
2120                            if *node_id as usize >= topo.nodes.len() {
2121                                bail!(
2122                                    "Spec {:?} has invalid NUMA node ID {} (available nodes: 0-{})",
2123                                    spec.name,
2124                                    node_id,
2125                                    topo.nodes.len() - 1
2126                                );
2127                            }
2128                            mt.kind = bpf_intf::layer_match_kind_MATCH_NUMA_NODE as i32;
2129                            mt.numa_node_id = *node_id;
2130                        }
2131                    }
2132                }
2133                layer.matches[or_i].nr_match_ands = or.len() as i32;
2134            }
2135
2136            layer.nr_match_ors = spec.matches.len() as u32;
2137            layer.kind = spec.kind.as_bpf_enum();
2138
2139            {
2140                let LayerCommon {
2141                    min_exec_us,
2142                    yield_ignore,
2143                    perf,
2144                    preempt,
2145                    preempt_first,
2146                    exclusive,
2147                    skip_remote_node,
2148                    prev_over_idle_core,
2149                    growth_algo,
2150                    slice_us,
2151                    fifo,
2152                    weight,
2153                    disallow_open_after_us,
2154                    disallow_preempt_after_us,
2155                    xllc_mig_min_us,
2156                    placement,
2157                    member_expire_ms,
2158                    ..
2159                } = spec.kind.common();
2160
2161                layer.slice_ns = *slice_us * 1000;
2162                layer.fifo.write(*fifo);
2163                layer.min_exec_ns = min_exec_us * 1000;
2164                layer.yield_step_ns = if *yield_ignore > 0.999 {
2165                    0
2166                } else if *yield_ignore < 0.001 {
2167                    layer.slice_ns
2168                } else {
2169                    (layer.slice_ns as f64 * (1.0 - *yield_ignore)) as u64
2170                };
2171                let mut layer_name: String = spec.name.clone();
2172                layer_name.truncate(MAX_LAYER_NAME);
2173                copy_into_cstr(&mut layer.name, layer_name.as_str());
2174                layer.preempt.write(*preempt);
2175                layer.preempt_first.write(*preempt_first);
2176                layer.excl.write(*exclusive);
2177                layer.skip_remote_node.write(*skip_remote_node);
2178                layer.prev_over_idle_core.write(*prev_over_idle_core);
2179                layer.growth_algo = growth_algo.as_bpf_enum();
2180                layer.weight = *weight;
2181                layer.member_expire_ms = *member_expire_ms;
2182                layer.disallow_open_after_ns = match disallow_open_after_us.unwrap() {
2183                    v if v == u64::MAX => v,
2184                    v => v * 1000,
2185                };
2186                layer.disallow_preempt_after_ns = match disallow_preempt_after_us.unwrap() {
2187                    v if v == u64::MAX => v,
2188                    v => v * 1000,
2189                };
2190                layer.xllc_mig_min_ns = (xllc_mig_min_us * 1000.0) as u64;
2191                layer_weights.push(layer.weight.try_into().unwrap());
2192                layer.perf = u32::try_from(*perf)?;
2193
2194                let task_place = |place: u32| crate::types::layer_task_place(place);
2195                layer.task_place = match placement {
2196                    LayerPlacement::Standard => {
2197                        task_place(bpf_intf::layer_task_place_PLACEMENT_STD)
2198                    }
2199                    LayerPlacement::Sticky => {
2200                        task_place(bpf_intf::layer_task_place_PLACEMENT_STICK)
2201                    }
2202                    LayerPlacement::Floating => {
2203                        task_place(bpf_intf::layer_task_place_PLACEMENT_FLOAT)
2204                    }
2205                };
2206            }
2207
2208            layer.is_protected.write(match spec.kind {
2209                LayerKind::Open { .. } => false,
2210                LayerKind::Confined { protected, .. } | LayerKind::Grouped { protected, .. } => {
2211                    protected
2212                }
2213            });
2214
2215            layer.idle_confined.write(match spec.kind {
2216                LayerKind::Grouped { idle_confined, .. } => idle_confined,
2217                _ => false,
2218            });
2219
2220            match &spec.cpuset {
2221                Some(mask) => {
2222                    Self::update_cpumask(mask, &mut layer.cpuset);
2223                    layer.has_cpuset.write(true);
2224                }
2225                None => {
2226                    for i in 0..layer.cpuset.len() {
2227                        layer.cpuset[i] = u8::MAX;
2228                    }
2229                    layer.has_cpuset.write(false);
2230                }
2231            };
2232
2233            perf_set |= layer.perf > 0;
2234        }
2235
2236        layer_iteration_order.sort_by(|i, j| layer_weights[*i].cmp(&layer_weights[*j]));
2237        for (idx, layer_idx) in layer_iteration_order.iter().enumerate() {
2238            skel.maps
2239                .rodata_data
2240                .as_mut()
2241                .unwrap()
2242                .layer_iteration_order[idx] = *layer_idx as u32;
2243        }
2244
2245        if perf_set && !compat::ksym_exists("scx_bpf_cpuperf_set")? {
2246            warn!("cpufreq support not available, ignoring perf configurations");
2247        }
2248
2249        Ok(cgroup_regexes)
2250    }
2251
2252    fn init_nodes(skel: &mut OpenBpfSkel, _opts: &Opts, topo: &Topology) {
2253        skel.maps.rodata_data.as_mut().unwrap().nr_nodes = topo.nodes.len() as u32;
2254        skel.maps.rodata_data.as_mut().unwrap().nr_llcs = 0;
2255
2256        for (&node_id, node) in &topo.nodes {
2257            debug!("configuring node {}, LLCs {:?}", node_id, node.llcs.len());
2258            skel.maps.rodata_data.as_mut().unwrap().nr_llcs += node.llcs.len() as u32;
2259            let raw_numa_slice = node.span.as_raw_slice();
2260            let node_cpumask_slice =
2261                &mut skel.maps.rodata_data.as_mut().unwrap().numa_cpumasks[node_id];
2262            let (left, _) = node_cpumask_slice.split_at_mut(raw_numa_slice.len());
2263            left.clone_from_slice(raw_numa_slice);
2264            debug!(
2265                "node {} mask: {:?}",
2266                node_id,
2267                skel.maps.rodata_data.as_ref().unwrap().numa_cpumasks[node_id]
2268            );
2269
2270            for llc in node.llcs.values() {
2271                debug!("configuring llc {:?} for node {:?}", llc.id, node_id);
2272                skel.maps.rodata_data.as_mut().unwrap().llc_numa_id_map[llc.id] = node_id as u32;
2273            }
2274        }
2275
2276        for cpu in topo.all_cpus.values() {
2277            skel.maps.rodata_data.as_mut().unwrap().cpu_llc_id_map[cpu.id] = cpu.llc_id as u32;
2278        }
2279    }
2280
2281    fn init_cpu_prox_map(topo: &Topology, cpu_ctxs: &mut [bpf_intf::cpu_ctx]) {
2282        let radiate = |mut vec: Vec<usize>, center_id: usize| -> Vec<usize> {
2283            vec.sort_by_key(|&id| (center_id as i32 - id as i32).abs());
2284            vec
2285        };
2286        let radiate_cpu =
2287            |mut vec: Vec<usize>, center_cpu: usize, center_core: usize| -> Vec<usize> {
2288                vec.sort_by_key(|&id| {
2289                    (
2290                        (center_core as i32 - topo.all_cpus.get(&id).unwrap().core_id as i32).abs(),
2291                        (center_cpu as i32 - id as i32).abs(),
2292                    )
2293                });
2294                vec
2295            };
2296
2297        for (&cpu_id, cpu) in &topo.all_cpus {
2298            // Collect the spans.
2299            let mut core_span = topo.all_cores[&cpu.core_id].span.clone();
2300            let llc_span = &topo.all_llcs[&cpu.llc_id].span;
2301            let node_span = &topo.nodes[&cpu.node_id].span;
2302            let sys_span = &topo.span;
2303
2304            // Make the spans exclusive.
2305            let sys_span = sys_span.and(&node_span.not());
2306            let node_span = node_span.and(&llc_span.not());
2307            let llc_span = llc_span.and(&core_span.not());
2308            core_span.clear_cpu(cpu_id).unwrap();
2309
2310            // Convert them into arrays.
2311            let mut sys_order: Vec<usize> = sys_span.iter().collect();
2312            let mut node_order: Vec<usize> = node_span.iter().collect();
2313            let mut llc_order: Vec<usize> = llc_span.iter().collect();
2314            let mut core_order: Vec<usize> = core_span.iter().collect();
2315
2316            // Shuffle them so that different CPUs follow different orders.
2317            // Each CPU radiates in both directions based on the cpu id and
2318            // radiates out to the closest cores based on core ids.
2319
2320            sys_order = radiate_cpu(sys_order, cpu_id, cpu.core_id);
2321            node_order = radiate(node_order, cpu.node_id);
2322            llc_order = radiate_cpu(llc_order, cpu_id, cpu.core_id);
2323            core_order = radiate_cpu(core_order, cpu_id, cpu.core_id);
2324
2325            // Concatenate and record the topology boundaries.
2326            let mut order: Vec<usize> = vec![];
2327            let mut idx: usize = 0;
2328
2329            idx += 1;
2330            order.push(cpu_id);
2331
2332            idx += core_order.len();
2333            order.append(&mut core_order);
2334            let core_end = idx;
2335
2336            idx += llc_order.len();
2337            order.append(&mut llc_order);
2338            let llc_end = idx;
2339
2340            idx += node_order.len();
2341            order.append(&mut node_order);
2342            let node_end = idx;
2343
2344            idx += sys_order.len();
2345            order.append(&mut sys_order);
2346            let sys_end = idx;
2347
2348            debug!(
2349                "CPU[{}] proximity map[{}/{}/{}/{}]: {:?}",
2350                cpu_id, core_end, llc_end, node_end, sys_end, &order
2351            );
2352
2353            // Record in cpu_ctx.
2354            let pmap = &mut cpu_ctxs[cpu_id].prox_map;
2355            for (i, &cpu) in order.iter().enumerate() {
2356                pmap.cpus[i] = cpu as u16;
2357            }
2358            pmap.core_end = core_end as u32;
2359            pmap.llc_end = llc_end as u32;
2360            pmap.node_end = node_end as u32;
2361            pmap.sys_end = sys_end as u32;
2362        }
2363    }
2364
2365    fn convert_cpu_ctxs(cpu_ctxs: Vec<bpf_intf::cpu_ctx>) -> Vec<Vec<u8>> {
2366        cpu_ctxs
2367            .iter()
2368            .map(|cpu_ctx| unsafe { plain::as_bytes(cpu_ctx) }.to_vec())
2369            .collect()
2370    }
2371
2372    fn init_cpus(skel: &BpfSkel, layer_specs: &[LayerSpec], topo: &Topology) -> Result<()> {
2373        let key = (0_u32).to_ne_bytes();
2374        let mut cpu_ctxs: Vec<bpf_intf::cpu_ctx> = vec![];
2375        let cpu_ctxs_vec = skel
2376            .maps
2377            .cpu_ctxs
2378            .lookup_percpu(&key, libbpf_rs::MapFlags::ANY)
2379            .context("Failed to lookup cpu_ctx")?
2380            .unwrap();
2381
2382        let op_layers: Vec<u32> = layer_specs
2383            .iter()
2384            .enumerate()
2385            .filter(|(_idx, spec)| match &spec.kind {
2386                LayerKind::Open { .. } => spec.kind.common().preempt,
2387                _ => false,
2388            })
2389            .map(|(idx, _)| idx as u32)
2390            .collect();
2391        let on_layers: Vec<u32> = layer_specs
2392            .iter()
2393            .enumerate()
2394            .filter(|(_idx, spec)| match &spec.kind {
2395                LayerKind::Open { .. } => !spec.kind.common().preempt,
2396                _ => false,
2397            })
2398            .map(|(idx, _)| idx as u32)
2399            .collect();
2400        let gp_layers: Vec<u32> = layer_specs
2401            .iter()
2402            .enumerate()
2403            .filter(|(_idx, spec)| match &spec.kind {
2404                LayerKind::Grouped { .. } => spec.kind.common().preempt,
2405                _ => false,
2406            })
2407            .map(|(idx, _)| idx as u32)
2408            .collect();
2409        let gn_layers: Vec<u32> = layer_specs
2410            .iter()
2411            .enumerate()
2412            .filter(|(_idx, spec)| match &spec.kind {
2413                LayerKind::Grouped { .. } => !spec.kind.common().preempt,
2414                _ => false,
2415            })
2416            .map(|(idx, _)| idx as u32)
2417            .collect();
2418
2419        // FIXME - this incorrectly assumes all possible CPUs are consecutive.
2420        for cpu in 0..*NR_CPUS_POSSIBLE {
2421            cpu_ctxs.push(
2422                *plain::from_bytes(cpu_ctxs_vec[cpu].as_slice())
2423                    .expect("cpu_ctx: short or misaligned buffer"),
2424            );
2425
2426            let topo_cpu = topo.all_cpus.get(&cpu).unwrap();
2427            let is_big = topo_cpu.core_type == CoreType::Big { turbo: true };
2428            cpu_ctxs[cpu].cpu = cpu as i32;
2429            cpu_ctxs[cpu].layer_id = MAX_LAYERS as u32;
2430            cpu_ctxs[cpu].is_big = is_big;
2431
2432            fastrand::seed(cpu as u64);
2433
2434            let mut ogp_order = op_layers.clone();
2435            ogp_order.append(&mut gp_layers.clone());
2436            fastrand::shuffle(&mut ogp_order);
2437
2438            let mut ogn_order = on_layers.clone();
2439            ogn_order.append(&mut gn_layers.clone());
2440            fastrand::shuffle(&mut ogn_order);
2441
2442            let mut op_order = op_layers.clone();
2443            fastrand::shuffle(&mut op_order);
2444
2445            let mut on_order = on_layers.clone();
2446            fastrand::shuffle(&mut on_order);
2447
2448            let mut gp_order = gp_layers.clone();
2449            fastrand::shuffle(&mut gp_order);
2450
2451            let mut gn_order = gn_layers.clone();
2452            fastrand::shuffle(&mut gn_order);
2453
2454            for i in 0..MAX_LAYERS {
2455                cpu_ctxs[cpu].ogp_layer_order[i] =
2456                    ogp_order.get(i).cloned().unwrap_or(MAX_LAYERS as u32);
2457                cpu_ctxs[cpu].ogn_layer_order[i] =
2458                    ogn_order.get(i).cloned().unwrap_or(MAX_LAYERS as u32);
2459
2460                cpu_ctxs[cpu].op_layer_order[i] =
2461                    op_order.get(i).cloned().unwrap_or(MAX_LAYERS as u32);
2462                cpu_ctxs[cpu].on_layer_order[i] =
2463                    on_order.get(i).cloned().unwrap_or(MAX_LAYERS as u32);
2464                cpu_ctxs[cpu].gp_layer_order[i] =
2465                    gp_order.get(i).cloned().unwrap_or(MAX_LAYERS as u32);
2466                cpu_ctxs[cpu].gn_layer_order[i] =
2467                    gn_order.get(i).cloned().unwrap_or(MAX_LAYERS as u32);
2468            }
2469        }
2470
2471        Self::init_cpu_prox_map(topo, &mut cpu_ctxs);
2472
2473        skel.maps
2474            .cpu_ctxs
2475            .update_percpu(
2476                &key,
2477                &Self::convert_cpu_ctxs(cpu_ctxs),
2478                libbpf_rs::MapFlags::ANY,
2479            )
2480            .context("Failed to update cpu_ctx")?;
2481
2482        Ok(())
2483    }
2484
2485    fn init_single_prox_map_per_llc(
2486        skel: &mut BpfSkel,
2487        topo: &Topology,
2488        prox_map_idx: &usize,
2489    ) -> Result<()> {
2490        for (&llc_id, llc) in &topo.all_llcs {
2491            // Collect the orders.
2492            let mut node_order: Vec<usize> =
2493                topo.nodes[&llc.node_id].llcs.keys().cloned().collect();
2494            let mut sys_order: Vec<usize> = topo.all_llcs.keys().cloned().collect();
2495
2496            // Make the orders exclusive.
2497            sys_order.retain(|id| !node_order.contains(id));
2498            node_order.retain(|&id| id != llc_id);
2499
2500            // Shuffle so that different LLCs follow different orders. See
2501            // init_cpu_prox_map().
2502            fastrand::seed((*prox_map_idx as u64) << 32 | llc_id as u64);
2503            fastrand::shuffle(&mut sys_order);
2504            fastrand::shuffle(&mut node_order);
2505
2506            // Concatenate and record the node boundary.
2507            let mut order: Vec<usize> = vec![];
2508            let mut idx: usize = 0;
2509
2510            idx += 1;
2511            order.push(llc_id);
2512
2513            idx += node_order.len();
2514            order.append(&mut node_order);
2515            let node_end = idx;
2516
2517            idx += sys_order.len();
2518            order.append(&mut sys_order);
2519            let sys_end = idx;
2520
2521            debug!(
2522                "LLC[{}] proximity map {}[{}/{}]: {:?}",
2523                llc_id, prox_map_idx, node_end, sys_end, &order
2524            );
2525
2526            // Record in llc_ctx.
2527            //
2528            // XXX - This would be a lot easier if llc_ctx were in the bss.
2529            // See BpfStats::read().
2530            let key = (llc_id as u32).to_ne_bytes();
2531            let v = skel
2532                .maps
2533                .llc_data
2534                .lookup(&key, libbpf_rs::MapFlags::ANY)
2535                .unwrap()
2536                .unwrap();
2537            let mut llcc: bpf_intf::llc_ctx =
2538                *plain::from_bytes(v.as_slice()).expect("llc_ctx: short or misaligned buffer");
2539
2540            let pmap = &mut llcc.prox_maps[*prox_map_idx];
2541            for (i, &llc_id) in order.iter().enumerate() {
2542                pmap.llcs[i] = llc_id as u16;
2543            }
2544            pmap.node_end = node_end as u32;
2545            pmap.sys_end = sys_end as u32;
2546
2547            skel.maps.llc_data.update(
2548                &key,
2549                unsafe { plain::as_bytes(&llcc) },
2550                libbpf_rs::MapFlags::ANY,
2551            )?
2552        }
2553
2554        Ok(())
2555    }
2556
2557    fn init_llc_prox_map(skel: &mut BpfSkel, topo: &Topology) -> Result<()> {
2558        let num_proximity_maps = bpf_intf::consts_NUM_PROXIMITY_MAPS as usize;
2559        for prox_map_idx in 0..num_proximity_maps {
2560            Self::init_single_prox_map_per_llc(skel, topo, &prox_map_idx)?;
2561        }
2562
2563        Ok(())
2564    }
2565
2566    fn init_node_prox_map(skel: &mut BpfSkel, topo: &Topology) -> Result<()> {
2567        for (&node_id, node) in &topo.nodes {
2568            let mut order: Vec<(usize, usize)> = node
2569                .distance
2570                .iter()
2571                .enumerate()
2572                .filter(|&(nid, _)| nid != node_id)
2573                .map(|(nid, &dist)| (nid, dist))
2574                .collect();
2575            order.sort_by_key(|&(_, dist)| dist);
2576
2577            let key = (node_id as u32).to_ne_bytes();
2578
2579            // The map entry may not exist yet — create a zeroed one.
2580            let v = skel.maps.node_data.lookup(&key, libbpf_rs::MapFlags::ANY);
2581            let mut nodec: bpf_intf::node_ctx = match v {
2582                Ok(Some(v)) => {
2583                    *plain::from_bytes(v.as_slice()).expect("node_ctx: short or misaligned buffer")
2584                }
2585                _ => unsafe { MaybeUninit::zeroed().assume_init() },
2586            };
2587
2588            let pmap = &mut nodec.prox_map;
2589            for (i, &(nid, _)) in order.iter().enumerate() {
2590                pmap.nodes[i] = nid as u16;
2591            }
2592            pmap.sys_end = order.len() as u32;
2593
2594            let order_vec: Vec<_> = order.iter().map(|(n, d)| (*n, *d)).collect();
2595            debug!(
2596                "NODE[{}] prox_map[{}]: {:?}",
2597                node_id, pmap.sys_end, order_vec
2598            );
2599
2600            skel.maps.node_data.update(
2601                &key,
2602                unsafe { plain::as_bytes(&nodec) },
2603                libbpf_rs::MapFlags::ANY,
2604            )?;
2605        }
2606        Ok(())
2607    }
2608
2609    fn init_node_ctx(skel: &mut BpfSkel, topo: &Topology, nr_layers: usize) -> Result<()> {
2610        let all_layers: Vec<u32> = (0..nr_layers as u32).collect();
2611        let node_empty_layers: Vec<Vec<u32>> =
2612            (0..topo.nodes.len()).map(|_| all_layers.clone()).collect();
2613        Self::refresh_node_ctx(skel, topo, &node_empty_layers, true);
2614        Ok(())
2615    }
2616
2617    fn init(
2618        opts: &'a Opts,
2619        layer_specs: &[LayerSpec],
2620        open_object: &'a mut MaybeUninit<OpenObject>,
2621        hint_to_layer_map: &HashMap<u64, HintLayerInfo>,
2622        membw_tracking: bool,
2623    ) -> Result<Self> {
2624        let nr_layers = layer_specs.len();
2625        let mut disable_topology = opts.disable_topology.unwrap_or(false);
2626
2627        let topo = Arc::new(if disable_topology {
2628            Topology::with_flattened_llc_node()?
2629        } else if opts.topology.virt_llc.is_some() {
2630            Topology::with_args(&opts.topology)?
2631        } else {
2632            Topology::new()?
2633        });
2634
2635        /*
2636         * FIXME: scx_layered incorrectly assumes that node, LLC and CPU IDs
2637         * are consecutive. Verify that they are on this system and bail if
2638         * not. It's lucky that core ID is not used anywhere as core IDs are
2639         * not consecutive on some Ryzen CPUs.
2640         */
2641        if topo.nodes.keys().enumerate().any(|(i, &k)| i != k) {
2642            bail!("Holes in node IDs detected: {:?}", topo.nodes.keys());
2643        }
2644        if topo.all_llcs.keys().enumerate().any(|(i, &k)| i != k) {
2645            bail!("Holes in LLC IDs detected: {:?}", topo.all_llcs.keys());
2646        }
2647        if topo.all_cpus.keys().enumerate().any(|(i, &k)| i != k) {
2648            bail!("Holes in CPU IDs detected: {:?}", topo.all_cpus.keys());
2649        }
2650
2651        let netdevs = if opts.netdev_irq_balance {
2652            let devs = read_netdevs()?;
2653            let total_irqs: usize = devs.values().map(|d| d.irqs.len()).sum();
2654            let breakdown = devs
2655                .iter()
2656                .map(|(iface, d)| format!("{iface}={}", d.irqs.len()))
2657                .collect::<Vec<_>>()
2658                .join(", ");
2659            info!(
2660                "Netdev IRQ balancing enabled: overriding {total_irqs} IRQ{} \
2661                 across {} interface{} [{breakdown}]",
2662                if total_irqs == 1 { "" } else { "s" },
2663                devs.len(),
2664                if devs.len() == 1 { "" } else { "s" },
2665            );
2666            devs
2667        } else {
2668            BTreeMap::new()
2669        };
2670
2671        if !disable_topology {
2672            if topo.nodes.len() == 1 && topo.nodes[&0].llcs.len() == 1 {
2673                disable_topology = true;
2674            };
2675            info!(
2676                "Topology awareness not specified, selecting {} based on hardware",
2677                if disable_topology {
2678                    "disabled"
2679                } else {
2680                    "enabled"
2681                }
2682            );
2683        };
2684
2685        let cpu_pool = CpuPool::new(topo.clone(), opts.allow_partial_core)?;
2686
2687        // If disabling topology awareness clear out any set NUMA/LLC configs and
2688        // it will fallback to using all cores.
2689        let layer_specs: Vec<_> = if disable_topology {
2690            info!("Disabling topology awareness");
2691            layer_specs
2692                .iter()
2693                .cloned()
2694                .map(|mut s| {
2695                    s.kind.common_mut().nodes.clear();
2696                    s.kind.common_mut().llcs.clear();
2697                    s
2698                })
2699                .collect()
2700        } else {
2701            layer_specs.to_vec()
2702        };
2703
2704        // Validate that spec node/LLC references exist in the topology.
2705        for spec in layer_specs.iter() {
2706            let mut seen = BTreeSet::new();
2707            for &node_id in spec.nodes().iter() {
2708                if !topo.nodes.contains_key(&node_id) {
2709                    bail!(
2710                        "layer {:?}: nodes references node {} which does not \
2711                         exist in the topology (available: {:?})",
2712                        spec.name,
2713                        node_id,
2714                        topo.nodes.keys().collect::<Vec<_>>()
2715                    );
2716                }
2717                if !seen.insert(node_id) {
2718                    bail!(
2719                        "layer {:?}: nodes contains duplicate node {}",
2720                        spec.name,
2721                        node_id
2722                    );
2723                }
2724            }
2725
2726            seen.clear();
2727            for &llc_id in spec.llcs().iter() {
2728                if !topo.all_llcs.contains_key(&llc_id) {
2729                    bail!(
2730                        "layer {:?}: llcs references LLC {} which does not \
2731                         exist in the topology (available: {:?})",
2732                        spec.name,
2733                        llc_id,
2734                        topo.all_llcs.keys().collect::<Vec<_>>()
2735                    );
2736                }
2737                if !seen.insert(llc_id) {
2738                    bail!(
2739                        "layer {:?}: llcs contains duplicate LLC {}",
2740                        spec.name,
2741                        llc_id
2742                    );
2743                }
2744            }
2745        }
2746
2747        for spec in layer_specs.iter() {
2748            let has_numa_node_match = spec
2749                .matches
2750                .iter()
2751                .flatten()
2752                .any(|m| matches!(m, LayerMatch::NumaNode(_)));
2753            let has_node_spread_algo = matches!(
2754                spec.kind.common().growth_algo,
2755                LayerGrowthAlgo::NodeSpread
2756                    | LayerGrowthAlgo::NodeSpreadReverse
2757                    | LayerGrowthAlgo::NodeSpreadRandom
2758            );
2759            if has_numa_node_match && has_node_spread_algo {
2760                bail!(
2761                    "layer {:?}: NumaNode matcher cannot be combined with {:?} \
2762                     growth algorithm. NodeSpread* allocates CPUs equally across \
2763                     ALL NUMA nodes, but NumaNode restricts tasks to one node's \
2764                     CPUs — CPUs on other nodes are wasted and utilization \
2765                     will never exceed 1/numa_nodes. Use a non-spread algorithm \
2766                     (e.g. Linear, Topo) instead.",
2767                    spec.name,
2768                    spec.kind.common().growth_algo
2769                );
2770            }
2771        }
2772
2773        // Check kernel features
2774        init_libbpf_logging(None);
2775        let kfuncs_in_syscall = scx_bpf_compat::kfuncs_supported_in_syscall()?;
2776        if !kfuncs_in_syscall {
2777            warn!("Using slow path: kfuncs not supported in syscall programs (a8e03b6bbb2c ∉ ker)");
2778        }
2779
2780        // Open the BPF prog first for verification.
2781        let debug_level = if opts.log_level.contains("trace") {
2782            2
2783        } else if opts.log_level.contains("debug") {
2784            1
2785        } else {
2786            0
2787        };
2788        let mut skel_builder = BpfSkelBuilder::default();
2789        skel_builder.obj_builder.debug(debug_level > 1);
2790
2791        info!(
2792            "Running scx_layered (build ID: {})",
2793            build_id::full_version(env!("CARGO_PKG_VERSION"))
2794        );
2795        let open_opts = opts.libbpf.clone().into_bpf_open_opts();
2796        let mut skel = scx_ops_open!(skel_builder, open_object, layered, open_opts)?;
2797
2798        // No memory BW tracking by default
2799        skel.progs.scx_pmu_switch_tc.set_autoload(membw_tracking);
2800        skel.progs.scx_pmu_tick_tc.set_autoload(membw_tracking);
2801
2802        let mut loaded_kprobes = HashSet::new();
2803
2804        // enable autoloads for conditionally loaded things
2805        // immediately after creating skel (because this is always before loading)
2806        if opts.enable_gpu_support {
2807            // by default, enable open if gpu support is enabled.
2808            // open has been observed to be relatively cheap to kprobe.
2809            if opts.gpu_kprobe_level >= 1 {
2810                compat::cond_kprobe_load("nvidia_open", &skel.progs.kprobe_nvidia_open)?;
2811                loaded_kprobes.insert("nvidia_open");
2812            }
2813            // enable the rest progressively based upon how often they are called
2814            // for observed workloads
2815            if opts.gpu_kprobe_level >= 2 {
2816                compat::cond_kprobe_load("nvidia_mmap", &skel.progs.kprobe_nvidia_mmap)?;
2817                loaded_kprobes.insert("nvidia_mmap");
2818            }
2819            if opts.gpu_kprobe_level >= 3 {
2820                compat::cond_kprobe_load("nvidia_poll", &skel.progs.kprobe_nvidia_poll)?;
2821                loaded_kprobes.insert("nvidia_poll");
2822            }
2823        }
2824
2825        let ext_sched_class_addr = get_kallsyms_addr("ext_sched_class");
2826        let idle_sched_class_addr = get_kallsyms_addr("idle_sched_class");
2827
2828        let event = if membw_tracking {
2829            setup_membw_tracking(&mut skel)?
2830        } else {
2831            0
2832        };
2833
2834        let rodata = skel.maps.rodata_data.as_mut().unwrap();
2835
2836        if let (Ok(ext_addr), Ok(idle_addr)) = (ext_sched_class_addr, idle_sched_class_addr) {
2837            rodata.ext_sched_class_addr = ext_addr;
2838            rodata.idle_sched_class_addr = idle_addr;
2839        } else {
2840            warn!(
2841                "Unable to get sched_class addresses from /proc/kallsyms, disabling skip_preempt."
2842            );
2843        }
2844
2845        rodata.slice_ns = scx_enums.SCX_SLICE_DFL;
2846        rodata.max_exec_ns = 20 * scx_enums.SCX_SLICE_DFL;
2847
2848        // Initialize skel according to @opts.
2849        skel.struct_ops.layered_mut().exit_dump_len = opts.exit_dump_len;
2850
2851        if !opts.disable_queued_wakeup {
2852            match *compat::SCX_OPS_ALLOW_QUEUED_WAKEUP {
2853                0 => info!("Kernel does not support queued wakeup optimization"),
2854                v => skel.struct_ops.layered_mut().flags |= v,
2855            }
2856        }
2857
2858        rodata.percpu_kthread_preempt = !opts.disable_percpu_kthread_preempt;
2859        rodata.percpu_kthread_preempt_all =
2860            !opts.disable_percpu_kthread_preempt && opts.percpu_kthread_preempt_all;
2861        rodata.debug = debug_level as u32;
2862        rodata.slice_ns = opts.slice_us * 1000;
2863        rodata.max_exec_ns = if opts.max_exec_us > 0 {
2864            opts.max_exec_us * 1000
2865        } else {
2866            opts.slice_us * 1000 * 20
2867        };
2868        rodata.nr_cpu_ids = *NR_CPU_IDS as u32;
2869        rodata.nr_possible_cpus = *NR_CPUS_POSSIBLE as u32;
2870        rodata.smt_enabled = topo.smt_enabled;
2871        rodata.has_little_cores = topo.has_little_cores();
2872        rodata.antistall_sec = opts.antistall_sec;
2873        rodata.monitor_disable = opts.monitor_disable;
2874        rodata.lo_fb_wait_ns = opts.lo_fb_wait_us * 1000;
2875        rodata.lo_fb_share_ppk = ((opts.lo_fb_share * 1024.0) as u32).clamp(1, 1024);
2876        rodata.enable_antistall = !opts.disable_antistall;
2877        rodata.enable_match_debug = opts.enable_match_debug;
2878        rodata.enable_gpu_support = opts.enable_gpu_support;
2879        rodata.kfuncs_supported_in_syscall = kfuncs_in_syscall;
2880
2881        for (cpu, sib) in topo.sibling_cpus().iter().enumerate() {
2882            rodata.__sibling_cpu[cpu] = *sib;
2883        }
2884        for cpu in topo.all_cpus.keys() {
2885            rodata.all_cpus[cpu / 8] |= 1 << (cpu % 8);
2886        }
2887
2888        rodata.nr_op_layers = layer_specs
2889            .iter()
2890            .filter(|spec| match &spec.kind {
2891                LayerKind::Open { .. } => spec.kind.common().preempt,
2892                _ => false,
2893            })
2894            .count() as u32;
2895        rodata.nr_on_layers = layer_specs
2896            .iter()
2897            .filter(|spec| match &spec.kind {
2898                LayerKind::Open { .. } => !spec.kind.common().preempt,
2899                _ => false,
2900            })
2901            .count() as u32;
2902        rodata.nr_gp_layers = layer_specs
2903            .iter()
2904            .filter(|spec| match &spec.kind {
2905                LayerKind::Grouped { .. } => spec.kind.common().preempt,
2906                _ => false,
2907            })
2908            .count() as u32;
2909        rodata.nr_gn_layers = layer_specs
2910            .iter()
2911            .filter(|spec| match &spec.kind {
2912                LayerKind::Grouped { .. } => !spec.kind.common().preempt,
2913                _ => false,
2914            })
2915            .count() as u32;
2916        rodata.nr_excl_layers = layer_specs
2917            .iter()
2918            .filter(|spec| spec.kind.common().exclusive)
2919            .count() as u32;
2920
2921        let mut min_open = u64::MAX;
2922        let mut min_preempt = u64::MAX;
2923
2924        for spec in layer_specs.iter() {
2925            if let LayerKind::Open { common, .. } = &spec.kind {
2926                min_open = min_open.min(common.disallow_open_after_us.unwrap());
2927                min_preempt = min_preempt.min(common.disallow_preempt_after_us.unwrap());
2928            }
2929        }
2930
2931        rodata.min_open_layer_disallow_open_after_ns = match min_open {
2932            u64::MAX => *DFL_DISALLOW_OPEN_AFTER_US,
2933            v => v,
2934        };
2935        rodata.min_open_layer_disallow_preempt_after_ns = match min_preempt {
2936            u64::MAX => *DFL_DISALLOW_PREEMPT_AFTER_US,
2937            v => v,
2938        };
2939
2940        // We set the pin path before loading the skeleton. This will ensure
2941        // libbpf creates and pins the map, or reuses the pinned map fd for us,
2942        // so that we can keep reusing the older map already pinned on scheduler
2943        // restarts.
2944        let layered_task_hint_map_path = &opts.task_hint_map;
2945        let hint_map = &mut skel.maps.scx_layered_task_hint_map;
2946        // Only set pin path if a path is provided.
2947        if !layered_task_hint_map_path.is_empty() {
2948            hint_map.set_pin_path(layered_task_hint_map_path).unwrap();
2949            rodata.task_hint_map_enabled = true;
2950        }
2951
2952        if !opts.hi_fb_thread_name.is_empty() {
2953            let bpf_hi_fb_thread_name = &mut rodata.hi_fb_thread_name;
2954            copy_into_cstr(bpf_hi_fb_thread_name, opts.hi_fb_thread_name.as_str());
2955            rodata.enable_hi_fb_thread_name_match = true;
2956        }
2957
2958        let cgroup_regexes = Self::init_layers(&mut skel, &layer_specs, &topo)?;
2959        skel.maps.rodata_data.as_mut().unwrap().nr_cgroup_regexes = cgroup_regexes.len() as u32;
2960        Self::init_nodes(&mut skel, opts, &topo);
2961
2962        let mut skel = scx_ops_load!(skel, layered, uei)?;
2963
2964        // Populate the mapping of hints to layer IDs for faster lookups
2965        if !hint_to_layer_map.is_empty() {
2966            for (k, v) in hint_to_layer_map.iter() {
2967                let key: u32 = *k as u32;
2968
2969                // Create hint_layer_info struct
2970                let mut info_bytes = vec![0u8; std::mem::size_of::<bpf_intf::hint_layer_info>()];
2971                let info_ptr = info_bytes.as_mut_ptr() as *mut bpf_intf::hint_layer_info;
2972                unsafe {
2973                    (*info_ptr).layer_id = v.layer_id as u32;
2974                    (*info_ptr).system_cpu_util_below = match v.system_cpu_util_below {
2975                        Some(threshold) => (threshold * 10000.0) as u64,
2976                        None => u64::MAX, // disabled sentinel
2977                    };
2978                    (*info_ptr).dsq_insert_below = match v.dsq_insert_below {
2979                        Some(threshold) => (threshold * 10000.0) as u64,
2980                        None => u64::MAX, // disabled sentinel
2981                    };
2982                }
2983
2984                skel.maps.hint_to_layer_id_map.update(
2985                    &key.to_ne_bytes(),
2986                    &info_bytes,
2987                    libbpf_rs::MapFlags::ANY,
2988                )?;
2989            }
2990        }
2991
2992        if membw_tracking {
2993            create_perf_fds(&mut skel, event)?;
2994        }
2995
2996        let mut layers = vec![];
2997        let layer_growth_orders =
2998            LayerGrowthAlgo::layer_core_orders(&cpu_pool, &layer_specs, &topo)?;
2999        for (idx, spec) in layer_specs.iter().enumerate() {
3000            let growth_order = layer_growth_orders
3001                .get(&idx)
3002                .with_context(|| "layer has no growth order".to_string())?;
3003            layers.push(Layer::new(spec, &topo, growth_order)?);
3004        }
3005
3006        let mut idle_qos_enabled = layers
3007            .iter()
3008            .any(|layer| layer.kind.common().idle_resume_us.unwrap_or(0) > 0);
3009        if idle_qos_enabled && !cpu_idle_resume_latency_supported() {
3010            warn!("idle_resume_us not supported, ignoring");
3011            idle_qos_enabled = false;
3012        }
3013
3014        Self::init_cpus(&skel, &layer_specs, &topo)?;
3015        Self::init_llc_prox_map(&mut skel, &topo)?;
3016        Self::init_node_prox_map(&mut skel, &topo)?;
3017        Self::init_node_ctx(&mut skel, &topo, nr_layers)?;
3018
3019        // Other stuff.
3020        let proc_reader = fb_procfs::ProcReader::new();
3021
3022        // Handle setup if layered is running in a pid namespace.
3023        let input = ProgramInput {
3024            ..Default::default()
3025        };
3026        let prog = &mut skel.progs.initialize_pid_namespace;
3027
3028        let _ = prog.test_run(input);
3029
3030        // XXX If we try to refresh the cpumasks here before attaching, we
3031        // sometimes (non-deterministically) don't see the updated values in
3032        // BPF. It would be better to update the cpumasks here before we
3033        // attach, but the value will quickly converge anyways so it's not a
3034        // huge problem in the interim until we figure it out.
3035
3036        // Allow all tasks to open and write to BPF task hint map, now that
3037        // we should have it pinned at the desired location.
3038        if !layered_task_hint_map_path.is_empty() {
3039            let path = CString::new(layered_task_hint_map_path.as_bytes()).unwrap();
3040            let mode: libc::mode_t = 0o666;
3041            unsafe {
3042                if libc::chmod(path.as_ptr(), mode) != 0 {
3043                    trace!("'chmod' to 666 of task hint map failed, continuing...");
3044                }
3045            }
3046        }
3047
3048        // Attach.
3049        let struct_ops = scx_ops_attach!(skel, layered)?;
3050
3051        // Turn on installed kprobes
3052        if opts.enable_gpu_support {
3053            if loaded_kprobes.contains("nvidia_open") {
3054                compat::cond_kprobe_attach("nvidia_open", &skel.progs.kprobe_nvidia_open)?;
3055            }
3056            if loaded_kprobes.contains("nvidia_mmap") {
3057                compat::cond_kprobe_attach("nvidia_mmap", &skel.progs.kprobe_nvidia_mmap)?;
3058            }
3059            if loaded_kprobes.contains("nvidia_poll") {
3060                compat::cond_kprobe_attach("nvidia_poll", &skel.progs.kprobe_nvidia_poll)?;
3061            }
3062        }
3063
3064        let stats_server = StatsServer::new(stats::server_data()).launch()?;
3065        let mut gpu_task_handler =
3066            GpuTaskAffinitizer::new(opts.gpu_affinitize_secs, opts.enable_gpu_affinitize);
3067        gpu_task_handler.init(topo.clone());
3068
3069        let sched = Self {
3070            struct_ops: Some(struct_ops),
3071            layer_specs,
3072
3073            sched_intv: Duration::from_secs_f64(opts.interval),
3074            layer_refresh_intv: Duration::from_millis(opts.layer_refresh_ms_avgruntime),
3075
3076            cpu_pool,
3077            layers,
3078            idle_qos_enabled,
3079
3080            sched_stats: Stats::new(
3081                &mut skel,
3082                &proc_reader,
3083                topo.clone(),
3084                &gpu_task_handler,
3085                opts.util_compensation,
3086            )?,
3087            layer_peak_utils: vec![0.0; nr_layers],
3088
3089            cgroup_regexes: Some(cgroup_regexes),
3090            nr_layer_cpus_ranges: vec![(0, 0); nr_layers],
3091            xnuma_mig_src: vec![vec![false; topo.nodes.len()]; nr_layers],
3092            growth_denied: vec![vec![false; topo.nodes.len()]; nr_layers],
3093            processing_dur: Default::default(),
3094
3095            proc_reader,
3096            skel,
3097
3098            topo,
3099            netdevs,
3100            stats_server,
3101            gpu_task_handler,
3102        };
3103
3104        info!("Layered Scheduler Attached. Run `scx_layered --monitor` for metrics.");
3105
3106        Ok(sched)
3107    }
3108
3109    fn update_cpumask(mask: &Cpumask, bpfmask: &mut [u8]) {
3110        for cpu in 0..mask.len() {
3111            if mask.test_cpu(cpu) {
3112                bpfmask[cpu / 8] |= 1 << (cpu % 8);
3113            } else {
3114                bpfmask[cpu / 8] &= !(1 << (cpu % 8));
3115            }
3116        }
3117    }
3118
3119    fn update_bpf_layer_cpumask(layer: &Layer, bpf_layer: &mut types::layer) {
3120        trace!("[{}] Updating BPF CPUs: {}", layer.name, &layer.cpus);
3121        Self::update_cpumask(&layer.cpus, &mut bpf_layer.cpus);
3122
3123        bpf_layer.nr_cpus = layer.nr_cpus as u32;
3124        for (llc_id, &nr_llc_cpus) in layer.nr_llc_cpus.iter().enumerate() {
3125            bpf_layer.nr_llc_cpus[llc_id] = nr_llc_cpus as u32;
3126        }
3127        for (node_id, &nr_node_cpus) in layer.nr_node_cpus.iter().enumerate() {
3128            bpf_layer.node[node_id].nr_cpus = nr_node_cpus as u32;
3129        }
3130
3131        bpf_layer.refresh_cpus = 1;
3132    }
3133
3134    fn update_netdev_cpumasks(&mut self) -> Result<()> {
3135        let available_cpus = self.cpu_pool.available_cpus();
3136        if available_cpus.is_empty() {
3137            return Ok(());
3138        }
3139
3140        for (iface, netdev) in self.netdevs.iter_mut() {
3141            let node = self
3142                .topo
3143                .nodes
3144                .values()
3145                .find(|n| n.id == netdev.node())
3146                .ok_or_else(|| anyhow!("Failed to get netdev node"))?;
3147            let node_cpus = node.span.clone();
3148            for (irq, irqmask) in netdev.irqs.iter_mut() {
3149                irqmask.clear_all();
3150                for cpu in available_cpus.iter() {
3151                    if !node_cpus.test_cpu(cpu) {
3152                        continue;
3153                    }
3154                    let _ = irqmask.set_cpu(cpu);
3155                }
3156                // If no CPUs are available in the node then spread the load across the node
3157                if irqmask.weight() == 0 {
3158                    for cpu in node_cpus.iter() {
3159                        let _ = irqmask.set_cpu(cpu);
3160                    }
3161                }
3162                trace!("{} updating irq {} cpumask {:?}", iface, irq, irqmask);
3163            }
3164            netdev.apply_cpumasks()?;
3165            debug!(
3166                "{iface}: applied affinity override to {} IRQ{}",
3167                netdev.irqs.len(),
3168                if netdev.irqs.len() == 1 { "" } else { "s" },
3169            );
3170        }
3171
3172        Ok(())
3173    }
3174
3175    fn clamp_target_by_membw(
3176        &self,
3177        layer: &Layer,
3178        membw_limit: f64,
3179        membw: f64,
3180        curtarget: u64,
3181    ) -> usize {
3182        let ncpu: u64 = layer.cpus.weight() as u64;
3183        let membw = (membw * 1024_f64.powf(3.0)).round() as u64;
3184        let membw_limit = (membw_limit * 1024_f64.powf(3.0)).round() as u64;
3185        let last_membw_percpu = membw.checked_div(ncpu).unwrap_or(0);
3186
3187        // Either there is no memory bandwidth limit set, or the counters
3188        // are not fully initialized yet. Just return the current target.
3189        if membw_limit == 0 || last_membw_percpu == 0 {
3190            return curtarget as usize;
3191        }
3192
3193        (membw_limit / last_membw_percpu) as usize
3194    }
3195
3196    /// Decompose per-layer CPU targets into per-node pinned demand and
3197    /// unpinned demand for unified_alloc(). Uses layer_node_pinned_utils
3198    /// to split each layer's target. All outputs are in alloc units.
3199    fn calc_raw_demands(&self, targets: &[(usize, usize)]) -> Vec<LayerDemand> {
3200        let au = self.cpu_pool.alloc_unit();
3201        let pinned_utils = &self.sched_stats.layer_node_pinned_utils;
3202        let nr_nodes = self.topo.nodes.len();
3203
3204        targets
3205            .iter()
3206            .enumerate()
3207            .map(|(idx, &(target, _min))| {
3208                let layer = &self.layers[idx];
3209                let weight = layer.kind.common().weight as usize;
3210
3211                // Open layers don't participate in allocation.
3212                if matches!(layer.kind, LayerKind::Open { .. }) {
3213                    return LayerDemand {
3214                        raw_pinned: vec![0; nr_nodes],
3215                        raw_unpinned: 0,
3216                        weight,
3217                        spread: false,
3218                    };
3219                }
3220
3221                // Only NodeSpread* keeps strict equal-per-node placement
3222                // via the spread:bool path.  RoundRobin migrated to a
3223                // single-tier node_groups (balanced via cap-weighted
3224                // water_fill in place_unpinned, no bottleneck cap).
3225                let spread = matches!(
3226                    layer.growth_algo,
3227                    LayerGrowthAlgo::NodeSpread
3228                        | LayerGrowthAlgo::NodeSpreadReverse
3229                        | LayerGrowthAlgo::NodeSpreadRandom
3230                );
3231
3232                let util_high = match &layer.kind {
3233                    LayerKind::Confined { util_range, .. }
3234                    | LayerKind::Grouped { util_range, .. } => util_range.1,
3235                    _ => 1.0,
3236                };
3237
3238                // Convert per-node pinned utilization to CPU demand.
3239                let mut raw_pinned = vec![0usize; nr_nodes];
3240                for n in 0..nr_nodes {
3241                    let pu = pinned_utils[idx][n];
3242                    if pu < 0.01 {
3243                        continue;
3244                    }
3245                    // Check this layer has allowed_cpus on this node.
3246                    let node_span = &self.topo.nodes[&n].span;
3247                    if layer.allowed_cpus.and(node_span).is_empty() {
3248                        continue;
3249                    }
3250                    let cpus = (pu / util_high).ceil() as usize;
3251                    // Round up to alloc units.
3252                    let units = cpus.div_ceil(au);
3253                    raw_pinned[n] = units;
3254                }
3255
3256                // Unpinned = remainder of the target.
3257                let target_units = target.div_ceil(au);
3258                let pinned_units: usize = raw_pinned.iter().sum();
3259                let raw_unpinned = target_units.saturating_sub(pinned_units);
3260
3261                LayerDemand {
3262                    raw_pinned,
3263                    raw_unpinned,
3264                    weight,
3265                    spread,
3266                }
3267            })
3268            .collect()
3269    }
3270
3271    /// Calculate how many CPUs each layer would like to have if there were
3272    /// no competition. When util_compensation is enabled, compensated
3273    /// utilization (scaled for irq/softirq/stolen overhead) is used
3274    /// instead of raw utilization. The CPU range is determined by
3275    /// applying the inverse of util_range and capping by cpus_range.
3276    /// If the current allocation is within the acceptable range, no
3277    /// change is made. Returns (target, min) pair for each layer.
3278    fn calc_target_nr_cpus(&mut self) -> Vec<(usize, usize)> {
3279        let nr_cpus = self.cpu_pool.topo.all_cpus.len();
3280        let membws = &self.sched_stats.layer_membws;
3281
3282        let mut records: Vec<(u64, u64, u64, u64, usize, usize, usize)> = vec![];
3283        let mut targets: Vec<(usize, usize)> = vec![];
3284
3285        for (idx, layer) in self.layers.iter().enumerate() {
3286            targets.push(match &layer.kind {
3287                LayerKind::Confined {
3288                    util_range,
3289                    cpus_range,
3290                    cpus_range_frac,
3291                    membw_gb,
3292                    ..
3293                }
3294                | LayerKind::Grouped {
3295                    util_range,
3296                    cpus_range,
3297                    cpus_range_frac,
3298                    membw_gb,
3299                    ..
3300                } => {
3301                    let cpus_range =
3302                        resolve_cpus_pct_range(cpus_range, cpus_range_frac, nr_cpus).unwrap();
3303
3304                    // A grouped layer can choose to include open cputime
3305                    // for sizing. Also, as an empty layer can only get CPU
3306                    // time through fallback (counted as owned) or open
3307                    // execution, add open cputime for empty layers.
3308                    let (owned, open) = {
3309                        let utils = if self.sched_stats.util_compensation {
3310                            &self.sched_stats.layer_utils_compensated[idx]
3311                        } else {
3312                            &self.sched_stats.layer_utils[idx]
3313                        };
3314                        (utils[LAYER_USAGE_OWNED], utils[LAYER_USAGE_OPEN])
3315                    };
3316
3317                    let membw_owned = membws[idx][LAYER_USAGE_OWNED];
3318                    let membw_open = membws[idx][LAYER_USAGE_OPEN];
3319
3320                    let mut util = owned;
3321                    let mut membw = membw_owned;
3322                    if layer.kind.util_includes_open_cputime() || layer.nr_cpus == 0 {
3323                        util += open;
3324                        membw += membw_open;
3325                    }
3326
3327                    let util = if util < 0.01 { 0.0 } else { util };
3328
3329                    let peak_util = if layer.kind.common().util_peak_half_life_ms == 0 {
3330                        util
3331                    } else {
3332                        update_peak_util(
3333                            self.layer_peak_utils[idx],
3334                            util,
3335                            Duration::from_millis(layer.kind.common().util_peak_half_life_ms),
3336                            self.sched_stats.elapsed,
3337                        )
3338                    };
3339                    self.layer_peak_utils[idx] = peak_util;
3340                    let (low, high) = calc_peak_aware_target_range(util, peak_util, *util_range);
3341
3342                    let membw_limit = match membw_gb {
3343                        Some(membw_limit) => *membw_limit,
3344                        None => 0.0,
3345                    };
3346
3347                    trace!(
3348                        "layer {0} (membw, membw_limit): ({membw} gi_b, {membw_limit} gi_b)",
3349                        layer.name
3350                    );
3351
3352                    let target = layer.cpus.weight().clamp(low, high);
3353
3354                    records.push((
3355                        (owned * 100.0) as u64,
3356                        (open * 100.0) as u64,
3357                        (util * 100.0) as u64,
3358                        (peak_util * 100.0) as u64,
3359                        low,
3360                        high,
3361                        target,
3362                    ));
3363
3364                    let target = target.clamp(cpus_range.0, cpus_range.1);
3365                    let membw_target =
3366                        self.clamp_target_by_membw(layer, membw_limit, membw, target as u64);
3367
3368                    trace!("CPU target pre- and post-membw adjustment: {target} -> {membw_target}");
3369
3370                    // If there's no way to drop our memory usage down enough,
3371                    // pin the target CPUs to low and drop a warning.
3372                    if membw_target < cpus_range.0 {
3373                        warn!("cannot satisfy memory bw limit for layer {}", layer.name);
3374                        warn!("membw_target {membw_target} low {}", cpus_range.0);
3375                    };
3376
3377                    // Memory bandwidth target cannot override imposed limits or bump
3378                    // the target above what CPU usage-based throttling requires.
3379                    let target = membw_target.clamp(cpus_range.0, target);
3380
3381                    (target, cpus_range.0)
3382                }
3383                LayerKind::Open { .. } => (0, 0),
3384            });
3385        }
3386
3387        trace!(
3388            "(owned, open, util, peak, low, high, target): {:?}",
3389            &records
3390        );
3391        targets
3392    }
3393
3394    // Figure out a tuple (LLCs, extra_cpus) in terms of the target CPUs
3395    // computed by weighted_target_nr_cpus. Returns the number of full LLCs
3396    // occupied by a layer, and any extra cores needed beyond those LLCs.
3397    fn compute_target_llcs(target: usize, node_llcs: &BTreeMap<usize, Arc<Llc>>) -> (usize, usize) {
3398        let mut remaining = target;
3399        let mut full = 0;
3400        for llc in node_llcs.values() {
3401            let cpus_in_llc = llc.span.weight();
3402            if remaining >= cpus_in_llc {
3403                full += 1;
3404                remaining -= cpus_in_llc;
3405            } else {
3406                let mut extra = 0;
3407                for core in llc.cores.values() {
3408                    if remaining == 0 {
3409                        break;
3410                    }
3411                    extra += 1;
3412                    remaining = remaining.saturating_sub(core.cpus.len());
3413                }
3414                return (full, extra);
3415            }
3416        }
3417        (full, 0)
3418    }
3419
3420    // Recalculate the core order for layers using StickyDynamic growth
3421    // algorithm. Uses per-node targets from unified_alloc() to decide how
3422    // many LLCs each layer gets on each node, then builds core_order and
3423    // applies CPU changes.
3424    fn recompute_layer_core_order(
3425        &mut self,
3426        layer_targets: &[(usize, usize)],
3427        layer_allocs: &[LayerAlloc],
3428        au: usize,
3429    ) -> Result<bool> {
3430        let nr_nodes = self.topo.nodes.len();
3431
3432        // Phase 1 — Free per-node: return excess LLCs to cpu_pool.
3433        debug!(
3434            " free: before pass: free_llcs={:?}",
3435            self.cpu_pool.free_llcs
3436        );
3437        for &(idx, _) in layer_targets.iter().rev() {
3438            let layer = &mut self.layers[idx];
3439
3440            if layer.growth_algo != LayerGrowthAlgo::StickyDynamic {
3441                continue;
3442            }
3443
3444            let alloc = &layer_allocs[idx];
3445
3446            for n in 0..nr_nodes {
3447                let assigned_on_n = layer.assigned_llcs[n].len();
3448                let node = self.topo.nodes.get(&n).unwrap();
3449                let target_full_n =
3450                    Self::compute_target_llcs(alloc.node_target(n) * au, &node.llcs).0;
3451                let mut to_free = assigned_on_n.saturating_sub(target_full_n);
3452
3453                debug!(
3454                    " free: layer={} node={} assigned={} target_full={} to_free={}",
3455                    layer.name, n, assigned_on_n, target_full_n, to_free,
3456                );
3457
3458                while to_free > 0 {
3459                    if let Some(llc) = layer.assigned_llcs[n].pop() {
3460                        self.cpu_pool.return_llc(llc);
3461                        to_free -= 1;
3462                        debug!(" layer={} freed_llc={} from node={}", layer.name, llc, n);
3463                    } else {
3464                        break;
3465                    }
3466                }
3467            }
3468        }
3469        debug!(" free: after pass: free_llcs={:?}", self.cpu_pool.free_llcs);
3470
3471        // Phase 2 — Acquire per-node: claim LLCs from cpu_pool.
3472        for &(idx, _) in layer_targets.iter().rev() {
3473            let layer = &mut self.layers[idx];
3474
3475            if layer.growth_algo != LayerGrowthAlgo::StickyDynamic {
3476                continue;
3477            }
3478
3479            let alloc = &layer_allocs[idx];
3480
3481            for n in 0..nr_nodes {
3482                let cur_on_n = layer.assigned_llcs[n].len();
3483                let node = self.topo.nodes.get(&n).unwrap();
3484                let target_full_n =
3485                    Self::compute_target_llcs(alloc.node_target(n) * au, &node.llcs).0;
3486                let mut to_alloc = target_full_n.saturating_sub(cur_on_n);
3487
3488                debug!(
3489                    " alloc: layer={} node={} cur={} target_full={} to_alloc={} free={}",
3490                    layer.name,
3491                    n,
3492                    cur_on_n,
3493                    target_full_n,
3494                    to_alloc,
3495                    self.cpu_pool.free_llcs.get(&n).map_or(0, |v| v.len()),
3496                );
3497
3498                while to_alloc > 0 {
3499                    if let Some(llc) = self.cpu_pool.take_llc_from_node(n) {
3500                        layer.assigned_llcs[n].push(llc);
3501                        to_alloc -= 1;
3502                        debug!(" layer={} alloc_llc={} on node={}", layer.name, llc, n);
3503                    } else {
3504                        break;
3505                    }
3506                }
3507            }
3508
3509            debug!(
3510                " alloc: layer={} assigned_llcs={:?}",
3511                layer.name, layer.assigned_llcs
3512            );
3513        }
3514
3515        // Phase 3 — Spillover per-node: consume extra cores from free LLCs.
3516        for &(idx, _) in layer_targets.iter() {
3517            let layer = &mut self.layers[idx];
3518
3519            if layer.growth_algo != LayerGrowthAlgo::StickyDynamic {
3520                continue;
3521            }
3522
3523            layer.core_order = vec![Vec::new(); nr_nodes];
3524            let alloc = &layer_allocs[idx];
3525
3526            for n in 0..nr_nodes {
3527                let node = self.topo.nodes.get(&n).unwrap();
3528                let mut extra = Self::compute_target_llcs(alloc.node_target(n) * au, &node.llcs).1;
3529
3530                if let Some(node_llcs) = self.cpu_pool.free_llcs.get_mut(&n) {
3531                    for entry in node_llcs.iter_mut() {
3532                        if extra == 0 {
3533                            break;
3534                        }
3535                        let llc_id = entry.0;
3536                        let llc = self.topo.all_llcs.get(&llc_id).unwrap();
3537                        let avail = llc.cores.len() - entry.1;
3538                        let mut used = extra.min(avail);
3539                        let cores_to_add = used;
3540
3541                        let shift = entry.1;
3542                        entry.1 += used;
3543
3544                        for core in llc.cores.iter().skip(shift) {
3545                            if used == 0 {
3546                                break;
3547                            }
3548                            layer.core_order[n].push(core.1.id);
3549                            used -= 1;
3550                        }
3551
3552                        extra -= cores_to_add;
3553                    }
3554                }
3555            }
3556
3557            for node_cores in &mut layer.core_order {
3558                node_cores.reverse();
3559            }
3560        }
3561
3562        // Reset consumed entries in free LLCs.
3563        for node_llcs in self.cpu_pool.free_llcs.values_mut() {
3564            for entry in node_llcs.iter_mut() {
3565                entry.1 = 0;
3566            }
3567        }
3568
3569        // Phase 4 — Build core_order: append cores from assigned LLCs.
3570        for &(idx, _) in layer_targets.iter() {
3571            let layer = &mut self.layers[idx];
3572
3573            if layer.growth_algo != LayerGrowthAlgo::StickyDynamic {
3574                continue;
3575            }
3576
3577            let all_assigned: HashSet<usize> =
3578                layer.assigned_llcs.iter().flatten().copied().collect();
3579
3580            for core in self.topo.all_cores.iter() {
3581                let llc_id = core.1.llc_id;
3582                if all_assigned.contains(&llc_id) {
3583                    let nid = core.1.node_id;
3584                    layer.core_order[nid].push(core.1.id);
3585                }
3586            }
3587            for node_cores in &mut layer.core_order {
3588                node_cores.reverse();
3589            }
3590
3591            debug!(
3592                " alloc: layer={} core_order={:?}",
3593                layer.name, layer.core_order
3594            );
3595        }
3596
3597        // Phase 5 — Apply CPU changes for StickyDynamic layers.
3598        // Two phases: first free per-node, then allocate per-node.
3599        let mut updated = false;
3600
3601        // Free excess CPUs per-node.
3602        for &(idx, _) in layer_targets.iter() {
3603            let layer = &mut self.layers[idx];
3604
3605            if layer.growth_algo != LayerGrowthAlgo::StickyDynamic {
3606                continue;
3607            }
3608
3609            for n in 0..nr_nodes {
3610                let mut node_target = Cpumask::new();
3611                for &core_id in &layer.core_order[n] {
3612                    if let Some(core) = self.topo.all_cores.get(&core_id) {
3613                        node_target |= &core.span;
3614                    }
3615                }
3616                node_target &= &layer.allowed_cpus;
3617
3618                let node_span = &self.topo.nodes[&n].span;
3619                let node_cur = layer.cpus.and(node_span);
3620                let cpus_to_free = node_cur.and(&node_target.not());
3621
3622                if cpus_to_free.weight() > 0 {
3623                    debug!(
3624                        " apply: layer={} freeing CPUs on node {}: {}",
3625                        layer.name, n, cpus_to_free
3626                    );
3627                    layer.cpus &= &cpus_to_free.not();
3628                    layer.nr_cpus -= cpus_to_free.weight();
3629                    for cpu in cpus_to_free.iter() {
3630                        layer.nr_llc_cpus[self.cpu_pool.topo.all_cpus[&cpu].llc_id] -= 1;
3631                        layer.nr_node_cpus[n] -= 1;
3632                    }
3633                    self.cpu_pool.free(&cpus_to_free)?;
3634                    updated = true;
3635                }
3636            }
3637        }
3638
3639        // Allocate needed CPUs per-node.
3640        for &(idx, _) in layer_targets.iter() {
3641            let layer = &mut self.layers[idx];
3642
3643            if layer.growth_algo != LayerGrowthAlgo::StickyDynamic {
3644                continue;
3645            }
3646
3647            for n in 0..nr_nodes {
3648                let mut node_target = Cpumask::new();
3649                for &core_id in &layer.core_order[n] {
3650                    if let Some(core) = self.topo.all_cores.get(&core_id) {
3651                        node_target |= &core.span;
3652                    }
3653                }
3654                node_target &= &layer.allowed_cpus;
3655
3656                let available_cpus = self.cpu_pool.available_cpus();
3657                let desired_to_alloc = node_target.and(&layer.cpus.clone().not());
3658                let cpus_to_alloc = desired_to_alloc.clone().and(&available_cpus);
3659
3660                if desired_to_alloc.weight() > cpus_to_alloc.weight() {
3661                    debug!(
3662                        " apply: layer={} node {} wanted to alloc {} CPUs but only {} available",
3663                        layer.name,
3664                        n,
3665                        desired_to_alloc.weight(),
3666                        cpus_to_alloc.weight()
3667                    );
3668                }
3669
3670                if cpus_to_alloc.weight() > 0 {
3671                    debug!(
3672                        " apply: layer={} allocating CPUs on node {}: {}",
3673                        layer.name, n, cpus_to_alloc
3674                    );
3675                    layer.cpus |= &cpus_to_alloc;
3676                    layer.nr_cpus += cpus_to_alloc.weight();
3677                    for cpu in cpus_to_alloc.iter() {
3678                        layer.nr_llc_cpus[self.cpu_pool.topo.all_cpus[&cpu].llc_id] += 1;
3679                        layer.nr_node_cpus[n] += 1;
3680                    }
3681                    self.cpu_pool.mark_allocated(&cpus_to_alloc)?;
3682                    updated = true;
3683                }
3684            }
3685
3686            debug!(
3687                " apply: layer={} final cpus.weight()={} nr_cpus={}",
3688                layer.name,
3689                layer.cpus.weight(),
3690                layer.nr_cpus
3691            );
3692        }
3693
3694        Ok(updated)
3695    }
3696
3697    fn refresh_node_ctx(
3698        skel: &mut BpfSkel,
3699        topo: &Topology,
3700        node_empty_layers: &[Vec<u32>],
3701        init: bool,
3702    ) {
3703        for &nid in topo.nodes.keys() {
3704            let mut arg: bpf_intf::refresh_node_ctx_arg =
3705                unsafe { MaybeUninit::zeroed().assume_init() };
3706            arg.node_id = nid as u32;
3707            arg.init = init as u32;
3708
3709            let empty = &node_empty_layers[nid];
3710            arg.nr_empty_layer_ids = empty.len() as u32;
3711            for (i, &lid) in empty.iter().enumerate() {
3712                arg.empty_layer_ids[i] = lid;
3713            }
3714            for i in empty.len()..MAX_LAYERS {
3715                arg.empty_layer_ids[i] = MAX_LAYERS as u32;
3716            }
3717
3718            if init {
3719                // Per-node LLC list — static topology, only set during init.
3720                let node = &topo.nodes[&nid];
3721                let llcs: Vec<u32> = node.llcs.keys().map(|&id| id as u32).collect();
3722                arg.nr_llcs = llcs.len() as u32;
3723                for (i, &llc_id) in llcs.iter().enumerate() {
3724                    arg.llcs[i] = llc_id;
3725                }
3726            }
3727
3728            let input = ProgramInput {
3729                context_in: Some(unsafe { plain::as_mut_bytes(&mut arg) }),
3730                ..Default::default()
3731            };
3732            let _ = skel.progs.refresh_node_ctx.test_run(input);
3733        }
3734    }
3735
3736    fn refresh_cpumasks(&mut self) -> Result<()> {
3737        let layer_is_open = |layer: &Layer| matches!(layer.kind, LayerKind::Open { .. });
3738
3739        let mut updated = false;
3740        let raw_targets = self.calc_target_nr_cpus();
3741        let au = self.cpu_pool.alloc_unit();
3742        let total_cpus = self.cpu_pool.topo.all_cpus.len();
3743
3744        // Dampen shrink: only drop halfway per cycle to avoid unnecessary
3745        // changes. There's some dampening built into util metrics but slow
3746        // down freeing further. This is solely based on intuition. Drop or
3747        // update according to real-world behavior.
3748        let targets: Vec<(usize, usize)> = raw_targets
3749            .iter()
3750            .enumerate()
3751            .map(|(idx, &(target, min))| {
3752                let cur = self.layers[idx].nr_cpus;
3753                if target < cur {
3754                    let dampened = cur - (cur - target).div_ceil(2);
3755                    (dampened.max(min), min)
3756                } else {
3757                    (target, min)
3758                }
3759            })
3760            .collect();
3761
3762        // Build demands for unified_alloc and compute per-node allocations.
3763        let demands = self.calc_raw_demands(&targets);
3764        let nr_nodes = self.topo.nodes.len();
3765        let node_caps: Vec<usize> = self
3766            .topo
3767            .nodes
3768            .values()
3769            .map(|n| n.span.weight() / au)
3770            .collect();
3771        let all_layer_nodes: Vec<&[usize]> = self
3772            .layer_specs
3773            .iter()
3774            .map(|s| s.nodes().as_slice())
3775            .collect();
3776        let norders: Vec<Vec<usize>> = (0..self.layers.len())
3777            .map(|idx| {
3778                layer_core_growth::node_order(
3779                    self.layer_specs[idx].nodes(),
3780                    &self.topo,
3781                    idx,
3782                    &all_layer_nodes,
3783                )
3784            })
3785            .collect();
3786        let node_groups: Vec<Vec<Vec<usize>>> = (0..self.layers.len())
3787            .map(|idx| {
3788                layer_core_growth::node_groups(
3789                    self.layer_specs[idx].nodes(),
3790                    &self.topo,
3791                    idx,
3792                    &all_layer_nodes,
3793                    &self.layer_specs[idx].kind.common().growth_algo,
3794                )
3795            })
3796            .collect();
3797        let layer_allocs = unified_alloc(total_cpus / au, &node_caps, &demands, &node_groups);
3798
3799        // Convert allocations back to CPU counts. Shrink dampening is
3800        // already applied to the targets fed into unified_alloc above.
3801        let cpu_targets: Vec<usize> = layer_allocs.iter().map(|a| a.total() * au).collect();
3802
3803        // Snapshot per-layer CPU counts for ALLOC debug logging.
3804        let prev_nr_cpus: Vec<usize> = self.layers.iter().map(|l| l.nr_cpus).collect();
3805
3806        let mut ascending: Vec<(usize, usize)> = cpu_targets.iter().copied().enumerate().collect();
3807        ascending.sort_by_key(|a| a.1);
3808
3809        // Snapshot per-layer per-node CPU counts before allocation changes.
3810        let prev_node_cpus: Vec<Vec<usize>> =
3811            self.layers.iter().map(|l| l.nr_node_cpus.clone()).collect();
3812
3813        // Per-node SD allocation requires multiple LLCs. On flat topologies
3814        // (single LLC, e.g. VMs with topology disabled), SD layers fall through
3815        // to the non-SD grow/shrink loops below.
3816        let use_sd_alloc = self.topo.all_llcs.len() > 1;
3817        let sticky_dynamic_updated = if use_sd_alloc {
3818            self.recompute_layer_core_order(&ascending, &layer_allocs, au)?
3819        } else {
3820            false
3821        };
3822        updated |= sticky_dynamic_updated;
3823
3824        // Update BPF cpumasks for StickyDynamic layers if they were updated
3825        if sticky_dynamic_updated {
3826            for (idx, layer) in self.layers.iter().enumerate() {
3827                if layer.growth_algo == LayerGrowthAlgo::StickyDynamic {
3828                    Self::update_bpf_layer_cpumask(
3829                        layer,
3830                        &mut self.skel.maps.bss_data.as_mut().unwrap().layers[idx],
3831                    );
3832                }
3833            }
3834        }
3835
3836        // Shrink per-node: free excess CPUs from each node.
3837        for &(idx, _target) in ascending.iter().rev() {
3838            let layer = &mut self.layers[idx];
3839            if layer_is_open(layer) {
3840                continue;
3841            }
3842
3843            // Skip StickyDynamic layers when per-node SD allocation is active.
3844            if layer.growth_algo == LayerGrowthAlgo::StickyDynamic && use_sd_alloc {
3845                continue;
3846            }
3847
3848            let alloc = &layer_allocs[idx];
3849            let mut freed = false;
3850
3851            for n in 0..nr_nodes {
3852                let desired = alloc.node_target(n) * au;
3853                let mut to_free = layer.nr_node_cpus[n].saturating_sub(desired);
3854                let node_span = &self.topo.nodes[&n].span;
3855
3856                while to_free > 0 {
3857                    let node_cands = layer.cpus.and(node_span);
3858                    let cpus_to_free = match self
3859                        .cpu_pool
3860                        .next_to_free(&node_cands, layer.core_order[n].iter().rev())?
3861                    {
3862                        Some(ret) => ret,
3863                        None => break,
3864                    };
3865                    let nr = cpus_to_free.weight();
3866                    trace!(
3867                        "[{}] freeing CPUs on node {}: {}",
3868                        layer.name, n, &cpus_to_free
3869                    );
3870                    layer.cpus &= &cpus_to_free.not();
3871                    layer.nr_cpus -= nr;
3872                    for cpu in cpus_to_free.iter() {
3873                        let node_id = self.cpu_pool.topo.all_cpus[&cpu].node_id;
3874                        layer.nr_llc_cpus[self.cpu_pool.topo.all_cpus[&cpu].llc_id] -= 1;
3875                        layer.nr_node_cpus[node_id] -= 1;
3876                        layer.nr_pinned_cpus[node_id] =
3877                            layer.nr_pinned_cpus[node_id].min(layer.nr_node_cpus[node_id]);
3878                    }
3879                    self.cpu_pool.free(&cpus_to_free)?;
3880                    to_free = to_free.saturating_sub(nr);
3881                    freed = true;
3882                }
3883            }
3884
3885            if freed {
3886                Self::update_bpf_layer_cpumask(
3887                    layer,
3888                    &mut self.skel.maps.bss_data.as_mut().unwrap().layers[idx],
3889                );
3890                updated = true;
3891            }
3892        }
3893
3894        // Grow layers per-node using allocations from unified_alloc.
3895        for &(idx, _target) in &ascending {
3896            let layer = &mut self.layers[idx];
3897
3898            if layer_is_open(layer) {
3899                continue;
3900            }
3901
3902            // Skip StickyDynamic layers when per-node SD allocation is active.
3903            if layer.growth_algo == LayerGrowthAlgo::StickyDynamic && use_sd_alloc {
3904                continue;
3905            }
3906
3907            let alloc = &layer_allocs[idx];
3908            let norder = &norders[idx];
3909            let mut alloced = false;
3910
3911            for &node_id in norder.iter() {
3912                let node_target = alloc.node_target(node_id) * au;
3913                let cur_node = layer.nr_node_cpus[node_id];
3914                if node_target <= cur_node {
3915                    continue;
3916                }
3917                let pinned_target = alloc.pinned[node_id] * au;
3918                let mut nr_to_alloc = node_target - cur_node;
3919                let node_span = &self.topo.nodes[&node_id].span;
3920                let node_allowed = layer.allowed_cpus.and(node_span);
3921
3922                while nr_to_alloc > 0 {
3923                    let nr_alloced = match self.cpu_pool.alloc_cpus(
3924                        &node_allowed,
3925                        &layer.core_order[node_id],
3926                        nr_to_alloc,
3927                    ) {
3928                        Some(new_cpus) => {
3929                            let nr = new_cpus.weight();
3930                            layer.cpus |= &new_cpus;
3931                            layer.nr_cpus += nr;
3932                            for cpu in new_cpus.iter() {
3933                                layer.nr_llc_cpus[self.cpu_pool.topo.all_cpus[&cpu].llc_id] += 1;
3934                                let nid = self.cpu_pool.topo.all_cpus[&cpu].node_id;
3935                                layer.nr_node_cpus[nid] += 1;
3936                                if layer.nr_pinned_cpus[nid] < pinned_target {
3937                                    layer.nr_pinned_cpus[nid] += 1;
3938                                }
3939                            }
3940                            nr
3941                        }
3942                        None => 0,
3943                    };
3944                    if nr_alloced == 0 {
3945                        break;
3946                    }
3947                    alloced = true;
3948                    nr_to_alloc -= nr_alloced.min(nr_to_alloc);
3949                }
3950            }
3951
3952            if alloced {
3953                Self::update_bpf_layer_cpumask(
3954                    layer,
3955                    &mut self.skel.maps.bss_data.as_mut().unwrap().layers[idx],
3956                );
3957                updated = true;
3958            }
3959        }
3960
3961        // Tell BPF whether all CPUs are allocated.  When the system is
3962        // fully allocated, idle_confined layers have nowhere to grow, so
3963        // they should be allowed to use open (unprotected) idle CPUs from
3964        // other layers.
3965        let total_allocated: usize = self.layers.iter().map(|l| l.nr_cpus).sum();
3966        let fully_allocated = total_allocated >= total_cpus;
3967        for (idx, layer) in self.layers.iter().enumerate() {
3968            if !layer_is_open(layer) {
3969                self.skel.maps.bss_data.as_mut().unwrap().layers[idx]
3970                    .fully_allocated
3971                    .write(fully_allocated);
3972            }
3973        }
3974
3975        // Recompute growth_denied: for each node, check whether the
3976        // unpinned portion of the layer warranted growth (unpinned util /
3977        // util_high > unpinned CPUs) but the node didn't gain CPUs.  Only
3978        // unpinned matters because pinned tasks can't migrate cross-node.
3979        let node_utils = &self.sched_stats.layer_node_utils;
3980        let pinned_utils = &self.sched_stats.layer_node_pinned_utils;
3981        for (idx, layer) in self.layers.iter().enumerate() {
3982            self.growth_denied[idx].fill(false);
3983            let util_high = match layer.kind.util_range() {
3984                Some((_, high)) => high,
3985                None => continue,
3986            };
3987            for n in 0..nr_nodes {
3988                let unpinned_util = (node_utils[idx][n] - pinned_utils[idx][n]).max(0.0);
3989                let unpinned_cpus_needed = unpinned_util / util_high;
3990                let unpinned_cpus_have =
3991                    layer.nr_node_cpus[n].saturating_sub(layer.nr_pinned_cpus[n]) as f64;
3992                let wanted = unpinned_cpus_needed > unpinned_cpus_have;
3993                let got = layer.nr_node_cpus[n] > prev_node_cpus[idx][n];
3994                if wanted && !got {
3995                    self.growth_denied[idx][n] = true;
3996                }
3997            }
3998        }
3999
4000        // Log per-layer allocation changes.
4001        if updated {
4002            for (idx, layer) in self.layers.iter().enumerate() {
4003                if layer_is_open(layer) {
4004                    continue;
4005                }
4006                let prev = prev_nr_cpus[idx];
4007                let cur = layer.nr_cpus;
4008                if prev != cur {
4009                    debug!(
4010                        "ALLOC {} algo={:?} cpus:{}→{} mask={:x}",
4011                        layer.name, layer.growth_algo, prev, cur, layer.cpus,
4012                    );
4013                }
4014            }
4015            debug!(
4016                "ALLOC pool_available={}",
4017                self.cpu_pool.available_cpus().weight()
4018            );
4019        }
4020
4021        // Log allocation changes.
4022        if updated {
4023            let nr_nodes = self.topo.nodes.len();
4024            for (idx, layer) in self.layers.iter().enumerate() {
4025                if layer_is_open(layer) {
4026                    continue;
4027                }
4028                let prev = &prev_node_cpus[idx];
4029                let cur = &layer.nr_node_cpus;
4030                if prev == cur {
4031                    continue;
4032                }
4033                let per_node: String = (0..nr_nodes)
4034                    .map(|n| format!("n{}:{}→{}", n, prev[n], cur[n]))
4035                    .collect::<Vec<_>>()
4036                    .join(" ");
4037                let prev_total: usize = prev.iter().sum();
4038                let cur_total: usize = cur[..nr_nodes].iter().sum();
4039                let target: String = (0..nr_nodes)
4040                    .map(|n| format!("n{}:{}", n, layer_allocs[idx].node_target(n) * au))
4041                    .collect::<Vec<_>>()
4042                    .join(" ");
4043                debug!(
4044                    "ALLOC {} algo={:?} {} total:{}→{} target:[{}] mask={:x}",
4045                    layer.name,
4046                    layer.growth_algo,
4047                    per_node,
4048                    prev_total,
4049                    cur_total,
4050                    target,
4051                    layer.cpus,
4052                );
4053            }
4054            debug!(
4055                "ALLOC pool_available={}",
4056                self.cpu_pool.available_cpus().weight()
4057            );
4058        }
4059
4060        // Give the rest to the open layers.
4061        if updated {
4062            for (idx, layer) in self.layers.iter_mut().enumerate() {
4063                if !layer_is_open(layer) {
4064                    continue;
4065                }
4066
4067                let bpf_layer = &mut self.skel.maps.bss_data.as_mut().unwrap().layers[idx];
4068                let available_cpus = self.cpu_pool.available_cpus().and(&layer.allowed_cpus);
4069                let nr_available_cpus = available_cpus.weight();
4070
4071                // Open layers need the intersection of allowed cpus and
4072                // available cpus. Recompute per-LLC and per-node counts
4073                // since open layers bypass alloc/free.
4074                layer.cpus = available_cpus;
4075                layer.nr_cpus = nr_available_cpus;
4076                for llc in self.cpu_pool.topo.all_llcs.values() {
4077                    layer.nr_llc_cpus[llc.id] = layer.cpus.and(&llc.span).weight();
4078                }
4079                for node in self.cpu_pool.topo.nodes.values() {
4080                    layer.nr_node_cpus[node.id] = layer.cpus.and(&node.span).weight();
4081                    layer.nr_pinned_cpus[node.id] = 0;
4082                }
4083                Self::update_bpf_layer_cpumask(layer, bpf_layer);
4084            }
4085
4086            for (&node_id, &cpu) in &self.cpu_pool.fallback_cpus {
4087                self.skel.maps.bss_data.as_mut().unwrap().fallback_cpus[node_id] = cpu as u32;
4088            }
4089
4090            for (lidx, layer) in self.layers.iter().enumerate() {
4091                self.nr_layer_cpus_ranges[lidx] = (
4092                    self.nr_layer_cpus_ranges[lidx].0.min(layer.nr_cpus),
4093                    self.nr_layer_cpus_ranges[lidx].1.max(layer.nr_cpus),
4094                );
4095            }
4096
4097            // Trigger updates on the BPF side.
4098            let input = ProgramInput {
4099                ..Default::default()
4100            };
4101            let prog = &mut self.skel.progs.refresh_layer_cpumasks;
4102            let _ = prog.test_run(input);
4103
4104            // Update per-node empty layer IDs via BPF prog.
4105            let nr_nodes = self.topo.nodes.len();
4106            let node_empty_layers: Vec<Vec<u32>> = (0..nr_nodes)
4107                .map(|nid| {
4108                    self.layers
4109                        .iter()
4110                        .enumerate()
4111                        .filter(|(_lidx, layer)| layer.nr_node_cpus[nid] == 0)
4112                        .map(|(lidx, _)| lidx as u32)
4113                        .collect()
4114                })
4115                .collect();
4116            Self::refresh_node_ctx(&mut self.skel, &self.topo, &node_empty_layers, false);
4117        }
4118
4119        if let Err(e) = self.update_netdev_cpumasks() {
4120            warn!("Failed to update netdev IRQ cpumasks: {:#}", e);
4121        }
4122        Ok(())
4123    }
4124
4125    fn refresh_idle_qos(&mut self) -> Result<()> {
4126        if !self.idle_qos_enabled {
4127            return Ok(());
4128        }
4129
4130        let mut cpu_idle_qos = vec![0; *NR_CPU_IDS];
4131        for layer in self.layers.iter() {
4132            let idle_resume_us = layer.kind.common().idle_resume_us.unwrap_or(0) as i32;
4133            for cpu in layer.cpus.iter() {
4134                cpu_idle_qos[cpu] = idle_resume_us;
4135            }
4136        }
4137
4138        for (cpu, idle_resume_usec) in cpu_idle_qos.iter().enumerate() {
4139            update_cpu_idle_resume_latency(cpu, *idle_resume_usec)?;
4140        }
4141
4142        Ok(())
4143    }
4144
4145    fn refresh_xnuma(&mut self) {
4146        let nr_nodes = self.topo.nodes.len();
4147        if nr_nodes <= 1 {
4148            return;
4149        }
4150
4151        let duty_sums = &self.sched_stats.layer_node_duty_sums;
4152
4153        for (layer_idx, spec) in self.layer_specs.iter().enumerate() {
4154            let common = spec.kind.common();
4155            let threshold = common.xnuma_threshold;
4156            let threshold_delta = common.xnuma_threshold_delta;
4157            let bpf_layer = &mut self.skel.maps.bss_data.as_mut().unwrap().layers[layer_idx];
4158
4159            if threshold.0 <= 0.0 && threshold.1 <= 0.0 {
4160                // Off — all open, infinite budget
4161                for src in 0..nr_nodes {
4162                    bpf_layer.node[src].xnuma_is_mig_src.write(true);
4163                    for dst in 0..nr_nodes {
4164                        bpf_layer.node[src].xnuma[dst].rate = u64::MAX;
4165                    }
4166                }
4167                self.xnuma_mig_src[layer_idx].fill(false);
4168                continue;
4169            }
4170
4171            let layer = &self.layers[layer_idx];
4172            let is_mig_src = xnuma_check_active(
4173                &duty_sums[layer_idx],
4174                &layer.nr_node_cpus,
4175                threshold,
4176                threshold_delta,
4177                &self.growth_denied[layer_idx],
4178                &self.xnuma_mig_src[layer_idx],
4179            );
4180
4181            self.xnuma_mig_src[layer_idx] = is_mig_src.clone();
4182
4183            let result = xnuma_compute_rates(&duty_sums[layer_idx], &layer.nr_node_cpus);
4184
4185            // Write rates before is_mig_src so BPF sees valid rates
4186            // when the gate activates.
4187            for src in 0..nr_nodes {
4188                for dst in 0..nr_nodes {
4189                    bpf_layer.node[src].xnuma[dst].rate = result.rates[src][dst];
4190                }
4191            }
4192            for (nid, is_src) in is_mig_src.iter().enumerate().take(nr_nodes) {
4193                bpf_layer.node[nid].xnuma_is_mig_src.write(*is_src);
4194            }
4195        }
4196    }
4197
4198    fn step(&mut self) -> Result<()> {
4199        let started_at = Instant::now();
4200        self.sched_stats.refresh(
4201            &mut self.skel,
4202            &self.proc_reader,
4203            started_at,
4204            self.processing_dur,
4205            &self.gpu_task_handler,
4206        )?;
4207
4208        // Update BPF with EWMA values
4209        self.skel
4210            .maps
4211            .bss_data
4212            .as_mut()
4213            .unwrap()
4214            .system_cpu_util_ewma = (self.sched_stats.system_cpu_util_ewma * 10000.0) as u64;
4215
4216        for layer_id in 0..self.sched_stats.nr_layers {
4217            self.skel
4218                .maps
4219                .bss_data
4220                .as_mut()
4221                .unwrap()
4222                .layer_dsq_insert_ewma[layer_id] =
4223                (self.sched_stats.layer_dsq_insert_ewma[layer_id] * 10000.0) as u64;
4224        }
4225
4226        self.refresh_cpumasks()?;
4227        self.refresh_xnuma();
4228        self.refresh_idle_qos()?;
4229        self.gpu_task_handler.maybe_affinitize();
4230        self.processing_dur += Instant::now().duration_since(started_at);
4231        Ok(())
4232    }
4233
4234    fn generate_sys_stats(
4235        &mut self,
4236        stats: &Stats,
4237        cpus_ranges: &mut [(usize, usize)],
4238    ) -> Result<SysStats> {
4239        let bstats = &stats.bpf_stats;
4240        let mut sys_stats = SysStats::new(stats, bstats, &self.cpu_pool.fallback_cpus)?;
4241
4242        for (lidx, (spec, layer)) in self.layer_specs.iter().zip(self.layers.iter()).enumerate() {
4243            let layer_stats = LayerStats::new(
4244                lidx,
4245                layer,
4246                stats,
4247                bstats,
4248                cpus_ranges[lidx],
4249                self.xnuma_mig_src[lidx].iter().any(|&a| a),
4250            );
4251            sys_stats.layers.insert(spec.name.to_string(), layer_stats);
4252            cpus_ranges[lidx] = (layer.nr_cpus, layer.nr_cpus);
4253        }
4254
4255        Ok(sys_stats)
4256    }
4257
4258    // Helper function to process a cgroup creation (common logic for walkdir and inotify)
4259    fn process_cgroup_creation(
4260        path: &Path,
4261        cgroup_regexes: &HashMap<u32, Regex>,
4262        cgroup_path_to_id: &mut HashMap<String, u64>,
4263        sender: &crossbeam::channel::Sender<CgroupEvent>,
4264    ) {
4265        let path_str = path.to_string_lossy().to_string();
4266
4267        // Get cgroup ID (inode number)
4268        let cgroup_id = std::fs::metadata(path)
4269            .map(|metadata| {
4270                use std::os::unix::fs::MetadataExt;
4271                metadata.ino()
4272            })
4273            .unwrap_or(0);
4274
4275        // Build match bitmap by testing against CgroupRegex rules
4276        let mut match_bitmap = 0u64;
4277        for (rule_id, regex) in cgroup_regexes {
4278            if regex.is_match(&path_str) {
4279                match_bitmap |= 1u64 << rule_id;
4280            }
4281        }
4282
4283        // Store in hash
4284        cgroup_path_to_id.insert(path_str.clone(), cgroup_id);
4285
4286        // Send event
4287        if let Err(e) = sender.send(CgroupEvent::Created {
4288            path: path_str,
4289            cgroup_id,
4290            match_bitmap,
4291        }) {
4292            error!("Failed to send cgroup creation event: {}", e);
4293        }
4294    }
4295
4296    fn start_cgroup_watcher(
4297        shutdown: Arc<AtomicBool>,
4298        cgroup_regexes: HashMap<u32, Regex>,
4299    ) -> Result<Receiver<CgroupEvent>> {
4300        let mut inotify = Inotify::init().context("Failed to initialize inotify")?;
4301        let mut wd_to_path = HashMap::new();
4302
4303        // Create crossbeam channel for cgroup events (bounded to prevent memory issues)
4304        let (sender, receiver) = crossbeam::channel::bounded::<CgroupEvent>(1024);
4305
4306        // Watch for directory creation and deletion events
4307        let root_wd = inotify
4308            .watches()
4309            .add("/sys/fs/cgroup", WatchMask::CREATE | WatchMask::DELETE)
4310            .context("Failed to add watch for /sys/fs/cgroup")?;
4311        wd_to_path.insert(root_wd, PathBuf::from("/sys/fs/cgroup"));
4312
4313        // Also recursively watch existing directories for new subdirectories
4314        Self::add_recursive_watches(&mut inotify, &mut wd_to_path, Path::new("/sys/fs/cgroup"))?;
4315
4316        // Spawn watcher thread
4317        std::thread::spawn(move || {
4318            let mut buffer = [0; 4096];
4319            let inotify_fd = inotify.as_raw_fd();
4320            // Maintain hash of cgroup path -> cgroup ID (inode number)
4321            let mut cgroup_path_to_id = HashMap::<String, u64>::new();
4322
4323            // Populate existing cgroups
4324            for entry in WalkDir::new("/sys/fs/cgroup")
4325                .into_iter()
4326                .filter_map(|e| e.ok())
4327                .filter(|e| e.file_type().is_dir())
4328            {
4329                let path = entry.path();
4330                Self::process_cgroup_creation(
4331                    path,
4332                    &cgroup_regexes,
4333                    &mut cgroup_path_to_id,
4334                    &sender,
4335                );
4336            }
4337
4338            while !shutdown.load(Ordering::Relaxed) {
4339                // Use select to wait for events with a 100ms timeout
4340                let ready = unsafe {
4341                    let mut read_fds: libc::fd_set = std::mem::zeroed();
4342                    libc::FD_ZERO(&mut read_fds);
4343                    libc::FD_SET(inotify_fd, &mut read_fds);
4344
4345                    let mut timeout = libc::timeval {
4346                        tv_sec: 0,
4347                        tv_usec: 100_000, // 100ms
4348                    };
4349
4350                    libc::select(
4351                        inotify_fd + 1,
4352                        &mut read_fds,
4353                        std::ptr::null_mut(),
4354                        std::ptr::null_mut(),
4355                        &mut timeout,
4356                    )
4357                };
4358
4359                if ready <= 0 {
4360                    // Timeout or error, continue loop to check shutdown
4361                    continue;
4362                }
4363
4364                // Read events non-blocking
4365                let events = match inotify.read_events(&mut buffer) {
4366                    Ok(events) => events,
4367                    Err(e) => {
4368                        error!("Error reading inotify events: {}", e);
4369                        break;
4370                    }
4371                };
4372
4373                for event in events {
4374                    if !event.mask.contains(inotify::EventMask::CREATE)
4375                        && !event.mask.contains(inotify::EventMask::DELETE)
4376                    {
4377                        continue;
4378                    }
4379
4380                    let name = match event.name {
4381                        Some(name) => name,
4382                        None => continue,
4383                    };
4384
4385                    let parent_path = match wd_to_path.get(&event.wd) {
4386                        Some(parent) => parent,
4387                        None => {
4388                            warn!("Unknown watch descriptor: {:?}", event.wd);
4389                            continue;
4390                        }
4391                    };
4392
4393                    let path = parent_path.join(name.to_string_lossy().as_ref());
4394
4395                    if event.mask.contains(inotify::EventMask::CREATE) {
4396                        if !path.is_dir() {
4397                            continue;
4398                        }
4399
4400                        Self::process_cgroup_creation(
4401                            &path,
4402                            &cgroup_regexes,
4403                            &mut cgroup_path_to_id,
4404                            &sender,
4405                        );
4406
4407                        // Add watch for this new directory
4408                        match inotify
4409                            .watches()
4410                            .add(&path, WatchMask::CREATE | WatchMask::DELETE)
4411                        {
4412                            Ok(wd) => {
4413                                wd_to_path.insert(wd, path.clone());
4414                            }
4415                            Err(e) => {
4416                                warn!(
4417                                    "Failed to add watch for new cgroup {}: {}",
4418                                    path.display(),
4419                                    e
4420                                );
4421                            }
4422                        }
4423                    } else if event.mask.contains(inotify::EventMask::DELETE) {
4424                        let path_str = path.to_string_lossy().to_string();
4425
4426                        // Get cgroup ID from our hash (since the directory is gone, we can't stat it)
4427                        let cgroup_id = cgroup_path_to_id.remove(&path_str).unwrap_or(0);
4428
4429                        // Send removal event to main thread
4430                        if let Err(e) = sender.send(CgroupEvent::Removed {
4431                            path: path_str,
4432                            cgroup_id,
4433                        }) {
4434                            error!("Failed to send cgroup removal event: {}", e);
4435                        }
4436
4437                        // Find and remove the watch descriptor for this path
4438                        let wd_to_remove = wd_to_path.iter().find_map(|(wd, watched_path)| {
4439                            if watched_path == &path {
4440                                Some(wd.clone())
4441                            } else {
4442                                None
4443                            }
4444                        });
4445                        if let Some(wd) = wd_to_remove {
4446                            wd_to_path.remove(&wd);
4447                        }
4448                    }
4449                }
4450            }
4451        });
4452
4453        Ok(receiver)
4454    }
4455
4456    fn add_recursive_watches(
4457        inotify: &mut Inotify,
4458        wd_to_path: &mut HashMap<inotify::WatchDescriptor, PathBuf>,
4459        path: &Path,
4460    ) -> Result<()> {
4461        for entry in WalkDir::new(path)
4462            .into_iter()
4463            .filter_map(|e| e.ok())
4464            .filter(|e| e.file_type().is_dir())
4465            .skip(1)
4466        {
4467            let entry_path = entry.path();
4468            // Add watch for this directory
4469            match inotify
4470                .watches()
4471                .add(entry_path, WatchMask::CREATE | WatchMask::DELETE)
4472            {
4473                Ok(wd) => {
4474                    wd_to_path.insert(wd, entry_path.to_path_buf());
4475                }
4476                Err(e) => {
4477                    debug!("Failed to add watch for {}: {}", entry_path.display(), e);
4478                }
4479            }
4480        }
4481        Ok(())
4482    }
4483
4484    fn run(&mut self, shutdown: Arc<AtomicBool>) -> Result<UserExitInfo> {
4485        let (res_ch, req_ch) = self.stats_server.channels();
4486        let mut next_sched_at = Instant::now() + self.sched_intv;
4487        let enable_layer_refresh = !self.layer_refresh_intv.is_zero();
4488        let mut next_layer_refresh_at = Instant::now() + self.layer_refresh_intv;
4489        let mut cpus_ranges = HashMap::<ThreadId, Vec<(usize, usize)>>::new();
4490
4491        // Start the cgroup watcher only if there are CgroupRegex rules
4492        let cgroup_regexes = self.cgroup_regexes.take().unwrap();
4493        let cgroup_event_rx = if !cgroup_regexes.is_empty() {
4494            Some(Self::start_cgroup_watcher(
4495                shutdown.clone(),
4496                cgroup_regexes,
4497            )?)
4498        } else {
4499            None
4500        };
4501
4502        while !shutdown.load(Ordering::Relaxed) && !uei_exited!(&self.skel, uei) {
4503            let now = Instant::now();
4504
4505            if now >= next_sched_at {
4506                self.step()?;
4507                while next_sched_at < now {
4508                    next_sched_at += self.sched_intv;
4509                }
4510            }
4511
4512            if enable_layer_refresh && now >= next_layer_refresh_at {
4513                self.skel
4514                    .maps
4515                    .bss_data
4516                    .as_mut()
4517                    .unwrap()
4518                    .layer_refresh_seq_avgruntime += 1;
4519                while next_layer_refresh_at < now {
4520                    next_layer_refresh_at += self.layer_refresh_intv;
4521                }
4522            }
4523
4524            // Handle both stats requests and cgroup events with timeout
4525            let timeout_duration = next_sched_at.saturating_duration_since(Instant::now());
4526            let never_rx = crossbeam::channel::never();
4527            let cgroup_rx = cgroup_event_rx.as_ref().unwrap_or(&never_rx);
4528
4529            select! {
4530                recv(req_ch) -> msg => match msg {
4531                    Ok(StatsReq::Hello(tid)) => {
4532                        cpus_ranges.insert(
4533                            tid,
4534                            self.layers.iter().map(|l| (l.nr_cpus, l.nr_cpus)).collect(),
4535                        );
4536                        let stats =
4537                            Stats::new(&mut self.skel, &self.proc_reader, self.topo.clone(), &self.gpu_task_handler, self.sched_stats.util_compensation)?;
4538                        res_ch.send(StatsRes::Hello(Box::new(stats)))?;
4539                    }
4540                    Ok(StatsReq::Refresh(tid, mut stats)) => {
4541                        // Propagate self's layer cpu ranges into each stat's.
4542                        for i in 0..self.nr_layer_cpus_ranges.len() {
4543                            for ranges in cpus_ranges.values_mut() {
4544                                ranges[i] = (
4545                                    ranges[i].0.min(self.nr_layer_cpus_ranges[i].0),
4546                                    ranges[i].1.max(self.nr_layer_cpus_ranges[i].1),
4547                                );
4548                            }
4549                            self.nr_layer_cpus_ranges[i] =
4550                                (self.layers[i].nr_cpus, self.layers[i].nr_cpus);
4551                        }
4552
4553                        stats.refresh(
4554                            &mut self.skel,
4555                            &self.proc_reader,
4556                            now,
4557                            self.processing_dur,
4558                            &self.gpu_task_handler,
4559                        )?;
4560                        let sys_stats =
4561                            self.generate_sys_stats(&stats, cpus_ranges.get_mut(&tid).unwrap())?;
4562                        res_ch.send(StatsRes::Refreshed(Box::new((*stats, sys_stats))))?;
4563                    }
4564                    Ok(StatsReq::Bye(tid)) => {
4565                        cpus_ranges.remove(&tid);
4566                        res_ch.send(StatsRes::Bye)?;
4567                    }
4568                    Err(e) => Err(e)?,
4569                },
4570
4571                recv(cgroup_rx) -> event => match event {
4572                    Ok(CgroupEvent::Created { path, cgroup_id, match_bitmap }) => {
4573                        // Insert into BPF map
4574                        self.skel.maps.cgroup_match_bitmap.update(
4575                            &cgroup_id.to_ne_bytes(),
4576                            &match_bitmap.to_ne_bytes(),
4577                            libbpf_rs::MapFlags::ANY,
4578                        ).with_context(|| format!(
4579                            "Failed to insert cgroup {}({}) into BPF map. Cgroup map may be full \
4580                             (max 16384 entries). Aborting.",
4581                            cgroup_id, path
4582                        ))?;
4583
4584                        debug!("Added cgroup {} to BPF map with bitmap 0x{:x}", cgroup_id, match_bitmap);
4585                    }
4586                    Ok(CgroupEvent::Removed { path, cgroup_id }) => {
4587                        // Delete from BPF map
4588                        if let Err(e) = self.skel.maps.cgroup_match_bitmap.delete(&cgroup_id.to_ne_bytes()) {
4589                            warn!("Failed to delete cgroup {} from BPF map: {}", cgroup_id, e);
4590                        } else {
4591                            debug!("Removed cgroup {}({}) from BPF map", cgroup_id, path);
4592                        }
4593                    }
4594                    Err(e) => {
4595                        error!("Error receiving cgroup event: {}", e);
4596                    }
4597                },
4598
4599                recv(crossbeam::channel::after(timeout_duration)) -> _ => {
4600                    // Timeout - continue main loop
4601                }
4602            }
4603        }
4604
4605        let _ = self.struct_ops.take();
4606        uei_report!(&self.skel, uei)
4607    }
4608}
4609
4610impl Drop for Scheduler<'_> {
4611    fn drop(&mut self) {
4612        info!("Unregister {SCHEDULER_NAME} scheduler");
4613
4614        if !self.netdevs.is_empty() {
4615            for (iface, netdev) in &self.netdevs {
4616                if let Err(e) = netdev.restore_cpumasks() {
4617                    warn!("Failed to restore {iface} IRQ affinity: {e}");
4618                }
4619            }
4620            info!("Restored original netdev IRQ affinity");
4621        }
4622
4623        if let Some(struct_ops) = self.struct_ops.take() {
4624            drop(struct_ops);
4625        }
4626    }
4627}
4628
4629fn write_example_file(path: &str) -> Result<()> {
4630    let mut f = fs::OpenOptions::new()
4631        .create_new(true)
4632        .write(true)
4633        .open(path)?;
4634    Ok(f.write_all(serde_json::to_string_pretty(&*EXAMPLE_CONFIG)?.as_bytes())?)
4635}
4636
4637struct HintLayerInfo {
4638    layer_id: usize,
4639    system_cpu_util_below: Option<f64>,
4640    dsq_insert_below: Option<f64>,
4641}
4642
4643fn verify_layer_specs(specs: &[LayerSpec]) -> Result<HashMap<u64, HintLayerInfo>> {
4644    let mut hint_to_layer_map = HashMap::<u64, (usize, String, Option<f64>, Option<f64>)>::new();
4645
4646    let nr_specs = specs.len();
4647    if nr_specs == 0 {
4648        bail!("No layer spec");
4649    }
4650    if nr_specs > MAX_LAYERS {
4651        bail!("Too many layer specs");
4652    }
4653
4654    for (idx, spec) in specs.iter().enumerate() {
4655        if idx < nr_specs - 1 {
4656            if spec.matches.is_empty() {
4657                bail!("Non-terminal spec {:?} has NULL matches", spec.name);
4658            }
4659        } else if spec.matches.len() != 1 || !spec.matches[0].is_empty() {
4660            bail!("Terminal spec {:?} must have an empty match", spec.name);
4661        }
4662
4663        if spec.matches.len() > MAX_LAYER_MATCH_ORS {
4664            bail!(
4665                "Spec {:?} has too many ({}) OR match blocks",
4666                spec.name,
4667                spec.matches.len()
4668            );
4669        }
4670
4671        for (ands_idx, ands) in spec.matches.iter().enumerate() {
4672            if ands.len() > NR_LAYER_MATCH_KINDS {
4673                bail!(
4674                    "Spec {:?}'s {}th OR block has too many ({}) match conditions",
4675                    spec.name,
4676                    ands_idx,
4677                    ands.len()
4678                );
4679            }
4680            let mut hint_equals_cnt = 0;
4681            let mut system_cpu_util_below_cnt = 0;
4682            let mut dsq_insert_below_cnt = 0;
4683            let mut hint_value: Option<u64> = None;
4684            let mut system_cpu_util_threshold: Option<f64> = None;
4685            let mut dsq_insert_threshold: Option<f64> = None;
4686            for one in ands.iter() {
4687                match one {
4688                    LayerMatch::CgroupPrefix(prefix) => {
4689                        if prefix.len() > MAX_PATH {
4690                            bail!("Spec {:?} has too long a cgroup prefix", spec.name);
4691                        }
4692                    }
4693                    LayerMatch::CgroupSuffix(suffix) => {
4694                        if suffix.len() > MAX_PATH {
4695                            bail!("Spec {:?} has too long a cgroup suffix", spec.name);
4696                        }
4697                    }
4698                    LayerMatch::CgroupContains(substr) => {
4699                        if substr.len() > MAX_PATH {
4700                            bail!("Spec {:?} has too long a cgroup substr", spec.name);
4701                        }
4702                    }
4703                    LayerMatch::CommPrefix(prefix) => {
4704                        if prefix.len() > MAX_COMM {
4705                            bail!("Spec {:?} has too long a comm prefix", spec.name);
4706                        }
4707                    }
4708                    LayerMatch::PcommPrefix(prefix) => {
4709                        if prefix.len() > MAX_COMM {
4710                            bail!("Spec {:?} has too long a process name prefix", spec.name);
4711                        }
4712                    }
4713                    LayerMatch::SystemCpuUtilBelow(threshold) => {
4714                        if *threshold < 0.0 || *threshold > 1.0 {
4715                            bail!(
4716                                "Spec {:?} has SystemCpuUtilBelow threshold outside the range [0.0, 1.0]",
4717                                spec.name
4718                            );
4719                        }
4720                        system_cpu_util_threshold = Some(*threshold);
4721                        system_cpu_util_below_cnt += 1;
4722                    }
4723                    LayerMatch::DsqInsertBelow(threshold) => {
4724                        if *threshold < 0.0 || *threshold > 1.0 {
4725                            bail!(
4726                                "Spec {:?} has DsqInsertBelow threshold outside the range [0.0, 1.0]",
4727                                spec.name
4728                            );
4729                        }
4730                        dsq_insert_threshold = Some(*threshold);
4731                        dsq_insert_below_cnt += 1;
4732                    }
4733                    LayerMatch::HintEquals(hint) => {
4734                        if *hint > 1024 {
4735                            bail!(
4736                                "Spec {:?} has hint value outside the range [0, 1024]",
4737                                spec.name
4738                            );
4739                        }
4740                        hint_value = Some(*hint);
4741                        hint_equals_cnt += 1;
4742                    }
4743                    _ => {}
4744                }
4745            }
4746            if hint_equals_cnt > 1 {
4747                bail!("Only 1 HintEquals match permitted per AND block");
4748            }
4749            let high_freq_matcher_cnt = system_cpu_util_below_cnt + dsq_insert_below_cnt;
4750            if high_freq_matcher_cnt > 0 {
4751                if hint_equals_cnt != 1 {
4752                    bail!(
4753                        "High-frequency matchers (SystemCpuUtilBelow, DsqInsertBelow) must be used with one HintEquals"
4754                    );
4755                }
4756                if system_cpu_util_below_cnt > 1 {
4757                    bail!("Only 1 SystemCpuUtilBelow match permitted per AND block");
4758                }
4759                if dsq_insert_below_cnt > 1 {
4760                    bail!("Only 1 DsqInsertBelow match permitted per AND block");
4761                }
4762                if ands.len() != hint_equals_cnt + system_cpu_util_below_cnt + dsq_insert_below_cnt
4763                {
4764                    bail!(
4765                        "High-frequency matchers must be used only with HintEquals (no other matchers)"
4766                    );
4767                }
4768            } else if hint_equals_cnt == 1 && ands.len() != 1 {
4769                bail!("HintEquals match cannot be in conjunction with other matches");
4770            }
4771
4772            // Insert hint into map if present
4773            if let Some(hint) = hint_value {
4774                if let Some((layer_id, name, _, _)) = hint_to_layer_map.get(&hint) {
4775                    if *layer_id != idx {
4776                        bail!(
4777                            "Spec {:?} has hint value ({}) that is already mapped to Spec {:?}",
4778                            spec.name,
4779                            hint,
4780                            name
4781                        );
4782                    }
4783                } else {
4784                    hint_to_layer_map.insert(
4785                        hint,
4786                        (
4787                            idx,
4788                            spec.name.clone(),
4789                            system_cpu_util_threshold,
4790                            dsq_insert_threshold,
4791                        ),
4792                    );
4793                }
4794            }
4795        }
4796
4797        match spec.kind {
4798            LayerKind::Confined {
4799                cpus_range,
4800                util_range,
4801                ..
4802            }
4803            | LayerKind::Grouped {
4804                cpus_range,
4805                util_range,
4806                ..
4807            } => {
4808                if let Some((cpus_min, cpus_max)) = cpus_range
4809                    && cpus_min > cpus_max
4810                {
4811                    bail!(
4812                        "Spec {:?} has invalid cpus_range({}, {})",
4813                        spec.name,
4814                        cpus_min,
4815                        cpus_max
4816                    );
4817                }
4818                if util_range.0 >= util_range.1 {
4819                    bail!(
4820                        "Spec {:?} has invalid util_range ({}, {})",
4821                        spec.name,
4822                        util_range.0,
4823                        util_range.1
4824                    );
4825                }
4826            }
4827            _ => {}
4828        }
4829    }
4830
4831    Ok(hint_to_layer_map
4832        .into_iter()
4833        .map(|(k, v)| {
4834            (
4835                k,
4836                HintLayerInfo {
4837                    layer_id: v.0,
4838                    system_cpu_util_below: v.2,
4839                    dsq_insert_below: v.3,
4840                },
4841            )
4842        })
4843        .collect())
4844}
4845
4846fn name_suffix(cgroup: &str, len: usize) -> String {
4847    let suffixlen = std::cmp::min(len, cgroup.len());
4848    let suffixrev: String = cgroup.chars().rev().take(suffixlen).collect();
4849
4850    suffixrev.chars().rev().collect()
4851}
4852
4853fn traverse_sysfs(dir: &Path) -> Result<Vec<PathBuf>> {
4854    let mut paths = vec![];
4855
4856    if !dir.is_dir() {
4857        panic!("path {:?} does not correspond to directory", dir);
4858    }
4859
4860    let direntries = fs::read_dir(dir)?;
4861
4862    for entry in direntries {
4863        let path = entry?.path();
4864        if path.is_dir() {
4865            paths.append(&mut traverse_sysfs(&path)?);
4866            paths.push(path);
4867        }
4868    }
4869
4870    Ok(paths)
4871}
4872
4873fn find_cpumask(cgroup: &str) -> Cpumask {
4874    let mut path = String::from(cgroup);
4875    path.push_str("/cpuset.cpus.effective");
4876
4877    let description = fs::read_to_string(&mut path).unwrap();
4878
4879    Cpumask::from_cpulist(&description).unwrap()
4880}
4881
4882fn expand_template(rule: &LayerMatch) -> Result<Vec<(LayerMatch, Cpumask)>> {
4883    match rule {
4884        LayerMatch::CgroupSuffix(suffix) => Ok(traverse_sysfs(Path::new("/sys/fs/cgroup"))?
4885            .into_iter()
4886            .map(|cgroup| String::from(cgroup.to_str().expect("could not parse cgroup path")))
4887            .filter(|cgroup| cgroup.ends_with(suffix))
4888            .map(|cgroup| {
4889                (
4890                    {
4891                        let mut slashterminated = cgroup.clone();
4892                        slashterminated.push('/');
4893                        LayerMatch::CgroupSuffix(name_suffix(&slashterminated, 64))
4894                    },
4895                    find_cpumask(&cgroup),
4896                )
4897            })
4898            .collect()),
4899        LayerMatch::CgroupRegex(expr) => Ok(traverse_sysfs(Path::new("/sys/fs/cgroup"))?
4900            .into_iter()
4901            .map(|cgroup| String::from(cgroup.to_str().expect("could not parse cgroup path")))
4902            .filter(|cgroup| {
4903                let re = Regex::new(expr).unwrap();
4904                re.is_match(cgroup)
4905            })
4906            .map(|cgroup| {
4907                (
4908                    // Here we convert the regex match into a suffix match because we still need to
4909                    // do the matching on the bpf side and doing a regex match in bpf isn't
4910                    // easily done.
4911                    {
4912                        let mut slashterminated = cgroup.clone();
4913                        slashterminated.push('/');
4914                        LayerMatch::CgroupSuffix(name_suffix(&slashterminated, 64))
4915                    },
4916                    find_cpumask(&cgroup),
4917                )
4918            })
4919            .collect()),
4920        _ => panic!("Unimplemented template enum {:?}", rule),
4921    }
4922}
4923
4924fn create_perf_fds(skel: &mut BpfSkel, event: u64) -> Result<()> {
4925    let mut attr = perf::bindings::perf_event_attr {
4926        size: std::mem::size_of::<perf::bindings::perf_event_attr>() as u32,
4927        type_: perf::bindings::PERF_TYPE_RAW,
4928        config: event,
4929        sample_type: 0u64,
4930        ..Default::default()
4931    };
4932    attr.__bindgen_anon_1.sample_period = 0u64;
4933    attr.set_disabled(0);
4934
4935    let perf_events_map = &skel.maps.scx_pmu_map;
4936    let map_fd = unsafe { libbpf_sys::bpf_map__fd(perf_events_map.as_libbpf_object().as_ptr()) };
4937
4938    let mut failures = 0u64;
4939
4940    for cpu in 0..*NR_CPUS_POSSIBLE {
4941        let fd = unsafe { perf::perf_event_open(&mut attr as *mut _, -1, cpu as i32, -1, 0) };
4942        if fd < 0 {
4943            failures += 1;
4944            trace!(
4945                "perf_event_open failed cpu={cpu} errno={}",
4946                std::io::Error::last_os_error()
4947            );
4948            continue;
4949        }
4950
4951        let key = cpu as u32;
4952        let val = fd as u32;
4953        let ret = unsafe {
4954            libbpf_sys::bpf_map_update_elem(
4955                map_fd,
4956                &key as *const _ as *const _,
4957                &val as *const _ as *const _,
4958                0,
4959            )
4960        };
4961        if ret != 0 {
4962            trace!("bpf_map_update_elem failed cpu={cpu} fd={fd} ret={ret}");
4963        } else {
4964            trace!("mapped cpu={cpu} -> fd={fd}");
4965        }
4966    }
4967
4968    if failures > 0 {
4969        println!("membw tracking: failed to install {failures} counters");
4970        // Keep going, do not fail the scheduler for this
4971    }
4972
4973    Ok(())
4974}
4975
4976// Set up the counters
4977fn setup_membw_tracking(skel: &mut OpenBpfSkel) -> Result<u64> {
4978    let pmumanager = PMUManager::new()?;
4979    let codename = &pmumanager.codename as &str;
4980
4981    let pmuspec = match codename {
4982        "amdzen1" | "amdzen2" | "amdzen3" => {
4983            trace!("found AMD codename {codename}");
4984            pmumanager.pmus.get("ls_any_fills_from_sys.mem_io_local")
4985        }
4986        "amdzen4" | "amdzen5" => {
4987            trace!("found AMD codename {codename}");
4988            pmumanager.pmus.get("ls_any_fills_from_sys.dram_io_all")
4989        }
4990
4991        "haswell" | "broadwell" | "broadwellde" | "broadwellx" | "skylake" | "skylakex"
4992        | "cascadelakex" | "arrowlake" | "meteorlake" | "sapphirerapids" | "emeraldrapids"
4993        | "graniterapids" => {
4994            trace!("found Intel codename {codename}");
4995            pmumanager.pmus.get("LONGEST_LAT_CACHE.MISS")
4996        }
4997
4998        _ => {
4999            trace!("found unknown codename {codename}");
5000            None
5001        }
5002    };
5003
5004    let spec = pmuspec.ok_or("not_found").unwrap();
5005    let config = (spec.umask << 8) | spec.event[0];
5006
5007    // Install the counter in the BPF map
5008    skel.maps.rodata_data.as_mut().unwrap().membw_event = config;
5009
5010    Ok(config)
5011}
5012
5013#[clap_main::clap_main]
5014fn main(opts: Opts) -> Result<()> {
5015    if opts.version {
5016        println!(
5017            "scx_layered {}",
5018            build_id::full_version(env!("CARGO_PKG_VERSION"))
5019        );
5020        return Ok(());
5021    }
5022
5023    if opts.help_stats {
5024        stats::server_data().describe_meta(&mut std::io::stdout(), None)?;
5025        return Ok(());
5026    }
5027
5028    let env_filter = EnvFilter::try_from_default_env()
5029        .or_else(|_| match EnvFilter::try_new(&opts.log_level) {
5030            Ok(filter) => Ok(filter),
5031            Err(e) => {
5032                eprintln!(
5033                    "invalid log envvar: {}, using info, err is: {}",
5034                    opts.log_level, e
5035                );
5036                EnvFilter::try_new("info")
5037            }
5038        })
5039        .unwrap_or_else(|_| EnvFilter::new("info"));
5040
5041    match tracing_subscriber::fmt()
5042        .with_env_filter(env_filter)
5043        .with_target(true)
5044        .with_thread_ids(true)
5045        .with_file(true)
5046        .with_line_number(true)
5047        .try_init()
5048    {
5049        Ok(()) => {}
5050        Err(e) => eprintln!("failed to init logger: {}", e),
5051    }
5052
5053    if opts.verbose > 0 {
5054        warn!("Setting verbose via -v is deprecated and will be an error in future releases.");
5055    }
5056
5057    if opts.no_load_frac_limit {
5058        warn!("--no-load-frac-limit is deprecated and noop");
5059    }
5060    if opts.layer_preempt_weight_disable != 0.0 {
5061        warn!("--layer-preempt-weight-disable is deprecated and noop");
5062    }
5063    if opts.layer_growth_weight_disable != 0.0 {
5064        warn!("--layer-growth-weight-disable is deprecated and noop");
5065    }
5066    if opts.local_llc_iteration {
5067        warn!("--local_llc_iteration is deprecated and noop");
5068    }
5069
5070    debug!("opts={:?}", &opts);
5071
5072    if let Some(run_id) = opts.run_id {
5073        info!("scx_layered run_id: {}", run_id);
5074    }
5075
5076    let shutdown = Arc::new(AtomicBool::new(false));
5077    let shutdown_clone = shutdown.clone();
5078    ctrlc::set_handler(move || {
5079        shutdown_clone.store(true, Ordering::Relaxed);
5080    })
5081    .context("Error setting Ctrl-C handler")?;
5082
5083    if let Some(intv) = opts.monitor.or(opts.stats) {
5084        let shutdown_copy = shutdown.clone();
5085        let stats_columns = opts.stats_columns;
5086        let stats_no_llc = opts.stats_no_llc;
5087        let jh = std::thread::spawn(move || {
5088            match stats::monitor(
5089                Duration::from_secs_f64(intv),
5090                shutdown_copy,
5091                stats_columns,
5092                stats_no_llc,
5093            ) {
5094                Ok(_) => {
5095                    debug!("stats monitor thread finished successfully")
5096                }
5097                Err(error_object) => {
5098                    warn!(
5099                        "stats monitor thread finished because of an error {}",
5100                        error_object
5101                    )
5102                }
5103            }
5104        });
5105        if opts.monitor.is_some() {
5106            let _ = jh.join();
5107            return Ok(());
5108        }
5109    }
5110
5111    if let Some(path) = &opts.example {
5112        write_example_file(path)?;
5113        return Ok(());
5114    }
5115
5116    let mut layer_config = match opts.run_example {
5117        true => EXAMPLE_CONFIG.clone(),
5118        false => LayerConfig { specs: vec![] },
5119    };
5120
5121    for (idx, input) in opts.specs.iter().enumerate() {
5122        let specs = LayerSpec::parse(input)
5123            .context(format!("Failed to parse specs[{}] ({:?})", idx, input))?;
5124
5125        for spec in specs {
5126            match spec.template {
5127                Some(ref rule) => {
5128                    let matches = expand_template(rule)?;
5129                    // in the absence of matching cgroups, have template layers
5130                    // behave as non-template layers do.
5131                    if matches.is_empty() {
5132                        layer_config.specs.push(spec);
5133                    } else {
5134                        for (mt, mask) in matches {
5135                            let mut genspec = spec.clone();
5136
5137                            genspec.cpuset = Some(mask);
5138
5139                            // Push the new "and" rule into each "or" term.
5140                            for orterm in &mut genspec.matches {
5141                                orterm.push(mt.clone());
5142                            }
5143
5144                            match &mt {
5145                                LayerMatch::CgroupSuffix(cgroup) => genspec.name.push_str(cgroup),
5146                                _ => bail!("Template match has unexpected type"),
5147                            }
5148
5149                            // Push the generated layer into the config
5150                            layer_config.specs.push(genspec);
5151                        }
5152                    }
5153                }
5154
5155                None => {
5156                    layer_config.specs.push(spec);
5157                }
5158            }
5159        }
5160    }
5161
5162    for spec in layer_config.specs.iter_mut() {
5163        let common = spec.kind.common_mut();
5164
5165        if common.slice_us == 0 {
5166            common.slice_us = opts.slice_us;
5167        }
5168
5169        if common.weight == 0 {
5170            common.weight = DEFAULT_LAYER_WEIGHT;
5171        }
5172        common.weight = common.weight.clamp(MIN_LAYER_WEIGHT, MAX_LAYER_WEIGHT);
5173
5174        if common.preempt {
5175            if common.disallow_open_after_us.is_some() {
5176                warn!(
5177                    "Preempt layer {} has non-null disallow_open_after_us, ignored",
5178                    &spec.name
5179                );
5180            }
5181            if common.disallow_preempt_after_us.is_some() {
5182                warn!(
5183                    "Preempt layer {} has non-null disallow_preempt_after_us, ignored",
5184                    &spec.name
5185                );
5186            }
5187            common.disallow_open_after_us = Some(u64::MAX);
5188            common.disallow_preempt_after_us = Some(u64::MAX);
5189        } else {
5190            if common.disallow_open_after_us.is_none() {
5191                common.disallow_open_after_us = Some(*DFL_DISALLOW_OPEN_AFTER_US);
5192            }
5193
5194            if common.disallow_preempt_after_us.is_none() {
5195                common.disallow_preempt_after_us = Some(*DFL_DISALLOW_PREEMPT_AFTER_US);
5196            }
5197        }
5198
5199        if common.idle_smt.is_some() {
5200            warn!("Layer {} has deprecated flag \"idle_smt\"", &spec.name);
5201        }
5202
5203        if common.allow_node_aligned.is_some() {
5204            warn!(
5205                "Layer {} has deprecated flag \"allow_node_aligned\", node-aligned tasks are now always dispatched on layer DSQs",
5206                &spec.name
5207            );
5208        }
5209    }
5210
5211    let membw_required = layer_config.specs.iter().any(|spec| match spec.kind {
5212        LayerKind::Confined { membw_gb, .. } | LayerKind::Grouped { membw_gb, .. } => {
5213            membw_gb.is_some()
5214        }
5215        LayerKind::Open { .. } => false,
5216    });
5217
5218    if opts.print_and_exit {
5219        println!("specs={}", serde_json::to_string_pretty(&layer_config)?);
5220        return Ok(());
5221    }
5222
5223    debug!("specs={}", serde_json::to_string_pretty(&layer_config)?);
5224    let hint_to_layer_map = verify_layer_specs(&layer_config.specs)?;
5225
5226    let mut open_object = MaybeUninit::uninit();
5227    loop {
5228        let mut sched = Scheduler::init(
5229            &opts,
5230            &layer_config.specs,
5231            &mut open_object,
5232            &hint_to_layer_map,
5233            membw_required,
5234        )?;
5235        if !sched.run(shutdown.clone())?.should_restart() {
5236            break;
5237        }
5238    }
5239
5240    Ok(())
5241}
5242
5243#[cfg(test)]
5244mod peak_util_tests {
5245    use super::*;
5246
5247    #[test]
5248    fn test_peak_pins_recent_burst() {
5249        let elapsed = Duration::from_millis(100);
5250        let half_life = Duration::from_secs(1);
5251        let mut peak = 0.0;
5252
5253        for _ in 0..5 {
5254            peak = update_peak_util(peak, 0.1, half_life, elapsed);
5255        }
5256
5257        peak = update_peak_util(peak, 2.0, half_life, elapsed);
5258
5259        for _ in 0..5 {
5260            peak = update_peak_util(peak, 0.1, half_life, elapsed);
5261        }
5262
5263        assert!(
5264            peak > 1.3,
5265            "peak should still hold the recent burst, got {peak}"
5266        );
5267    }
5268
5269    #[test]
5270    fn test_peak_decays_after_quiet_period() {
5271        let elapsed = Duration::from_millis(100);
5272        let half_life = Duration::from_secs(1);
5273        let mut peak = 2.0;
5274
5275        for _ in 0..50 {
5276            peak = update_peak_util(peak, 0.0, half_life, elapsed);
5277        }
5278
5279        assert!(
5280            peak < 0.1,
5281            "peak should decay after a long quiet period, got {peak}"
5282        );
5283    }
5284
5285    #[test]
5286    fn test_peak_matches_steady_input() {
5287        let elapsed = Duration::from_millis(100);
5288        let half_life = Duration::from_secs(1);
5289        let mut peak = 0.0;
5290
5291        for _ in 0..30 {
5292            peak = update_peak_util(peak, 1.0, half_life, elapsed);
5293        }
5294
5295        assert!(
5296            (peak - 1.0).abs() < 1e-9,
5297            "peak should match steady-state input, got {peak}"
5298        );
5299    }
5300
5301    #[test]
5302    fn test_peak_only_changes_shrink_target() {
5303        let util_range = (0.7, 0.9);
5304        let (low_no_peak, high_no_peak) = calc_peak_aware_target_range(1.8, 1.8, util_range);
5305        let (low_with_peak, high_with_peak) = calc_peak_aware_target_range(1.8, 3.1, util_range);
5306
5307        assert_eq!(low_no_peak, 2);
5308        assert_eq!(
5309            low_with_peak, low_no_peak,
5310            "peak should not affect growth target"
5311        );
5312        assert_eq!(high_no_peak, 2);
5313        assert_eq!(
5314            high_with_peak, 4,
5315            "peak should only widen the shrink boundary"
5316        );
5317    }
5318}
5319
5320#[cfg(test)]
5321mod xnuma_tests {
5322    use super::*;
5323
5324    // Default thresholds for tests
5325    const THRESH: (f64, f64) = (0.6, 0.7);
5326    const DELTA: (f64, f64) = (0.2, 0.3);
5327
5328    // =====================================================================
5329    // xnuma_check_active tests — per-(layer, node) two-threshold
5330    // =====================================================================
5331
5332    #[test]
5333    fn test_activation_below_threshold() {
5334        // Both nodes below threshold — both closed
5335        let duty = vec![40.0, 40.0];
5336        let allocs = vec![96, 96];
5337        let gd = vec![true, true];
5338        let cur = vec![false, false];
5339        let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5340        // ratio = 0.42 < 0.6 (low) → deactivate
5341        assert!(!result[0]);
5342        assert!(!result[1]);
5343    }
5344
5345    #[test]
5346    fn test_activation_above_all_thresholds() {
5347        // N0 overloaded + imbalanced + growth denied → open
5348        let duty = vec![90.0, 20.0];
5349        let allocs = vec![96, 96];
5350        let gd = vec![true, false];
5351        let cur = vec![false, false];
5352        let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5353        // N0: load=0.94>0.7, eq=110/192=0.573, surplus=90-55=35, ratio=0.365>0.3, gd=true → open
5354        assert!(result[0]);
5355        // N1: load=0.21<0.6 → closed
5356        assert!(!result[1]);
5357    }
5358
5359    #[test]
5360    fn test_activation_requires_growth_denied() {
5361        // N0 overloaded + imbalanced but growth NOT denied → closed
5362        let duty = vec![90.0, 20.0];
5363        let allocs = vec![96, 96];
5364        let gd = vec![false, false]; // growth succeeded
5365        let cur = vec![false, false];
5366        let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5367        assert!(!result[0]); // !growth_denied → deactivate
5368    }
5369
5370    #[test]
5371    fn test_symmetric_high_load_stays_closed() {
5372        // Both nodes above threshold but balanced (delta=0) → closed
5373        let duty = vec![80.0, 80.0];
5374        let allocs = vec![96, 96];
5375        let gd = vec![true, true];
5376        let cur = vec![false, false];
5377        let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5378        // load=0.83>0.7 but surplus=0, surplus_ratio=0 < 0.3 → no activation
5379        assert!(!result[0]);
5380        assert!(!result[1]);
5381    }
5382
5383    #[test]
5384    fn test_hysteresis_stays_active() {
5385        // N0 was active, now between thresholds → stays active
5386        let duty = vec![75.0, 20.0];
5387        let allocs = vec![96, 96];
5388        let gd = vec![true, false];
5389        let cur = vec![true, false]; // N0 was active
5390        let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5391        // N0: load=0.78 > 0.6(lo) and < 0.7(hi), surplus=(75-20)/2=27.5, ratio=0.286 > 0.2(lo)
5392        // gd=true. activate: 0.78<0.7 NO. deactivate: 0.78>0.6 NO, 0.286>0.2 NO, gd=true NO
5393        // → hysteresis preserves true
5394        assert!(result[0]);
5395    }
5396
5397    #[test]
5398    fn test_hysteresis_stays_inactive() {
5399        // N0 was inactive, in hysteresis band → stays inactive
5400        let duty = vec![75.0, 20.0];
5401        let allocs = vec![96, 96];
5402        let gd = vec![true, false];
5403        let cur = vec![false, false]; // N0 was inactive
5404        let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5405        // Same conditions but currently inactive → stays inactive
5406        assert!(!result[0]);
5407    }
5408
5409    #[test]
5410    fn test_deactivation_load_drops() {
5411        // N0 was active, load drops below low → closes
5412        let duty = vec![50.0, 50.0];
5413        let allocs = vec![96, 96];
5414        let gd = vec![true, true];
5415        let cur = vec![true, false];
5416        let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5417        // N0: load=0.52 < 0.6(lo) → deactivate
5418        assert!(!result[0]);
5419    }
5420
5421    #[test]
5422    fn test_deactivation_growth_succeeds() {
5423        // N0 was active, growth now succeeds → closes
5424        let duty = vec![90.0, 20.0];
5425        let allocs = vec![96, 96];
5426        let gd = vec![false, false]; // growth succeeded
5427        let cur = vec![true, false];
5428        let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5429        // !growth_denied → deactivate regardless of load/delta
5430        assert!(!result[0]);
5431    }
5432
5433    #[test]
5434    fn test_zero_alloc_with_duty_and_growth_denied() {
5435        // Zero alloc but duty > 0 and growth denied → open
5436        let duty = vec![50.0, 0.0];
5437        let allocs = vec![0, 96];
5438        let gd = vec![true, false];
5439        let cur = vec![false, false];
5440        let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5441        assert!(result[0]); // zero alloc + duty + gd → source
5442    }
5443
5444    #[test]
5445    fn test_all_zero_alloc() {
5446        let duty = vec![0.0, 0.0];
5447        let allocs = vec![0, 0];
5448        let gd = vec![true, true];
5449        let cur = vec![true, true];
5450        let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5451        assert!(!result[0]);
5452        assert!(!result[1]);
5453    }
5454
5455    #[test]
5456    fn test_three_nodes_mixed() {
5457        // 3 nodes: N0 overloaded+imbalanced+gd, N1 moderate, N2 low
5458        let duty = vec![90.0, 40.0, 10.0];
5459        let allocs = vec![96, 96, 96];
5460        let gd = vec![true, true, false];
5461        let cur = vec![false, false, false];
5462        let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5463        // eq_ratio = 140/288 ≈ 0.486
5464        // N0: load=0.94>0.7, surplus=90-46.7=43.3, ratio=0.451>0.3, gd=true → open
5465        assert!(result[0]);
5466        // N1: load=0.42<0.6 → deactivate
5467        assert!(!result[1]);
5468        // N2: load=0.10<0.6 → deactivate
5469        assert!(!result[2]);
5470    }
5471
5472    #[test]
5473    fn test_per_node_independence() {
5474        // N0 active, N1 deactivates independently
5475        let duty = vec![90.0, 50.0];
5476        let allocs = vec![96, 96];
5477        let gd = vec![true, true];
5478        let cur = vec![true, true];
5479        let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5480        // N0: load=0.94>0.7, eq=140/192=0.729, surplus=90-70=20, ratio=0.208>0.2(lo)
5481        //   activate: 0.94>0.7 YES, 0.208>0.3 NO → no activate
5482        //   deactivate: 0.94>0.6 NO, 0.208>0.2 NO, gd=true NO → no deactivate → preserve true
5483        assert!(result[0]);
5484        // N1: load=0.52<0.6 → deactivate
5485        assert!(!result[1]);
5486    }
5487
5488    // =====================================================================
5489    // xnuma_compute_rates tests — basic
5490    // =====================================================================
5491
5492    #[test]
5493    fn test_rates_balanced_load() {
5494        // Equal load per CPU on both nodes — no migration needed
5495        let duty = vec![48.0, 48.0];
5496        let allocs = vec![96, 96];
5497        let result = xnuma_compute_rates(&duty, &allocs);
5498
5499        // No surplus/deficit → all rates zero
5500        for src in 0..2 {
5501            for dst in 0..2 {
5502                assert_eq!(result.rates[src][dst], 0);
5503            }
5504        }
5505    }
5506
5507    #[test]
5508    fn test_rates_one_overloaded() {
5509        // N0 has 80 CPUs of duty, N1 has 40. Both have 96 CPUs.
5510        // eq_ratio = 120/192 = 0.625
5511        // N0 expected = 0.625 * 96 = 60, surplus = 80 - 60 = 20
5512        // N1 expected = 0.625 * 96 = 60, deficit = 60 - 40 = 20
5513        let duty = vec![80.0, 40.0];
5514        let allocs = vec![96, 96];
5515        let result = xnuma_compute_rates(&duty, &allocs);
5516
5517        // rate[0][1] should be positive (migrate from N0 to N1)
5518        assert!(result.rates[0][1] > 0);
5519        // rate[1][0] should be 0 (N1 has no surplus)
5520        assert_eq!(result.rates[1][0], 0);
5521        // Self-rates always 0
5522        assert_eq!(result.rates[0][0], 0);
5523        assert_eq!(result.rates[1][1], 0);
5524
5525        // Verify the rate magnitude: migration = 20.0 * DAMPEN, scaled
5526        let expected_rate = (20.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64;
5527        assert_eq!(result.rates[0][1], expected_rate);
5528    }
5529
5530    #[test]
5531    fn test_rates_asymmetric_allocation() {
5532        // N0 has 48 CPUs, N1 has 144 CPUs. Duty 48 each.
5533        // eq_ratio = 96/192 = 0.5
5534        // N0 expected = 0.5 * 48 = 24, surplus = 48 - 24 = 24
5535        // N1 expected = 0.5 * 144 = 72, deficit = 72 - 48 = 24
5536        let duty = vec![48.0, 48.0];
5537        let allocs = vec![48, 144];
5538        let result = xnuma_compute_rates(&duty, &allocs);
5539
5540        let expected_rate = (24.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64;
5541        assert_eq!(result.rates[0][1], expected_rate);
5542        assert_eq!(result.rates[1][0], 0);
5543    }
5544
5545    // =====================================================================
5546    // xnuma_compute_rates tests — multi-node
5547    // =====================================================================
5548
5549    #[test]
5550    fn test_rates_three_nodes_one_source() {
5551        // N0 overloaded, N1 and N2 are deficit
5552        // allocs: 96 each, duty: N0=120, N1=30, N2=30
5553        // eq_ratio = 180/288 = 0.625
5554        // N0 expected = 60, surplus = 60
5555        // N1 expected = 60, deficit = 30
5556        // N2 expected = 60, deficit = 30
5557        let duty = vec![120.0, 30.0, 30.0];
5558        let allocs = vec![96, 96, 96];
5559        let result = xnuma_compute_rates(&duty, &allocs);
5560
5561        // Total deficit = 60. N1 gets 30/60 = 50%, N2 gets 30/60 = 50%
5562        // rate[0][1] = 60 * 30/60 * DAMPEN = 15
5563        // rate[0][2] = 60 * 30/60 * DAMPEN = 15
5564        let rate_01 = (30.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64;
5565        let rate_02 = (30.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64;
5566        assert_eq!(result.rates[0][1], rate_01);
5567        assert_eq!(result.rates[0][2], rate_02);
5568
5569        // No reverse flow
5570        assert_eq!(result.rates[1][0], 0);
5571        assert_eq!(result.rates[2][0], 0);
5572        assert_eq!(result.rates[1][2], 0);
5573        assert_eq!(result.rates[2][1], 0);
5574    }
5575
5576    #[test]
5577    fn test_rates_three_nodes_unequal_deficit() {
5578        // N0 overloaded. N1 slight deficit, N2 large deficit.
5579        // allocs: 96 each, duty: N0=120, N1=50, N2=10
5580        // eq_ratio = 180/288 = 0.625
5581        // N0: expected=60, surplus=60
5582        // N1: expected=60, deficit=10
5583        // N2: expected=60, deficit=50
5584        let duty = vec![120.0, 50.0, 10.0];
5585        let allocs = vec![96, 96, 96];
5586        let result = xnuma_compute_rates(&duty, &allocs);
5587
5588        // Total deficit = 60. N1 share = 10/60, N2 share = 50/60
5589        // rate[0][1] = 60 * 10/60 * DAMPEN = 5
5590        // rate[0][2] = 60 * 50/60 * DAMPEN = 25
5591        let rate_01 = (10.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64;
5592        let rate_02 = (50.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64;
5593        assert_eq!(result.rates[0][1], rate_01);
5594        assert_eq!(result.rates[0][2], rate_02);
5595    }
5596
5597    #[test]
5598    fn test_rates_two_sources_one_sink() {
5599        // N0 and N1 overloaded, N2 deficit
5600        // allocs: 96 each, duty: N0=80, N1=70, N2=30
5601        // total = 180, eq_ratio = 180/288 = 0.625
5602        // N0: expected=60, surplus=20
5603        // N1: expected=60, surplus=10
5604        // N2: expected=60, deficit=30
5605        let duty = vec![80.0, 70.0, 30.0];
5606        let allocs = vec![96, 96, 96];
5607        let result = xnuma_compute_rates(&duty, &allocs);
5608
5609        // Total deficit = 30. Only N2 is deficit, so deficit share = 1.0
5610        // rate[0][2] = 20 * 1.0 * DAMPEN = 10
5611        // rate[1][2] = 10 * 1.0 * DAMPEN = 5
5612        let rate_02 = (20.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64;
5613        let rate_12 = (10.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64;
5614        assert_eq!(result.rates[0][2], rate_02);
5615        assert_eq!(result.rates[1][2], rate_12);
5616
5617        // No flow between sources
5618        assert_eq!(result.rates[0][1], 0);
5619        assert_eq!(result.rates[1][0], 0);
5620    }
5621
5622    // =====================================================================
5623    // Conservation invariants
5624    // =====================================================================
5625
5626    #[test]
5627    fn test_conservation_per_source_outbound() {
5628        // Each source's total outbound should equal its surplus * DAMPEN (scaled).
5629        // Total outbound from src = surplus[src] * DAMPEN * DUTY_CYCLE_SCALE
5630        let duty = vec![100.0, 30.0, 50.0, 20.0];
5631        let allocs = vec![96, 96, 96, 96];
5632        let nr = 4;
5633        let result = xnuma_compute_rates(&duty, &allocs);
5634
5635        let total_duty: f64 = duty.iter().sum();
5636        let total_alloc: f64 = allocs.iter().map(|&a| a as f64).sum();
5637        let eq_ratio = total_duty / total_alloc;
5638
5639        for src in 0..nr {
5640            let expected = eq_ratio * allocs[src] as f64;
5641            let surplus = (duty[src] - expected).max(0.0);
5642            let total_outbound: u64 = (0..nr).map(|dst| result.rates[src][dst]).sum();
5643            let expected_rate = (surplus * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64;
5644            assert_eq!(
5645                total_outbound, expected_rate,
5646                "node {} outbound mismatch",
5647                src
5648            );
5649        }
5650    }
5651
5652    #[test]
5653    fn test_conservation_surplus_equals_deficit() {
5654        // Mathematical invariant: total surplus == total deficit in water-fill
5655        let duty = [100.0, 30.0, 50.0];
5656        let allocs = [96, 96, 96];
5657
5658        let total_duty: f64 = duty.iter().sum();
5659        let total_alloc: f64 = allocs.iter().map(|&a| a as f64).sum();
5660        let eq_ratio = total_duty / total_alloc;
5661
5662        let mut total_surplus = 0.0f64;
5663        let mut total_deficit = 0.0f64;
5664        for i in 0..3 {
5665            let expected = eq_ratio * allocs[i] as f64;
5666            let delta = duty[i] - expected;
5667            if delta > 0.0 {
5668                total_surplus += delta;
5669            } else {
5670                total_deficit += -delta;
5671            }
5672        }
5673        assert!((total_surplus - total_deficit).abs() < 1e-10);
5674    }
5675
5676    #[test]
5677    fn test_self_rates_always_zero() {
5678        let duty = vec![100.0, 30.0, 50.0];
5679        let allocs = vec![96, 96, 96];
5680        let result = xnuma_compute_rates(&duty, &allocs);
5681
5682        for nid in 0..3 {
5683            assert_eq!(result.rates[nid][nid], 0);
5684        }
5685    }
5686
5687    #[test]
5688    fn test_deficit_nodes_have_zero_outbound() {
5689        // Deficit nodes should have zero outbound rates
5690        let duty = vec![100.0, 20.0, 30.0, 50.0];
5691        let allocs = vec![96, 96, 96, 96];
5692        let result = xnuma_compute_rates(&duty, &allocs);
5693
5694        // eq_ratio = 200/384 ≈ 0.521. N0: surplus=50, N3: surplus=0
5695        // N1: deficit=30, N2: deficit=20
5696        // Deficit nodes (outbound sum == 0) should have all zero rates
5697        for nid in 0..4 {
5698            let outbound: u64 = (0..4).map(|dst| result.rates[nid][dst]).sum();
5699            if outbound == 0 {
5700                for dst in 0..4 {
5701                    assert_eq!(
5702                        result.rates[nid][dst], 0,
5703                        "deficit node {} has non-zero rate to {}",
5704                        nid, dst
5705                    );
5706                }
5707            }
5708        }
5709    }
5710
5711    // =====================================================================
5712    // Edge cases
5713    // =====================================================================
5714
5715    #[test]
5716    fn test_rates_zero_duty_everywhere() {
5717        let duty = vec![0.0, 0.0];
5718        let allocs = vec![96, 96];
5719        let result = xnuma_compute_rates(&duty, &allocs);
5720
5721        for src in 0..2 {
5722            for dst in 0..2 {
5723                assert_eq!(result.rates[src][dst], 0);
5724            }
5725        }
5726    }
5727
5728    #[test]
5729    fn test_rates_zero_alloc() {
5730        let duty = vec![50.0, 50.0];
5731        let allocs = vec![0, 0];
5732        let result = xnuma_compute_rates(&duty, &allocs);
5733
5734        // total_alloc = 0 → early return with all zeros
5735        for src in 0..2 {
5736            for dst in 0..2 {
5737                assert_eq!(result.rates[src][dst], 0);
5738            }
5739        }
5740    }
5741
5742    #[test]
5743    fn test_rates_all_load_one_node() {
5744        // All duty on N0, nothing on N1
5745        // allocs: 96 each, duty: N0=96, N1=0
5746        // eq_ratio = 96/192 = 0.5
5747        // N0: expected=48, surplus=48
5748        // N1: expected=48, deficit=48
5749        let duty = vec![96.0, 0.0];
5750        let allocs = vec![96, 96];
5751        let result = xnuma_compute_rates(&duty, &allocs);
5752
5753        let expected_rate = (48.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64;
5754        assert_eq!(result.rates[0][1], expected_rate);
5755    }
5756
5757    #[test]
5758    fn test_rates_single_node() {
5759        // Single node — balanced by definition
5760        let duty = vec![96.0];
5761        let allocs = vec![96];
5762        let result = xnuma_compute_rates(&duty, &allocs);
5763
5764        assert_eq!(result.rates[0][0], 0);
5765    }
5766
5767    #[test]
5768    fn test_rates_one_node_zero_alloc() {
5769        // N0 has allocation, N1 has zero
5770        // eq_ratio = 80/96
5771        // N0: expected=80, surplus=0 → balanced
5772        // N1: expected=0, deficit=0 → zero alloc skipped effectively
5773        let duty = vec![80.0, 0.0];
5774        let allocs = vec![96, 0];
5775        let result = xnuma_compute_rates(&duty, &allocs);
5776
5777        // N0: surplus = 80 - (80/96 * 96) = 0
5778        // N1: deficit = (80/96 * 0) - 0 = 0
5779        // All zeros — nothing to migrate
5780        assert_eq!(result.rates[0][1], 0);
5781        assert_eq!(result.rates[1][0], 0);
5782    }
5783
5784    // =====================================================================
5785    // Rate magnitude and scaling
5786    // =====================================================================
5787
5788    #[test]
5789    fn test_rate_scaling() {
5790        // Verify rates are in DUTY_CYCLE_SCALE units
5791        let duty = vec![80.0, 40.0];
5792        let allocs = vec![96, 96];
5793        let result = xnuma_compute_rates(&duty, &allocs);
5794
5795        // surplus = 20, deficit = 20 → migration = 20 * DAMPEN = 10
5796        // rate = 10 * (1 << 20)
5797        let expected = (20.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64;
5798        assert_eq!(result.rates[0][1], expected);
5799        assert_eq!(expected, 10 * (1 << 20));
5800    }
5801
5802    #[test]
5803    fn test_rates_tiny_imbalance() {
5804        // Very small imbalance — should still produce non-zero rate
5805        let duty = vec![48.001, 47.999];
5806        let allocs = vec![96, 96];
5807        let result = xnuma_compute_rates(&duty, &allocs);
5808
5809        // surplus ≈ 0.001, rate ≈ 0.001 * 2^20 ≈ 1048
5810        assert!(result.rates[0][1] > 0);
5811        assert!(result.rates[0][1] < (1 << 20)); // Less than 1.0 CPU worth
5812    }
5813
5814    #[test]
5815    fn test_rates_large_values() {
5816        // Large system: 8 nodes, 96 CPUs each
5817        let duty = vec![300.0, 100.0, 50.0, 50.0, 50.0, 50.0, 50.0, 50.0];
5818        let allocs = vec![96, 96, 96, 96, 96, 96, 96, 96];
5819        let result = xnuma_compute_rates(&duty, &allocs);
5820
5821        // eq_ratio = 700/768 ≈ 0.911. N0 surplus=212.5, N1 surplus=12.5
5822        // N2-N7 deficit=37.5 each. N0 and N1 have surplus.
5823        assert!(result.rates[0][2] > 0); // N0 → N2 (deficit node)
5824        assert!(result.rates[1][2] > 0); // N1 → N2 (N1 also has surplus)
5825
5826        // Verify conservation: per-source outbound ≈ surplus * DAMPEN (within
5827        // truncation tolerance — each `as u64` can lose up to 1 per cell)
5828        let total_duty: f64 = duty.iter().sum();
5829        let total_alloc: f64 = allocs.iter().map(|&a| a as f64).sum();
5830        let eq_ratio = total_duty / total_alloc;
5831        let nr = 8;
5832        for src in 0..nr {
5833            let surplus = (duty[src] - eq_ratio * allocs[src] as f64).max(0.0);
5834            let outbound: u64 = (0..nr).map(|dst| result.rates[src][dst]).sum();
5835            let expected = (surplus * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64;
5836            let tolerance = nr as u64; // up to 1 per destination cell
5837            assert!(
5838                outbound.abs_diff(expected) <= tolerance,
5839                "node {} outbound {} vs expected {}, diff {}",
5840                src,
5841                outbound,
5842                expected,
5843                outbound.abs_diff(expected)
5844            );
5845        }
5846    }
5847
5848    // =====================================================================
5849    // Hysteresis integration
5850    // =====================================================================
5851
5852    #[test]
5853    fn test_hysteresis_cycle() {
5854        let allocs = vec![96, 96];
5855        let gd = vec![true, false]; // N0 growth denied
5856
5857        // Start: all closed, low load
5858        let active =
5859            xnuma_check_active(&[40.0, 40.0], &allocs, THRESH, DELTA, &gd, &[false, false]);
5860        assert!(!active[0]); // load 0.42 < 0.6(lo) → deactivate
5861
5862        // N0 overloaded + imbalanced + gd → opens
5863        let active =
5864            xnuma_check_active(&[90.0, 20.0], &allocs, THRESH, DELTA, &gd, &[false, false]);
5865        assert!(active[0]); // load 0.94>0.7, surplus=35, ratio=0.365>0.3, gd=true
5866
5867        // Load decreases but in hysteresis band → stays open
5868        let active = xnuma_check_active(&[75.0, 20.0], &allocs, THRESH, DELTA, &gd, &[true, false]);
5869        assert!(active[0]); // load 0.78 between 0.6 and 0.7, surplus ok, gd → preserve
5870
5871        // Load drops below low threshold → closes
5872        let active = xnuma_check_active(&[50.0, 50.0], &allocs, THRESH, DELTA, &gd, &[true, false]);
5873        assert!(!active[0]); // load 0.52 < 0.6(lo) → deactivate
5874    }
5875
5876    #[test]
5877    fn test_hysteresis_growth_toggle() {
5878        // N0 was open, then growth succeeds → closes
5879        let allocs = vec![96, 96];
5880        let gd_denied = vec![true, false];
5881        let gd_ok = vec![false, false];
5882
5883        // Active with growth denied
5884        let active = xnuma_check_active(
5885            &[90.0, 20.0],
5886            &allocs,
5887            THRESH,
5888            DELTA,
5889            &gd_denied,
5890            &[false, false],
5891        );
5892        assert!(active[0]);
5893
5894        // Growth succeeds → !gd → deactivate
5895        let active = xnuma_check_active(
5896            &[90.0, 20.0],
5897            &allocs,
5898            THRESH,
5899            DELTA,
5900            &gd_ok,
5901            &[true, false],
5902        );
5903        assert!(!active[0]); // !growth_denied → deactivate
5904    }
5905
5906    // =====================================================================
5907    // Proportional distribution to multiple sinks
5908    // =====================================================================
5909
5910    #[test]
5911    fn test_proportional_sink_distribution() {
5912        // 4 nodes: N0 source, N1-N3 sinks with different deficits
5913        // allocs: 96 each, duty: N0=180, N1=20, N2=40, N3=0
5914        // total = 240, eq_ratio = 240/384 = 0.625
5915        // N0: expected=60, surplus=120
5916        // N1: expected=60, deficit=40
5917        // N2: expected=60, deficit=20
5918        // N3: expected=60, deficit=60
5919        // total_deficit = 120
5920        let duty = vec![180.0, 20.0, 40.0, 0.0];
5921        let allocs = vec![96, 96, 96, 96];
5922        let result = xnuma_compute_rates(&duty, &allocs);
5923
5924        // rate[0][1] = 120 * 40/120 * DAMPEN = 20
5925        // rate[0][2] = 120 * 20/120 * DAMPEN = 10
5926        // rate[0][3] = 120 * 60/120 * DAMPEN = 30
5927        assert_eq!(
5928            result.rates[0][1],
5929            (40.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64
5930        );
5931        assert_eq!(
5932            result.rates[0][2],
5933            (20.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64
5934        );
5935        assert_eq!(
5936            result.rates[0][3],
5937            (60.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64
5938        );
5939
5940        // Total outbound from N0 = 20 + 10 + 30 = 60 (half of surplus, dampened)
5941        let total_from_n0: u64 = (0..4).map(|dst| result.rates[0][dst]).sum();
5942        assert_eq!(
5943            total_from_n0,
5944            (120.0 * XNUMA_RATE_DAMPEN * DUTY_CYCLE_SCALE) as u64
5945        );
5946    }
5947}
5948
5949#[cfg(test)]
5950mod util_compensation_tests {
5951    /// Compute the scale factor: s = Δt / (Δt - irq - softirq - stolen)
5952    fn compute_scale(delta_total: u64, overhead: u64) -> f64 {
5953        let available = delta_total.saturating_sub(overhead);
5954        if available > 0 {
5955            (delta_total as f64 / available as f64).clamp(1.0, 20.0)
5956        } else {
5957            1.0
5958        }
5959    }
5960
5961    /// Simulate the per-CPU-scaled aggregation that refresh() does.
5962    fn scaled_aggregate(
5963        deltas: &[Vec<u64>],
5964        scales: &[f64],
5965        nr_layers: usize,
5966        elapsed: f64,
5967    ) -> Vec<f64> {
5968        (0..nr_layers)
5969            .map(|layer| {
5970                let mut sum = 0.0f64;
5971                for (cpu, cpu_deltas) in deltas.iter().enumerate() {
5972                    sum += cpu_deltas[layer] as f64 * scales[cpu];
5973                }
5974                sum / 1_000_000_000.0 / elapsed
5975            })
5976            .collect()
5977    }
5978
5979    #[test]
5980    fn test_scale_no_overhead() {
5981        // No irq/softirq/stolen → scale = 1.0
5982        assert!((compute_scale(1000, 0) - 1.0).abs() < 0.01);
5983    }
5984
5985    #[test]
5986    fn test_scale_half_overhead() {
5987        // 50% overhead → s = 1000/500 = 2.0
5988        assert!((compute_scale(1000, 500) - 2.0).abs() < 0.01);
5989    }
5990
5991    #[test]
5992    fn test_scale_high_overhead() {
5993        // 90% overhead → s = 1000/100 = 10.0
5994        assert!((compute_scale(1000, 900) - 10.0).abs() < 0.01);
5995    }
5996
5997    #[test]
5998    fn test_scale_very_high_overhead() {
5999        // 95% overhead → s = 1000/50 = 20.0 (at clamp boundary)
6000        assert!((compute_scale(1000, 950) - 20.0).abs() < 0.01);
6001    }
6002
6003    #[test]
6004    fn test_scale_clamped_at_max() {
6005        // 98% overhead → ratio = 50.0, clamped to 20.0
6006        assert!((compute_scale(1000, 980) - 20.0).abs() < 0.01);
6007    }
6008
6009    #[test]
6010    fn test_scale_all_overhead() {
6011        // 100% overhead → available = 0, returns 1.0
6012        assert!((compute_scale(1000, 1000) - 1.0).abs() < 0.01);
6013    }
6014
6015    #[test]
6016    fn test_scale_idle_cpu() {
6017        // Idle CPU: delta_total = 0 → returns 1.0
6018        assert!((compute_scale(0, 0) - 1.0).abs() < 0.01);
6019    }
6020
6021    #[test]
6022    fn test_scale_small_overhead_mostly_idle() {
6023        // 1% softirq + 1% user + 98% idle = total 10000, overhead 100
6024        // s = 10000 / 9900 ≈ 1.01 — nearly 1.0, not pathological
6025        let s = compute_scale(10000, 100);
6026        assert!((s - 1.01).abs() < 0.01, "expected ~1.01, got {}", s);
6027    }
6028
6029    #[test]
6030    fn test_uniform_scale_matches_unscaled() {
6031        let deltas = vec![vec![1_000_000_000u64; 2]; 4];
6032        let scales = vec![1.0; 4];
6033        let result = scaled_aggregate(&deltas, &scales, 2, 1.0);
6034        assert!((result[0] - 4.0).abs() < 0.01);
6035        assert!((result[1] - 4.0).abs() < 0.01);
6036    }
6037
6038    #[test]
6039    fn test_hot_cpu_weighted_more() {
6040        // CPU 0: scale=2.0, CPU 1: scale=1.0
6041        // Layer 0: 900ms on CPU 0, 100ms on CPU 1
6042        // Compensated: 900*2.0 + 100*1.0 = 1900ms = 1.9
6043        let deltas = vec![vec![900_000_000, 0], vec![100_000_000, 0]];
6044        let scales = vec![2.0, 1.0];
6045        let result = scaled_aggregate(&deltas, &scales, 2, 1.0);
6046        assert!(
6047            (result[0] - 1.9).abs() < 0.01,
6048            "expected 1.9, got {}",
6049            result[0]
6050        );
6051    }
6052
6053    #[test]
6054    fn test_cold_cpu_weighted_less() {
6055        let deltas = vec![vec![100_000_000, 0], vec![900_000_000, 0]];
6056        let scales = vec![2.0, 1.0];
6057        let result = scaled_aggregate(&deltas, &scales, 2, 1.0);
6058        assert!(
6059            (result[0] - 1.1).abs() < 0.01,
6060            "expected 1.1, got {}",
6061            result[0]
6062        );
6063    }
6064
6065    #[test]
6066    fn test_no_usage_no_compensation() {
6067        let deltas = vec![vec![0u64; 2]; 4];
6068        let scales = vec![5.0; 4];
6069        let result = scaled_aggregate(&deltas, &scales, 2, 1.0);
6070        assert_eq!(result[0], 0.0);
6071        assert_eq!(result[1], 0.0);
6072    }
6073
6074    #[test]
6075    fn test_multilayer_independent_scaling() {
6076        let deltas = vec![
6077            vec![800_000_000, 200_000_000],
6078            vec![200_000_000, 800_000_000],
6079        ];
6080        let scales = vec![3.0, 1.0];
6081        let result = scaled_aggregate(&deltas, &scales, 2, 1.0);
6082        assert!((result[0] - 2.6).abs() < 0.01);
6083        assert!((result[1] - 1.4).abs() < 0.01);
6084    }
6085
6086    #[test]
6087    fn test_elapsed_time_normalization() {
6088        let deltas = vec![vec![500_000_000u64; 1]; 1];
6089        let scales = vec![1.0];
6090        let result = scaled_aggregate(&deltas, &scales, 1, 2.0);
6091        assert!((result[0] - 0.25).abs() < 0.01);
6092    }
6093
6094    #[test]
6095    fn test_many_cpus_mixed_scales() {
6096        let deltas = vec![vec![1_000_000_000u64; 1]; 8];
6097        let scales = vec![2.5, 1.0, 2.5, 1.0, 2.5, 1.0, 2.5, 1.0];
6098        let result = scaled_aggregate(&deltas, &scales, 1, 1.0);
6099        assert!((result[0] - 14.0).abs() < 0.01);
6100    }
6101
6102    #[test]
6103    fn test_compensated_ge_raw() {
6104        // Since scale >= 1.0 always, compensated >= raw for any input.
6105        let deltas = vec![
6106            vec![500_000_000u64; 3],
6107            vec![300_000_000; 3],
6108            vec![200_000_000; 3],
6109        ];
6110        let scales_raw = vec![1.0; 3];
6111        let scales_comp = vec![1.5, 2.0, 1.0];
6112        let raw = scaled_aggregate(&deltas, &scales_raw, 3, 1.0);
6113        let comp = scaled_aggregate(&deltas, &scales_comp, 3, 1.0);
6114        for i in 0..3 {
6115            assert!(
6116                comp[i] >= raw[i] - 0.001,
6117                "layer {}: comp {} < raw {}",
6118                i,
6119                comp[i],
6120                raw[i]
6121            );
6122        }
6123    }
6124}