use super::MetricsError;
use super::labels::{ComponentLabels, OwnedGauge, PartitionGauges};
use super::names;
use super::ownership::{SeriesClaim, series_key};
use crate::record::PartitionId;
use metrics::{Counter, Histogram};
use std::collections::HashMap;
use std::sync::Mutex;
use std::time::Duration;
#[derive(Debug)]
pub struct SourceMetrics {
records: Counter,
bytes: Counter,
poll_duration: Histogram,
rebalance_assign: Counter,
rebalance_revoke: Counter,
lanes_active: OwnedGauge,
partition_lag: PartitionGauges,
_claim: Option<SeriesClaim>,
}
impl SourceMetrics {
pub fn new(labels: &ComponentLabels) -> Self {
Self::build(labels, SeriesClaim::claim_or_shadow(Self::key(labels)))
}
pub fn try_new(labels: &ComponentLabels) -> Result<Self, MetricsError> {
let claim = SeriesClaim::try_claim(Self::key(labels))?;
Ok(Self::build(labels, Some(claim)))
}
#[must_use]
pub fn shadow(labels: &ComponentLabels) -> Self {
Self::build(labels, None)
}
fn key(labels: &ComponentLabels) -> String {
series_key("source", labels, "")
}
fn build(labels: &ComponentLabels, claim: Option<SeriesClaim>) -> Self {
let owned = claim.is_some();
SourceMetrics {
records: labels.counter(names::SOURCE_RECORDS_TOTAL),
bytes: labels.counter(names::SOURCE_BYTES_TOTAL),
poll_duration: labels.histogram(names::SOURCE_POLL_DURATION_SECONDS),
rebalance_assign: labels.counter1(
names::SOURCE_REBALANCES_TOTAL,
names::L_EVENT,
"assign",
),
rebalance_revoke: labels.counter1(
names::SOURCE_REBALANCES_TOTAL,
names::L_EVENT,
"revoke",
),
lanes_active: OwnedGauge::new(labels.gauge(names::SOURCE_LANES_ACTIVE), owned),
partition_lag: PartitionGauges {
name: names::SOURCE_LAG_RECORDS,
labels: labels.clone(),
gauges: Mutex::new(HashMap::new()),
owned,
},
_claim: claim,
}
}
#[inline]
pub fn batch(&self, records: u64, bytes: u64) {
self.records.increment(records);
self.bytes.increment(bytes);
}
#[inline]
pub fn poll_duration(&self, d: Duration) {
self.poll_duration.record(d.as_secs_f64());
}
pub fn set_partition_lag(&self, partition: PartitionId, lag: u64) {
self.partition_lag.set(partition, lag as f64);
}
pub fn retain_partitions(&self, keep: &[PartitionId]) {
self.partition_lag.retain(keep);
}
pub fn rebalance_assigned(&self) {
self.rebalance_assign.increment(1);
}
pub fn rebalance_revoked(&self) {
self.rebalance_revoke.increment(1);
}
pub fn set_lanes_active(&self, lanes: usize) {
self.lanes_active.set(lanes as f64);
}
}