use std::sync::Arc;
#[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,
}
}
}
pub trait ObservationSink: Send + Sync {
fn emit(&self, observation: &ReliabilityObservation);
}
#[derive(Debug, Default, Clone)]
pub struct NoopSink;
impl ObservationSink for NoopSink {
fn emit(&self, _observation: &ReliabilityObservation) {}
}
#[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);
}
}
}
}
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());
}
}