1use std::sync::atomic::{AtomicBool, Ordering};
15use std::time::Duration;
16
17use anyhow::Result;
18
19use crate::chaos::{self, RawWindow};
20use crate::procdb::ProcessDb;
21use crate::scheduler::{PandemoniumStats, Scheduler};
22use crate::tuning::{
23 self, detect_regime, scaled_regime_knobs, MwuController, MwuSignals, Regime, HIST_BUCKETS,
24};
25
26const CHAOS_WIN: usize = 16;
33
34const SLEEP_BUCKETS: usize = 4;
39
40pub fn monitor_loop(
46 sched: &mut Scheduler,
47 shutdown: &'static AtomicBool,
48 verbose: bool,
49 nr_cpus: u64,
50) -> Result<bool> {
51 let mut prev = PandemoniumStats::default();
52 let mut prev_hist = [[0u64; HIST_BUCKETS]; 3];
53 let mut prev_sleep = [0u64; SLEEP_BUCKETS];
54 let mut regime = Regime::Mixed;
55
56 let mut idle_win: RawWindow<CHAOS_WIN> = RawWindow::new();
60 let mut wake_win: RawWindow<CHAOS_WIN> = RawWindow::new();
61 let mut prev_bp_h: f64 = 0.0;
62 let chaos_count = chaos::ChaosCounter::new();
63 let mut prev_lambda_above: bool = false;
64 let mut tau_ns = sched.read_tuning_knobs().topology_tau_ns;
68 let mut mwu = MwuController::new(scaled_regime_knobs(regime, nr_cpus, tau_ns));
69 let mut pending_regime = regime;
70 let mut regime_hold: u32 = 0;
71 let mut light_ticks: u64 = 0;
72 let mut mixed_ticks: u64 = 0;
73 let mut heavy_ticks: u64 = 0;
74 let mut stability_score: u32 = 0;
79 let mut tick_counter: u64 = 0;
80
81 let mut quiesce = tuning::QuiescenceState::new();
86 let mut retune_interval: u32 = tuning::RETUNE_INTERVAL_BASE;
87 let mut ticks_since_retune: u32 = 0;
88 let mut frozen_ticks: u64 = 0;
89
90 let mut procdb = match ProcessDb::new() {
91 Ok(db) => Some(db),
92 Err(e) => {
93 log_warn!("PROCDB INIT FAILED: {}", e);
94 None
95 }
96 };
97
98 let live = sched.read_tuning_knobs();
102 let mut rk = scaled_regime_knobs(regime, nr_cpus, tau_ns);
103 rk.topology_tau_ns = tau_ns;
104 rk.codel_eq_ns = live.codel_eq_ns;
105 sched.write_tuning_knobs(&rk)?;
106 mwu.set_baseline(rk);
113 let mut last_written_knobs = rk;
117
118 while !shutdown.load(Ordering::Relaxed) && !sched.exited() {
119 crate::watchdog::LOOP_HEARTBEAT.fetch_add(1, Ordering::Relaxed);
120 std::thread::sleep(Duration::from_secs(1));
121
122 let stats = sched.read_stats();
123 let cur_hist = sched.read_wake_lat_hist();
124 let cur_sleep = sched.read_sleep_hist();
125
126 let mut wrapped = stats.nr_dispatches < prev.nr_dispatches;
131 if !wrapped {
132 'wrap: for tier in 0..3 {
133 for b in 0..HIST_BUCKETS {
134 if cur_hist[tier][b] < prev_hist[tier][b] {
135 wrapped = true;
136 break 'wrap;
137 }
138 }
139 }
140 }
141 if !wrapped {
142 for i in 0..SLEEP_BUCKETS {
143 if cur_sleep[i] < prev_sleep[i] {
144 wrapped = true;
145 break;
146 }
147 }
148 }
149 if wrapped {
150 log_warn!("WRAP DETECTED: BASELINE RESET, SKIPPING ADAPTIVE UPDATE");
151 prev = stats;
152 prev_hist = cur_hist;
153 prev_sleep = cur_sleep;
154 continue;
155 }
156
157 let delta_d = stats.nr_dispatches.wrapping_sub(prev.nr_dispatches);
159 let delta_idle = stats.nr_idle_hits.wrapping_sub(prev.nr_idle_hits);
160 let delta_shared = stats.nr_shared.wrapping_sub(prev.nr_shared);
161 let delta_preempt = stats.nr_preempt.wrapping_sub(prev.nr_preempt);
162 let delta_keep = stats.nr_keep_running.wrapping_sub(prev.nr_keep_running);
163 let delta_parks = stats.nr_osc_park.wrapping_sub(prev.nr_osc_park);
164 let delta_wake_sum = stats.wake_lat_sum.wrapping_sub(prev.wake_lat_sum);
165 let delta_wake_samples = stats.wake_lat_samples.wrapping_sub(prev.wake_lat_samples);
166 let delta_hard = stats.nr_hard_kicks.wrapping_sub(prev.nr_hard_kicks);
167 let delta_soft = stats.nr_soft_kicks.wrapping_sub(prev.nr_soft_kicks);
168 let delta_enq_wake = stats.nr_enq_wakeup.wrapping_sub(prev.nr_enq_wakeup);
169 let delta_enq_requeue = stats.nr_enq_requeue.wrapping_sub(prev.nr_enq_requeue);
170 let delta_rescue = stats
171 .nr_overflow_rescue
172 .wrapping_sub(prev.nr_overflow_rescue);
173 let scatter_now: u64 = stats.nr_cross_domain[0..6].iter().sum();
179 let scatter_prev: u64 = prev.nr_cross_domain[0..6].iter().sum();
180 let delta_scatter = scatter_now.saturating_sub(scatter_prev);
181 let scatter_pct = if delta_d > 0 {
182 delta_scatter * 100 / delta_d
183 } else {
184 0
185 };
186 let wake_avg_us = if delta_wake_samples > 0 {
187 delta_wake_sum / delta_wake_samples / 1000
188 } else {
189 0
190 };
191
192 let d_idle_sum = stats.wake_lat_idle_sum.wrapping_sub(prev.wake_lat_idle_sum);
194 let d_idle_cnt = stats.wake_lat_idle_cnt.wrapping_sub(prev.wake_lat_idle_cnt);
195 let d_kick_sum = stats.wake_lat_kick_sum.wrapping_sub(prev.wake_lat_kick_sum);
196 let d_kick_cnt = stats.wake_lat_kick_cnt.wrapping_sub(prev.wake_lat_kick_cnt);
197 let lat_idle_us = if d_idle_cnt > 0 {
198 d_idle_sum / d_idle_cnt / 1000
199 } else {
200 0
201 };
202 let lat_kick_us = if d_kick_cnt > 0 {
203 d_kick_sum / d_kick_cnt / 1000
204 } else {
205 0
206 };
207 let delta_reenq = stats.nr_reenqueue.wrapping_sub(prev.nr_reenqueue);
208
209 let dl2_hb = stats.nr_l2_hit_batch.wrapping_sub(prev.nr_l2_hit_batch);
211 let dl2_mb = stats.nr_l2_miss_batch.wrapping_sub(prev.nr_l2_miss_batch);
212 let dl2_hi = stats
213 .nr_l2_hit_interactive
214 .wrapping_sub(prev.nr_l2_hit_interactive);
215 let dl2_mi = stats
216 .nr_l2_miss_interactive
217 .wrapping_sub(prev.nr_l2_miss_interactive);
218 let dl2_hl = stats
219 .nr_l2_hit_lat_crit
220 .wrapping_sub(prev.nr_l2_hit_lat_crit);
221 let dl2_ml = stats
222 .nr_l2_miss_lat_crit
223 .wrapping_sub(prev.nr_l2_miss_lat_crit);
224 let l2_pct_b = if dl2_hb + dl2_mb > 0 {
225 dl2_hb * 100 / (dl2_hb + dl2_mb)
226 } else {
227 0
228 };
229 let l2_pct_i = if dl2_hi + dl2_mi > 0 {
230 dl2_hi * 100 / (dl2_hi + dl2_mi)
231 } else {
232 0
233 };
234 let l2_pct_l = if dl2_hl + dl2_ml > 0 {
235 dl2_hl * 100 / (dl2_hl + dl2_ml)
236 } else {
237 0
238 };
239
240 let idle_pct = if delta_d > 0 {
241 delta_idle * 100 / delta_d
242 } else {
243 0
244 };
245
246 let mut delta_hist = [[0u64; HIST_BUCKETS]; 3];
248 for tier in 0..3 {
249 for b in 0..HIST_BUCKETS {
250 delta_hist[tier][b] = cur_hist[tier][b] - prev_hist[tier][b];
251 }
252 }
253
254 let tp99_b_ns = tuning::compute_p99_from_histogram(&delta_hist[0]);
256 let tp99_i_ns = tuning::compute_p99_from_histogram(&delta_hist[1]);
257 let tp99_l_ns = tuning::compute_p99_from_histogram(&delta_hist[2]);
258
259 let mut agg = [0u64; HIST_BUCKETS];
261 for t in 0..3 {
262 for b in 0..HIST_BUCKETS {
263 agg[b] += delta_hist[t][b];
264 }
265 }
266 let p99_ns = tuning::compute_p99_from_histogram(&agg);
267
268 let mut delta_sleep = [0u64; SLEEP_BUCKETS];
270 for i in 0..SLEEP_BUCKETS {
271 delta_sleep[i] = cur_sleep[i] - prev_sleep[i];
272 }
273 let sleep_total: u64 = delta_sleep.iter().sum();
274 let io_pct = if sleep_total > 0 {
275 (delta_sleep[0] + delta_sleep[1]) * 100 / sleep_total
276 } else {
277 0
278 };
279
280 idle_win.push(idle_pct as f64);
284 wake_win.push(delta_enq_wake as f64);
285
286 let (idle_lambda, _idle_hvg_s) = chaos::hvg_stats(&idle_win);
291 let wake_bp_h = chaos::bandt_pompe_d3(&wake_win);
292 let bp_delta = wake_bp_h - prev_bp_h;
293 let mean_idle = chaos::mean(&idle_win);
294 let rqa = chaos::rqa_det(&idle_win);
295
296 let lambda_above = idle_lambda >= chaos::HVG_LAMBDA_CHAOTIC_MIN;
301 let chaos_crossing = (lambda_above && !prev_lambda_above) || bp_delta > 0.10;
302 if chaos_crossing {
303 chaos_count.bump();
304 }
305 prev_lambda_above = lambda_above;
306
307 let detected = detect_regime(mean_idle, idle_lambda, wake_bp_h);
309
310 let mut regime_changed_this_tick = false;
311 if detected != regime {
312 if detected == pending_regime {
313 regime_hold += 1;
314 } else {
315 pending_regime = detected;
316 regime_hold = 1;
317 }
318 if regime_hold >= 2 {
319 regime = detected;
320 let live = sched.read_tuning_knobs();
324 tau_ns = live.topology_tau_ns;
325 let mut rk = scaled_regime_knobs(regime, nr_cpus, tau_ns);
326 rk.topology_tau_ns = tau_ns;
327 rk.codel_eq_ns = live.codel_eq_ns;
328 sched.write_tuning_knobs(&rk)?;
329 last_written_knobs = rk;
330 regime_changed_this_tick = true;
331 mwu.set_baseline(rk);
332 mwu.reset_regime(regime);
336 retune_interval = tuning::RETUNE_INTERVAL_BASE;
337 ticks_since_retune = 0;
338 }
339 } else {
340 pending_regime = regime;
341 regime_hold = 0;
342 }
343
344 let mwu_converged = mwu.converged(regime);
352 let frozen = quiesce.update(idle_lambda, rqa, mwu_converged);
353 if frozen {
354 frozen_ticks += 1;
355 }
356
357 if !regime_changed_this_tick && !frozen {
363 ticks_since_retune += 1;
364 if ticks_since_retune >= retune_interval {
365 ticks_since_retune = 0;
366 let signals = MwuSignals {
367 p99_ns,
368 interactive_p99_ns: tp99_i_ns,
369 io_pct,
370 rescue_count: delta_rescue,
371 wakeup_rate: delta_enq_wake,
376 scatter_pct,
377 hvg_lambda: idle_lambda,
378 bp_h_delta: bp_delta,
379 rqa_det: rqa,
383 };
384 let osc_state = sched.read_oscillator_state();
391 let mut knobs = mwu.update(
392 &signals,
393 regime.p99_ceiling(),
394 nr_cpus,
395 tau_ns,
396 &osc_state,
397 regime,
398 );
399 let live = sched.read_tuning_knobs();
403 knobs.topology_tau_ns = live.topology_tau_ns;
404 knobs.codel_eq_ns = live.codel_eq_ns;
405 let changed = tuning::knobs_differ(&knobs, &last_written_knobs);
411 if changed {
412 sched.write_tuning_knobs(&knobs)?;
413 last_written_knobs = knobs;
414 }
415 let disturbed = mwu.had_losses() || chaos_crossing;
416 retune_interval =
417 tuning::next_retune_interval(retune_interval, !changed, disturbed);
418 }
419 }
420
421 let tighten_delta = if mwu.had_losses() { 1u64 } else { 0u64 };
423 stability_score = tuning::compute_stability_score(
424 stability_score,
425 regime_changed_this_tick,
426 tighten_delta,
427 p99_ns,
428 regime.p99_ceiling(),
429 );
430
431 let (db_total, db_confident) = if let Some(ref mut db) = procdb {
433 db.ingest();
434 db.flush_predictions();
435 db.tick();
436 db.summary()
437 } else {
438 (0, 0)
439 };
440
441 let p99_us = p99_ns / 1000;
442 let tp99_b = tp99_b_ns / 1000;
443 let tp99_i = tp99_i_ns / 1000;
444 let tp99_l = tp99_l_ns / 1000;
445 let knobs = sched.read_tuning_knobs();
446
447 let sojourn_ms = stats.batch_sojourn_ns / 1_000_000;
448 let sojourn_thresh_ms = knobs.sojourn_thresh_ns / 1_000_000;
449 let longrun_label = if stats.longrun_mode_active > 0 {
450 " LONGRUN"
451 } else {
452 ""
453 };
454
455 if verbose && tuning::should_print_telemetry(tick_counter, stability_score) {
456 let rqa_disp = rqa.unwrap_or(-1.0);
457 let frozen_disp = if frozen { 1 } else { 0 };
458 println!(
459 "d/s: {:<8} idle: {}% shared: {:<6} preempt: {:<4} keep: {:<4} kick: H={:<4} S={:<4} enq: W={:<4} R={:<4} wake: {}us p99: {}us [B:{} I:{} L:{}] lat_idle: {}us lat_kick: {}us procdb: {}/{} sleep: io={}% slice: {}us batch: {}us reenq: {} sjrn: {}ms/{}ms rescue: {} l2: B={}% I={}% L={}% chaos: lam={:.2} H={:.2} det={:.2} x={} frozen: {} (n={}) retune_iv: {} [{}{}]",
460 delta_d, idle_pct, delta_shared, delta_preempt, delta_keep,
461 delta_hard, delta_soft, delta_enq_wake, delta_enq_requeue,
462 wake_avg_us, p99_us, tp99_b, tp99_i, tp99_l,
463 lat_idle_us, lat_kick_us,
464 db_total, db_confident,
465 io_pct, knobs.slice_ns / 1000, knobs.batch_slice_ns / 1000,
466 delta_reenq, sojourn_ms, sojourn_thresh_ms,
467 delta_rescue,
468 l2_pct_b, l2_pct_i, l2_pct_l,
469 idle_lambda, wake_bp_h, rqa_disp, chaos_count.load(),
470 frozen_disp, frozen_ticks, retune_interval,
471 regime.label(), longrun_label,
472 );
473 }
474
475 sched.log.snapshot(
476 delta_d,
477 delta_idle,
478 delta_shared,
479 delta_preempt,
480 delta_keep,
481 delta_parks,
482 wake_avg_us,
483 delta_hard,
484 delta_soft,
485 lat_idle_us,
486 lat_kick_us,
487 );
488
489 match regime {
490 Regime::Light => light_ticks += 1,
491 Regime::Mixed => mixed_ticks += 1,
492 Regime::Heavy => heavy_ticks += 1,
493 }
494
495 tick_counter += 1;
496 prev_hist = cur_hist;
497 prev_sleep = cur_sleep;
498 prev = stats;
499 prev_bp_h = wake_bp_h;
500 }
501
502 if let Some(ref db) = procdb {
504 let path = ProcessDb::default_path();
505 match db.save(&path) {
506 Ok(()) => {
507 let (total, confident) = db.summary();
508 log_info!(
509 "PROCDB: SAVED {}/{} PROFILES TO {}",
510 confident,
511 total,
512 path.display()
513 );
514 }
515 Err(e) => log_warn!("PROCDB SAVE FAILED: {}", e),
516 }
517 }
518
519 let final_knobs = sched.read_tuning_knobs();
521 let final_stats = sched.read_stats();
522 let l2_total_b = final_stats.nr_l2_hit_batch + final_stats.nr_l2_miss_batch;
523 let l2_total_i = final_stats.nr_l2_hit_interactive + final_stats.nr_l2_miss_interactive;
524 let l2_total_l = final_stats.nr_l2_hit_lat_crit + final_stats.nr_l2_miss_lat_crit;
525 let l2_cum_b = if l2_total_b > 0 {
526 final_stats.nr_l2_hit_batch * 100 / l2_total_b
527 } else {
528 0
529 };
530 let l2_cum_i = if l2_total_i > 0 {
531 final_stats.nr_l2_hit_interactive * 100 / l2_total_i
532 } else {
533 0
534 };
535 let l2_cum_l = if l2_total_l > 0 {
536 final_stats.nr_l2_hit_lat_crit * 100 / l2_total_l
537 } else {
538 0
539 };
540 let x = &final_stats.nr_cross_domain;
544 let x_scatter: u64 = x[0..6].iter().sum();
545 let x_scatter_pct = if final_stats.nr_dispatches > 0 {
546 x_scatter * 100 / final_stats.nr_dispatches
547 } else {
548 0
549 };
550 println!(
551 "[KNOBS] regime={} slice_ns={} batch_ns={} preempt_ns={} mwu={:.3} ticks=L:{}/M:{}/H:{} frozen={} 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={}",
552 regime.label(), final_knobs.slice_ns, final_knobs.batch_slice_ns,
553 final_knobs.preempt_thresh_ns,
554 mwu.scale(regime),
555 light_ticks, mixed_ticks, heavy_ticks, frozen_ticks,
556 l2_cum_b, l2_cum_i, l2_cum_l,
557 x_scatter_pct, x[0], x[1], x[2], x[3], x[4], x[5], x[6], x[7],
558 );
559
560 let should_restart = sched.read_exit_info();
562 Ok(should_restart)
563}