use std::sync::LazyLock;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Instant;
use crate::render::NumberFormat;
pub(crate) static ENABLED : LazyLock<bool> =
LazyLock::new(|| is_enabled_by(std::env::var_os("MEZURA_PHASE_TIMING").as_deref()));
static OPEN_NANOS : AtomicU64 = AtomicU64::new(0);
static READ_NANOS : AtomicU64 = AtomicU64::new(0);
static PARSE_NANOS : AtomicU64 = AtomicU64::new(0);
static BYTES : AtomicU64 = AtomicU64::new(0);
static FILES : AtomicU64 = AtomicU64::new(0);
static STARVED : AtomicU64 = AtomicU64::new(0);
static STARVED_NANOS : AtomicU64 = AtomicU64::new(0);
#[derive(Debug, Default)]
pub(crate) struct Totals {
pub open_nanos: u64,
pub read_nanos: u64,
pub parse_nanos: u64,
pub bytes: u64,
pub files: u64,
pub starved: u64,
pub starved_nanos: u64,
}
impl Totals {
pub(crate) fn publish(&self) {
OPEN_NANOS.fetch_add(self.open_nanos, Ordering::Relaxed);
READ_NANOS.fetch_add(self.read_nanos, Ordering::Relaxed);
PARSE_NANOS.fetch_add(self.parse_nanos, Ordering::Relaxed);
BYTES.fetch_add(self.bytes, Ordering::Relaxed);
FILES.fetch_add(self.files, Ordering::Relaxed);
STARVED.fetch_add(self.starved, Ordering::Relaxed);
STARVED_NANOS.fetch_add(self.starved_nanos, Ordering::Relaxed);
}
}
pub(crate) fn nanos_since(start: Instant) -> u64 {
start.elapsed().as_nanos() as u64
}
pub(crate) fn report(consumers: usize, run_millis: u128) -> String {
format_report(&Totals {
open_nanos: OPEN_NANOS.load(Ordering::Relaxed),
read_nanos: READ_NANOS.load(Ordering::Relaxed),
parse_nanos: PARSE_NANOS.load(Ordering::Relaxed),
bytes: BYTES.load(Ordering::Relaxed),
files: FILES.load(Ordering::Relaxed),
starved: STARVED.load(Ordering::Relaxed),
starved_nanos: STARVED_NANOS.load(Ordering::Relaxed)
}, consumers, run_millis, std::thread::available_parallelism().map_or(0, |x| x.get()))
}
fn is_enabled_by(value: Option<&std::ffi::OsStr>) -> bool {
value.is_some_and(|value| !matches!(value.to_string_lossy().trim().to_ascii_lowercase().as_str(),
"" | "0" | "no" | "false" | "off"))
}
fn format_report(totals: &Totals, consumers: usize, run_millis: u128, hardware_threads: usize) -> String {
let consumers = consumers.max(1);
let thread_nanos = (consumers as f64 * run_millis as f64 * 1_000_000.0).max(1.0);
let share = |x: u64| 100.0 * x as f64 / thread_nanos;
let read_gb_per_second = if totals.read_nanos > 0 {totals.bytes as f64 / totals.read_nanos as f64} else {0.0};
let consumer_word = if consumers == 1 {"consumer"} else {"consumers"};
let wait_word = if totals.starved == 1 {"wait"} else {"waits"};
let queueing = if hardware_threads > 0 && consumers > hardware_threads {
let hardware_word = if hardware_threads == 1 {"hardware thread"} else {"hardware threads"};
format!("\n[phase] shares are of elapsed time over {consumers} {consumer_word} on \
{hardware_threads} {hardware_word}, so a thread waiting for a core counts in the phase it is in")
} else {
String::new()
};
format!("[phase] {consumers} {consumer_word}: starved {:.1}% ({} {wait_word}) | open {:.1}% ({}) | read {:.1}% ({}) | parse {:.1}% ({})\n\
[phase] {} files, {:.1} MB, read at {:.2} GB/s per thread{queueing}",
share(totals.starved_nanos), totals.starved,
share(totals.open_nanos), format_summed_millis(totals.open_nanos),
share(totals.read_nanos), format_summed_millis(totals.read_nanos),
share(totals.parse_nanos), format_summed_millis(totals.parse_nanos),
totals.files, totals.bytes as f64 / 1_048_576.0, read_gb_per_second)
}
fn format_summed_millis(nanos: u64) -> String {
format!("{} ms", NumberFormat::new(Some(','), '.').integer((nanos / 1_000_000) as usize))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_four_shares_are_of_the_consumers_own_time_and_close_on_a_hundred() {
let totals = Totals {
open_nanos: 30_228_000_000, read_nanos: 9_763_000_000, parse_nanos: 14_442_000_000,
bytes: 7_633_832_346, files: 308_744,
starved: 131_681, starved_nanos: 165_012_000_000
};
let report = format_report(&totals, 64, 3482, 64);
assert!(report.contains("64 consumers: starved 74.0% (131681 waits) | open 13.6% (30,228 ms) \
| read 4.4% (9,763 ms) | parse 6.5% (14,442 ms)"), "{report}");
assert!(report.contains("308744 files, 7280.2 MB, read at 0.78 GB/s per thread"), "{report}");
}
#[test]
fn the_share_and_the_count_of_waits_are_two_different_statements() {
let trickling = Totals {starved: 131_681, starved_nanos: 165_012_000_000, ..Default::default()};
let stalled = Totals {starved: 12, starved_nanos: 165_012_000_000, ..Default::default()};
assert!(format_report(&trickling, 64, 3482, 64).contains("starved 74.0% (131681 waits)"));
assert!(format_report(&stalled, 64, 3482, 64).contains("starved 74.0% (12 waits)"));
}
#[test]
fn more_consumers_than_cores_is_said_out_loud_and_a_core_each_is_not() {
let totals = Totals::default();
let crowded = format_report(&totals, 64, 3482, 16);
assert!(crowded.contains("[phase] shares are of elapsed time over 64 consumers on 16 hardware \
threads, so a thread waiting for a core counts in the phase it is in"), "{crowded}");
for hardware_threads in [64, 128] {
let roomy = format_report(&totals, 64, 3482, hardware_threads);
assert!(!roomy.contains("hardware thread"),
"{hardware_threads} hardware threads for 64 consumers was still called crowded:\n{roomy}");
}
assert!(!format_report(&totals, 64, 3482, 0).contains("hardware thread"));
}
#[test]
fn the_value_decides_and_not_merely_the_presence_of_the_name() {
let set_to = |value: &str| is_enabled_by(Some(std::ffi::OsStr::new(value)));
assert!(!is_enabled_by(None), "the report is on without the variable at all");
for off in ["0", "no", "false", "off", "", " ", "OFF", "False"] {
assert!(!set_to(off), "'{off}' left the report on");
}
for on in ["1", "yes", "true", "on", "please"] {
assert!(set_to(on), "'{on}' did not turn the report on");
}
}
#[test]
fn a_run_with_nothing_in_it_reports_zeroes_rather_than_dividing_by_zero() {
let report = format_report(&Totals::default(), 0, 0, 0);
assert!(report.contains("1 consumer: starved 0.0% (0 waits) | open 0.0% (0 ms) | read 0.0% (0 ms) \
| parse 0.0% (0 ms)"), "{report}");
assert!(report.contains("0 files, 0.0 MB, read at 0.00 GB/s per thread"), "{report}");
}
}