Skip to main content

scx_pandemonium/
main.rs

1// PANDEMONIUM -- SCHED_EXT KERNEL SCHEDULER
2// ADAPTIVE DESKTOP SCHEDULING FOR LINUX
3//
4// SCHEDULING DECISIONS HAPPEN IN BPF (ZERO KERNEL-USERSPACE ROUND TRIPS)
5// RUST USERSPACE HANDLES: ADAPTIVE CONTROL LOOP, MONITORING, BENCHMARKING
6
7#[allow(non_upper_case_globals)]
8#[allow(non_camel_case_types)]
9#[allow(non_snake_case)]
10#[allow(dead_code)]
11mod bpf_skel;
12
13mod bpf_intf;
14
15#[macro_use]
16mod log;
17mod adaptive;
18mod chaos;
19mod cli;
20mod procdb;
21mod scheduler;
22mod topology;
23mod tuning;
24mod watchdog;
25
26use std::mem::MaybeUninit;
27use std::sync::atomic::{AtomicBool, Ordering};
28use std::time::Duration;
29
30use anyhow::Result;
31use clap::{Parser, Subcommand};
32
33use scheduler::Scheduler;
34use scx_utils::build_id;
35
36static SHUTDOWN: AtomicBool = AtomicBool::new(false);
37
38#[derive(Parser)]
39#[command(name = "scx_pandemonium")]
40#[command(
41    version,
42    disable_version_flag = true,
43    about = "PANDEMONIUM -- ADAPTIVE LINUX SCHEDULER"
44)]
45struct Cli {
46    #[command(subcommand)]
47    command: Option<SubCmd>,
48
49    #[arg(short, long)]
50    verbose: bool,
51
52    /// Print scheduler version and exit.
53    #[arg(long)]
54    version: bool,
55
56    /// Internal: dump in-memory ring log on shutdown
57    #[arg(long, hide = true)]
58    dump_log: bool,
59
60    /// Internal: override CPU count for scaling formulas (test harness use)
61    #[arg(long, hide = true)]
62    nr_cpus: Option<u64>,
63
64    /// Run BPF scheduler only, disable Rust adaptive control loop
65    #[arg(long)]
66    no_adaptive: bool,
67
68    /// Override the topology-derived Phi distance scale (phi_dist_scale_q16).
69    /// 0 disables the Phi steal-resist (flat CoDel target); omit for the
70    /// topology value. Test/bench use -- the override holds across both the
71    /// adaptive and --no-adaptive paths.
72    #[arg(long)]
73    phi_scale: Option<u64>,
74}
75
76#[derive(Subcommand)]
77enum SubCmd {
78    /// Internal: interactive wakeup probe (Python test harness use)
79    #[command(hide = true)]
80    Probe,
81
82    /// Internal: CPU-pinned stress worker (Python test harness use)
83    #[command(hide = true)]
84    StressWorker(StressWorkerArgs),
85}
86
87#[derive(Parser)]
88struct StressWorkerArgs {
89    /// CPU to pin the stress worker to
90    #[arg(long)]
91    cpu: u32,
92}
93
94fn main() -> Result<()> {
95    let cli = Cli::parse();
96
97    let verbose = cli.verbose;
98    let dump_log = cli.dump_log;
99    let nr_cpus = cli.nr_cpus;
100    let no_adaptive = cli.no_adaptive;
101    let phi_scale = cli.phi_scale;
102
103    if cli.version {
104        println!(
105            "scx_pandemonium {}",
106            build_id::full_version(env!("CARGO_PKG_VERSION"))
107        );
108        return Ok(());
109    }
110
111    match cli.command {
112        None => run_scheduler(verbose, dump_log, nr_cpus, no_adaptive, phi_scale),
113        Some(SubCmd::Probe) => {
114            cli::probe::run_probe();
115            Ok(())
116        }
117        Some(SubCmd::StressWorker(args)) => {
118            cli::stress::run_stress_worker(args.cpu);
119            Ok(())
120        }
121    }
122}
123
124fn run_scheduler(
125    verbose: bool,
126    dump_log: bool,
127    nr_cpus: Option<u64>,
128    no_adaptive: bool,
129    phi_scale: Option<u64>,
130) -> Result<()> {
131    ctrlc::set_handler(move || {
132        SHUTDOWN.store(true, Ordering::Relaxed);
133    })?;
134
135    // WATCHDOG: ABORTS IF THE CONTROL LOOP STALLS FOR MORE THAN 10 SECONDS.
136    // LIBBPF MAP OPERATIONS CAN HANG ON KERNEL STALL / VERIFIER RELOAD /
137    // PERCPU CONTENTION; WITHOUT THIS, TELEMETRY AND KNOB WRITES STOP SILENTLY.
138    watchdog::spawn(&SHUTDOWN, Duration::from_secs(10));
139
140    let nr_cpus_display =
141        nr_cpus.unwrap_or_else(|| libbpf_rs::num_possible_cpus().unwrap_or(1) as u64);
142    let governor = std::fs::read_to_string("/sys/devices/system/cpu/cpu0/cpufreq/scaling_governor")
143        .unwrap_or_default()
144        .trim()
145        .to_string();
146
147    let smt_on = std::fs::read_to_string("/sys/devices/system/cpu/smt/active")
148        .map(|s| s.trim() == "1")
149        .unwrap_or(false);
150
151    log_info!(
152        "scx_pandemonium {} SMT {}",
153        build_id::full_version(env!("CARGO_PKG_VERSION")),
154        if smt_on { "on" } else { "off" }
155    );
156    log_info!(
157        "CPUS: {} (governor: {})",
158        nr_cpus_display,
159        if governor.is_empty() {
160            "unknown"
161        } else {
162            &governor
163        }
164    );
165    log_info!("VERBOSE: {}", verbose);
166
167    let mut is_restart = false;
168    loop {
169        // ON RESTART, WAIT FOR KERNEL STRUCT_OPS CLEANUP.
170        // DETACH IS ASYNCHRONOUS -- UNDER HEAVY LOAD (12C SATURATED),
171        // THE KERNEL NEEDS TIME TO FULLY UNREGISTER THE OLD SCHEDULER.
172        if is_restart {
173            std::thread::sleep(Duration::from_secs(2));
174        }
175
176        let mut open_object = MaybeUninit::uninit();
177        let mut sched = Scheduler::init(&mut open_object, nr_cpus)?;
178
179        // POPULATE CACHE TOPOLOGY MAP AT STARTUP
180        match topology::CpuTopology::detect(nr_cpus_display as usize) {
181            Ok(topo) => {
182                topo.log_summary();
183                if let Err(e) = topo.populate_bpf_map(&mut sched) {
184                    log_warn!("CACHE TOPOLOGY MAP WRITE FAILED: {}", e);
185                }
186                if let Err(e) = topo.populate_l2_siblings_map(&sched) {
187                    log_warn!("L2 SIBLINGS MAP WRITE FAILED: {}", e);
188                }
189                // RESISTANCE AFFINITY: COMPUTE R_EFF VIA LAPLACIAN PSEUDOINVERSE
190                // AND POPULATE BPF AFFINITY RANK MAP. SPECTRUM CARRIES lambda_2
191                // AND tau_ns FOR UNIVERSAL TOPOLOGY-DERIVED SCALING.
192                let (reff, rank, mut spectrum) = topo.compute_resistance_affinity();
193                if let Some(pv) = phi_scale {
194                    log_info!(
195                        "PHI OVERRIDE: phi_dist_scale_q16 {} -> {} (--phi-scale)",
196                        spectrum.phi_dist_scale_q16,
197                        pv
198                    );
199                    spectrum.phi_dist_scale_q16 = pv;
200                }
201                topo.log_resistance_affinity(&reff, &rank, spectrum);
202                // T2: derive the emergent de-facto-NUMA domain tree from the cache
203                // graph (min-conductance cuts) and log it. T3's bounded steal reads it.
204                let domains = topo.compute_domain_tree();
205                topo.log_domains(&domains);
206                // T3b.1: flatten the tree to the per-CPU-pair crossing-price matrix
207                // and write it 1:1 with the affinity rank (domain_phi map).
208                let domain_phi = topo.domain_cross_phi_matrix(&domains);
209                if let Err(e) = topo.populate_affinity_rank_map(
210                    &sched,
211                    &reff,
212                    &rank,
213                    spectrum.phi_dist_scale_q16,
214                    &domain_phi,
215                ) {
216                    log_warn!("AFFINITY RANK MAP WRITE FAILED: {}", e);
217                }
218                // T3b.2: partition the tree into emergent overflow domains (L3
219                // granularity) and write cpu_domain -- the the discrete domain map replacement.
220                let ov_domains = topo.overflow_domain_count();
221                let cpu_dom = topo.domain_partition(&domains, ov_domains);
222                for (cpu, &d) in cpu_dom.iter().enumerate() {
223                    if let Err(e) = sched.write_cpu_domain(cpu as u32, d) {
224                        log_warn!("CPU DOMAIN MAP WRITE FAILED (cpu {}): {}", cpu, e);
225                        break;
226                    }
227                }
228                log_info!(
229                    "OVERFLOW DOMAINS: {} (emergent, cpu_domain populated)",
230                    ov_domains
231                );
232                // WRITE tau_ns + codel_eq_ns INTO tuning_knobs. BPF'S tick() ON
233                // CPU 0 PICKS THESE UP AND DERIVES THE TAU-SCALED TIMING STATICS
234                // AND THE R_eff-DERIVED CODEL EQUILIBRIUM TARGET.
235                if let Err(e) = sched.write_topology_fields(spectrum.tau_ns, spectrum.codel_eq_ns) {
236                    log_warn!("TOPOLOGY KNOB WRITE FAILED: {}", e);
237                }
238            }
239            Err(e) => log_warn!("CACHE TOPOLOGY DETECT FAILED: {}", e),
240        }
241
242        let should_restart = if no_adaptive {
243            // BPF-ONLY MODE: SCHEDULER RUNS WITH DEFAULT KNOBS, NO RUST TUNING
244            // STILL PRINTS STATS SO BENCHMARKS GET TELEMETRY FOR BOTH PHASES
245            log_info!("PANDEMONIUM IS ACTIVE (BPF ONLY, CTRL+C TO EXIT)");
246            // ONE-SHOT PROCDB WARM-START. BPF-only mode has no adaptive loop,
247            // so without this every app launch re-learns task classes from cold
248            // (12C BPF app-launch 16ms vs ADAPTIVE ~2ms). ProcessDb::new() loads
249            // the persisted profiles and flush_predictions() populates
250            // task_class_init, which enable() reads on every spawn. Construct,
251            // log, drop -- no loop, no 1Hz tax. Stale-but-warm beats cold.
252            match crate::procdb::ProcessDb::new() {
253                Ok(db) => {
254                    let (total, confident) = db.summary();
255                    log_info!(
256                        "PROCDB: BPF-mode warm-start {}/{} confident profiles",
257                        confident,
258                        total
259                    );
260                }
261                Err(e) => log_warn!("PROCDB WARM-START FAILED: {}", e),
262            }
263            let mut prev = scheduler::PandemoniumStats::default();
264            while !SHUTDOWN.load(Ordering::Relaxed) && !sched.exited() {
265                watchdog::LOOP_HEARTBEAT.fetch_add(1, Ordering::Relaxed);
266                std::thread::sleep(Duration::from_secs(1));
267
268                let stats = sched.read_stats();
269
270                let delta_d = stats.nr_dispatches.wrapping_sub(prev.nr_dispatches);
271                let delta_idle = stats.nr_idle_hits.wrapping_sub(prev.nr_idle_hits);
272                let delta_shared = stats.nr_shared.wrapping_sub(prev.nr_shared);
273                let delta_preempt = stats.nr_preempt.wrapping_sub(prev.nr_preempt);
274                let delta_keep = stats.nr_keep_running.wrapping_sub(prev.nr_keep_running);
275                let delta_parks = stats.nr_osc_park.wrapping_sub(prev.nr_osc_park);
276                let delta_wake_sum = stats.wake_lat_sum.wrapping_sub(prev.wake_lat_sum);
277                let delta_wake_samples = stats.wake_lat_samples.wrapping_sub(prev.wake_lat_samples);
278                let delta_hard = stats.nr_hard_kicks.wrapping_sub(prev.nr_hard_kicks);
279                let delta_soft = stats.nr_soft_kicks.wrapping_sub(prev.nr_soft_kicks);
280                let delta_enq_wake = stats.nr_enq_wakeup.wrapping_sub(prev.nr_enq_wakeup);
281                let delta_enq_requeue = stats.nr_enq_requeue.wrapping_sub(prev.nr_enq_requeue);
282                let wake_avg_us = if delta_wake_samples > 0 {
283                    delta_wake_sum / delta_wake_samples / 1000
284                } else {
285                    0
286                };
287
288                let d_idle_sum = stats.wake_lat_idle_sum.wrapping_sub(prev.wake_lat_idle_sum);
289                let d_idle_cnt = stats.wake_lat_idle_cnt.wrapping_sub(prev.wake_lat_idle_cnt);
290                let d_kick_sum = stats.wake_lat_kick_sum.wrapping_sub(prev.wake_lat_kick_sum);
291                let d_kick_cnt = stats.wake_lat_kick_cnt.wrapping_sub(prev.wake_lat_kick_cnt);
292                let lat_idle_us = if d_idle_cnt > 0 {
293                    d_idle_sum / d_idle_cnt / 1000
294                } else {
295                    0
296                };
297                let lat_kick_us = if d_kick_cnt > 0 {
298                    d_kick_sum / d_kick_cnt / 1000
299                } else {
300                    0
301                };
302                let delta_reenq = stats.nr_reenqueue.wrapping_sub(prev.nr_reenqueue);
303
304                // L2 CACHE AFFINITY DELTAS
305                let dl2_hb = stats.nr_l2_hit_batch.wrapping_sub(prev.nr_l2_hit_batch);
306                let dl2_mb = stats.nr_l2_miss_batch.wrapping_sub(prev.nr_l2_miss_batch);
307                let dl2_hi = stats
308                    .nr_l2_hit_interactive
309                    .wrapping_sub(prev.nr_l2_hit_interactive);
310                let dl2_mi = stats
311                    .nr_l2_miss_interactive
312                    .wrapping_sub(prev.nr_l2_miss_interactive);
313                let dl2_hl = stats
314                    .nr_l2_hit_lat_crit
315                    .wrapping_sub(prev.nr_l2_hit_lat_crit);
316                let dl2_ml = stats
317                    .nr_l2_miss_lat_crit
318                    .wrapping_sub(prev.nr_l2_miss_lat_crit);
319                let l2_pct_b = if dl2_hb + dl2_mb > 0 {
320                    dl2_hb * 100 / (dl2_hb + dl2_mb)
321                } else {
322                    0
323                };
324                let l2_pct_i = if dl2_hi + dl2_mi > 0 {
325                    dl2_hi * 100 / (dl2_hi + dl2_mi)
326                } else {
327                    0
328                };
329                let l2_pct_l = if dl2_hl + dl2_ml > 0 {
330                    dl2_hl * 100 / (dl2_hl + dl2_ml)
331                } else {
332                    0
333                };
334
335                let idle_pct = if delta_d > 0 {
336                    delta_idle * 100 / delta_d
337                } else {
338                    0
339                };
340
341                let sojourn_ms = stats.batch_sojourn_ns / 1_000_000;
342                let longrun_label = if stats.longrun_mode_active > 0 {
343                    " LONGRUN"
344                } else {
345                    ""
346                };
347
348                if verbose {
349                    println!(
350                        "d/s: {:<8} idle: {}% shared: {:<6} preempt: {:<4} keep: {:<4} kick: H={:<4} S={:<4} enq: W={:<4} R={:<4} wake: {}us lat_idle: {}us lat_kick: {}us reenq: {} sjrn: {}ms l2: B={}% I={}% L={}% [BPF{}]",
351                        delta_d, idle_pct, delta_shared, delta_preempt, delta_keep,
352                        delta_hard, delta_soft, delta_enq_wake, delta_enq_requeue,
353                        wake_avg_us, lat_idle_us, lat_kick_us,
354                        delta_reenq, sojourn_ms, l2_pct_b, l2_pct_i, l2_pct_l,
355                        longrun_label,
356                    );
357                }
358
359                sched.log.snapshot(
360                    delta_d,
361                    delta_idle,
362                    delta_shared,
363                    delta_preempt,
364                    delta_keep,
365                    delta_parks,
366                    wake_avg_us,
367                    delta_hard,
368                    delta_soft,
369                    lat_idle_us,
370                    lat_kick_us,
371                );
372
373                prev = stats;
374            }
375
376            // KNOBS SUMMARY: CAPTURED BY TEST HARNESS FOR ARCHIVE
377            let knobs = sched.read_tuning_knobs();
378            let final_stats = sched.read_stats();
379            let l2_total_b = final_stats.nr_l2_hit_batch + final_stats.nr_l2_miss_batch;
380            let l2_total_i = final_stats.nr_l2_hit_interactive + final_stats.nr_l2_miss_interactive;
381            let l2_total_l = final_stats.nr_l2_hit_lat_crit + final_stats.nr_l2_miss_lat_crit;
382            let l2_cum_b = if l2_total_b > 0 {
383                final_stats.nr_l2_hit_batch * 100 / l2_total_b
384            } else {
385                0
386            };
387            let l2_cum_i = if l2_total_i > 0 {
388                final_stats.nr_l2_hit_interactive * 100 / l2_total_i
389            } else {
390                0
391            };
392            let l2_cum_l = if l2_total_l > 0 {
393                final_stats.nr_l2_hit_lat_crit * 100 / l2_total_l
394            } else {
395                0
396            };
397            // CROSS-DOMAIN SCATTER ATTRIBUTION (PER XDOM_* PATH), ON THE [KNOBS]
398            // LINE SO THE BENCH SUITE CAPTURES IT UNIFORMLY ACROSS BPF/ADAPTIVE
399            // (LETS THE SUITE COMPARE SCATTER BETWEEN MODES). scatter_pct IS THE
400            // PLACEMENT-SIDE FRACTION (idx 0..6).
401            let x = &final_stats.nr_cross_domain;
402            let x_scatter: u64 = x[0..6].iter().sum();
403            let x_scatter_pct = if final_stats.nr_dispatches > 0 {
404                x_scatter * 100 / final_stats.nr_dispatches
405            } else {
406                0
407            };
408            println!(
409                "[KNOBS] regime=BPF slice_ns={} batch_ns={} preempt_ns={} l2_hit=B:{}%/I:{}%/L:{}% cross_domain_scatter_pct={} cross_domain_sel_tight={} cross_domain_sel_sync={} cross_domain_sel_normal={} cross_domain_sel_dfl={} cross_domain_enq_t1={} cross_domain_enq_t2={} cross_domain_steal={} cross_domain_step5={}",
410                knobs.slice_ns, knobs.batch_slice_ns,
411                knobs.preempt_thresh_ns,
412                l2_cum_b, l2_cum_i, l2_cum_l,
413                x_scatter_pct, x[0], x[1], x[2], x[3], x[4], x[5], x[6], x[7],
414            );
415
416            sched.read_exit_info()
417        } else {
418            // ADAPTIVE MODE: BPF + SINGLE-THREAD MONITOR LOOP
419            log_info!("PANDEMONIUM IS ACTIVE (CTRL+C TO EXIT)");
420            adaptive::monitor_loop(&mut sched, &SHUTDOWN, verbose, nr_cpus_display)?
421        };
422
423        log_info!("PANDEMONIUM IS SHUTTING DOWN");
424
425        if dump_log {
426            sched.log.dump();
427        }
428        sched.log.summary();
429
430        if !should_restart || SHUTDOWN.load(Ordering::Relaxed) {
431            break;
432        }
433
434        // RESET SHUTDOWN FOR RESTART
435        SHUTDOWN.store(false, Ordering::Relaxed);
436        log_info!("RESTARTING PANDEMONIUM...");
437        is_restart = true;
438    }
439
440    log_info!("Shutdown complete");
441    Ok(())
442}