use std::sync::atomic::{AtomicU64, Ordering};
use serde_json::json;
pub(crate) const STATS_ENV: &str = "ONEPIPELINE_LOOP_STATS";
pub(crate) const STATS_FILE: &str = "loop-stats.json";
static PASSES: AtomicU64 = AtomicU64::new(0);
static STATUSES: AtomicU64 = AtomicU64::new(0);
static PUBLICATIONS: AtomicU64 = AtomicU64::new(0);
static UPSTREAM_READS: AtomicU64 = AtomicU64::new(0);
static RELEASE_ASKS: AtomicU64 = AtomicU64::new(0);
static STORE_BYTES: AtomicU64 = AtomicU64::new(0);
pub(crate) fn pass() {
PASSES.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn statuses_derived() {
STATUSES.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn published() {
PUBLICATIONS.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn upstream_read() {
UPSTREAM_READS.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn release_asked() {
RELEASE_ASKS.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn store_read(bytes: u64) {
STORE_BYTES.fetch_add(bytes, Ordering::Relaxed);
}
#[cfg(test)]
pub(crate) fn store_bytes() -> u64 {
STORE_BYTES.load(Ordering::Relaxed)
}
fn asked() -> bool {
std::env::var_os(STATS_ENV).is_some_and(|value| !value.is_empty())
}
pub(crate) fn flush(paths: &crate::ledger::RunPaths) -> crate::error::Result<()> {
if !asked() {
return Ok(());
}
let document = json!({
"passes": PASSES.load(Ordering::Relaxed),
"statuses": STATUSES.load(Ordering::Relaxed),
"publications": PUBLICATIONS.load(Ordering::Relaxed),
"upstream_reads": UPSTREAM_READS.load(Ordering::Relaxed),
"release_asks": RELEASE_ASKS.load(Ordering::Relaxed),
"store_bytes": STORE_BYTES.load(Ordering::Relaxed),
});
crate::ledger::write_json(&paths.dir.join(STATS_FILE), &document)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn nothing_is_written_when_nobody_asked_for_it() {
std::env::remove_var(STATS_ENV);
assert!(!asked());
std::env::set_var(STATS_ENV, "");
assert!(!asked(), "an empty setting asks for nothing");
std::env::remove_var(STATS_ENV);
let paths = crate::ledger::RunPaths::under(&std::env::temp_dir(), "nobody");
flush(&paths).expect("an unmeasured run writes nothing and cannot fail");
assert!(!paths.dir.join(STATS_FILE).exists());
}
#[test]
fn every_count_a_journey_reads_is_written_under_its_own_name() {
let root =
std::env::temp_dir().join(format!("onepipeline-loopstats-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
let paths = crate::ledger::RunPaths::under(&root, "measured");
paths.create().expect("the run directory");
std::env::set_var(STATS_ENV, "1");
pass();
statuses_derived();
published();
upstream_read();
release_asked();
store_read(7);
flush(&paths).expect("the counts are written");
std::env::remove_var(STATS_ENV);
let written: serde_json::Value = crate::ledger::read_json_opt(&paths.dir.join(STATS_FILE))
.expect("the counts are written");
for name in [
"passes",
"statuses",
"publications",
"upstream_reads",
"release_asks",
"store_bytes",
] {
assert!(
written[name].as_u64().is_some_and(|count| count > 0),
"{name} is not a count in {written}"
);
}
let _ = std::fs::remove_dir_all(&root);
}
}