1mod bpf_skel;
9pub use bpf_skel::*;
10pub mod bpf_intf;
11pub use bpf_intf::*;
12
13mod stats;
14
15use std::mem::MaybeUninit;
16use std::sync::atomic::AtomicBool;
17use std::sync::atomic::Ordering;
18use std::sync::Arc;
19use std::time::Duration;
20
21use anyhow::bail;
22use anyhow::Context;
23use anyhow::Result;
24use clap::Parser;
25use crossbeam::channel::RecvTimeoutError;
26use libbpf_rs::skel::Skel;
27use libbpf_rs::OpenObject;
28use libbpf_rs::ProgramInput;
29use log::debug;
30use log::info;
31use log::warn;
32use scx_arena::ArenaLib;
33use scx_stats::prelude::*;
34use scx_utils::build_id;
35use scx_utils::compat;
36use scx_utils::get_primary_cpus;
37use scx_utils::libbpf_clap_opts::LibbpfOpts;
38use scx_utils::scx_ops_attach;
39use scx_utils::scx_ops_cid_load;
40use scx_utils::scx_ops_cid_open;
41use scx_utils::try_set_rlimit_infinity;
42use scx_utils::uei_exited;
43use scx_utils::uei_report;
44use scx_utils::Powermode;
45use scx_utils::Topology;
46use scx_utils::UserExitInfo;
47use scx_utils::NR_CPUS_POSSIBLE;
48use scx_utils::NR_CPU_IDS;
49use stats::Metrics;
50
51const SCHEDULER_NAME: &str = "scx_cidland";
52
53fn run_syscall_prog<T>(prog: &libbpf_rs::ProgramMut<'_>, args: &mut T) -> Result<()> {
58 let input = ProgramInput {
59 context_in: Some(unsafe {
60 std::slice::from_raw_parts_mut(args as *mut T as *mut u8, std::mem::size_of::<T>())
61 }),
62 ..Default::default()
63 };
64
65 let output = prog.test_run(input)?;
66 if output.return_value != 0 {
67 bail!(
68 "{} returned {}",
69 prog.name().to_string_lossy(),
70 output.return_value as i32
71 );
72 }
73
74 Ok(())
75}
76
77#[derive(Debug, Parser)]
93struct Opts {
94 #[clap(short = 's', long, default_value = "1000")]
96 slice_us: u64,
97
98 #[clap(short = 'l', long, default_value = "20000")]
104 slice_lag_us: u64,
105
106 #[clap(short = 'm', long, value_name = "CPU_LIST")]
121 primary_domain: Option<String>,
122
123 #[clap(long, default_value = "0")]
125 exit_dump_len: u32,
126
127 #[clap(long)]
129 stats: Option<f64>,
130
131 #[clap(long)]
134 monitor: Option<f64>,
135
136 #[clap(short = 'v', long, action = clap::ArgAction::SetTrue)]
138 verbose: bool,
139
140 #[clap(short = 'V', long, action = clap::ArgAction::SetTrue)]
142 version: bool,
143
144 #[clap(long)]
146 help_stats: bool,
147
148 #[clap(flatten, next_help_heading = "Libbpf Options")]
149 pub libbpf: LibbpfOpts,
150}
151
152fn parse_primary_domain(arg: &str) -> Result<Vec<usize>> {
155 let mode = match arg {
156 "performance" => Some(Powermode::Performance),
157 "powersave" => Some(Powermode::Powersave),
158 "turbo" => Some(Powermode::Turbo),
159 "all" => Some(Powermode::Any),
160 _ => None,
161 };
162
163 let Some(mode) = mode else {
164 return parse_cpu_list(arg);
165 };
166
167 let mut cpus = get_primary_cpus(mode).context("detecting the primary CPUs")?;
168 if cpus.is_empty() {
169 bail!("no CPU matches \"{arg}\" on this system");
170 }
171 cpus.sort_unstable();
172 cpus.dedup();
173
174 Ok(cpus)
175}
176
177fn parse_cpu_list(arg: &str) -> Result<Vec<usize>> {
179 let mut cpus = Vec::new();
180
181 for token in arg.split(',') {
182 let token = token.trim();
183
184 if token.is_empty() {
185 continue;
186 }
187
188 if let Some((start, end)) = token.split_once('-') {
189 let start: usize = start
190 .trim()
191 .parse()
192 .with_context(|| format!("invalid range start in {token:?}"))?;
193 let end: usize = end
194 .trim()
195 .parse()
196 .with_context(|| format!("invalid range end in {token:?}"))?;
197 if start > end {
198 bail!("invalid range {token:?}");
199 }
200 cpus.extend(start..=end);
201 } else {
202 cpus.push(
203 token
204 .parse()
205 .with_context(|| format!("invalid cpu id {token:?}"))?,
206 );
207 }
208 }
209
210 if cpus.is_empty() {
211 bail!("no CPU specified");
212 }
213 cpus.sort_unstable();
214 cpus.dedup();
215
216 Ok(cpus)
217}
218
219struct Scheduler<'a> {
220 _arenalib: ArenaLib,
221 skel: BpfSkel<'a>,
222 struct_ops: Option<libbpf_rs::Link>,
223 stats_server: StatsServer<(), Metrics>,
224}
225
226impl<'a> Scheduler<'a> {
227 fn init(opts: &'a Opts, open_object: &'a mut MaybeUninit<OpenObject>) -> Result<Self> {
228 try_set_rlimit_infinity();
229
230 if opts.slice_us == 0 {
231 bail!("--slice-us must be greater than 0");
232 }
233
234 let topo = Topology::new().context("detecting system topology")?;
235 info!(
236 "{} {} ({} CPUs, {} LLCs)",
237 SCHEDULER_NAME,
238 build_id::full_version(env!("CARGO_PKG_VERSION")),
239 *NR_CPUS_POSSIBLE,
240 topo.all_llcs.len(),
241 );
242
243 let mut skel_builder = BpfSkelBuilder::default();
245 skel_builder.obj_builder.debug(opts.verbose);
246 let open_opts = opts.libbpf.clone().into_bpf_open_opts();
247 let mut skel = scx_ops_cid_open!(skel_builder, open_object, cidland_ops, open_opts)
248 .context("opening BPF skeleton (does this kernel support cid-form sched_ext?)")?;
249
250 skel.struct_ops.cidland_ops_mut().exit_dump_len = opts.exit_dump_len;
251
252 let rodata = skel
253 .maps
254 .rodata_data
255 .as_mut()
256 .expect("rodata_data missing after skel open");
257 rodata.slice_ns = opts.slice_us * 1000;
258 rodata.slice_lag = opts.slice_lag_us * 1000;
259
260 let mut primary_cpus: Vec<usize> = Vec::new();
264 if let Some(ref domain) = opts.primary_domain {
265 let cpus = parse_primary_domain(domain).context("parsing primary domain")?;
266
267 if let Some(cpu) = cpus.iter().find(|cpu| **cpu >= *NR_CPU_IDS) {
268 bail!(
269 "primary domain cpu {} exceeds nr_cpu_ids {}",
270 cpu,
271 *NR_CPU_IDS
272 );
273 }
274 if cpus.len() < *NR_CPU_IDS {
275 info!("primary domain: {:?}", cpus);
276 primary_cpus = cpus;
277 rodata.primary_all = false;
278 }
279 }
280
281 skel.struct_ops.cidland_ops_mut().flags = *compat::SCX_OPS_ENQ_EXITING
287 | *compat::SCX_OPS_ENQ_LAST
288 | *compat::SCX_OPS_ENQ_MIGRATION_DISABLED
289 | *compat::SCX_OPS_ALLOW_QUEUED_WAKEUP;
290 info!(
291 "scheduler flags: {:#x}",
292 skel.struct_ops.cidland_ops_mut().flags
293 );
294
295 let mut skel = scx_ops_cid_load!(skel, cidland_ops, uei).context("loading BPF skeleton")?;
297
298 let nr_cpus = (*NR_CPU_IDS).max(*NR_CPUS_POSSIBLE);
305 let mut args = types::cidland_arena_args {
306 nr_cpus: nr_cpus as u64,
307 };
308 run_syscall_prog(&skel.progs.cidland_arena_init, &mut args)
309 .context("running cidland_arena_init")?;
310
311 let mut words = vec![0u64; nr_cpus.div_ceil(64)];
314 for cpu in &primary_cpus {
315 words[cpu / 64] |= 1u64 << (cpu % 64);
316 }
317 for (idx, word) in words.iter().enumerate() {
318 if *word == 0 {
319 continue;
320 }
321 let mut args = types::cidland_primary_args {
322 idx: idx as u64,
323 word: *word,
324 };
325 run_syscall_prog(&skel.progs.cidland_set_primary_word, &mut args)
326 .context("running cidland_set_primary_word")?;
327 }
328
329 let arenalib =
334 ArenaLib::start(skel.object_mut()).context("starting arena userspace services")?;
335
336 let struct_ops = Some(scx_ops_attach!(skel, cidland_ops).context("attaching scheduler")?);
337 let stats_server = StatsServer::new(stats::server_data()).launch()?;
338
339 Ok(Self {
340 _arenalib: arenalib,
341 skel,
342 struct_ops,
343 stats_server,
344 })
345 }
346
347 fn get_metrics(&self) -> Metrics {
348 let bss_data = self
349 .skel
350 .maps
351 .bss_data
352 .as_ref()
353 .expect("bss_data missing after skel load");
354 Metrics {
355 nr_direct_dispatches: bss_data.nr_direct_dispatches,
356 nr_shared_enqueues: bss_data.nr_shared_enqueues,
357 nr_idle_kicks: bss_data.nr_idle_kicks,
358 nr_local_llc: bss_data.nr_local_llc,
359 nr_remote_llc: bss_data.nr_remote_llc,
360 }
361 }
362
363 fn exited(&mut self) -> bool {
364 uei_exited!(&self.skel, uei)
365 }
366
367 fn run(&mut self, shutdown: Arc<AtomicBool>) -> Result<UserExitInfo> {
368 let (res_ch, req_ch) = self.stats_server.channels();
369 while !shutdown.load(Ordering::Relaxed) && !self.exited() {
370 match req_ch.recv_timeout(Duration::from_secs(1)) {
371 Ok(()) => res_ch.send(self.get_metrics())?,
372 Err(RecvTimeoutError::Timeout) => {}
373 Err(e) => Err(e)?,
374 }
375 }
376
377 let _ = self.struct_ops.take();
378 uei_report!(&self.skel, uei)
379 }
380}
381
382impl Drop for Scheduler<'_> {
383 fn drop(&mut self) {
384 info!("Unregister {SCHEDULER_NAME} scheduler");
385 }
386}
387
388fn main() -> Result<()> {
389 let opts = Opts::parse();
390
391 if opts.version {
392 println!(
393 "{} {}",
394 SCHEDULER_NAME,
395 build_id::full_version(env!("CARGO_PKG_VERSION"))
396 );
397 return Ok(());
398 }
399
400 if opts.help_stats {
401 stats::server_data().describe_meta(&mut std::io::stdout(), None)?;
402 return Ok(());
403 }
404
405 let loglevel = simplelog::LevelFilter::Info;
406
407 let mut lcfg = simplelog::ConfigBuilder::new();
408 lcfg.set_time_offset_to_local()
409 .expect("Failed to set local time offset")
410 .set_time_level(simplelog::LevelFilter::Error)
411 .set_location_level(simplelog::LevelFilter::Off)
412 .set_target_level(simplelog::LevelFilter::Off)
413 .set_thread_level(simplelog::LevelFilter::Off);
414 simplelog::TermLogger::init(
415 loglevel,
416 lcfg.build(),
417 simplelog::TerminalMode::Stderr,
418 simplelog::ColorChoice::Auto,
419 )?;
420
421 let shutdown = Arc::new(AtomicBool::new(false));
422 let shutdown_clone = shutdown.clone();
423 ctrlc::set_handler(move || {
424 shutdown_clone.store(true, Ordering::Relaxed);
425 })
426 .context("Error setting Ctrl-C handler")?;
427
428 if let Some(intv) = opts.monitor.or(opts.stats) {
429 let shutdown_copy = shutdown.clone();
430 let jh = std::thread::spawn(move || {
431 match stats::monitor(Duration::from_secs_f64(intv), shutdown_copy) {
432 Ok(_) => debug!("stats monitor thread finished successfully"),
433 Err(error_object) => {
434 warn!("stats monitor thread finished because of an error {error_object}")
435 }
436 }
437 });
438 if opts.monitor.is_some() {
439 let _ = jh.join();
440 return Ok(());
441 }
442 }
443
444 let mut open_object = MaybeUninit::uninit();
445 loop {
446 let mut sched = Scheduler::init(&opts, &mut open_object)?;
447 if !sched.run(shutdown.clone())?.should_restart() {
448 break;
449 }
450 }
451
452 Ok(())
453}