azums_core/backend/
observability.rs1use 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 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}