1mod 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#[derive(Debug, Parser)]
470#[command(verbatim_doc_comment)]
471struct Opts {
472 #[clap(short = 'v', long, action = clap::ArgAction::Count)]
474 verbose: u8,
475
476 #[clap(short = 's', long, default_value = "20000")]
478 slice_us: u64,
479
480 #[clap(short = 'M', long, default_value = "0")]
485 max_exec_us: u64,
486
487 #[clap(short = 'i', long, default_value = "0.1")]
489 interval: f64,
490
491 #[clap(short = 'n', long, default_value = "false")]
494 no_load_frac_limit: bool,
495
496 #[clap(long, default_value = "0")]
498 exit_dump_len: u32,
499
500 #[clap(long, default_value = "info")]
503 log_level: String,
504
505 #[arg(short = 't', long, num_args = 0..=1, default_missing_value = "true", require_equals = true)]
509 disable_topology: Option<bool>,
510
511 #[clap(long)]
513 monitor_disable: bool,
514
515 #[clap(short = 'e', long)]
517 example: Option<String>,
518
519 #[clap(long, default_value = "0.0")]
523 layer_preempt_weight_disable: f64,
524
525 #[clap(long, default_value = "0.0")]
529 layer_growth_weight_disable: f64,
530
531 #[clap(long)]
533 stats: Option<f64>,
534
535 #[clap(long)]
538 monitor: Option<f64>,
539
540 #[clap(long, default_value = "95")]
542 stats_columns: usize,
543
544 #[clap(long)]
546 stats_no_llc: bool,
547
548 #[clap(long)]
550 run_example: bool,
551
552 #[clap(long)]
554 allow_partial_core: bool,
555
556 #[clap(long, default_value = "false")]
559 local_llc_iteration: bool,
560
561 #[clap(long, default_value = "10000")]
566 lo_fb_wait_us: u64,
567
568 #[clap(long, default_value = ".05")]
571 lo_fb_share: f64,
572
573 #[clap(long, default_value = "false")]
575 disable_antistall: bool,
576
577 #[clap(long, default_value = "false")]
579 enable_gpu_affinitize: bool,
580
581 #[clap(long, default_value = "900")]
584 gpu_affinitize_secs: u64,
585
586 #[clap(long, default_value = "false")]
591 enable_match_debug: bool,
592
593 #[clap(long, default_value = "3")]
595 antistall_sec: u64,
596
597 #[clap(long, default_value = "false")]
599 enable_gpu_support: bool,
600
601 #[clap(long, default_value = "3")]
609 gpu_kprobe_level: u64,
610
611 #[clap(long, default_value = "false")]
615 util_compensation: bool,
616
617 #[clap(long, default_value = "false")]
619 netdev_irq_balance: bool,
620
621 #[clap(long, default_value = "false")]
623 disable_queued_wakeup: bool,
624
625 #[clap(long, default_value = "false")]
627 disable_percpu_kthread_preempt: bool,
628
629 #[clap(long, default_value = "false")]
633 percpu_kthread_preempt_all: bool,
634
635 #[clap(short = 'V', long, action = clap::ArgAction::SetTrue)]
637 version: bool,
638
639 #[clap(long)]
641 run_id: Option<u64>,
642
643 #[clap(long)]
645 help_stats: bool,
646
647 specs: Vec<String>,
649
650 #[clap(long, default_value = "2000")]
653 layer_refresh_ms_avgruntime: u64,
654
655 #[clap(long, default_value = "")]
657 task_hint_map: String,
658
659 #[clap(long, default_value = "false")]
661 print_and_exit: bool,
662
663 #[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#[derive(Debug, Clone)]
676enum CgroupEvent {
677 Created {
678 path: String,
679 cgroup_id: u64, match_bitmap: u64, },
682 Removed {
683 path: String,
684 cgroup_id: u64, },
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>>>, }
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 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 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, 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>>, prev_layer_node_duty_raw: Vec<Vec<u64>>, layer_membws: Vec<Vec<f64>>, prev_layer_membw_agg: Vec<Vec<u64>>, cpu_busy: f64, prev_total_cpu: fb_procfs::CpuStat,
892 prev_pmu_resctrl_membw: (u64, u64), util_compensation: bool,
895 layer_utils_compensated: Vec<Vec<f64>>, prev_cpu_layer_usages: Vec<u64>, prev_per_cpu_stats: BTreeMap<u32, fb_procfs::CpuStat>,
898
899 system_cpu_util_ewma: f64, layer_dsq_insert_ewma: Vec<f64>, 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 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 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 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 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 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 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 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 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 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 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 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 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 for (node_id, node) in &topo.nodes {
1526 if nodes.contains(node_id) {
1528 for &id in node.all_cpus.keys() {
1529 allowed_cpus.set_cpu(id)?;
1530 }
1531 }
1532 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 for (node_id, node) in &topo.nodes {
1556 if nodes.contains(node_id) {
1558 for &id in node.all_cpus.keys() {
1559 allowed_cpus.set_cpu(id)?;
1560 }
1561 }
1562 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 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 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 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 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
1852struct XnumaRates {
1854 rates: Vec<Vec<u64>>,
1856}
1857
1858fn 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
1920fn 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 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 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 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 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 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 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 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 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 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 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 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 sys_order.retain(|id| !node_order.contains(id));
2498 node_order.retain(|&id| id != llc_id);
2499
2500 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 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 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 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 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 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 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 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 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 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 if opts.enable_gpu_support {
2807 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 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 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 let layered_task_hint_map_path = &opts.task_hint_map;
2945 let hint_map = &mut skel.maps.scx_layered_task_hint_map;
2946 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 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 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, };
2978 (*info_ptr).dsq_insert_below = match v.dsq_insert_below {
2979 Some(threshold) => (threshold * 10000.0) as u64,
2980 None => u64::MAX, };
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 let proc_reader = fb_procfs::ProcReader::new();
3021
3022 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 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 let struct_ops = scx_ops_attach!(skel, layered)?;
3050
3051 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 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 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 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 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 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 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 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 let units = cpus.div_ceil(au);
3253 raw_pinned[n] = units;
3254 }
3255
3256 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 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 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 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 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 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 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 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 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 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 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 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 let mut updated = false;
3600
3601 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 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 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 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 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 let cpu_targets: Vec<usize> = layer_allocs.iter().map(|a| a.total() * au).collect();
3802
3803 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 let prev_node_cpus: Vec<Vec<usize>> =
3811 self.layers.iter().map(|l| l.nr_node_cpus.clone()).collect();
3812
3813 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 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 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 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 for &(idx, _target) in &ascending {
3896 let layer = &mut self.layers[idx];
3897
3898 if layer_is_open(layer) {
3899 continue;
3900 }
3901
3902 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 cgroup_path_to_id.insert(path_str.clone(), cgroup_id);
4285
4286 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 let (sender, receiver) = crossbeam::channel::bounded::<CgroupEvent>(1024);
4305
4306 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 Self::add_recursive_watches(&mut inotify, &mut wd_to_path, Path::new("/sys/fs/cgroup"))?;
4315
4316 std::thread::spawn(move || {
4318 let mut buffer = [0; 4096];
4319 let inotify_fd = inotify.as_raw_fd();
4320 let mut cgroup_path_to_id = HashMap::<String, u64>::new();
4322
4323 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 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, };
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 continue;
4362 }
4363
4364 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 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 let cgroup_id = cgroup_path_to_id.remove(&path_str).unwrap_or(0);
4428
4429 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 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 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 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 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 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 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 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 }
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 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 {
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 }
4972
4973 Ok(())
4974}
4975
4976fn 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 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 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 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 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 const THRESH: (f64, f64) = (0.6, 0.7);
5326 const DELTA: (f64, f64) = (0.2, 0.3);
5327
5328 #[test]
5333 fn test_activation_below_threshold() {
5334 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 assert!(!result[0]);
5342 assert!(!result[1]);
5343 }
5344
5345 #[test]
5346 fn test_activation_above_all_thresholds() {
5347 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 assert!(result[0]);
5355 assert!(!result[1]);
5357 }
5358
5359 #[test]
5360 fn test_activation_requires_growth_denied() {
5361 let duty = vec![90.0, 20.0];
5363 let allocs = vec![96, 96];
5364 let gd = vec![false, false]; let cur = vec![false, false];
5366 let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5367 assert!(!result[0]); }
5369
5370 #[test]
5371 fn test_symmetric_high_load_stays_closed() {
5372 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 assert!(!result[0]);
5380 assert!(!result[1]);
5381 }
5382
5383 #[test]
5384 fn test_hysteresis_stays_active() {
5385 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]; let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5391 assert!(result[0]);
5395 }
5396
5397 #[test]
5398 fn test_hysteresis_stays_inactive() {
5399 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]; let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5405 assert!(!result[0]);
5407 }
5408
5409 #[test]
5410 fn test_deactivation_load_drops() {
5411 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 assert!(!result[0]);
5419 }
5420
5421 #[test]
5422 fn test_deactivation_growth_succeeds() {
5423 let duty = vec![90.0, 20.0];
5425 let allocs = vec![96, 96];
5426 let gd = vec![false, false]; let cur = vec![true, false];
5428 let result = xnuma_check_active(&duty, &allocs, THRESH, DELTA, &gd, &cur);
5429 assert!(!result[0]);
5431 }
5432
5433 #[test]
5434 fn test_zero_alloc_with_duty_and_growth_denied() {
5435 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]); }
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 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 assert!(result[0]);
5466 assert!(!result[1]);
5468 assert!(!result[2]);
5470 }
5471
5472 #[test]
5473 fn test_per_node_independence() {
5474 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 assert!(result[0]);
5484 assert!(!result[1]);
5486 }
5487
5488 #[test]
5493 fn test_rates_balanced_load() {
5494 let duty = vec![48.0, 48.0];
5496 let allocs = vec![96, 96];
5497 let result = xnuma_compute_rates(&duty, &allocs);
5498
5499 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 let duty = vec![80.0, 40.0];
5514 let allocs = vec![96, 96];
5515 let result = xnuma_compute_rates(&duty, &allocs);
5516
5517 assert!(result.rates[0][1] > 0);
5519 assert_eq!(result.rates[1][0], 0);
5521 assert_eq!(result.rates[0][0], 0);
5523 assert_eq!(result.rates[1][1], 0);
5524
5525 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 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 #[test]
5550 fn test_rates_three_nodes_one_source() {
5551 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 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 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 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 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 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 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 assert_eq!(result.rates[0][1], 0);
5619 assert_eq!(result.rates[1][0], 0);
5620 }
5621
5622 #[test]
5627 fn test_conservation_per_source_outbound() {
5628 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 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 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 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 #[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 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 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 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 let duty = vec![80.0, 0.0];
5774 let allocs = vec![96, 0];
5775 let result = xnuma_compute_rates(&duty, &allocs);
5776
5777 assert_eq!(result.rates[0][1], 0);
5781 assert_eq!(result.rates[1][0], 0);
5782 }
5783
5784 #[test]
5789 fn test_rate_scaling() {
5790 let duty = vec![80.0, 40.0];
5792 let allocs = vec![96, 96];
5793 let result = xnuma_compute_rates(&duty, &allocs);
5794
5795 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 let duty = vec![48.001, 47.999];
5806 let allocs = vec![96, 96];
5807 let result = xnuma_compute_rates(&duty, &allocs);
5808
5809 assert!(result.rates[0][1] > 0);
5811 assert!(result.rates[0][1] < (1 << 20)); }
5813
5814 #[test]
5815 fn test_rates_large_values() {
5816 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 assert!(result.rates[0][2] > 0); assert!(result.rates[1][2] > 0); 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; 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 #[test]
5853 fn test_hysteresis_cycle() {
5854 let allocs = vec![96, 96];
5855 let gd = vec![true, false]; let active =
5859 xnuma_check_active(&[40.0, 40.0], &allocs, THRESH, DELTA, &gd, &[false, false]);
5860 assert!(!active[0]); let active =
5864 xnuma_check_active(&[90.0, 20.0], &allocs, THRESH, DELTA, &gd, &[false, false]);
5865 assert!(active[0]); let active = xnuma_check_active(&[75.0, 20.0], &allocs, THRESH, DELTA, &gd, &[true, false]);
5869 assert!(active[0]); let active = xnuma_check_active(&[50.0, 50.0], &allocs, THRESH, DELTA, &gd, &[true, false]);
5873 assert!(!active[0]); }
5875
5876 #[test]
5877 fn test_hysteresis_growth_toggle() {
5878 let allocs = vec![96, 96];
5880 let gd_denied = vec![true, false];
5881 let gd_ok = vec![false, false];
5882
5883 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 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]); }
5905
5906 #[test]
5911 fn test_proportional_sink_distribution() {
5912 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 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 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 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 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 assert!((compute_scale(1000, 0) - 1.0).abs() < 0.01);
5983 }
5984
5985 #[test]
5986 fn test_scale_half_overhead() {
5987 assert!((compute_scale(1000, 500) - 2.0).abs() < 0.01);
5989 }
5990
5991 #[test]
5992 fn test_scale_high_overhead() {
5993 assert!((compute_scale(1000, 900) - 10.0).abs() < 0.01);
5995 }
5996
5997 #[test]
5998 fn test_scale_very_high_overhead() {
5999 assert!((compute_scale(1000, 950) - 20.0).abs() < 0.01);
6001 }
6002
6003 #[test]
6004 fn test_scale_clamped_at_max() {
6005 assert!((compute_scale(1000, 980) - 20.0).abs() < 0.01);
6007 }
6008
6009 #[test]
6010 fn test_scale_all_overhead() {
6011 assert!((compute_scale(1000, 1000) - 1.0).abs() < 0.01);
6013 }
6014
6015 #[test]
6016 fn test_scale_idle_cpu() {
6017 assert!((compute_scale(0, 0) - 1.0).abs() < 0.01);
6019 }
6020
6021 #[test]
6022 fn test_scale_small_overhead_mostly_idle() {
6023 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 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 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}