1use metrics::{Label, gauge};
10use serde::{Deserialize, Serialize};
11
12#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)]
14pub struct SourceLag {
15 #[serde(default, skip_serializing_if = "Option::is_none")]
17 pub bytes: Option<u64>,
18 #[serde(default, skip_serializing_if = "Option::is_none")]
20 pub events: Option<u64>,
21 #[serde(default, skip_serializing_if = "Option::is_none")]
23 pub seconds: Option<f64>,
24}
25
26impl SourceLag {
27 pub fn bytes(n: u64) -> Self {
29 Self {
30 bytes: Some(n),
31 ..Default::default()
32 }
33 }
34
35 pub fn events(n: u64) -> Self {
37 Self {
38 events: Some(n),
39 ..Default::default()
40 }
41 }
42
43 pub fn seconds(s: f64) -> Self {
45 Self {
46 seconds: Some(if s.is_finite() { s.max(0.0) } else { 0.0 }),
47 ..Default::default()
48 }
49 }
50
51 pub fn is_empty(&self) -> bool {
53 self.bytes.is_none() && self.events.is_none() && self.seconds.is_none()
54 }
55
56 pub fn human(&self) -> String {
58 let mut parts = Vec::new();
59 if let Some(b) = self.bytes {
60 parts.push(human_bytes(b));
61 }
62 if let Some(e) = self.events {
63 parts.push(format!("{e} event{}", if e == 1 { "" } else { "s" }));
64 }
65 if let Some(s) = self.seconds {
66 parts.push(human_seconds(s));
67 }
68 parts.join(" · ")
69 }
70}
71
72fn human_bytes(b: u64) -> String {
73 const UNITS: [&str; 5] = ["B", "KiB", "MiB", "GiB", "TiB"];
74 let mut v = b as f64;
75 let mut i = 0;
76 while v >= 1024.0 && i + 1 < UNITS.len() {
77 v /= 1024.0;
78 i += 1;
79 }
80 if i == 0 {
81 format!("{b} B")
82 } else {
83 format!("{v:.0} {}", UNITS[i])
84 }
85}
86
87fn human_seconds(s: f64) -> String {
88 let s = s.max(0.0);
89 if s < 60.0 {
90 return format!("{s:.0}s");
91 }
92 let total = s as u64;
93 let (h, m, sec) = (total / 3600, (total % 3600) / 60, total % 60);
94 if h > 0 {
95 format!("{h}h {m}m")
96 } else {
97 format!("{m}m {sec}s")
98 }
99}
100
101#[derive(Debug, Default)]
104pub struct LagObserver {
105 last: std::sync::Mutex<Option<SourceLag>>,
106}
107
108impl LagObserver {
109 pub fn new() -> Self {
111 Self::default()
112 }
113
114 pub fn record(&self, lag: SourceLag) {
116 if let Ok(mut g) = self.last.lock() {
117 *g = Some(lag);
118 }
119 }
120
121 pub fn last(&self) -> Option<SourceLag> {
123 self.last.lock().ok().and_then(|g| *g)
124 }
125}
126
127pub fn record_lag_gauges(labels: &[Label], lag: &SourceLag) {
129 if let Some(b) = lag.bytes {
130 gauge!("faucet_source_lag_bytes", labels.to_vec()).set(b as f64);
131 }
132 if let Some(e) = lag.events {
133 gauge!("faucet_source_lag_events", labels.to_vec()).set(e as f64);
134 }
135 if let Some(s) = lag.seconds {
136 gauge!("faucet_source_lag_seconds", labels.to_vec()).set(s);
137 }
138}
139
140pub const LAG_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_secs(15);
142
143pub(crate) struct LagPoller<'a> {
147 source: &'a dyn crate::Source,
148 labels: Vec<Label>,
149 observer: Option<std::sync::Arc<LagObserver>>,
150 interval: std::time::Duration,
151 last_poll: std::sync::Mutex<Option<std::time::Instant>>,
152 warned: std::sync::atomic::AtomicBool,
153}
154
155impl<'a> LagPoller<'a> {
156 pub(crate) fn new(
157 source: &'a dyn crate::Source,
158 pipeline: &str,
159 row: &str,
160 observer: Option<std::sync::Arc<LagObserver>>,
161 ) -> Self {
162 use metrics::SharedString;
163 Self {
164 labels: vec![
165 Label::new("pipeline", SharedString::from(pipeline.to_string())),
166 Label::new("row", SharedString::from(row.to_string())),
167 Label::new(
168 "connector",
169 SharedString::from(source.connector_name().to_string()),
170 ),
171 ],
172 source,
173 observer,
174 interval: LAG_POLL_INTERVAL,
175 last_poll: std::sync::Mutex::new(None),
176 warned: std::sync::atomic::AtomicBool::new(false),
177 }
178 }
179
180 #[cfg(test)]
181 fn with_interval(mut self, interval: std::time::Duration) -> Self {
182 self.interval = interval;
183 self
184 }
185
186 pub(crate) async fn poll(&self, force: bool) {
189 let now = std::time::Instant::now();
190 {
191 let Ok(mut last) = self.last_poll.lock() else {
192 return;
193 };
194 if !force && last.is_some_and(|t| now.duration_since(t) < self.interval) {
195 return;
196 }
197 *last = Some(now);
198 }
199 match self.source.lag().await {
200 Ok(Some(lag)) if !lag.is_empty() => {
201 record_lag_gauges(&self.labels, &lag);
202 if let Some(o) = &self.observer {
203 o.record(lag);
204 }
205 }
206 Ok(_) => {}
207 Err(e) => {
208 if !self.warned.swap(true, std::sync::atomic::Ordering::Relaxed) {
209 tracing::warn!(
210 connector = self.source.connector_name(),
211 error = %e,
212 "source lag query failed; lag is not reported for this run (logged once)"
213 );
214 }
215 }
216 }
217 }
218}
219
220#[cfg(test)]
221mod tests {
222 use super::*;
223
224 #[test]
225 fn constructors_and_rendering() {
226 assert_eq!(SourceLag::bytes(10).bytes, Some(10));
227 assert_eq!(SourceLag::events(3).events, Some(3));
228 assert_eq!(SourceLag::seconds(-4.0).seconds, Some(0.0));
229 assert_eq!(SourceLag::seconds(f64::NAN).seconds, Some(0.0));
230 assert!(SourceLag::default().is_empty());
231 assert!(!SourceLag::bytes(0).is_empty());
232 assert_eq!(SourceLag::bytes(512).human(), "512 B");
233 assert_eq!(SourceLag::bytes(412 * 1024 * 1024).human(), "412 MiB");
234 assert_eq!(SourceLag::events(1).human(), "1 event");
235 assert_eq!(SourceLag::events(5).human(), "5 events");
236 assert_eq!(SourceLag::seconds(42.0).human(), "42s");
237 assert_eq!(SourceLag::seconds(200.0).human(), "3m 20s");
238 assert_eq!(SourceLag::seconds(7300.0).human(), "2h 1m");
239 let both = SourceLag {
240 bytes: Some(2048),
241 seconds: Some(5.0),
242 events: None,
243 };
244 assert_eq!(both.human(), "2 KiB · 5s");
245 }
246
247 struct LagSource {
248 calls: std::sync::atomic::AtomicUsize,
249 result: Result<Option<SourceLag>, ()>,
250 }
251
252 #[async_trait::async_trait]
253 impl crate::Source for LagSource {
254 async fn fetch_with_context(
255 &self,
256 _: &std::collections::HashMap<String, serde_json::Value>,
257 ) -> Result<Vec<serde_json::Value>, crate::FaucetError> {
258 Ok(Vec::new())
259 }
260 async fn lag(&self) -> Result<Option<SourceLag>, crate::FaucetError> {
261 self.calls
262 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
263 self.result
264 .map_err(|_| crate::FaucetError::Source("lag query failed".into()))
265 }
266 }
267
268 fn lag_source(result: Result<Option<SourceLag>, ()>) -> LagSource {
269 LagSource {
270 calls: std::sync::atomic::AtomicUsize::new(0),
271 result,
272 }
273 }
274
275 #[tokio::test]
276 async fn poller_throttles_records_and_forces() {
277 let src = lag_source(Ok(Some(SourceLag::bytes(7))));
278 let obs = std::sync::Arc::new(LagObserver::new());
279 assert_eq!(obs.last(), None);
280 let p = LagPoller::new(&src, "p", "r", Some(std::sync::Arc::clone(&obs)))
281 .with_interval(std::time::Duration::from_secs(3600));
282 p.poll(false).await;
283 p.poll(false).await;
284 assert_eq!(src.calls.load(std::sync::atomic::Ordering::Relaxed), 1);
285 p.poll(true).await;
286 assert_eq!(src.calls.load(std::sync::atomic::Ordering::Relaxed), 2);
287 assert_eq!(obs.last(), Some(SourceLag::bytes(7)));
288 }
289
290 #[tokio::test]
291 async fn poller_ignores_empty_and_failing_lag() {
292 let obs = std::sync::Arc::new(LagObserver::new());
293 let empty = lag_source(Ok(Some(SourceLag::default())));
294 LagPoller::new(&empty, "p", "r", Some(std::sync::Arc::clone(&obs)))
295 .poll(true)
296 .await;
297 let none = lag_source(Ok(None));
298 LagPoller::new(&none, "p", "r", Some(std::sync::Arc::clone(&obs)))
299 .poll(true)
300 .await;
301 let failing = lag_source(Err(()));
302 let p = LagPoller::new(&failing, "p", "r", Some(std::sync::Arc::clone(&obs)));
303 p.poll(true).await;
304 p.poll(true).await;
305 assert_eq!(failing.calls.load(std::sync::atomic::Ordering::Relaxed), 2);
306 assert_eq!(obs.last(), None);
307 LagPoller::new(&none, "p", "r", None).poll(true).await;
308 use crate::Source;
309 assert!(
310 none.fetch_with_context(&Default::default())
311 .await
312 .unwrap()
313 .is_empty()
314 );
315 }
316
317 #[test]
318 fn gauges_record_without_a_recorder() {
319 record_lag_gauges(
320 &[],
321 &SourceLag {
322 bytes: Some(1),
323 events: Some(2),
324 seconds: Some(3.0),
325 },
326 );
327 record_lag_gauges(&[], &SourceLag::default());
328 }
329}