use std::{sync::Arc, time::Duration};
use async_trait::async_trait;
use libdd_shared_runtime::{BasicRuntime, SharedRuntime, SharedRuntimeError, Worker};
use crate::{core::telemetry, span_exporter::QueueMetricsFetcher, TraceRegistry};
const EMIT_INTERVAL: Duration = Duration::from_secs(10);
pub struct TelemetryMetricsCollector {
registry: TraceRegistry,
exporter_queue_metrics: QueueMetricsFetcher,
interval: Option<tokio::time::Interval>,
}
impl std::fmt::Debug for TelemetryMetricsCollector {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TelemetryMetricsCollector")
.finish_non_exhaustive()
}
}
impl TelemetryMetricsCollector {
pub fn start(
registry: TraceRegistry,
exporter_queue_metrics: QueueMetricsFetcher,
shared_runtime: &Arc<BasicRuntime>,
) -> Result<(), SharedRuntimeError> {
let worker = Self {
registry,
exporter_queue_metrics,
interval: None,
};
shared_runtime.spawn_worker(worker, false).map(|_handle| ())
}
fn emit_metrics(&mut self) {
use telemetry::TelemetryMetric::*;
let registry_metrics = self.registry.get_metrics();
let exporter_queue_metrics = self.exporter_queue_metrics.get_metrics();
telemetry::add_points([
(registry_metrics.spans_created as f64, SpansCreated),
(registry_metrics.spans_finished as f64, SpansFinished),
(
registry_metrics.trace_segments_created as f64,
TraceSegmentsCreated,
),
(
registry_metrics.trace_segments_closed as f64,
TraceSegmentsClosed,
),
(
registry_metrics.trace_partial_flush_count as f64,
TracePartialFlushCount,
),
(
exporter_queue_metrics.spans_queued as f64,
SpansEnqueuedForSerialization,
),
(
exporter_queue_metrics.spans_dropped_full_buffer as f64,
SpansDroppedBufferFull,
),
]);
}
}
#[async_trait]
impl Worker for TelemetryMetricsCollector {
async fn run(&mut self) {
self.emit_metrics();
}
async fn trigger(&mut self) {
let interval = self.interval.get_or_insert_with(|| {
let mut interval = tokio::time::interval_at(
tokio::time::Instant::now() + EMIT_INTERVAL,
EMIT_INTERVAL,
);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval
});
interval.tick().await;
}
}