Skip to main content

etdl_core/
telemetry.rs

1use std::fmt;
2
3#[derive(Debug)]
4pub struct Error {
5    message: String,
6}
7
8impl Error {
9    pub fn new(msg: impl Into<String>) -> Self {
10        Error {
11            message: msg.into(),
12        }
13    }
14}
15
16impl fmt::Display for Error {
17    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
18        write!(f, "{}", self.message)
19    }
20}
21
22impl std::error::Error for Error {}
23
24impl From<String> for Error {
25    fn from(s: String) -> Self {
26        Error::new(s)
27    }
28}
29
30impl From<serde_json::Error> for Error {
31    fn from(e: serde_json::Error) -> Self {
32        Error::new(format!("serialization error: {}", e))
33    }
34}
35
36impl From<crate::publisher::PublishError> for Error {
37    fn from(e: crate::publisher::PublishError) -> Self {
38        Error::new(format!("publish error: {}", e))
39    }
40}
41
42pub type WorkflowError = Error;
43
44pub fn attach_node_span_attribute(node_id: &str) {
45    eprintln!("[etdl.telemetry] span attribute etdl.node.id={}", node_id);
46}
47
48pub fn emit_anomaly_event(
49    node_id: &str,
50    outcome: &str,
51    declared_probability: f64,
52    observed_frequency: f64,
53) {
54    eprintln!(
55        "[etdl.telemetry] SLA ANOMALY | node={} outcome={} declared={:.6} observed={:.6} deviation={:.6}",
56        node_id,
57        outcome,
58        declared_probability,
59        observed_frequency,
60        (observed_frequency - declared_probability).abs()
61    );
62}
63
64pub fn inject_traceparent(message_type: &str) -> String {
65    let trace_id = generate_trace_id();
66    let span_id = generate_span_id();
67    let traceparent = format!("00-{}-{}-01", trace_id, span_id);
68
69    eprintln!(
70        "[etdl.telemetry] inject traceparent into {} message: {}",
71        message_type, traceparent
72    );
73
74    traceparent
75}
76
77fn generate_trace_id() -> String {
78    let mut bytes = [0u8; 16];
79    fill_random(&mut bytes);
80    hex(&bytes)
81}
82
83fn generate_span_id() -> String {
84    let mut bytes = [0u8; 8];
85    fill_random(&mut bytes);
86    hex(&bytes)
87}
88
89/// Fill `buf` with OS randomness; fall back to a time+counter mix if the OS
90/// source is unavailable (rare; format remains valid and unique in-process).
91fn fill_random(buf: &mut [u8]) {
92    if getrandom::getrandom(buf).is_err() {
93        use std::sync::atomic::{AtomicU64, Ordering};
94        use std::time::{SystemTime, UNIX_EPOCH};
95        static COUNTER: AtomicU64 = AtomicU64::new(0);
96        let counter = COUNTER.fetch_add(1, Ordering::Relaxed);
97        let nanos = SystemTime::now()
98            .duration_since(UNIX_EPOCH)
99            .map(|d| d.as_nanos() as u64)
100            .unwrap_or(0);
101        let mut seed = nanos ^ counter.wrapping_mul(0x9E3779B97F4A7C15);
102        for slot in buf.iter_mut() {
103            seed ^= seed << 13;
104            seed ^= seed >> 7;
105            seed ^= seed << 17;
106            *slot = seed as u8;
107        }
108    }
109}
110
111fn hex(bytes: &[u8]) -> String {
112    let mut s = String::with_capacity(bytes.len() * 2);
113    for b in bytes {
114        s.push_str(&format!("{:02x}", b));
115    }
116    s
117}
118
119#[cfg(test)]
120mod tests {
121    use super::*;
122
123    #[test]
124    fn traceparent_is_well_formed() {
125        let tp = inject_traceparent("Test");
126        // version-trace-span-flags
127        let parts: Vec<&str> = tp.split('-').collect();
128        assert_eq!(parts.len(), 4, "got {}", tp);
129        assert_eq!(parts[0], "00");
130        assert_eq!(parts[1].len(), 32, "trace-id must be 32 hex, got {}", tp);
131        assert_eq!(parts[2].len(), 16, "span-id must be 16 hex, got {}", tp);
132        assert_eq!(parts[3], "01");
133        assert!(parts[1].chars().all(|c| c.is_ascii_hexdigit()));
134        assert!(parts[2].chars().all(|c| c.is_ascii_hexdigit()));
135    }
136
137    #[test]
138    fn trace_ids_are_not_all_zero() {
139        let a = generate_trace_id();
140        let b = generate_span_id();
141        assert_ne!(a, "00000000000000000000000000000000");
142        assert_ne!(b, "0000000000000000");
143    }
144
145    #[test]
146    fn hex_lengths() {
147        assert_eq!(hex(&[0u8; 16]).len(), 32);
148        assert_eq!(hex(&[0u8; 8]).len(), 16);
149    }
150}