etdl-core 0.2.1

ETDL runtime: BranchMonitor, retry policies, SLA anomaly detection, chaos injection, and telemetry for reliability-aware event-driven services
Documentation
//! Lightweight runtime evidence collection.
//!
//! The runtime collects **immutable observations** for later offline analysis.
//! It does NOT run Bayesian inference, query reliability databases, run Monte
//! Carlo, or call AI — those are analysis-time concerns (see the
//! `etdl-reliability` crate). This keeps the runtime service-local and
//! lightweight, per the ETDL architecture.

use std::sync::Arc;

/// An immutable reliability observation: what happened, when, and under what
/// conditions. No sensitive payload data by default.
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct ReliabilityObservation {
    pub id: String,
    pub event: String,
    pub timestamp: String,
    pub service: Option<String>,
    pub operation: Option<String>,
    pub environment: Option<String>,
    pub deployment: Option<String>,
    pub outcome: String,
    pub conditions: Vec<String>,
    pub duration_ms: Option<u64>,
    pub trace_id: Option<String>,
}

impl ReliabilityObservation {
    pub fn new(id: impl Into<String>, event: impl Into<String>) -> Self {
        ReliabilityObservation {
            id: id.into(),
            event: event.into(),
            timestamp: String::new(),
            service: None,
            operation: None,
            environment: None,
            deployment: None,
            outcome: String::new(),
            conditions: Vec::new(),
            duration_ms: None,
            trace_id: None,
        }
    }
}

/// A destination for observations. Implementations may write JSON Lines, CSV,
/// OpenTelemetry, a database adapter, or a message stream. These are optional
/// adapters; the runtime does not require any of them.
pub trait ObservationSink: Send + Sync {
    fn emit(&self, observation: &ReliabilityObservation);
}

/// A sink that drops observations (default). Enables "no telemetry configured".
#[derive(Debug, Default, Clone)]
pub struct NoopSink;

impl ObservationSink for NoopSink {
    fn emit(&self, _observation: &ReliabilityObservation) {}
}

/// A sink that writes observations as JSON Lines to a `Vec<String>` for tests
/// and simple capture.
#[derive(Debug, Default)]
pub struct CapturingSink {
    lines: std::sync::Mutex<Vec<String>>,
}

impl CapturingSink {
    pub fn new() -> Self {
        Self::default()
    }

    pub fn lines(&self) -> Vec<String> {
        self.lines.lock().map(|g| g.clone()).unwrap_or_default()
    }
}

impl ObservationSink for CapturingSink {
    fn emit(&self, observation: &ReliabilityObservation) {
        if let Ok(line) = serde_json::to_string(observation) {
            if let Ok(mut g) = self.lines.lock() {
                g.push(line);
            }
        }
    }
}

/// Shared sink handle used by the runtime.
pub type SharedSink = Arc<dyn ObservationSink>;

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

    #[test]
    fn capturing_sink_records() {
        let sink = CapturingSink::new();
        let obs = ReliabilityObservation::new("obs-1", "failure.network.timeout");
        sink.emit(&obs);
        let lines = sink.lines();
        assert_eq!(lines.len(), 1);
        assert!(lines[0].contains("failure.network.timeout"));
    }

    #[test]
    fn observation_is_plain_data() {
        let obs = ReliabilityObservation::new("obs-1", "failure.network.timeout");
        assert_eq!(obs.event, "failure.network.timeout");
        assert!(obs.duration_ms.is_none());
    }
}