espresso_telemetry/
push_task.rs1use 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
24fn 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#[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 ticker.tick().await;
70
71 loop {
72 tokio::select! {
73 _ = ticker.tick() => {
74 push_once(&client, &url, &jwt, ®istry, &external_labels, &rate_limit_warned, &telemetry_log_filter).await;
75 }
76 _ = &mut shutdown => {
77 push_once(&client, &url, &jwt, ®istry, &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}