improvements
This commit is contained in:
@@ -0,0 +1,110 @@
|
||||
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<MetricValue> 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<dyn std::error::Error>> {
|
||||
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<dyn std::error::Error>> {
|
||||
debug!(msg = ?msg, "Handling msg");
|
||||
match msg {
|
||||
ChimemonMessage::TimeReport(_tr) => Ok(()),
|
||||
ChimemonMessage::SourceReport(sr) => self.handle_source_report(sr).await,
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user