Skip to content

Commit 07b0533

Browse files
committed
feat(datadog): use ureq to send metrics to datadog
Signed-off-by: Jérémie Drouet <jeremie.drouet@gmail.com>
1 parent cf9a628 commit 07b0533

4 files changed

Lines changed: 211 additions & 44 deletions

File tree

Cargo.toml

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ readme = "README.md"
1313

1414
[features]
1515
default = ["datadog"]
16-
datadog = ["datadog-client"]
16+
datadog = ["ureq"]
1717

1818
[dependencies]
1919
loggerv = "0.7.2"
@@ -22,7 +22,6 @@ clap = "2.33.3"
2222
regex = "1"
2323
procfs = "0.8.1"
2424
actix-web = "3"
25-
futures = "0.3"
2625
riemann_client = "0.9.0"
2726
hostname = "0.3.1"
2827
protobuf = "2.20.0"
@@ -31,7 +30,7 @@ serde_json = "1.0"
3130
warp10 = "1.0.0"
3231
time = "0.2.25"
3332

34-
datadog-client = { version = "0.1", optional = true }
33+
ureq = { version = "2.0.2", features = ["json"], optional = true }
3534

3635
[profile.release]
3736
lto = true

src/exporters/datadog.rs

Lines changed: 163 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,165 @@
11
use crate::exporters::*;
22
use crate::sensors::{Sensor, Topology};
3-
use datadog_client::client::{Client, Config};
4-
use datadog_client::metrics::{Point, Serie, Type};
3+
use serde::ser::SerializeSeq;
4+
use serde::{Serialize, Serializer};
55
use std::collections::HashMap;
66
use std::thread;
77
use std::time::{Duration, Instant};
88

