bot-forge 1.0.2

Rust CLI for installing agent skills and developer tools from configurable forms.
Documentation
//! Persistent Cargo build telemetry and deterministic scheduling estimates.
//!
//! Execution records samples in memory and flushes them in batches. The CLI composes [`summary`]
//! with storage status, keeping persistent observability independent of state correctness and
//! execution policy.

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
}

/// Load per-fingerprint duration and memory hints for scheduler construction.
///
/// Missing or invalid telemetry returns an empty map so historical observations never become a
/// correctness prerequisite.
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,
}

/// Queue one successful artifact lookup or build sample for a later batch flush.
///
/// Lock poisoning drops the observation; telemetry failures do not fail installation.
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,
        });
    }
}

/// Persist all currently queued observations as one locked, atomic update.
///
/// The queue is drained before I/O to keep producers unblocked. Persistence failure therefore
/// loses that telemetry batch but cannot change the execution 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();
}

/// Retain new peaks immediately and decay old peaks by one eighth per lower sample.
fn decayed_peak(previous: u64, sample: u64) -> u64 {
    if previous == 0 || sample >= previous {
        sample
    } else {
        previous.saturating_mul(7).saturating_add(sample) / 8
    }
}

/// Add 25% headroom and round the result up to a 512 MiB scheduling quantum.
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
}

/// Update a bounded-history mean whose effective window grows to at most eight builds.
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);
    }
}