use async_trait::async_trait; use futures::stream; use influxdb2::{ Client, models::{DataPoint, FieldValue}, }; use tokio::{select, sync::broadcast, time::timeout}; use tokio_util::sync::CancellationToken; use tracing::{debug, error, info, instrument}; use crate::{ ChimemonMessage, ChimemonTarget, ChimemonTargetChannel, MetricValue, SourceReport, config::InfluxConfig, fatal, }; pub struct InfluxTarget { name: String, config: InfluxConfig, influx: Client, } impl From for FieldValue { fn from(value: MetricValue) -> Self { match value { MetricValue::Bool(b) => FieldValue::Bool(b), MetricValue::Float(f) => FieldValue::F64(f), MetricValue::Int(i) => FieldValue::I64(i), } } } #[async_trait] impl ChimemonTarget for InfluxTarget { type Config = InfluxConfig; const TASK_NAME: &'static str = "influx-task"; fn new(name: &str, config: Self::Config) -> Self { let influx = Client::new(&config.url, &config.org, &config.token); Self { name: name.to_owned(), config: config, influx, } } async fn run(self, mut chan: ChimemonTargetChannel, cancel: CancellationToken) { info!("Influx task started"); loop { let msg = select! { _ = cancel.cancelled() => { return }, msg = chan.recv() => msg }; debug!(msg = ?msg, "Got msg"); let msg = match msg { Ok(msg) => msg, Err(broadcast::error::RecvError::Closed) => { fatal!("Permanent channel closed, terminating") } Err(broadcast::error::RecvError::Lagged(_)) => { error!("Channel lagged"); continue; } }; if let Err(e) = self.handle_msg(&msg).await { error!(error = ?e, msg=?&msg, "Error handling message"); } } } } impl InfluxTarget { #[instrument(skip_all)] async fn handle_source_report( &self, sr: &SourceReport, ) -> Result<(), Box> { debug!("Handling source report {}", sr.name); let mut dps = Vec::new(); for metric_set in &sr.details.to_metrics() { let mut builder = DataPoint::builder(&sr.name); builder = self .config .tags .iter() .fold(builder, |builder, (k, v)| builder.tag(k, v)); builder = metric_set .tags .iter() .fold(builder, |builder, (k, v)| builder.tag(*k, v)); builder = metric_set.metrics.iter().fold(builder, |builder, metric| { builder.field(metric.name, metric.value) }); dps.push(builder.build()?); } debug!("Sending {} datapoints to influx", dps.len()); timeout( self.config.timeout, self.influx.write(&self.config.bucket, stream::iter(dps)), ) .await??; debug!("All datapoints sent"); Ok(()) } async fn handle_msg(&self, msg: &ChimemonMessage) -> Result<(), Box> { debug!(msg = ?msg, "Handling msg"); match msg { ChimemonMessage::TimeReport(_tr) => Ok(()), ChimemonMessage::SourceReport(sr) => self.handle_source_report(sr).await, } } }