use super::MetricsError;
use super::labels::{ComponentLabels, OwnedGauge};
use super::names;
use super::ownership::{SeriesClaim, series_key};
use metrics::{Counter, Histogram};
use std::sync::Arc;
use std::time::Duration;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum AcquireReason {
Create,
Reclaimed,
Expired,
Reassigned,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum RevocationOutcome {
Requested,
Drained,
Forced,
Cancelled,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum SplitLossReason {
Fenced,
Starved,
Revoked,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum WriteOutcome {
Ok,
Conflict,
Error,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum ReplanOutcome {
Ok,
Error,
Noop,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum StoreOp {
Get,
Put,
Delete,
List,
Watch,
}
#[derive(Clone, Debug)]
pub struct CoordinationMetrics {
splits_owned: OwnedGauge,
splits_completed: OwnedGauge,
splits_quarantined: OwnedGauge,
live_workers: OwnedGauge,
leader: OwnedGauge,
idle: OwnedGauge,
splits_draining: OwnedGauge,
acquired_create: Counter,
acquired_reclaimed: Counter,
acquired_expired: Counter,
acquired_reassigned: Counter,
lost_fenced: Counter,
lost_starved: Counter,
lost_revoked: Counter,
releases: Counter,
revocations_requested: Counter,
revocations_drained: Counter,
revocations_forced: Counter,
revocations_cancelled: Counter,
splits_planned: Counter,
replans_ok: Counter,
replans_error: Counter,
replans_noop: Counter,
split_failures: Counter,
quarantines: Counter,
writes_ok: Counter,
writes_conflict: Counter,
writes_error: Counter,
write_duration: Histogram,
replan_duration: Histogram,
reconcile_duration: Histogram,
store_op_get: Histogram,
store_op_put: Histogram,
store_op_delete: Histogram,
store_op_list: Histogram,
store_op_watch: Histogram,
drain_duration: Histogram,
assignment_latency: Histogram,
_claim: Option<Arc<SeriesClaim>>,
}
impl CoordinationMetrics {
pub fn new(labels: &ComponentLabels) -> Self {
let claim = SeriesClaim::claim_or_shadow(Self::key(labels));
Self::build(labels, claim.map(Arc::new))
}
pub fn try_new(labels: &ComponentLabels) -> Result<Self, MetricsError> {
let claim = SeriesClaim::try_claim(Self::key(labels))?;
Ok(Self::build(labels, Some(Arc::new(claim))))
}
fn key(labels: &ComponentLabels) -> String {
series_key("coordination", labels, "")
}
fn build(labels: &ComponentLabels, claim: Option<Arc<SeriesClaim>>) -> Self {
let owned = claim.is_some();
let gauge = |name| OwnedGauge::new(labels.gauge(name), owned);
let acquired = |reason| {
labels.counter1(
names::COORDINATION_ACQUISITIONS_TOTAL,
names::L_REASON,
reason,
)
};
let lost = |reason| {
labels.counter1(
names::COORDINATION_SPLIT_LOSSES_TOTAL,
names::L_REASON,
reason,
)
};
let replans =
|outcome| labels.counter1(names::COORDINATION_REPLANS_TOTAL, names::L_OUTCOME, outcome);
let revocations = |outcome| {
labels.counter1(
names::COORDINATION_REVOCATIONS_TOTAL,
names::L_OUTCOME,
outcome,
)
};
let writes =
|outcome| labels.counter1(names::COORDINATION_WRITES_TOTAL, names::L_OUTCOME, outcome);
let store_op = |op| {
labels.histogram1(
names::COORDINATION_STORE_OP_DURATION_SECONDS,
names::L_OP,
op,
)
};
CoordinationMetrics {
splits_owned: gauge(names::COORDINATION_SPLITS_OWNED),
splits_completed: gauge(names::COORDINATION_SPLITS_COMPLETED),
splits_quarantined: gauge(names::COORDINATION_SPLITS_QUARANTINED),
live_workers: gauge(names::COORDINATION_LIVE_WORKERS),
leader: gauge(names::COORDINATION_LEADER),
idle: gauge(names::COORDINATION_IDLE),
splits_draining: gauge(names::COORDINATION_SPLITS_DRAINING),
acquired_create: acquired("create"),
acquired_reclaimed: acquired("reclaimed"),
acquired_expired: acquired("expired"),
acquired_reassigned: acquired("reassigned"),
lost_fenced: lost("fenced"),
lost_starved: lost("starved"),
lost_revoked: lost("revoked"),
releases: labels.counter(names::COORDINATION_RELEASES_TOTAL),
revocations_requested: revocations("requested"),
revocations_drained: revocations("drained"),
revocations_forced: revocations("forced"),
revocations_cancelled: revocations("cancelled"),
splits_planned: labels.counter(names::COORDINATION_SPLITS_PLANNED_TOTAL),
replans_ok: replans("ok"),
replans_error: replans("error"),
replans_noop: replans("noop"),
split_failures: labels.counter(names::COORDINATION_SPLIT_FAILURES_TOTAL),
quarantines: labels.counter(names::COORDINATION_QUARANTINES_TOTAL),
writes_ok: writes("ok"),
writes_conflict: writes("conflict"),
writes_error: writes("error"),
write_duration: labels.histogram(names::COORDINATION_WRITE_DURATION_SECONDS),
replan_duration: labels.histogram(names::COORDINATION_REPLAN_DURATION_SECONDS),
reconcile_duration: labels.histogram(names::COORDINATION_RECONCILE_DURATION_SECONDS),
store_op_get: store_op("get"),
store_op_put: store_op("put"),
store_op_delete: store_op("delete"),
store_op_list: store_op("list"),
store_op_watch: store_op("watch"),
drain_duration: labels.histogram(names::COORDINATION_DRAIN_DURATION_SECONDS),
assignment_latency: labels.histogram(names::COORDINATION_ASSIGNMENT_LATENCY_SECONDS),
_claim: claim,
}
}
pub fn set_splits_owned(&self, owned: usize) {
self.splits_owned.set(owned as f64);
}
pub fn set_splits_completed(&self, completed: usize) {
self.splits_completed.set(completed as f64);
}
pub fn set_splits_quarantined(&self, quarantined: usize) {
self.splits_quarantined.set(quarantined as f64);
}
pub fn set_live_workers(&self, workers: usize) {
self.live_workers.set(workers as f64);
}
pub fn set_leader(&self, leader: bool) {
self.leader.set(if leader { 1.0 } else { 0.0 });
}
pub fn set_idle(&self, idle: bool) {
self.idle.set(if idle { 1.0 } else { 0.0 });
}
pub fn acquired(&self, reason: AcquireReason) {
match reason {
AcquireReason::Create => self.acquired_create.increment(1),
AcquireReason::Reclaimed => self.acquired_reclaimed.increment(1),
AcquireReason::Expired => self.acquired_expired.increment(1),
AcquireReason::Reassigned => self.acquired_reassigned.increment(1),
}
}
pub fn revocation(&self, outcome: RevocationOutcome) {
match outcome {
RevocationOutcome::Requested => self.revocations_requested.increment(1),
RevocationOutcome::Drained => self.revocations_drained.increment(1),
RevocationOutcome::Forced => self.revocations_forced.increment(1),
RevocationOutcome::Cancelled => self.revocations_cancelled.increment(1),
}
}
pub fn drain_duration(&self, d: Duration) {
self.drain_duration.record(d.as_secs_f64());
}
pub fn assignment_latency(&self, d: Duration) {
self.assignment_latency.record(d.as_secs_f64());
}
pub fn set_splits_draining(&self, draining: usize) {
self.splits_draining.set(draining as f64);
}
pub fn lost(&self, reason: SplitLossReason) {
match reason {
SplitLossReason::Fenced => self.lost_fenced.increment(1),
SplitLossReason::Starved => self.lost_starved.increment(1),
SplitLossReason::Revoked => self.lost_revoked.increment(1),
}
}
pub fn released(&self, splits: u64) {
self.releases.increment(splits);
}
pub fn planned(&self, splits: u64) {
self.splits_planned.increment(splits);
}
pub fn replan(&self, outcome: ReplanOutcome, d: Duration) {
match outcome {
ReplanOutcome::Ok => self.replans_ok.increment(1),
ReplanOutcome::Error => self.replans_error.increment(1),
ReplanOutcome::Noop => self.replans_noop.increment(1),
}
self.replan_duration.record(d.as_secs_f64());
}
pub fn failed(&self) {
self.split_failures.increment(1);
}
pub fn quarantined(&self) {
self.quarantines.increment(1);
}
pub fn write(&self, outcome: WriteOutcome, d: Duration) {
match outcome {
WriteOutcome::Ok => self.writes_ok.increment(1),
WriteOutcome::Conflict => self.writes_conflict.increment(1),
WriteOutcome::Error => self.writes_error.increment(1),
}
self.write_duration.record(d.as_secs_f64());
}
pub fn reconcile(&self, d: Duration) {
self.reconcile_duration.record(d.as_secs_f64());
}
pub fn store_op(&self, op: StoreOp, d: Duration) {
let h = match op {
StoreOp::Get => &self.store_op_get,
StoreOp::Put => &self.store_op_put,
StoreOp::Delete => &self.store_op_delete,
StoreOp::List => &self.store_op_list,
StoreOp::Watch => &self.store_op_watch,
};
h.record(d.as_secs_f64());
}
}