rename prom-client crate to prometheus-http-client
This commit is contained in:
@@ -0,0 +1,29 @@
|
||||
[package]
|
||||
name = "prometheus-http-client"
|
||||
version = "0.1.0"
|
||||
edition.workspace = true
|
||||
|
||||
[lints]
|
||||
workspace = true
|
||||
|
||||
[dependencies]
|
||||
async-trait = { workspace = true }
|
||||
chrono = { workspace = true }
|
||||
displaydoc = { workspace = true }
|
||||
reqwest = { workspace = true }
|
||||
serde = { workspace = true }
|
||||
tracing = { workspace = true }
|
||||
url = { workspace = true }
|
||||
|
||||
plotters = { workspace = true, optional = true }
|
||||
|
||||
[dev-dependencies]
|
||||
conf = { workspace = true }
|
||||
conf-extra = { workspace = true }
|
||||
serde_json = { workspace = true }
|
||||
tokio = { workspace = true }
|
||||
tracing-subscriber = { workspace = true }
|
||||
|
||||
[features]
|
||||
default = ["plot"]
|
||||
plot = ["dep:plotters"]
|
||||
@@ -0,0 +1,79 @@
|
||||
//! Builder for QueryRangeRequest
|
||||
|
||||
use super::QueryRangeRequest;
|
||||
use chrono::{DateTime, TimeDelta, Utc};
|
||||
use std::{ops::Range, time::Duration};
|
||||
|
||||
/// Builder for constructing a QueryRangeRequest
|
||||
pub struct QueryRangeRequestBuilder {
|
||||
query: String,
|
||||
range: Option<Range<DateTime<Utc>>>,
|
||||
step: Option<Duration>,
|
||||
count: usize,
|
||||
}
|
||||
|
||||
impl QueryRangeRequestBuilder {
|
||||
/// Create a new builder with the given PromQL query
|
||||
pub fn new(query: String) -> Self {
|
||||
Self {
|
||||
query,
|
||||
range: None,
|
||||
step: None,
|
||||
count: 256,
|
||||
}
|
||||
}
|
||||
|
||||
/// Set the time range for the query
|
||||
pub fn range(mut self, range: Range<DateTime<Utc>>) -> Self {
|
||||
if self.range.is_some() {
|
||||
panic!("already set range: {:?}", self.range);
|
||||
}
|
||||
self.range = Some(range);
|
||||
self
|
||||
}
|
||||
|
||||
/// Set the time range to be from `time` ago until now
|
||||
pub fn since(mut self, time: Duration) -> Self {
|
||||
if self.range.is_some() {
|
||||
panic!("already set range: {:?}", self.range);
|
||||
}
|
||||
let end = Utc::now();
|
||||
let start = end - TimeDelta::from_std(time).unwrap();
|
||||
|
||||
self.range = Some(start..end);
|
||||
self
|
||||
}
|
||||
|
||||
/// Set the step interval between data points
|
||||
pub fn step(mut self, step: Duration) -> Self {
|
||||
if self.step.is_some() {
|
||||
panic!("already set step: {:?}", self.step);
|
||||
}
|
||||
self.step = Some(step);
|
||||
self
|
||||
}
|
||||
|
||||
/// Set the target number of data points (used to compute step if not set)
|
||||
pub fn count(mut self, count: usize) -> Self {
|
||||
self.count = count;
|
||||
self
|
||||
}
|
||||
|
||||
/// Build the QueryRangeRequest
|
||||
pub fn build(self) -> QueryRangeRequest {
|
||||
let query = self.query;
|
||||
let range = self.range.unwrap();
|
||||
|
||||
let step = self.step.unwrap_or_else(|| {
|
||||
let delta = range.end - range.start;
|
||||
(delta / (self.count as i32)).to_std().unwrap()
|
||||
});
|
||||
|
||||
QueryRangeRequest {
|
||||
query,
|
||||
start: range.start,
|
||||
end: range.end,
|
||||
step: step.as_secs_f64(),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
//! Error types for prom-client
|
||||
|
||||
use displaydoc::Display;
|
||||
use url::ParseError;
|
||||
|
||||
/// Errors that can occur when making Prometheus API requests
|
||||
#[derive(Debug, Display)]
|
||||
pub enum Error {
|
||||
/// URL: {0}
|
||||
Url(ParseError),
|
||||
/// Reqwest: {0}
|
||||
Reqwest(reqwest::Error),
|
||||
/// API: {0}: {1}
|
||||
API(String, String),
|
||||
/// Unexpected Result Type: {0}
|
||||
UnexpectedResultType(String),
|
||||
/// Missing data on success response
|
||||
MissingData,
|
||||
}
|
||||
|
||||
impl From<reqwest::Error> for Error {
|
||||
fn from(src: reqwest::Error) -> Self {
|
||||
Self::Reqwest(src)
|
||||
}
|
||||
}
|
||||
|
||||
impl From<ParseError> for Error {
|
||||
fn from(src: ParseError) -> Self {
|
||||
Self::Url(src)
|
||||
}
|
||||
}
|
||||
|
||||
impl std::error::Error for Error {}
|
||||
@@ -0,0 +1,84 @@
|
||||
//! Helpers for extracting common labels from metrics before presenting them
|
||||
use std::{
|
||||
collections::BTreeMap,
|
||||
fmt::{Debug, Display},
|
||||
};
|
||||
|
||||
/// A set of labels extracted by ExtractLabels
|
||||
pub type Labels = BTreeMap<String, String>;
|
||||
|
||||
/// Common labels extracted from a very generic sequence of metric labels (key value pairs)
|
||||
pub struct ExtractLabels {
|
||||
/// The metric name (from __name__ label)
|
||||
pub name: String,
|
||||
/// Labels that are common to all metrics in the set
|
||||
pub common_labels: Labels,
|
||||
/// Labels specific to each metric (excluding common labels)
|
||||
pub specific_labels: Vec<Labels>,
|
||||
}
|
||||
|
||||
impl ExtractLabels {
|
||||
/// Extract common and specific labels from an iterator of metric label sets.
|
||||
///
|
||||
/// Labels in `skip_labels` are excluded from both common and specific labels.
|
||||
pub fn new<'a, I, KV, K, V>(src: I, skip_labels: &[String]) -> Self
|
||||
where
|
||||
I: Iterator<Item = &'a KV>,
|
||||
KV: Clone + Debug + 'a,
|
||||
K: Display + 'a,
|
||||
V: Display + 'a,
|
||||
&'a KV: IntoIterator<Item = (&'a K, &'a V)>,
|
||||
{
|
||||
// Common labels needs to be String -> Option<String>, because when we find a conflict,
|
||||
// we need to poison that key (by putting None)
|
||||
let mut common_labels = BTreeMap::<String, Option<String>>::default();
|
||||
let mut specific_labels: Vec<Labels> = src
|
||||
.map(|kv| {
|
||||
let labels: Labels = kv
|
||||
.into_iter()
|
||||
.map(|(k, v)| {
|
||||
let k = k.to_string();
|
||||
let v = v.to_string();
|
||||
if let Some(existing_val) = common_labels.get_mut(&k) {
|
||||
// If the existing value is None, we already eliminated this label,
|
||||
// so don't add it back.
|
||||
if let Some(common_val) = existing_val
|
||||
&& common_val != &v
|
||||
{
|
||||
*existing_val = None;
|
||||
}
|
||||
} else {
|
||||
common_labels.insert(k.clone(), Some(v.clone()));
|
||||
}
|
||||
(k, v)
|
||||
})
|
||||
.collect();
|
||||
labels
|
||||
})
|
||||
.collect();
|
||||
|
||||
let mut common_labels = common_labels
|
||||
.into_iter()
|
||||
.filter_map(|(k, maybe_v)| maybe_v.map(|v| (k, v)))
|
||||
.collect::<Labels>();
|
||||
|
||||
for sl in skip_labels {
|
||||
common_labels.remove(sl);
|
||||
}
|
||||
for sp in &mut specific_labels {
|
||||
for sl in skip_labels {
|
||||
sp.remove(sl);
|
||||
}
|
||||
for cl in common_labels.keys() {
|
||||
sp.remove(cl);
|
||||
}
|
||||
}
|
||||
let name = common_labels.remove("__name__").unwrap_or_default();
|
||||
|
||||
ExtractLabels {
|
||||
name,
|
||||
common_labels,
|
||||
specific_labels,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,149 @@
|
||||
#![deny(missing_docs)]
|
||||
|
||||
//! Minimal API for getting time series data from prometheus
|
||||
//!
|
||||
//! To use it, instantiate one of the request objects,
|
||||
//! e.g. QueryRequest or QueryRangeRequest. When it's helpful a builder is provided.
|
||||
//!
|
||||
//! Then use `PromRequest` trait and call `send` or `send_with_client`.
|
||||
//! This takes the prometheus url, and optionally a reqwest client to use.
|
||||
//!
|
||||
//! On success, the result is `PromData`. One would usually call `into_matrix()?`
|
||||
//! or `into_vector()?` as expected for the request that is made.
|
||||
//!
|
||||
//! In prometheus, metric labels are just a set of key-value pairs. However,
|
||||
//! if you are expecting certain structure, you may use any KV object that implements
|
||||
//! `serde::Deserialize` in the `PromData` that results from the call.
|
||||
//!
|
||||
//! When `plot` feature is active, the plot module can be used to plot the timeseries data.
|
||||
//! This is mostly intended to be used as previews in alert messages.
|
||||
|
||||
use chrono::{DateTime, Utc};
|
||||
use serde::{Serialize, de::DeserializeOwned};
|
||||
use std::fmt::Debug;
|
||||
|
||||
mod builders;
|
||||
pub use builders::QueryRangeRequestBuilder;
|
||||
|
||||
mod error;
|
||||
pub use error::Error;
|
||||
|
||||
mod labels;
|
||||
pub use labels::{ExtractLabels, Labels};
|
||||
|
||||
mod messages;
|
||||
use messages::PromResponse;
|
||||
pub use messages::{
|
||||
AlertInfo, AlertStatus, AlertsResponse, MetricTimeseries, MetricVal, MetricValue, PromData,
|
||||
};
|
||||
|
||||
#[cfg(feature = "plot")]
|
||||
pub mod plot;
|
||||
|
||||
mod traits;
|
||||
pub use traits::PromRequest;
|
||||
|
||||
/// Query parameters for /api/v1/query prometheus request
|
||||
#[derive(Clone, Debug, Serialize)]
|
||||
pub struct QueryRequest {
|
||||
/// The PromQL query string
|
||||
pub query: String,
|
||||
/// Optional evaluation timestamp (defaults to current time)
|
||||
pub time: Option<DateTime<Utc>>,
|
||||
}
|
||||
|
||||
impl PromRequest for QueryRequest {
|
||||
const PATH: &str = "/api/v1/query";
|
||||
type Output<KV: Clone + Debug + DeserializeOwned> = PromData<KV>;
|
||||
}
|
||||
|
||||
/// Query parameters for /api/v1/query_range prometheus request
|
||||
/// Use builder to populate it
|
||||
#[derive(Clone, Debug, Serialize)]
|
||||
pub struct QueryRangeRequest {
|
||||
query: String,
|
||||
start: DateTime<Utc>,
|
||||
end: DateTime<Utc>,
|
||||
step: f64,
|
||||
}
|
||||
|
||||
impl QueryRangeRequest {
|
||||
/// Get builder for query range request with given query
|
||||
pub fn builder(query: impl Into<String>) -> QueryRangeRequestBuilder {
|
||||
QueryRangeRequestBuilder::new(query.into())
|
||||
}
|
||||
}
|
||||
|
||||
impl PromRequest for QueryRangeRequest {
|
||||
const PATH: &str = "/api/v1/query_range";
|
||||
type Output<KV: Clone + Debug + DeserializeOwned> = PromData<KV>;
|
||||
}
|
||||
|
||||
/// Query parameters for /api/v1/series prometheus request
|
||||
#[derive(Clone, Debug, Serialize)]
|
||||
#[serde(transparent)]
|
||||
pub struct SeriesRequest {
|
||||
/// Series selector arguments
|
||||
pub matches: MatchList,
|
||||
}
|
||||
|
||||
impl PromRequest for SeriesRequest {
|
||||
const PATH: &str = "/api/v1/series";
|
||||
type Output<KV: Clone + Debug + DeserializeOwned> = Vec<KV>;
|
||||
}
|
||||
|
||||
/// Query parameters for /api/v1/labels prometheus request
|
||||
#[derive(Clone, Debug, Serialize)]
|
||||
#[serde(transparent)]
|
||||
pub struct LabelsRequest {
|
||||
/// Series selector arguments to filter which labels are returned
|
||||
pub matches: MatchList,
|
||||
}
|
||||
|
||||
impl PromRequest for LabelsRequest {
|
||||
const PATH: &str = "/api/v1/labels";
|
||||
type Output<KV: Clone + Debug + DeserializeOwned> = Vec<String>;
|
||||
}
|
||||
|
||||
/// Query parameters for /api/v1/alerts prometheus request
|
||||
#[derive(Clone, Debug, Serialize)]
|
||||
pub struct AlertsRequest {}
|
||||
|
||||
impl PromRequest for AlertsRequest {
|
||||
const PATH: &str = "/api/v1/alerts";
|
||||
type Output<KV: Clone + Debug + DeserializeOwned> = AlertsResponse<KV>;
|
||||
}
|
||||
|
||||
/// Represents a sequence of match[]=...,match[]=...
|
||||
/// query parameters required by parts of prometheus api
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct MatchList(Vec<String>);
|
||||
|
||||
impl Serialize for MatchList {
|
||||
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
|
||||
where
|
||||
S: serde::Serializer,
|
||||
{
|
||||
use serde::ser::SerializeMap;
|
||||
let mut map = serializer.serialize_map(Some(self.0.len()))?;
|
||||
for v in &self.0 {
|
||||
map.serialize_entry("match[]", v)?;
|
||||
}
|
||||
map.end()
|
||||
}
|
||||
}
|
||||
|
||||
impl From<Vec<String>> for MatchList {
|
||||
fn from(src: Vec<String>) -> Self {
|
||||
Self(src)
|
||||
}
|
||||
}
|
||||
|
||||
impl FromIterator<String> for MatchList {
|
||||
fn from_iter<T>(iter: T) -> Self
|
||||
where
|
||||
T: IntoIterator<Item = String>,
|
||||
{
|
||||
Self(iter.into_iter().collect())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,214 @@
|
||||
//! Message types for Prometheus API responses
|
||||
//!
|
||||
//! The prometheus http interface (at port 9090) renders graphs in JS and gets the raw data using API requests like this:
|
||||
//!
|
||||
//! GET http://localhost:9090/api/v1/query_range?query=tick_time{quantile="0.99"}&step=14&start=1762534433.802&end=1762538033.802
|
||||
//!
|
||||
//! Response is:
|
||||
//!
|
||||
//! {
|
||||
//! status: "success"
|
||||
//! data: {
|
||||
//! resultType: "matrix",
|
||||
//! result: [
|
||||
//! {
|
||||
//! metric: { __name__: "tick_time", instance: "x.y.z.w", job: "ec2", .. },
|
||||
//! values: [
|
||||
//! [ 1762534433.802, "1.8974293514080933" ],
|
||||
//! [ 1762534447.802, "2.029724353457351" ],
|
||||
//! ..
|
||||
//! ]
|
||||
//! }
|
||||
//! ]
|
||||
//! }
|
||||
//! }
|
||||
//!
|
||||
//! For more detail see:
|
||||
//! https://prometheus.io/docs/prometheus/latest/querying/api/
|
||||
|
||||
use crate::Error;
|
||||
use chrono::{DateTime, Utc};
|
||||
use serde::{Deserialize, Deserializer, de::DeserializeOwned};
|
||||
use std::{
|
||||
collections::HashMap,
|
||||
fmt::{self, Debug, Display},
|
||||
};
|
||||
use tracing::warn;
|
||||
|
||||
/// Wrapper for f64 that deserializes from a string (prometheus returns numeric values as strings)
|
||||
#[derive(Clone, Copy, Debug)]
|
||||
pub struct MetricVal(pub f64);
|
||||
|
||||
impl Display for MetricVal {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
Display::fmt(&self.0, f)
|
||||
}
|
||||
}
|
||||
|
||||
impl AsRef<f64> for MetricVal {
|
||||
fn as_ref(&self) -> &f64 {
|
||||
&self.0
|
||||
}
|
||||
}
|
||||
|
||||
impl<'de> Deserialize<'de> for MetricVal {
|
||||
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
|
||||
where
|
||||
D: Deserializer<'de>,
|
||||
{
|
||||
let s = <&str>::deserialize(deserializer)?;
|
||||
s.parse::<f64>()
|
||||
.map(MetricVal)
|
||||
.map_err(serde::de::Error::custom)
|
||||
}
|
||||
}
|
||||
|
||||
/// A single metric value from an instant query
|
||||
#[derive(Clone, Debug, Deserialize)]
|
||||
#[serde(bound = "KV: DeserializeOwned")]
|
||||
pub struct MetricValue<KV = HashMap<String, String>>
|
||||
where
|
||||
KV: Clone + Debug,
|
||||
{
|
||||
/// The metric labels
|
||||
pub metric: KV,
|
||||
/// The timestamp and value (if present)
|
||||
#[serde(default)]
|
||||
pub value: Option<(f64, MetricVal)>,
|
||||
// TODO: Include histograms
|
||||
}
|
||||
|
||||
/// A metric timeseries from a range query
|
||||
#[derive(Clone, Debug, Deserialize)]
|
||||
#[serde(bound = "KV: DeserializeOwned")]
|
||||
pub struct MetricTimeseries<KV = HashMap<String, String>>
|
||||
where
|
||||
KV: Clone + Debug,
|
||||
{
|
||||
/// The metric labels
|
||||
pub metric: KV,
|
||||
/// The timestamp/value pairs
|
||||
#[serde(default)]
|
||||
pub values: Vec<(f64, MetricVal)>,
|
||||
// TODO: Include histograms
|
||||
}
|
||||
|
||||
/// The data payload from a Prometheus query response
|
||||
#[derive(Clone, Debug, Deserialize)]
|
||||
#[serde(bound = "KV: DeserializeOwned")]
|
||||
#[serde(tag = "resultType", content = "result", rename_all = "camelCase")]
|
||||
pub enum PromData<KV = HashMap<String, String>>
|
||||
where
|
||||
KV: Clone + Debug,
|
||||
{
|
||||
/// Result from a range query (multiple values per series)
|
||||
Matrix(Vec<MetricTimeseries<KV>>),
|
||||
/// Result from an instant query (single value per series)
|
||||
Vector(Vec<MetricValue<KV>>),
|
||||
}
|
||||
|
||||
impl<KV> PromData<KV>
|
||||
where
|
||||
KV: Clone + Debug,
|
||||
{
|
||||
/// Convert to matrix result, returning error if it was a vector
|
||||
pub fn into_matrix(self) -> Result<Vec<MetricTimeseries<KV>>, Error> {
|
||||
match self {
|
||||
Self::Matrix(data) => Ok(data),
|
||||
_ => Err(Error::UnexpectedResultType(format!("{self:?}"))),
|
||||
}
|
||||
}
|
||||
|
||||
/// Convert to vector result, returning error if it was a matrix
|
||||
pub fn into_vector(self) -> Result<Vec<MetricValue<KV>>, Error> {
|
||||
match self {
|
||||
Self::Vector(data) => Ok(data),
|
||||
_ => Err(Error::UnexpectedResultType(format!("{self:?}"))),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Deserialize, Eq, PartialEq)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub(crate) enum Status {
|
||||
Success,
|
||||
Error,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Deserialize)]
|
||||
#[serde(bound = "T: DeserializeOwned", rename_all = "camelCase")]
|
||||
pub(crate) struct PromResponse<T>
|
||||
where
|
||||
T: Clone + Debug,
|
||||
{
|
||||
pub status: Status,
|
||||
#[serde(default)]
|
||||
pub data: Option<T>,
|
||||
#[serde(default)]
|
||||
pub error_type: Option<String>,
|
||||
#[serde(default)]
|
||||
pub error: Option<String>,
|
||||
#[serde(default)]
|
||||
pub warnings: Vec<String>,
|
||||
}
|
||||
|
||||
impl<T> PromResponse<T>
|
||||
where
|
||||
T: Clone + Debug,
|
||||
{
|
||||
pub fn into_result(self) -> Result<T, Error> {
|
||||
for warning in self.warnings {
|
||||
warn!("Prometheus API response: {warning}");
|
||||
}
|
||||
|
||||
match self.status {
|
||||
Status::Success => Ok(self.data.ok_or(Error::MissingData)?),
|
||||
Status::Error => Err(Error::API(
|
||||
self.error_type.unwrap_or_default(),
|
||||
self.error.unwrap_or_default(),
|
||||
)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Response from /api/v1/alerts endpoint
|
||||
#[derive(Clone, Debug, Deserialize)]
|
||||
#[serde(bound = "KV: DeserializeOwned")]
|
||||
pub struct AlertsResponse<KV = HashMap<String, String>>
|
||||
where
|
||||
KV: Clone + Debug,
|
||||
{
|
||||
/// List of alerts
|
||||
pub alerts: Vec<AlertInfo<KV>>,
|
||||
}
|
||||
|
||||
/// Information about a single alert
|
||||
#[derive(Clone, Debug, Deserialize)]
|
||||
#[serde(bound = "KV: DeserializeOwned", rename_all = "camelCase")]
|
||||
pub struct AlertInfo<KV = HashMap<String, String>>
|
||||
where
|
||||
KV: Clone + Debug,
|
||||
{
|
||||
/// When the alert became active
|
||||
pub active_at: DateTime<Utc>,
|
||||
/// Alert annotations
|
||||
pub annotations: KV,
|
||||
/// Alert labels
|
||||
pub labels: KV,
|
||||
/// Current state of the alert
|
||||
pub state: AlertStatus,
|
||||
/// The value that triggered the alert
|
||||
pub value: String,
|
||||
}
|
||||
|
||||
/// The state of an alert
|
||||
#[derive(Clone, Copy, Debug, Deserialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum AlertStatus {
|
||||
/// Alert condition met but for duration not yet satisfied
|
||||
Pending,
|
||||
/// Alert is actively firing
|
||||
Firing,
|
||||
/// Alert has been resolved
|
||||
Resolved,
|
||||
}
|
||||
@@ -0,0 +1,307 @@
|
||||
//! Make a simple plot of time-series data from prometheus
|
||||
|
||||
use crate::{ExtractLabels, MetricTimeseries};
|
||||
use chrono::{DateTime, FixedOffset, TimeZone, Utc};
|
||||
use plotters::prelude::*;
|
||||
use std::{
|
||||
borrow::Borrow,
|
||||
error::Error,
|
||||
fmt::{Debug, Display, Write},
|
||||
ops::Range,
|
||||
path::Path,
|
||||
};
|
||||
|
||||
/// Styling options for the plot
|
||||
pub struct PlotStyle {
|
||||
/// The pixel size of the plot
|
||||
pub drawing_area: (u32, u32),
|
||||
/// The background color
|
||||
pub background: RGBAColor,
|
||||
/// The grid color
|
||||
pub grid: RGBAColor,
|
||||
/// The axis color
|
||||
pub axis: RGBAColor,
|
||||
/// The text color
|
||||
pub text_color: RGBAColor,
|
||||
/// The text font
|
||||
pub text_font: String,
|
||||
/// The text size
|
||||
pub text_size: u32,
|
||||
/// The caption size
|
||||
pub caption_size: u32,
|
||||
/// The colors to use for lines. If there are more lines than this, then colors will be repeated.
|
||||
pub data_colors: Vec<RGBAColor>,
|
||||
/// The colors to use for a "threshold" such as used in a PromQL alerting rule
|
||||
pub threshold_color: RGBAColor,
|
||||
/// Labels to skip rendering of
|
||||
pub skip_labels: Vec<String>,
|
||||
/// UTC offset to use when labelling the timestamps being plotted
|
||||
pub utc_offset_hours: i32,
|
||||
/// Optional title override - used when prometheus aggregations remove __name__
|
||||
pub title: Option<String>,
|
||||
}
|
||||
|
||||
impl Default for PlotStyle {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
drawing_area: (1920, 1200),
|
||||
background: WHITE.into(),
|
||||
grid: RGBAColor(100, 100, 100, 0.5),
|
||||
axis: BLACK.into(),
|
||||
text_color: BLACK.into(),
|
||||
text_font: "sans-serif".into(),
|
||||
text_size: 18,
|
||||
caption_size: 36,
|
||||
data_colors: [
|
||||
GREEN,
|
||||
BLUE,
|
||||
full_palette::ORANGE,
|
||||
YELLOW,
|
||||
MAGENTA,
|
||||
full_palette::TEAL,
|
||||
full_palette::PURPLE,
|
||||
]
|
||||
.iter()
|
||||
.cloned()
|
||||
.map(Into::into)
|
||||
.collect(),
|
||||
threshold_color: RED.mix(0.2),
|
||||
skip_labels: vec!["job".into(), "instance".into()],
|
||||
utc_offset_hours: 0,
|
||||
title: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl PlotStyle {
|
||||
/// Set the drawing_area
|
||||
pub fn with_drawing_area(mut self, drawing_area: impl Into<(u32, u32)>) -> Self {
|
||||
self.drawing_area = drawing_area.into();
|
||||
self
|
||||
}
|
||||
|
||||
/// Override the title of the plot
|
||||
pub fn with_title(mut self, title: impl Into<String>) -> Self {
|
||||
self.title = Some(title.into());
|
||||
self
|
||||
}
|
||||
|
||||
/// Set the UTC offset (timezone) used in the plot, in hours
|
||||
pub fn with_utc_offset(mut self, offset: i32) -> Self {
|
||||
self.utc_offset_hours = offset;
|
||||
self
|
||||
}
|
||||
}
|
||||
|
||||
/// A shaded region appearing on the plot to indicate values that would trigger an alert
|
||||
pub enum PlotThreshold {
|
||||
/// Shade values greater than this threshold
|
||||
GreaterThan(f64),
|
||||
/// Shade values less than this threshold
|
||||
LessThan(f64),
|
||||
}
|
||||
|
||||
impl PlotStyle {
|
||||
/// Use a dark color scheme for the plot
|
||||
pub fn dark_mode(mut self) -> Self {
|
||||
self.background = BLACK.into();
|
||||
self.grid = RGBAColor(100, 100, 100, 0.5);
|
||||
self.axis = WHITE.into();
|
||||
self.text_color = WHITE.into();
|
||||
self
|
||||
}
|
||||
|
||||
/// Plot a collection of metric timeseries data from prometheus, and possibly a "threshold" defined in an alert.
|
||||
/// Write the result to a path. The file extension of the path will determine the format, e.g. png, gif, etc.
|
||||
pub fn plot_timeseries<KV, K, V>(
|
||||
&self,
|
||||
path: impl AsRef<Path>,
|
||||
mts: &[MetricTimeseries<KV>],
|
||||
plot_threshold: Option<PlotThreshold>,
|
||||
) -> Result<(), Box<dyn Error>>
|
||||
where
|
||||
KV: Clone + Debug,
|
||||
K: Display,
|
||||
V: Display,
|
||||
for<'a> &'a KV: IntoIterator<Item = (&'a K, &'a V)>,
|
||||
{
|
||||
// Prepare to plot by scanning the data, finding x and y bounds, common labels, etc.
|
||||
let ExtractLabels {
|
||||
name,
|
||||
common_labels,
|
||||
specific_labels,
|
||||
} = ExtractLabels::new(mts.iter().map(|mts| &mts.metric), &self.skip_labels);
|
||||
let PreparedPlot {
|
||||
x_range,
|
||||
y_range,
|
||||
ts,
|
||||
} = PreparedPlot::prepare(mts)?;
|
||||
|
||||
// Figure out the caption for formatting style for date-times
|
||||
// Use title override if provided, otherwise use extracted metric name
|
||||
// Only append common_labels if no title override (since title likely already has labels)
|
||||
let mut caption = if let Some(title) = &self.title {
|
||||
title.clone()
|
||||
} else {
|
||||
let mut c = name;
|
||||
if !common_labels.is_empty() {
|
||||
write!(&mut c, " {common_labels:?}")?;
|
||||
}
|
||||
c
|
||||
};
|
||||
|
||||
// Format date-times differently depending on the range of date-times being displayed.
|
||||
//
|
||||
// If they are all on the same day, then omit the day, and put it in the caption instead
|
||||
let timezone =
|
||||
FixedOffset::east_opt(3600 * self.utc_offset_hours).ok_or("invalid timezone")?;
|
||||
let start_date_naive = x_range.start.with_timezone(&timezone).date_naive();
|
||||
let date_format_str =
|
||||
if start_date_naive == x_range.end.with_timezone(&timezone).date_naive() {
|
||||
write!(&mut caption, " {start_date_naive}")?;
|
||||
"%H:%M:%S"
|
||||
} else {
|
||||
"%m/%d %H:%M:%S"
|
||||
};
|
||||
|
||||
// Add timezone offset to the caption
|
||||
write!(&mut caption, " UTC{o:+}", o = self.utc_offset_hours)?;
|
||||
|
||||
// Actually start writing the file
|
||||
let root_area = BitMapBackend::new(&path, self.drawing_area).into_drawing_area();
|
||||
root_area.fill(&self.background)?;
|
||||
|
||||
let mut ctx = ChartBuilder::on(&root_area)
|
||||
.set_label_area_size(LabelAreaPosition::Left, 100)
|
||||
.set_label_area_size(LabelAreaPosition::Bottom, 40)
|
||||
.caption(
|
||||
caption,
|
||||
(self.text_font.as_str(), self.caption_size, &self.text_color)
|
||||
.into_text_style(&root_area),
|
||||
)
|
||||
.build_cartesian_2d(x_range.clone(), y_range.clone())?;
|
||||
|
||||
let text_style =
|
||||
(self.text_font.as_str(), self.text_size, &self.text_color).into_text_style(&root_area);
|
||||
|
||||
ctx.configure_mesh()
|
||||
.light_line_style(self.grid) // Dark gray grid lines
|
||||
.axis_style(self.axis) // White axis lines
|
||||
.bold_line_style(self.axis) // White bold lines
|
||||
.label_style(text_style.clone())
|
||||
.x_label_formatter(&|x| {
|
||||
x.with_timezone(&timezone)
|
||||
.format(date_format_str)
|
||||
.to_string()
|
||||
})
|
||||
.draw()?;
|
||||
|
||||
if let Some(threshold) = plot_threshold {
|
||||
let (limit, baseline) = match threshold {
|
||||
PlotThreshold::GreaterThan(limit) => (limit, y_range.end),
|
||||
PlotThreshold::LessThan(limit) => (limit, y_range.start),
|
||||
};
|
||||
ctx.draw_series(AreaSeries::new(
|
||||
[(x_range.start, limit), (x_range.end, limit)],
|
||||
baseline,
|
||||
self.threshold_color,
|
||||
))?;
|
||||
}
|
||||
|
||||
// Only show legend if there are multiple series (single series doesn't need a legend)
|
||||
let show_legend = specific_labels.len() > 1;
|
||||
|
||||
for (idx, (mut metric, vals)) in specific_labels.into_iter().zip(ts.into_iter()).enumerate()
|
||||
{
|
||||
let color = &self.data_colors[idx % self.data_colors.len()];
|
||||
|
||||
let name = metric.remove("__name__").unwrap_or_default();
|
||||
let label = format!("{name} {metric:?}");
|
||||
|
||||
let series = ctx.draw_series(LineSeries::new(vals, color))?;
|
||||
if show_legend {
|
||||
series
|
||||
.label(label)
|
||||
.legend(move |(x, y)| PathElement::new(vec![(x, y), (x + 20, y)], color));
|
||||
}
|
||||
}
|
||||
|
||||
if show_legend {
|
||||
ctx.configure_series_labels()
|
||||
.position(SeriesLabelPosition::LowerLeft)
|
||||
.border_style(self.axis)
|
||||
.background_style(self.background.mix(0.8))
|
||||
.label_font(text_style)
|
||||
.draw()?;
|
||||
}
|
||||
|
||||
// Signal any errors that occurred when writing the file
|
||||
// https://github.com/plotters-rs/plotters?tab=readme-ov-file#faq-list
|
||||
root_area.present()?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
struct PreparedPlot {
|
||||
x_range: Range<DateTime<Utc>>,
|
||||
y_range: Range<f64>,
|
||||
ts: Vec<Vec<(DateTime<Utc>, f64)>>,
|
||||
}
|
||||
|
||||
impl PreparedPlot {
|
||||
fn prepare<KV>(data: &[impl Borrow<MetricTimeseries<KV>>]) -> Result<Self, &'static str>
|
||||
where
|
||||
KV: Clone + Debug,
|
||||
{
|
||||
let mut x_range = None;
|
||||
let mut y_range = None;
|
||||
|
||||
let ts: Vec<Vec<(_, f64)>> = data
|
||||
.iter()
|
||||
.map(|mts| {
|
||||
let mts = mts.borrow();
|
||||
|
||||
mts.values
|
||||
.iter()
|
||||
.filter_map(|(k, v)| {
|
||||
let x = f64_to_datetime(k)?;
|
||||
let y = *v.as_ref();
|
||||
extend_range(&mut x_range, &x);
|
||||
extend_range(&mut y_range, &y);
|
||||
|
||||
Some((x, y))
|
||||
})
|
||||
.collect::<Vec<_>>()
|
||||
})
|
||||
.collect();
|
||||
|
||||
let x_range = x_range.ok_or("No data")?;
|
||||
let y_range = y_range.ok_or("No data")?;
|
||||
|
||||
Ok(Self {
|
||||
x_range,
|
||||
y_range,
|
||||
ts,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
fn f64_to_datetime(t: &f64) -> Option<DateTime<Utc>> {
|
||||
let seconds = t.trunc() as i64;
|
||||
let nanoseconds = (t.fract() * 1_000_000_000.0) as u32;
|
||||
|
||||
Utc.timestamp_opt(seconds, nanoseconds).single()
|
||||
}
|
||||
|
||||
fn extend_range<T: PartialOrd + Clone>(range: &mut Option<Range<T>>, val: &T) {
|
||||
if let Some(range) = range.as_mut() {
|
||||
if range.start > *val {
|
||||
range.start = val.clone();
|
||||
} else if range.end < *val {
|
||||
range.end = val.clone();
|
||||
}
|
||||
} else {
|
||||
*range = Some(val.clone()..val.clone())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
//! Trait for Prometheus API requests
|
||||
|
||||
use crate::{Error, PromResponse};
|
||||
use reqwest::{Client, Url};
|
||||
use serde::{Serialize, de::DeserializeOwned};
|
||||
use std::fmt::Debug;
|
||||
|
||||
/// Trait for types that can be sent as Prometheus API requests
|
||||
#[async_trait::async_trait]
|
||||
pub trait PromRequest: Serialize {
|
||||
/// The API path for this request type (e.g., "/api/v1/query")
|
||||
const PATH: &str;
|
||||
/// The output type returned by this request
|
||||
type Output<KV: Clone + Debug + DeserializeOwned>: Clone + Debug + DeserializeOwned;
|
||||
|
||||
/// Send the request to the given prometheus host URL
|
||||
async fn send<KV>(&self, host: &str) -> Result<Self::Output<KV>, Error>
|
||||
where
|
||||
KV: Clone + Debug + DeserializeOwned,
|
||||
{
|
||||
self.send_with_client(&Client::new(), host).await
|
||||
}
|
||||
|
||||
/// Send the request using the provided reqwest client
|
||||
async fn send_with_client<KV>(
|
||||
&self,
|
||||
client: &Client,
|
||||
host: &str,
|
||||
) -> Result<Self::Output<KV>, Error>
|
||||
where
|
||||
KV: Clone + Debug + DeserializeOwned,
|
||||
{
|
||||
let url = Url::parse(host)?.join(Self::PATH)?;
|
||||
|
||||
let resp: PromResponse<Self::Output<KV>> =
|
||||
client.get(url).query(&self).send().await?.json().await?;
|
||||
|
||||
let data = resp.into_result()?;
|
||||
Ok(data)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user