9+
#[derive(Clone, Debug)]
10+
pub enum Type {
11+
Count,
12+
Gauge,
13+
Rate,
14+
}
15+
16+
impl Type {
17+
pub fn as_str(&self) -> &str {
18+
match self {
19+
Self::Count => "count",
20+
Self::Gauge => "gauge",
21+
Self::Rate => "rate",
22+
}
23+
}
24+
}
25+
26+
impl Serialize for Type {
27+
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
28+
where
29+
S: Serializer,
30+
{
31+
serializer.serialize_str(self.as_str())
32+
}
33+
}
34+
35+
#[derive(Clone, Debug)]
36+
pub struct Point {
37+
timestamp: u64,
38+
value: f64,
39+
}
40+
41+
impl Point {
42+
pub fn new(timestamp: u64, value: f64) -> Self {
43+
Self { timestamp, value }
44+
}
45+
}
46+
47+
impl Serialize for Point {
48+
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
49+
where
50+
S: Serializer,
51+
{
52+
let mut seq = serializer.serialize_seq(Some(2))?;
53+
seq.serialize_element(&self.timestamp)?;
54+
seq.serialize_element(&self.value)?;
55+
seq.end()
56+
}
57+
}
58+
59+
#[derive(Debug, Clone, Serialize)]
60+
pub struct Serie {
61+
// The name of the host that produced the metric.
62+
#[serde(skip_serializing_if = "Option::is_none")]
63+
host: Option<String>,
64+
// If the type of the metric is rate or count, define the corresponding interval.
65+
#[serde(skip_serializing_if = "Option::is_none")]
66+
interval: Option<i64>,
67+
// The name of the timeseries.
68+
metric: String,
69+
// Points relating to a metric. All points must be tuples with timestamp and a scalar value (cannot be a string).
70+
// Timestamps should be in POSIX time in seconds, and cannot be more than ten minutes in the future or more than one hour in the past.
71+
points: Vec<Point>,
72+
// A list of tags associated with the metric.
73+
tags: Vec<String>,
74+
// The type of the metric either count, gauge, or rate.
75+
#[serde(rename = "type")]
76+
dtype: Type,
77+
}
78+
79+
impl Serie {
80+
pub fn new(metric: &str, dtype: Type) -> Self {
81+
Self {
82+
host: None,
83+
interval: None,
84+
metric: metric.to_string(),
85+
points: Vec::new(),
86+
tags: Vec::new(),
87+
dtype,
88+
}
89+
}
90+
}
91+
92+
impl Serie {
93+
pub fn set_host(mut self, host: &str) -> Self {
94+
self.host = Some(host.to_string());
95+
self
96+
}
97+
98+
pub fn set_interval(mut self, interval: i64) -> Self {
99+
self.interval = Some(interval);
100+
self
101+
}
102+
103+
pub fn set_points(mut self, points: Vec<Point>) -> Self {
104+
self.points = points;
105+
self
106+
}
107+
108+
pub fn add_point(mut self, point: Point) -> Self {
109+
self.points.push(point);
110+
self
111+
}
112+
}
113+
114+
impl Serie {
115+
pub fn set_tags(mut self, tags: Vec<String>) -> Self {
116+
self.tags = tags;
117+
self
118+
}
119+
120+
pub fn add_tag(mut self, tag: String) -> Self {
121+
self.tags.push(tag);
122+
self
123+
}
124+
}
125+
126+
struct Client {
127+
host: String,
128+
api_key: String,
129+
}
130+
131+
impl Client {
132+
pub fn new(parameters: &ArgMatches) -> Self {
133+
Self {
134+
host: parameters.value_of("host").unwrap().to_string(),
135+
api_key: parameters.value_of("api_key").unwrap().to_string(),
136+
}
137+
}
138+
139+
pub fn send(&self, series: &[Serie]) {
140+
let url = format!("{}/api/v1/series", self.host);
141+
let request = ureq::post(url.as_str())
142+
.set("DD-API-KEY", self.api_key.as_str())
143+
.send_json(serde_json::json!({ "series": series }));
144+
match request {
145+
Ok(response) => {
146+
if response.status() >= 400 {
147+
log::warn!(
148+
"couldn't send metrics to datadog: status {}",
149+
response.status_text()
150+
);
151+
if let Ok(body) = response.into_string() {
152+
log::warn!("response from server: {}", body);
153+
}
154+
} else {
155+
log::info!("metrics sent with success");
156+
}
157+
}
158+
Err(err) => log::warn!("error while sending metrics: {}", err),
159+
};
160+
}
161+
}
162+
9163
fn merge<A>(first: Vec<A>, second: Vec<A>) -> Vec<A> {
10164
second.into_iter().fold(first, |mut res, item| {
11165
res.push(item);
@@ -79,15 +233,8 @@ impl DatadogExporter {
79233
}
80234
}
81235

82-
fn build_client(parameters: &ArgMatches) -> Client {
83-
let config = Config::new(
84-
parameters.value_of("host").unwrap().to_string(),
85-
parameters.value_of("api_key").unwrap().to_string(),
86-
);
87-
Client::new(config)
88-
}
89-
90-
fn runner(&mut self, parameters: &ArgMatches) {
236+
fn runner(&mut self, parameters: &ArgMatches<'_>) {
237+
let client = Client::new(parameters);
91238
if let Some(timeout) = parameters.value_of("timeout") {
92239
let now = Instant::now();
93240
let timeout = timeout
@@ -110,18 +257,18 @@ impl DatadogExporter {
110257
info!("Measurement step is: {}s", step_duration);
111258

112259
while now.elapsed().as_secs() <= timeout {
113-
self.iterate(parameters);
260+
self.iterate(&client);
114261
thread::sleep(Duration::new(step_duration, step_duration_nano));
115262
}
116263
} else {
117-
self.iterate(parameters);
264+
self.iterate(&client);
118265
}
119266
}
120267

121-
fn iterate(&mut self, parameters: &ArgMatches) {
268+
fn iterate(&mut self, client: &Client) {
122269
self.topology.refresh();
123-
let _series = self.collect_series();
124-
let _client = Self::build_client(parameters);
270+
let series = self.collect_series();
271+
client.send(&series);
125272
}
126273

127274
fn create_consumption_serie(&self) -> Serie {

src/lib.rs

Lines changed: 44 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -5,9 +5,9 @@ pub mod exporters;
55
pub mod sensors;
66
use clap::ArgMatches;
77
use exporters::{
8-
json::JSONExporter, prometheus::PrometheusExporter, qemu::QemuExporter,
9-
riemann::RiemannExporter, stdout::StdoutExporter, warpten::Warp10Exporter, Exporter,
10-
ExporterOption,
8+
datadog::DatadogExporter, json::JSONExporter, prometheus::PrometheusExporter,
9+
qemu::QemuExporter, riemann::RiemannExporter, stdout::StdoutExporter, warpten::Warp10Exporter,
10+
Exporter, ExporterOption,
1111
};
1212
use sensors::{powercap_rapl::PowercapRAPLSensor, Sensor};
1313
use std::collections::HashMap;
@@ -53,36 +53,50 @@ fn get_sensor(matches: &ArgMatches) -> Box<dyn Sensor> {
5353
pub fn run(matches: ArgMatches) {
5454
loggerv::init_with_verbosity(matches.occurrences_of("v")).unwrap();
5555

56-
let sensor_boxed = get_sensor(&matches);
57-
let exporter_parameters;
58-
5956
if let Some(stdout_exporter_parameters) = matches.subcommand_matches("stdout") {
60-
exporter_parameters = stdout_exporter_parameters.clone();
61-
let mut exporter = StdoutExporter::new(sensor_boxed);
57+
let exporter_parameters = stdout_exporter_parameters.clone();
58+
let mut exporter = StdoutExporter::new(get_sensor(&matches));
6259
exporter.run(exporter_parameters);
63-
} else if let Some(json_exporter_parameters) = matches.subcommand_matches("json") {
64-
exporter_parameters = json_exporter_parameters.clone();
65-
let mut exporter = JSONExporter::new(sensor_boxed);
60+
return;
61+
}
62+
if let Some(json_exporter_parameters) = matches.subcommand_matches("json") {
63+
let exporter_parameters = json_exporter_parameters.clone();
64+
let mut exporter = JSONExporter::new(get_sensor(&matches));
6665
exporter.run(exporter_parameters);
67-
} else if let Some(riemann_exporter_parameters) = matches.subcommand_matches("riemann") {
68-
exporter_parameters = riemann_exporter_parameters.clone();
69-
let mut exporter = RiemannExporter::new(sensor_boxed);
66+
return;
67+
}
68+
if let Some(riemann_exporter_parameters) = matches.subcommand_matches("riemann") {
69+
let exporter_parameters = riemann_exporter_parameters.clone();
70+
let mut exporter = RiemannExporter::new(get_sensor(&matches));
71+
exporter.run(exporter_parameters);
72+
return;
73+
}
74+
if let Some(prometheus_exporter_parameters) = matches.subcommand_matches("prometheus") {
75+
let exporter_parameters = prometheus_exporter_parameters.clone();
76+
let mut exporter = PrometheusExporter::new(get_sensor(&matches));
7077
exporter.run(exporter_parameters);
71-
} else if let Some(prometheus_exporter_parameters) = matches.subcommand_matches("prometheus") {
72-
exporter_parameters = prometheus_exporter_parameters.clone();
73-
let mut exporter = PrometheusExporter::new(sensor_boxed);
78+
return;
79+
}
80+
if let Some(qemu_exporter_parameters) = matches.subcommand_matches("qemu") {
81+
let exporter_parameters = qemu_exporter_parameters.clone();
82+
let mut exporter = QemuExporter::new(get_sensor(&matches));
7483
exporter.run(exporter_parameters);
75-
} else if let Some(qemu_exporter_parameters) = matches.subcommand_matches("qemu") {
76-
exporter_parameters = qemu_exporter_parameters.clone();
77-
let mut exporter = QemuExporter::new(sensor_boxed);
84+
return;
85+
}
86+
if let Some(warp10_exporter_parameters) = matches.subcommand_matches("warp10") {
87+
let exporter_parameters = warp10_exporter_parameters.clone();
88+
let mut exporter = Warp10Exporter::new(get_sensor(&matches));
7889
exporter.run(exporter_parameters);
79-
} else if let Some(warp10_exporter_parameters) = matches.subcommand_matches("warp10") {
80-
exporter_parameters = warp10_exporter_parameters.clone();
81-
let mut exporter = Warp10Exporter::new(sensor_boxed);
90+
return;
91+
}
92+
#[cfg(feature = "datadog")]
93+
if let Some(datadog_exporter_parameters) = matches.subcommand_matches("datadog") {
94+
let exporter_parameters = datadog_exporter_parameters.clone();
95+
let mut exporter = DatadogExporter::new(get_sensor(&matches));
8296
exporter.run(exporter_parameters);
83-
} else {
84-
error!("Couldn't determine which exporter has been chosen.");
97+
return;
8598
}
99+
error!("Couldn't determine which exporter has been chosen.");
86100
}
87101

88102
/// Returns options needed for each exporter as a HashMap.
@@ -113,6 +127,11 @@ pub fn get_exporters_options() -> HashMap<String, HashMap<String, ExporterOption
113127
String::from("warp10"),
114128
exporters::warpten::Warp10Exporter::get_options(),
115129
);
130+
#[cfg(feature = "datadog")]
131+
options.insert(
132+
String::from("datadog"),
133+
exporters::datadog::DatadogExporter::get_options(),
134+
);
116135
options
117136
}
118137

src/main.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,8 @@ fn main() {
6464
"riemann" => "Riemann exporter sends power consumption metrics to a Riemann server",
6565
"qemu" => "Qemu exporter watches all Qemu/KVM virtual machines running on the host and exposes metrics of each of them in a dedicated folder",
6666
"warp10" => "Warp10 exporter sends data to a Warp10 host, through HTTP",
67+
#[cfg(feature = "datadog")]
68+
"datadog" => "Datadog exporter sends power consumption metrics to Datadog",
6769
_ => "Unknown exporter",
6870
}
6971
);

0 commit comments

Comments
 (0)