Skip to main content

faucet_core/
lag.rs

1//! Source lag (#733): how far a CDC / streaming pipeline is behind its source.
2//!
3//! A source that has a notion of "the head" (a WAL position, a binlog end, a
4//! partition high watermark) reports the distance from where it has read to
5//! that head through [`Source::lag`](crate::Source::lag). The pipeline polls it
6//! at page boundaries and exports `faucet_source_lag_bytes`,
7//! `faucet_source_lag_events` and `faucet_source_lag_seconds`.
8
9use metrics::{Label, gauge};
10use serde::{Deserialize, Serialize};
11
12/// Distance from the source's head, in whichever units the source can measure.
13#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)]
14pub struct SourceLag {
15    /// Bytes of unread change log (Postgres WAL, MySQL binlog).
16    #[serde(default, skip_serializing_if = "Option::is_none")]
17    pub bytes: Option<u64>,
18    /// Unread events / messages (Kafka offsets, SQL Server change rows).
19    #[serde(default, skip_serializing_if = "Option::is_none")]
20    pub events: Option<u64>,
21    /// Age of the oldest unread change, in seconds.
22    #[serde(default, skip_serializing_if = "Option::is_none")]
23    pub seconds: Option<f64>,
24}
25
26impl SourceLag {
27    /// Lag measured in bytes.
28    pub fn bytes(n: u64) -> Self {
29        Self {
30            bytes: Some(n),
31            ..Default::default()
32        }
33    }
34
35    /// Lag measured in events.
36    pub fn events(n: u64) -> Self {
37        Self {
38            events: Some(n),
39            ..Default::default()
40        }
41    }
42
43    /// Lag measured in seconds (negative clock skew clamps to zero).
44    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    /// Whether no unit carries a value.
52    pub fn is_empty(&self) -> bool {
53        self.bytes.is_none() && self.events.is_none() && self.seconds.is_none()
54    }
55
56    /// A compact human rendering: `412 MiB`, `1,204 events`, `3m 20s`, joined.
57    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/// The latest lag sample of one run, shared with the caller (#733). Attach
102/// with [`Pipeline::with_lag_observer`](crate::Pipeline::with_lag_observer).
103#[derive(Debug, Default)]
104pub struct LagObserver {
105    last: std::sync::Mutex<Option<SourceLag>>,
106}
107
108impl LagObserver {
109    /// An observer with no sample yet.
110    pub fn new() -> Self {
111        Self::default()
112    }
113
114    /// Replace the latest sample.
115    pub fn record(&self, lag: SourceLag) {
116        if let Ok(mut g) = self.last.lock() {
117            *g = Some(lag);
118        }
119    }
120
121    /// The latest sample, if the source reported one.
122    pub fn last(&self) -> Option<SourceLag> {
123        self.last.lock().ok().and_then(|g| *g)
124    }
125}
126
127/// Export one lag sample as the `faucet_source_lag_*` gauges.
128pub 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
140/// How often the pipeline asks a source for its lag while pages flow.
141pub const LAG_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_secs(15);
142
143/// Polls [`Source::lag`](crate::Source::lag) for one run: on the first page,
144/// then at most every [`LAG_POLL_INTERVAL`], and once more when the run ends.
145/// A failing query is logged once and never fails the run.
146pub(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    /// Poll unless the last poll was under the interval ago (`force` skips
187    /// the throttle).
188    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}