Skip to main content

azums_core/backend/
observability.rs

1use crate::model::Job;
2use async_trait::async_trait;
3use chrono::{DateTime, Utc};
4use serde::{Deserialize, Serialize};
5use std::collections::BTreeMap;
6use uuid::Uuid;
7
8#[derive(Debug, Clone, Serialize, Deserialize)]
9pub struct JobObservationEvent {
10    pub at: DateTime<Utc>,
11    pub job_id: Uuid,
12    pub attempt: Option<i32>,
13    pub worker_id: Option<String>,
14    pub queue: String,
15    pub duration_ms: Option<i32>,
16    pub status: String,
17    pub retry_count: i32,
18    pub error: Option<String>,
19    pub trace_id: Option<String>,
20}
21
22impl JobObservationEvent {
23    /// Returns OpenTelemetry-compatible span attributes for this job lifecycle observation.
24    pub fn span_attributes(&self) -> BTreeMap<&'static str, String> {
25        let mut attrs = BTreeMap::new();
26        attrs.insert("azums.job_id", self.job_id.to_string());
27        attrs.insert("azums.queue", self.queue.clone());
28        attrs.insert("azums.status", self.status.clone());
29        attrs.insert("azums.retry_count", self.retry_count.to_string());
30
31        if let Some(attempt) = self.attempt {
32            attrs.insert("azums.attempt", attempt.to_string());
33        }
34        if let Some(worker_id) = &self.worker_id {
35            attrs.insert("azums.worker_id", worker_id.clone());
36        }
37        if let Some(duration_ms) = self.duration_ms {
38            attrs.insert("azums.duration_ms", duration_ms.to_string());
39        }
40        if let Some(error) = &self.error {
41            attrs.insert("azums.error", error.clone());
42        }
43        if let Some(trace_id) = &self.trace_id {
44            attrs.insert("trace_id", trace_id.clone());
45        }
46
47        attrs
48    }
49}
50
51#[derive(Debug, Clone, Serialize, Deserialize)]
52pub struct JobExplanation {
53    pub job_id: Uuid,
54    pub job_type: String,
55    pub queue: String,
56    pub status: String,
57    pub retry_count: i32,
58    pub last_worker_id: Option<String>,
59    pub last_error: Option<String>,
60    pub trace_id: Option<String>,
61    pub events: Vec<JobObservationEvent>,
62    pub summary: String,
63}
64
65#[derive(Debug, Clone, Serialize, Deserialize)]
66pub struct QueueMetrics {
67    pub at: DateTime<Utc>,
68    pub queue: String,
69    pub jobs_total: u64,
70    pub jobs_completed: u64,
71    pub jobs_failed: u64,
72    pub jobs_retried: u64,
73    pub jobs_dlq: u64,
74    pub queue_depth: u64,
75    pub execution_latency_ms_avg: f64,
76    pub claim_latency_ms_avg: f64,
77    pub retry_latency_ms_avg: f64,
78    pub worker_count: u64,
79}
80
81#[async_trait]
82pub trait ObservabilityBackend: Send + Sync {
83    async fn explain_job(&self, job_id: Uuid) -> anyhow::Result<Option<JobExplanation>>;
84
85    async fn queue_metrics(&self, queue: Option<&str>) -> anyhow::Result<Vec<QueueMetrics>>;
86}
87
88pub fn trace_id_from_job(job: &Job) -> Option<String> {
89    job.payload
90        .get("trace_id")
91        .and_then(|value| value.as_str())
92        .map(ToString::to_string)
93        .or_else(|| {
94            job.payload
95                .get("metadata")
96                .and_then(|value| value.get("trace_id"))
97                .and_then(|value| value.as_str())
98                .map(ToString::to_string)
99        })
100}