Skip to main content

espresso_telemetry/
push_task.rs

1//! Periodic push of `prometheus::Registry` snapshots to a remote-write endpoint.
2//! A final flush runs on shutdown. Errors are `warn!`-logged, never fatal; HTTP
3//! 429 triggers one shared ERROR via [`crate::rate_limit::log_rate_limit_once`].
4
5use std::{
6    sync::{Arc, atomic::AtomicBool},
7    time::Duration,
8};
9
10use prometheus::Registry;
11use reqwest::{
12    Client, StatusCode,
13    header::{AUTHORIZATION, CONTENT_TYPE},
14};
15use tokio::{sync::oneshot, time::MissedTickBehavior};
16use url::Url;
17
18use crate::{
19    build_write_request, encode_to_snappy,
20    rate_limit::log_rate_limit_once,
21    remote_write::{Label, WriteRequest},
22};
23
24/// Stamp push-time labels onto every TimeSeries (existing labels win), then
25/// re-sort to stay remote-write 1.0 compliant.
26fn apply_external_labels(request: &mut WriteRequest, external: &[Label]) {
27    if external.is_empty() {
28        return;
29    }
30    for series in &mut request.timeseries {
31        for label in external {
32            if !series.labels.iter().any(|l| l.name == label.name) {
33                series.labels.push(label.clone());
34            }
35        }
36        series.labels.sort_by(|a, b| a.name.cmp(&b.name));
37    }
38}
39
40/// Run the periodic push loop until `shutdown` resolves, driving one final flush
41/// on the way out. Returns early if the HTTP client can't be built.
42#[allow(clippy::too_many_arguments)]
43pub(crate) async fn run(
44    registry: Arc<Registry>,
45    endpoint: Url,
46    jwt: String,
47    interval: Duration,
48    external_labels: Vec<Label>,
49    rate_limit_warned: Arc<AtomicBool>,
50    telemetry_log_filter: Arc<String>,
51    mut shutdown: oneshot::Receiver<()>,
52) {
53    let url = format!("{}/api/v1/write", endpoint.as_str().trim_end_matches('/'));
54    let client = match Client::builder()
55        .connect_timeout(Duration::from_secs(5))
56        .timeout(Duration::from_secs(10))
57        .build()
58    {
59        Ok(c) => c,
60        Err(e) => {
61            tracing::warn!(error = %e, "telemetry: metrics http client init failed; metrics push disabled");
62            return;
63        },
64    };
65
66    let mut ticker = tokio::time::interval(interval);
67    ticker.set_missed_tick_behavior(MissedTickBehavior::Delay);
68    // Skip the immediate first tick so the first push isn't an empty scrape.
69    ticker.tick().await;
70
71    loop {
72        tokio::select! {
73            _ = ticker.tick() => {
74                push_once(&client, &url, &jwt, &registry, &external_labels, &rate_limit_warned, &telemetry_log_filter).await;
75            }
76            _ = &mut shutdown => {
77                push_once(&client, &url, &jwt, &registry, &external_labels, &rate_limit_warned, &telemetry_log_filter).await;
78                break;
79            }
80        }
81    }
82}
83
84async fn push_once(
85    client: &Client,
86    url: &str,
87    jwt: &str,
88    registry: &Registry,
89    external_labels: &[Label],
90    rate_limit_warned: &AtomicBool,
91    telemetry_log_filter: &str,
92) {
93    let families = registry.gather();
94    let mut request = match build_write_request(&families) {
95        Ok(r) => r,
96        Err(e) => {
97            tracing::warn!(error = %e, "telemetry: skipping metrics push: encode failed");
98            return;
99        },
100    };
101    apply_external_labels(&mut request, external_labels);
102    let body = match encode_to_snappy(&request) {
103        Ok(b) => b,
104        Err(e) => {
105            tracing::warn!(error = %e, "telemetry: skipping metrics push: snappy compress failed");
106            return;
107        },
108    };
109
110    let resp = client
111        .post(url)
112        .header(AUTHORIZATION, format!("Bearer {jwt}"))
113        .header(CONTENT_TYPE, "application/x-protobuf")
114        .header("Content-Encoding", "snappy")
115        .body(body)
116        .send()
117        .await;
118
119    match resp {
120        Ok(r) if r.status().is_success() => {},
121        Ok(r) if r.status() == StatusCode::TOO_MANY_REQUESTS => {
122            let retry_after = r
123                .headers()
124                .get(reqwest::header::RETRY_AFTER)
125                .and_then(|v| v.to_str().ok())
126                .and_then(|s| s.trim().parse::<u64>().ok());
127            log_rate_limit_once(rate_limit_warned, telemetry_log_filter, retry_after);
128        },
129        Ok(r) => tracing::warn!(status = %r.status(), "telemetry: metrics push non-2xx"),
130        Err(e) => tracing::warn!(error = %e, "telemetry: metrics push failed"),
131    }
132}