use crate::entity_context::NodeTelemetryHandle;
use otel_arrow_dfe_config::MetricLevel;
use otel_arrow_dfe_telemetry::error::Error as TelemetryError;
use otel_arrow_dfe_telemetry::instrument::Counter;
use otel_arrow_dfe_telemetry::metrics::{MetricSet, MetricSetSnapshot};
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
use otel_arrow_dfe_telemetry_macros::metric_set;
use std::sync::{Arc, Mutex};
#[metric_set(name = "node.completion_emission")]
#[derive(Debug, Default, Clone)]
pub(crate) struct CompletionEmissionMetrics {
#[metric(name = "notify_ack.routed", unit = "{message}")]
pub notify_ack_routed: Counter<u64>,
#[metric(name = "notify_nack.routed", unit = "{message}")]
pub notify_nack_routed: Counter<u64>,
}
pub(crate) struct CompletionEmissionMetricsState {
metrics: MetricSet<CompletionEmissionMetrics>,
}
impl CompletionEmissionMetricsState {
fn new(telemetry_handle: &NodeTelemetryHandle) -> Self {
Self {
metrics: telemetry_handle.register_metric_set::<CompletionEmissionMetrics>(),
}
}
pub(crate) fn record_notify_ack_routed(&mut self) {
self.metrics.notify_ack_routed.inc();
}
pub(crate) fn record_notify_nack_routed(&mut self) {
self.metrics.notify_nack_routed.inc();
}
pub(crate) fn report(
&mut self,
metrics_reporter: &mut MetricsReporter,
) -> Result<(), TelemetryError> {
metrics_reporter.report(&mut self.metrics)
}
pub(crate) fn snapshot(&self) -> Option<MetricSetSnapshot> {
(!self.metrics.is_empty()).then(|| self.metrics.snapshot())
}
#[cfg(test)]
pub(crate) fn counts(&self) -> (u64, u64) {
(
self.metrics.notify_ack_routed.get(),
self.metrics.notify_nack_routed.get(),
)
}
}
pub(crate) type CompletionEmissionMetricsHandle = Arc<Mutex<CompletionEmissionMetricsState>>;
pub(crate) fn make_completion_emission_metrics(
telemetry_handle: &Option<NodeTelemetryHandle>,
level: MetricLevel,
) -> Option<CompletionEmissionMetricsHandle> {
if level >= MetricLevel::Normal {
telemetry_handle
.as_ref()
.map(|telemetry| Arc::new(Mutex::new(CompletionEmissionMetricsState::new(telemetry))))
} else {
None
}
}