use std::collections::BTreeMap;
use std::fs::OpenOptions;
use std::path::PathBuf;
use std::sync::{Mutex, OnceLock};
use fs2::FileExt;
use serde::{Deserialize, Serialize};
use crate::error::ForgeError;
use crate::fsutil::{atomic_write_file, create_dir_all, lock_exclusive_cancellable};
use crate::paths::app_home;
use crate::util::now_secs;
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct CargoTelemetry {
records: BTreeMap<String, CargoRecord>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(default, deny_unknown_fields)]
struct CargoRecord {
component: String,
runs: u64,
build_runs: u64,
cache_hits: u64,
average_build_ms: u64,
last_build_ms: u64,
peak_rss_mib: u64,
last_cargo_jobs: usize,
updated_at: u64,
failures: u64,
cancellations: u64,
corruptions: u64,
estimated_saved_ms: u64,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
pub(crate) struct TelemetrySummary {
pub(crate) samples: u64,
pub(crate) builds: u64,
pub(crate) artifact_hits: u64,
pub(crate) misses: u64,
pub(crate) corruptions: u64,
pub(crate) failures: u64,
pub(crate) cancellations: u64,
pub(crate) estimated_saved_ms: u64,
pub(crate) hit_rate_percent: u64,
}
pub(crate) fn summary() -> TelemetrySummary {
let mut summary = TelemetrySummary::default();
for record in read_document().unwrap_or_default().records.into_values() {
summary.samples = summary.samples.saturating_add(record.runs);
summary.builds = summary.builds.saturating_add(record.build_runs);
summary.artifact_hits = summary.artifact_hits.saturating_add(record.cache_hits);
summary.corruptions = summary.corruptions.saturating_add(record.corruptions);
summary.failures = summary.failures.saturating_add(record.failures);
summary.cancellations = summary.cancellations.saturating_add(record.cancellations);
summary.estimated_saved_ms = summary
.estimated_saved_ms
.saturating_add(record.estimated_saved_ms);
}
summary.misses = summary.builds;
summary.hit_rate_percent = summary
.artifact_hits
.saturating_mul(100)
.checked_div(summary.samples)
.unwrap_or(0);
summary
}
pub(crate) fn estimates() -> BTreeMap<String, CargoEstimate> {
read_document()
.map(|document| {
document
.records
.into_iter()
.filter_map(|(fingerprint, record)| {
(record.average_build_ms > 0 || record.peak_rss_mib > 0).then_some((
fingerprint,
CargoEstimate {
average_build_ms: record.average_build_ms,
memory_mib: memory_claim(record.peak_rss_mib),
},
))
})
.collect()
})
.unwrap_or_default()
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct CargoEstimate {
pub(crate) average_build_ms: u64,
pub(crate) memory_mib: u32,
}
struct CargoSample {
fingerprint: String,
component: String,
build_ms: u64,
cache_hit: bool,
cargo_jobs: usize,
peak_rss_mib: u64,
}
enum CargoUpdate {
Sample(CargoSample),
Outcome {
fingerprint: String,
component: String,
outcome: CargoOutcome,
},
}
enum CargoOutcome {
Failure,
Cancellation,
Corruption,
}
pub(crate) fn record_cargo(
fingerprint: &str,
component: &str,
build_ms: u64,
cache_hit: bool,
cargo_jobs: usize,
peak_rss_mib: u64,
) {
if let Ok(mut pending) = pending_updates().lock() {
pending.push(CargoUpdate::Sample(CargoSample {
fingerprint: fingerprint.to_string(),
component: component.to_string(),
build_ms,
cache_hit,
cargo_jobs,
peak_rss_mib,
}));
}
}
pub(crate) fn record_failure(fingerprint: &str, component: &str) {
record_outcome(fingerprint, component, CargoOutcome::Failure);
}
pub(crate) fn record_cancellation(fingerprint: &str, component: &str) {
record_outcome(fingerprint, component, CargoOutcome::Cancellation);
}
pub(crate) fn record_corruption(fingerprint: &str, component: &str) {
record_outcome(fingerprint, component, CargoOutcome::Corruption);
}
fn record_outcome(fingerprint: &str, component: &str, outcome: CargoOutcome) {
if let Ok(mut pending) = pending_updates().lock() {
pending.push(CargoUpdate::Outcome {
fingerprint: fingerprint.to_string(),
component: component.to_string(),
outcome,
});
}
}
pub(crate) fn flush() {
let updates = pending_updates()
.lock()
.map(|mut pending| std::mem::take(&mut *pending))
.unwrap_or_default();
if updates.is_empty() {
return;
}
let _ = update_document(|document| {
for update in updates {
match update {
CargoUpdate::Sample(sample) => apply_sample(document, sample),
CargoUpdate::Outcome {
fingerprint,
component,
outcome,
} => apply_outcome(document, fingerprint, component, outcome),
}
}
});
}
fn pending_updates() -> &'static Mutex<Vec<CargoUpdate>> {
static PENDING: OnceLock<Mutex<Vec<CargoUpdate>>> = OnceLock::new();
PENDING.get_or_init(|| Mutex::new(Vec::new()))
}
fn apply_sample(document: &mut CargoTelemetry, sample: CargoSample) {
let record = document
.records
.entry(sample.fingerprint)
.or_insert_with(|| CargoRecord {
component: sample.component.clone(),
last_cargo_jobs: sample.cargo_jobs.max(1),
updated_at: now_secs(),
..CargoRecord::default()
});
record.component = sample.component;
record.runs = record.runs.saturating_add(1);
record.cache_hits = record
.cache_hits
.saturating_add(u64::from(sample.cache_hit));
if sample.cache_hit {
record.estimated_saved_ms = record
.estimated_saved_ms
.saturating_add(record.average_build_ms);
} else {
record.build_runs = record.build_runs.saturating_add(1);
record.average_build_ms =
rolling_average(record.average_build_ms, record.build_runs, sample.build_ms);
record.last_build_ms = sample.build_ms;
record.peak_rss_mib = decayed_peak(record.peak_rss_mib, sample.peak_rss_mib);
}
record.last_cargo_jobs = sample.cargo_jobs.max(1);
record.updated_at = now_secs();
}
fn apply_outcome(
document: &mut CargoTelemetry,
fingerprint: String,
component: String,
outcome: CargoOutcome,
) {
let record = document.records.entry(fingerprint).or_default();
record.component = component;
match outcome {
CargoOutcome::Failure => record.failures = record.failures.saturating_add(1),
CargoOutcome::Cancellation => {
record.cancellations = record.cancellations.saturating_add(1);
}
CargoOutcome::Corruption => {
record.corruptions = record.corruptions.saturating_add(1);
}
}
record.updated_at = now_secs();
}
fn decayed_peak(previous: u64, sample: u64) -> u64 {
if previous == 0 || sample >= previous {
sample
} else {
previous.saturating_mul(7).saturating_add(sample) / 8
}
}
fn memory_claim(peak_rss_mib: u64) -> u32 {
let with_headroom = peak_rss_mib.saturating_mul(5).div_ceil(4).max(512);
with_headroom
.div_ceil(512)
.saturating_mul(512)
.min(u64::from(u32::MAX)) as u32
}
pub(crate) fn rolling_average(previous: u64, runs: u64, sample: u64) -> u64 {
if previous == 0 {
return sample;
}
let weight = runs.min(8);
previous
.saturating_mul(weight.saturating_sub(1))
.saturating_add(sample)
/ weight
}
fn update_document(action: impl FnOnce(&mut CargoTelemetry)) -> Result<(), ForgeError> {
let path = telemetry_path();
let parent = path
.parent()
.ok_or_else(|| ForgeError::Config("telemetry path is missing a parent directory".into()))?;
create_dir_all(parent)?;
let lock_path = parent.join("cargo.lock");
let lock = OpenOptions::new()
.create(true)
.truncate(false)
.read(true)
.write(true)
.open(&lock_path)
.map_err(|source| ForgeError::Io {
path: lock_path.clone(),
source,
})?;
lock_exclusive_cancellable(&lock, &lock_path, "telemetry lock")?;
let mut document = read_document().unwrap_or_default();
action(&mut document);
let bytes = serde_json::to_vec_pretty(&document).map_err(|error| {
ForgeError::Parse(format!("failed to serialize Cargo telemetry: {error}"))
})?;
atomic_write_file(&path, &bytes)?;
let _ = FileExt::unlock(&lock);
Ok(())
}
fn read_document() -> Result<CargoTelemetry, ForgeError> {
let path = telemetry_path();
if !path.is_file() {
return Ok(CargoTelemetry::default());
}
let bytes = std::fs::read(&path).map_err(|source| ForgeError::Io {
path: path.clone(),
source,
})?;
serde_json::from_slice(&bytes).map_err(|error| {
ForgeError::Parse(format!(
"Cargo telemetry {} is invalid: {error}",
path.display()
))
})
}
fn telemetry_path() -> PathBuf {
app_home().join("telemetry").join("cargo.json")
}
#[cfg(test)]
mod tests {
use crate::telemetry::{
CargoOutcome, CargoTelemetry, apply_outcome, memory_claim, rolling_average,
};
#[test]
fn rolling_average_adapts_without_unbounded_history() {
assert_eq!(rolling_average(0, 1, 100), 100);
assert_eq!(rolling_average(100, 2, 300), 200);
assert_eq!(rolling_average(800, 99, 0), 700);
assert_eq!(memory_claim(0), 512);
assert_eq!(memory_claim(900), 1536);
}
#[test]
fn outcome_counters_are_persisted_independently_of_success_samples() {
let mut document = CargoTelemetry::default();
apply_outcome(
&mut document,
"fingerprint".into(),
"demo".into(),
CargoOutcome::Failure,
);
apply_outcome(
&mut document,
"fingerprint".into(),
"demo".into(),
CargoOutcome::Cancellation,
);
apply_outcome(
&mut document,
"fingerprint".into(),
"demo".into(),
CargoOutcome::Corruption,
);
let record = &document.records["fingerprint"];
assert_eq!(record.failures, 1);
assert_eq!(record.cancellations, 1);
assert_eq!(record.corruptions, 1);
assert_eq!(record.runs, 0);
}
}