use std::sync::atomic::AtomicU64;
use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering;
use std::time::Instant;
pub const MAX_WORKERS: usize = 72;
pub const N_STAGE: usize = 8;
pub const STAGE_NAMES: [&str; N_STAGE] = [
"tile_entropy",
"tile_recon",
"deblock_cols",
"deblock_rows",
"cdef",
"superres",
"loop_restore",
"other",
];
pub const FILTER_STAGES: [usize; 5] = [2, 3, 4, 5, 6];
#[allow(clippy::declare_interior_mutable_const)]
const ZERO: AtomicU64 = AtomicU64::new(0);
static BUSY_NS: [AtomicU64; N_STAGE * MAX_WORKERS] = [ZERO; N_STAGE * MAX_WORKERS];
static BUSY_CNT: [AtomicU64; N_STAGE * MAX_WORKERS] = [ZERO; N_STAGE * MAX_WORKERS];
#[inline]
const fn ix(w: usize, s: usize) -> usize {
w * N_STAGE + s
}
static PARK_NS: [AtomicU64; MAX_WORKERS] = [ZERO; MAX_WORKERS];
static PARK_CNT: [AtomicU64; MAX_WORKERS] = [ZERO; MAX_WORKERS];
static ACTIVE: AtomicUsize = AtomicUsize::new(0);
static FILT_ACTIVE: AtomicUsize = AtomicUsize::new(0);
static FILT_CONC: [AtomicU64; MAX_WORKERS + 1] = [ZERO; MAX_WORKERS + 1];
static TILE_ACTIVE: AtomicUsize = AtomicUsize::new(0);
static TAIL_CONC: [AtomicU64; MAX_WORKERS + 1] = [ZERO; MAX_WORKERS + 1];
static TAIL_SAMPLES: AtomicU64 = AtomicU64::new(0);
pub const N_DEFER: usize = 5;
pub const DEFER_NAMES: [&str; N_DEFER] = [
"own_progress",
"pass2_progress",
"deblock_barrier",
"ref_progress",
"admitted",
];
static DEFER: [AtomicU64; N_DEFER] = [ZERO; N_DEFER];
#[inline]
pub fn defer(kind: usize) {
DEFER[kind].fetch_add(1, Ordering::Relaxed);
}
static CONC: [AtomicU64; MAX_WORKERS + 1] = [ZERO; MAX_WORKERS + 1];
static SAMPLES: AtomicU64 = AtomicU64::new(0);
static NEXT_SLOT: AtomicUsize = AtomicUsize::new(0);
thread_local! {
static SLOT: usize = NEXT_SLOT.fetch_add(1, Ordering::Relaxed).min(MAX_WORKERS - 1);
}
#[inline]
pub fn slot() -> usize {
SLOT.with(|s| *s)
}
#[inline]
const fn is_filter(stage: usize) -> bool {
stage >= 2 && stage <= 6
}
#[inline]
pub fn stage_begin_of(stage: usize) -> Instant {
ACTIVE.fetch_add(1, Ordering::Relaxed);
if is_filter(stage) {
FILT_ACTIVE.fetch_add(1, Ordering::Relaxed);
} else if stage < 2 {
TILE_ACTIVE.fetch_add(1, Ordering::Relaxed);
}
Instant::now()
}
#[inline]
pub fn stage_end(t0: Instant, stage: usize) {
let ns = t0.elapsed().as_nanos() as u64;
ACTIVE.fetch_sub(1, Ordering::Relaxed);
if is_filter(stage) {
FILT_ACTIVE.fetch_sub(1, Ordering::Relaxed);
} else if stage < 2 {
TILE_ACTIVE.fetch_sub(1, Ordering::Relaxed);
}
let w = slot();
BUSY_NS[ix(w, stage)].fetch_add(ns, Ordering::Relaxed);
BUSY_CNT[ix(w, stage)].fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn park_begin() -> Instant {
Instant::now()
}
#[inline]
pub fn park_end(t0: Instant) {
let ns = t0.elapsed().as_nanos() as u64;
let w = slot();
PARK_NS[w].fetch_add(ns, Ordering::Relaxed);
PARK_CNT[w].fetch_add(1, Ordering::Relaxed);
}
const SAMPLE_US: u64 = 50;
pub fn start_monitor() {
std::thread::spawn(|| {
loop {
let a = ACTIVE.load(Ordering::Relaxed).min(MAX_WORKERS);
CONC[a].fetch_add(1, Ordering::Relaxed);
let fa = FILT_ACTIVE.load(Ordering::Relaxed).min(MAX_WORKERS);
FILT_CONC[fa].fetch_add(1, Ordering::Relaxed);
if TILE_ACTIVE.load(Ordering::Relaxed) == 0 && fa > 0 {
TAIL_CONC[a].fetch_add(1, Ordering::Relaxed);
TAIL_SAMPLES.fetch_add(1, Ordering::Relaxed);
}
SAMPLES.fetch_add(1, Ordering::Relaxed);
std::thread::sleep(std::time::Duration::from_micros(SAMPLE_US));
}
});
}
pub fn reset() {
for w in 0..MAX_WORKERS {
for s in 0..N_STAGE {
BUSY_NS[ix(w, s)].store(0, Ordering::Relaxed);
BUSY_CNT[ix(w, s)].store(0, Ordering::Relaxed);
}
PARK_NS[w].store(0, Ordering::Relaxed);
PARK_CNT[w].store(0, Ordering::Relaxed);
}
for k in 0..=MAX_WORKERS {
CONC[k].store(0, Ordering::Relaxed);
FILT_CONC[k].store(0, Ordering::Relaxed);
TAIL_CONC[k].store(0, Ordering::Relaxed);
}
for k in 0..N_DEFER {
DEFER[k].store(0, Ordering::Relaxed);
}
SAMPLES.store(0, Ordering::Relaxed);
TAIL_SAMPLES.store(0, Ordering::Relaxed);
}
pub fn report(frames: u64) {
assert!(
NEXT_SLOT.load(Ordering::Relaxed) <= MAX_WORKERS,
"task probe worker slot overflow"
);
let f = frames.max(1) as f64;
let mut per_stage_total = [0u64; N_STAGE];
let mut per_worker_total = [0u64; MAX_WORKERS];
let mut used = 0usize;
for w in 0..MAX_WORKERS {
let mut wt = 0u64;
for s in 0..N_STAGE {
let ns = BUSY_NS[ix(w, s)].load(Ordering::Relaxed);
per_stage_total[s] += ns;
wt += ns;
}
per_worker_total[w] = wt;
if wt > 0 || PARK_NS[w].load(Ordering::Relaxed) > 0 {
used = used.max(w + 1);
}
}
println!("PROBE frames {frames}");
for s in 0..N_STAGE {
let cnt: u64 = (0..MAX_WORKERS)
.map(|w| BUSY_CNT[ix(w, s)].load(Ordering::Relaxed))
.sum();
println!(
"PROBE stage_ms_per_frame {} {:.3} count_per_frame {:.2}",
STAGE_NAMES[s],
per_stage_total[s] as f64 / 1e6 / f,
cnt as f64 / f
);
}
let filter_ns: u64 = FILTER_STAGES.iter().map(|&s| per_stage_total[s]).sum();
let tile_ns: u64 = per_stage_total[0] + per_stage_total[1];
println!(
"PROBE filter_chain_ms_per_frame {:.3}",
filter_ns as f64 / 1e6 / f
);
println!("PROBE tile_ms_per_frame {:.3}", tile_ns as f64 / 1e6 / f);
for w in 0..used {
println!(
"PROBE worker {w} busy_ms_per_frame {:.3} park_ms_per_frame {:.3} park_count_per_frame {:.2}",
per_worker_total[w] as f64 / 1e6 / f,
PARK_NS[w].load(Ordering::Relaxed) as f64 / 1e6 / f,
PARK_CNT[w].load(Ordering::Relaxed) as f64 / f
);
}
let samples = SAMPLES.load(Ordering::Relaxed).max(1);
let mut weighted = 0f64;
for k in 0..=MAX_WORKERS {
let c = CONC[k].load(Ordering::Relaxed);
if c > 0 {
println!(
"PROBE conc {k} samples {c} frac {:.4}",
c as f64 / samples as f64
);
weighted += (k * c as usize) as f64;
}
}
println!("PROBE mean_active {:.3}", weighted / samples as f64);
for k in 0..=MAX_WORKERS {
let c = FILT_CONC[k].load(Ordering::Relaxed);
if c > 0 {
println!(
"PROBE filtconc {k} samples {c} frac {:.4}",
c as f64 / samples as f64
);
}
}
let tail = TAIL_SAMPLES.load(Ordering::Relaxed);
let mut tail_weighted = 0f64;
for k in 0..=MAX_WORKERS {
let c = TAIL_CONC[k].load(Ordering::Relaxed);
if c > 0 {
println!(
"PROBE tailconc {k} samples {c} frac_of_all {:.4}",
c as f64 / samples as f64
);
tail_weighted += (k * c as usize) as f64;
}
}
println!(
"PROBE tail_frac_of_wall {:.4} tail_mean_active {:.3}",
tail as f64 / samples as f64,
if tail > 0 {
tail_weighted / tail as f64
} else {
0.0
}
);
for k in 0..N_DEFER {
println!(
"PROBE defer {} per_frame {:.2}",
DEFER_NAMES[k],
DEFER[k].load(Ordering::Relaxed) as f64 / f
);
}
let busy_samples: u64 = (1..=MAX_WORKERS)
.map(|k| CONC[k].load(Ordering::Relaxed))
.sum();
if busy_samples > 0 {
println!(
"PROBE mean_active_when_busy {:.3}",
weighted / busy_samples as f64
);
}
}