- 1
//! Observability, metrics collection, distributed trace propagation, and Dead-Letter Queue (DLQ) tracking. - 2
- 3
use std::sync::Arc; - 4
use std::sync::atomic::{AtomicU64, Ordering}; - 5
- 6
use chrono::{DateTime, Utc}; - 7
use serde::{Deserialize, Serialize}; - 8
- 9
use crate::envelope::{MessageEnvelope, TraceContext}; - 10
- 11
/// Snapshot of bus operational metrics. - 12
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] - 13
pub struct BusMetricsSnapshot { - 14
pub published_count: u64, - 15
pub received_count: u64, - 16
pub bytes_published: u64, - 17
pub bytes_received: u64, - 18
pub dead_letter_count: u64, - 19
pub active_queue_lag: u64, - 20
} - 21
- 22
/// Thread-safe lock-free metrics counters for the messaging fabric. - 23
#[derive(Debug, Default)] - 24
pub struct BusMetrics { - 25
published_count: AtomicU64, - 26
received_count: AtomicU64, - 27
bytes_published: AtomicU64, - 28
bytes_received: AtomicU64, - 29
dead_letter_count: AtomicU64, - 30
active_queue_lag: AtomicU64, - 31
} - 32
- 33
impl BusMetrics { - 34
pub fn new() -> Arc<Self> { - 35
Arc::new(Self::default()) - 36
} - 37
- 38
pub fn record_published(&self, bytes: usize) { - 39
self.published_count.fetch_add(1, Ordering::Relaxed); - 40
self.bytes_published - 41
.fetch_add(bytes as u64, Ordering::Relaxed); - 42
} - 43
- 44
pub fn record_received(&self, bytes: usize) { - 45
self.received_count.fetch_add(1, Ordering::Relaxed); - 46
self.bytes_received - 47
.fetch_add(bytes as u64, Ordering::Relaxed); - 48
} - 49
- 50
pub fn record_dead_letter(&self) { - 51
self.dead_letter_count.fetch_add(1, Ordering::Relaxed); - 52
} - 53
- 54
pub fn update_queue_lag(&self, lag: u64) { - 55
self.active_queue_lag.store(lag, Ordering::Relaxed); - 56
} - 57
- 58
pub fn snapshot(&self) -> BusMetricsSnapshot { - 59
BusMetricsSnapshot { - 60
published_count: self.published_count.load(Ordering::Relaxed), - 61
received_count: self.received_count.load(Ordering::Relaxed), - 62
bytes_published: self.bytes_published.load(Ordering::Relaxed), - 63
bytes_received: self.bytes_received.load(Ordering::Relaxed), - 64
dead_letter_count: self.dead_letter_count.load(Ordering::Relaxed), - 65
active_queue_lag: self.active_queue_lag.load(Ordering::Relaxed), - 66
} - 67
} - 68
} - 69
- 70
/// Dead-Letter Queue event representation capturing poison pills and exhausted retries. - 71
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] - 72
pub struct DeadLetterEvent { - 73
pub dlq_id: String, - 74
pub envelope_id: String, - 75
pub subject: String, - 76
pub worker_id: String, - 77
pub attempts: u32, - 78
pub failure_reason: String, - 79
pub failed_at: DateTime<Utc>, - 80
pub envelope: MessageEnvelope, - 81
} - 82
- 83
impl DeadLetterEvent { - 84
pub fn new( - 85
envelope: MessageEnvelope, - 86
subject: impl Into<String>, - 87
worker_id: impl Into<String>, - 88
attempts: u32, - 89
failure_reason: impl Into<String>, - 90
) -> Self { - 91
Self { - 92
dlq_id: uuid::Uuid::now_v7().to_string(), - 93
envelope_id: envelope.id.clone(), - 94
subject: subject.into(), - 95
worker_id: worker_id.into(), - 96
attempts, - 97
failure_reason: failure_reason.into(), - 98
failed_at: Utc::now(), - 99
envelope, - 100
} - 101
} - 102
} - 103
- 104
/// Inject or propagate W3C distributed trace context across execution boundaries. - 105
pub fn inject_trace_context(envelope: &mut MessageEnvelope, parent_context: Option<&TraceContext>) { - 106
if let Some(parent) = parent_context { - 107
envelope.trace = parent.child_span(); - 108
} - 109
} - 110
- 111
/// Extract trace context from an incoming message envelope. - 112
pub fn extract_trace_context(envelope: &MessageEnvelope) -> TraceContext { - 113
envelope.trace.clone() - 114
} - 115
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.