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
89fn 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 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}