Skip to main content

spate_core/metrics/
queue.rs

1//! Queue-edge handles (`spate_queue_*`).
2
3use super::MetricsError;
4use super::labels::{ComponentLabels, OwnedGauge};
5use super::names;
6use super::ownership::{SeriesClaim, series_key};
7use metrics::{Counter, SharedString};
8
9/// Queue-edge handles (`spate_queue_*`).
10#[derive(Debug)]
11pub struct QueueMetrics {
12    depth: OwnedGauge,
13    full_events: Counter,
14    _claim: Option<SeriesClaim>,
15}
16
17impl QueueMetrics {
18    /// Resolve handles for one queue edge (e.g. `chain->sink/shard-3`) and
19    /// publish its configured capacity.
20    ///
21    /// Claims this edge's series (the labels plus the queue name); a second
22    /// live handle set for the same edge logs and becomes a shadow, counting
23    /// but publishing neither depth nor capacity.
24    pub fn new(labels: &ComponentLabels, queue: &str, capacity: usize) -> Self {
25        let claim = SeriesClaim::claim_or_shadow(Self::key(labels, queue));
26        Self::build(labels, queue, capacity, claim)
27    }
28
29    /// Resolve handles for one queue edge, failing when another live handle
30    /// set already owns it. The pipeline builder's path.
31    ///
32    /// # Errors
33    ///
34    /// [`MetricsError::DuplicateSeries`] on a collision.
35    pub fn try_new(
36        labels: &ComponentLabels,
37        queue: &str,
38        capacity: usize,
39    ) -> Result<Self, MetricsError> {
40        let claim = SeriesClaim::try_claim(Self::key(labels, queue))?;
41        Ok(Self::build(labels, queue, capacity, Some(claim)))
42    }
43
44    fn key(labels: &ComponentLabels, queue: &str) -> String {
45        series_key("queue", labels, &format!("queue={queue}"))
46    }
47
48    fn build(
49        labels: &ComponentLabels,
50        queue: &str,
51        capacity: usize,
52        claim: Option<SeriesClaim>,
53    ) -> Self {
54        let owned = claim.is_some();
55        let queue: SharedString = queue.to_owned().into();
56        OwnedGauge::new(
57            labels.gauge1(names::QUEUE_CAPACITY, names::L_QUEUE, queue.clone()),
58            owned,
59        )
60        .set(capacity as f64);
61        QueueMetrics {
62            depth: OwnedGauge::new(
63                labels.gauge1(names::QUEUE_DEPTH, names::L_QUEUE, queue.clone()),
64                owned,
65            ),
66            full_events: labels.counter1(names::QUEUE_FULL_EVENTS_TOTAL, names::L_QUEUE, queue),
67            _claim: claim,
68        }
69    }
70
71    /// Set the current queue depth.
72    #[inline]
73    pub fn set_depth(&self, depth: usize) {
74        self.depth.set(depth as f64);
75    }
76
77    /// Count `try_send` rejections.
78    #[inline]
79    pub fn full_events(&self, n: u64) {
80        self.full_events.increment(n);
81    }
82}