af-workflow 0.4.0

Spec-driven workflow chassis: typed node expressions composed into a branched DAG. Port of agent_core/workflow.
Documentation
//! Dependency-free scheduler metrics and readiness.

use std::sync::atomic::{AtomicI64, AtomicU64, Ordering};

const READY_MISSED_PASSES: i64 = 3;

/// Bounded-label Prometheus metrics and readiness for a supervisor.
pub struct RuntimeMetrics {
    prefix: String,
    passes_total: AtomicU64,
    pass_failures_total: AtomicU64,
    claimed_total: AtomicU64,
    evaluation_failures_total: AtomicU64,
    overruns_total: AtomicU64,
    backlog_events_total: AtomicU64,
    last_pass_unix_ms: AtomicI64,
    last_pass_duration_ms: AtomicU64,
    active_instances: AtomicI64,
    due_instances: AtomicI64,
    stalled_instances: AtomicI64,
    started_unix_ms: AtomicI64,
}

impl RuntimeMetrics {
    /// Metrics named with `prefix`, initialised at `now_ms`.
    pub fn new(now_ms: i64, prefix: impl Into<String>) -> Self {
        Self {
            prefix: prefix.into(),
            passes_total: AtomicU64::new(0),
            pass_failures_total: AtomicU64::new(0),
            claimed_total: AtomicU64::new(0),
            evaluation_failures_total: AtomicU64::new(0),
            overruns_total: AtomicU64::new(0),
            backlog_events_total: AtomicU64::new(0),
            last_pass_unix_ms: AtomicI64::new(0),
            last_pass_duration_ms: AtomicU64::new(0),
            active_instances: AtomicI64::new(0),
            due_instances: AtomicI64::new(0),
            stalled_instances: AtomicI64::new(0),
            started_unix_ms: AtomicI64::new(now_ms),
        }
    }

    /// Metric name prefix.
    pub fn prefix(&self) -> &str {
        &self.prefix
    }

    /// Record one completed supervisor pass.
    pub fn record_pass(&self, now_ms: i64, duration_ms: u64, claimed: u64, failed: u64) {
        self.passes_total.fetch_add(1, Ordering::Relaxed);
        self.claimed_total.fetch_add(claimed, Ordering::Relaxed);
        self.evaluation_failures_total
            .fetch_add(failed, Ordering::Relaxed);
        self.last_pass_unix_ms.store(now_ms, Ordering::Relaxed);
        self.last_pass_duration_ms
            .store(duration_ms, Ordering::Relaxed);
    }

    /// Count a pass that failed as a whole.
    pub fn record_pass_failure(&self) {
        self.pass_failures_total.fetch_add(1, Ordering::Relaxed);
    }

    /// Count a pass that exceeded its interval.
    pub fn record_overrun(&self) {
        self.overruns_total.fetch_add(1, Ordering::Relaxed);
    }

    /// Count a pass that left due work unclaimed.
    pub fn record_backlog(&self) {
        self.backlog_events_total.fetch_add(1, Ordering::Relaxed);
    }

    /// Set the instance gauges.
    pub fn set_gauges(&self, active: i64, due: i64, stalled: i64) {
        self.active_instances.store(active, Ordering::Relaxed);
        self.due_instances.store(due, Ordering::Relaxed);
        self.stalled_instances.store(stalled, Ordering::Relaxed);
    }

    /// Readiness derived from the age of the last pass.
    pub fn readiness(&self, now_ms: i64, interval_secs: u64) -> Readiness {
        let active = self.active_instances.load(Ordering::Relaxed);
        let last = self.last_pass_unix_ms.load(Ordering::Relaxed);
        let reference = if last > 0 {
            last
        } else {
            self.started_unix_ms.load(Ordering::Relaxed)
        };
        let age_ms = now_ms.saturating_sub(reference).max(0);
        let budget_ms = (interval_secs as i64)
            .saturating_mul(1_000)
            .saturating_mul(READY_MISSED_PASSES);
        let ready = active == 0 || age_ms <= budget_ms;
        Readiness {
            ready,
            reason: (!ready).then(|| {
                format!(
                    "{active} active instance(s) but no completed pass in {}s",
                    age_ms / 1_000
                )
            }),
            last_pass_age_secs: age_ms / 1_000,
            active,
        }
    }

    /// Text exposition format.
    pub fn render_prometheus(&self, now_ms: i64) -> String {
        let last = self.last_pass_unix_ms.load(Ordering::Relaxed);
        let age = if last > 0 {
            now_ms.saturating_sub(last).max(0) / 1_000
        } else {
            -1
        };
        let values = [
            (
                "passes_total",
                "counter",
                self.passes_total.load(Ordering::Relaxed) as i64,
            ),
            (
                "pass_failures_total",
                "counter",
                self.pass_failures_total.load(Ordering::Relaxed) as i64,
            ),
            (
                "claimed_total",
                "counter",
                self.claimed_total.load(Ordering::Relaxed) as i64,
            ),
            (
                "evaluation_failures_total",
                "counter",
                self.evaluation_failures_total.load(Ordering::Relaxed) as i64,
            ),
            (
                "overruns_total",
                "counter",
                self.overruns_total.load(Ordering::Relaxed) as i64,
            ),
            (
                "backlog_events_total",
                "counter",
                self.backlog_events_total.load(Ordering::Relaxed) as i64,
            ),
            (
                "last_pass_duration_ms",
                "gauge",
                self.last_pass_duration_ms.load(Ordering::Relaxed) as i64,
            ),
            ("last_pass_age_seconds", "gauge", age),
            (
                "active_instances",
                "gauge",
                self.active_instances.load(Ordering::Relaxed),
            ),
            (
                "due_instances",
                "gauge",
                self.due_instances.load(Ordering::Relaxed),
            ),
            (
                "stalled_instances",
                "gauge",
                self.stalled_instances.load(Ordering::Relaxed),
            ),
        ];
        values
            .into_iter()
            .map(|(suffix, kind, value)| {
                let name = format!("{}_{suffix}", self.prefix);
                format!("# TYPE {name} {kind}\n{name} {value}\n")
            })
            .collect()
    }
}

/// Readiness probe result.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Readiness {
    /// Whether the supervisor passed recently enough.
    pub ready: bool,
    /// Human-readable reason.
    pub reason: Option<String>,
    /// Seconds since the last pass.
    pub last_pass_age_secs: i64,
    /// Active instances.
    pub active: i64,
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn readiness_detects_an_idle_active_runtime() {
        let metrics = RuntimeMetrics::new(0, "jobs");
        metrics.set_gauges(2, 2, 0);
        assert!(!metrics.readiness(31_000, 10).ready);
        metrics.record_pass(30_000, 25, 2, 1);
        assert!(metrics.readiness(31_000, 10).ready);
    }

    #[test]
    fn prometheus_output_is_namespaced_and_complete() {
        let metrics = RuntimeMetrics::new(0, "jobs");
        metrics.set_gauges(2, 1, 1);
        metrics.record_pass_failure();
        metrics.record_overrun();
        metrics.record_backlog();
        let output = metrics.render_prometheus(5_000);
        assert!(output.contains("jobs_active_instances 2"));
        assert!(output.contains("jobs_pass_failures_total 1"));
        assert!(output.contains("jobs_overruns_total 1"));
        assert!(output.contains("jobs_backlog_events_total 1"));
        assert_eq!(metrics.prefix(), "jobs");
    }
}