Skip to main content

faucet_core/
window.rs

1//! In-run datetime window slicing for forward incremental (#527).
2//!
3//! [`ReplicationBind`](crate::ReplicationBind) (#513) binds only a single *lower*
4//! bound (`?since=<bookmark>`); [`faucet backfill`](https://…) (#282) windows only
5//! a *bounded historical* range. Neither bounds each request of the ordinary
6//! forward-incremental run.
7//!
8//! Many APIs — analytics / ads / reporting feeds especially — require **both** a
9//! lower and an upper bound and **cap the span** (e.g. reject a range over 30 or
10//! 90 days). Against those, an unbounded `?since=<bookmark>` either errors or,
11//! worse, silently truncates. Window slicing bounds each request to a rolling
12//! `[start, end)` interval between the persisted bookmark and `now`, iterating the
13//! windows within a single run and persisting the window boundary as the bookmark
14//! at each step so the run is resumable mid-sweep. This is parity with Airbyte's
15//! `DatetimeBasedCursor` (`start_datetime` / `end_datetime` / `step` /
16//! `cursor_granularity` / `lookback_window`).
17//!
18//! The enumeration is a pure function ([`enumerate_windows`]); a source injects
19//! the rendered boundaries into its requests via the [`WindowBind`]s (the window
20//! analogue of [`ReplicationBind`](crate::ReplicationBind)).
21
22use crate::FaucetError;
23use crate::replication::{BindFormat, BindTarget, format_instant};
24use chrono::{DateTime, Duration, Utc};
25use schemars::JsonSchema;
26use serde::{Deserialize, Serialize};
27
28/// The placeholder replaced by the formatted window boundary inside a
29/// [`WindowBind::template`]. The `lower` bind renders the window **start**, the
30/// `upper` bind renders the window **end**.
31pub const WINDOW_PLACEHOLDER: &str = "${window}";
32
33/// The placeholder replaced by the formatted window **start** inside any
34/// [`WindowBind::template`], so one bind can carry both bounds (#772).
35pub const WINDOW_START_PLACEHOLDER: &str = "${window.start}";
36
37/// The placeholder replaced by the formatted window **end** (minus
38/// [`WindowSpec::granularity`]) inside any [`WindowBind::template`] (#772).
39pub const WINDOW_END_PLACEHOLDER: &str = "${window.end}";
40
41fn default_window_template() -> String {
42    WINDOW_PLACEHOLDER.to_owned()
43}
44
45/// Default [`WindowSpec::max_windows`]: a runaway backstop, far above any real
46/// sweep. On a first run against years of history at a small `step`, the sweep is
47/// truncated here and the next run resumes from the last window.
48pub const DEFAULT_MAX_WINDOWS: usize = 10_000;
49
50fn default_max_windows() -> usize {
51    DEFAULT_MAX_WINDOWS
52}
53
54/// One half-open `[start, end)` slice of the replication timeline.
55#[derive(Debug, Clone, Copy, PartialEq, Eq)]
56pub struct Window {
57    /// Inclusive lower bound.
58    pub start: DateTime<Utc>,
59    /// Exclusive upper bound.
60    pub end: DateTime<Utc>,
61}
62
63/// Injects a rendered window boundary into the outgoing request — the window
64/// analogue of [`ReplicationBind`](crate::ReplicationBind). Reuses the same
65/// [`BindTarget`] placement and [`BindFormat`] formatting.
66#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
67#[serde(deny_unknown_fields)]
68pub struct WindowBind {
69    /// Where to place the rendered boundary (query param / header / body field /
70    /// path placeholder).
71    #[serde(default)]
72    pub into: BindTarget,
73    /// The parameter / header / body-field / path-placeholder name. Optional
74    /// only for `into: body` with a `path`.
75    #[serde(default)]
76    pub name: String,
77    /// `into: body` only: an RFC 6901 JSON Pointer into the configured `body`
78    /// (`/dateRanges/0/startDate`) instead of a top-level `name` (#748).
79    #[serde(default, skip_serializing_if = "Option::is_none")]
80    pub path: Option<String>,
81    /// JSON type written by a body bind: `string` (default) or `number`.
82    #[serde(default)]
83    pub value_type: crate::replication::BindValueType,
84    /// Template rendered with [`WINDOW_PLACEHOLDER`] (`${window}`) replaced by the
85    /// formatted boundary. Defaults to the bare `${window}`; set e.g.
86    /// `"gte|${window}"` or `"[${window} TO *]"`. `${window.start}` and
87    /// `${window.end}` render the window's start and (granularity-adjusted) end
88    /// in any bind, so one bind can carry both bounds
89    /// (`"segments.date BETWEEN '${window.start}' AND '${window.end}'"`). The
90    /// rendered value is not escaped for any query language.
91    #[serde(default = "default_window_template")]
92    pub template: String,
93    /// How to format the boundary before substitution.
94    #[serde(default)]
95    pub format: BindFormat,
96}
97
98impl Default for WindowBind {
99    fn default() -> Self {
100        Self {
101            into: BindTarget::default(),
102            name: String::new(),
103            path: None,
104            value_type: crate::replication::BindValueType::default(),
105            template: default_window_template(),
106            format: BindFormat::default(),
107        }
108    }
109}
110
111impl WindowBind {
112    /// Whether this bind is the unset default (an omitted [`WindowSpec::upper`]).
113    pub fn is_unset(&self) -> bool {
114        *self == Self::default()
115    }
116
117    /// Validate the bind at config-load time. `side` names the field for errors
118    /// (`"lower"` / `"upper"`).
119    pub fn validate(&self, side: &str) -> Result<(), FaucetError> {
120        crate::replication::validate_bind_placement(
121            &format!("window slicing `{side}`"),
122            self.into,
123            &self.name,
124            self.path.as_deref(),
125        )?;
126        if ![
127            WINDOW_PLACEHOLDER,
128            WINDOW_START_PLACEHOLDER,
129            WINDOW_END_PLACEHOLDER,
130        ]
131        .iter()
132        .any(|p| self.template.contains(p))
133        {
134            return Err(FaucetError::Config(format!(
135                "window slicing: `{side}.template` must contain a `{WINDOW_PLACEHOLDER}`, \
136                 `{WINDOW_START_PLACEHOLDER}` or `{WINDOW_END_PLACEHOLDER}` placeholder"
137            )));
138        }
139        Ok(())
140    }
141
142    /// Render the bind for a concrete boundary instant.
143    pub fn render(&self, boundary: DateTime<Utc>) -> String {
144        let formatted = format_instant(boundary, self.format);
145        self.template.replace(WINDOW_PLACEHOLDER, &formatted)
146    }
147
148    /// Render the bind for a window: `${window}` becomes `own`, `${window.start}`
149    /// the window start and `${window.end}` the rendered end.
150    pub fn render_window(
151        &self,
152        own: DateTime<Utc>,
153        start: DateTime<Utc>,
154        end: DateTime<Utc>,
155    ) -> String {
156        self.template
157            .replace(
158                WINDOW_START_PLACEHOLDER,
159                &format_instant(start, self.format),
160            )
161            .replace(WINDOW_END_PLACEHOLDER, &format_instant(end, self.format))
162            .replace(WINDOW_PLACEHOLDER, &format_instant(own, self.format))
163    }
164}
165
166/// Declarative in-run datetime window slicing (#527).
167#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
168#[serde(deny_unknown_fields)]
169pub struct WindowSpec {
170    /// Window size — `45s` / `30m` / `6h` / `30d`, or a bare integer (= seconds).
171    /// Absolute UTC durations (`d` = 24h); calendar/DST-correct windows are a
172    /// [`faucet backfill`] concern, not the incremental cursor.
173    pub step: String,
174    /// Lower-bound bind, rendered with the window **start**.
175    pub lower: WindowBind,
176    /// Upper-bound bind, rendered with the window **end**. Optional: omit it when
177    /// `lower.template` renders both bounds through `${window.end}`; only the
178    /// `lower` bind is then applied.
179    #[serde(default, skip_serializing_if = "WindowBind::is_unset")]
180    pub upper: WindowBind,
181    /// Subtract this from each window's *rendered* upper bound so `[start, end]`
182    /// is non-overlapping for inclusive-inclusive APIs (Airbyte
183    /// `cursor_granularity`). Same grammar as `step`. The **persisted bookmark is
184    /// always the true half-open boundary**, so resume never gaps or overlaps —
185    /// only the value sent to the server is adjusted.
186    #[serde(default, skip_serializing_if = "Option::is_none")]
187    pub granularity: Option<String>,
188    /// Re-scan this much *before* the bookmark on the first window, to catch
189    /// late-arriving updates without a full replay. Same grammar as `step`.
190    #[serde(default, skip_serializing_if = "Option::is_none")]
191    pub lookback: Option<String>,
192    /// Safety cap on the number of windows enumerated in one run. On overflow the
193    /// sweep is truncated (logged, never silently), and the next run resumes from
194    /// the last window's end.
195    #[serde(default = "default_max_windows")]
196    pub max_windows: usize,
197}
198
199/// Parse a [`WindowSpec`] duration string (`step` / `granularity` / `lookback`)
200/// into an absolute [`chrono::Duration`]: `45s` / `30m` / `6h` / `30d` (`d` =
201/// 24h), or a bare integer (= seconds). Must be positive.
202pub fn parse_step(s: &str) -> Result<Duration, FaucetError> {
203    let s = s.trim();
204    let err = || {
205        FaucetError::Config(format!(
206            "window slicing: '{s}' is not a valid duration — use e.g. 45s, 30m, 6h, 30d"
207        ))
208    };
209    let (num, unit) = match s.chars().last() {
210        Some(c) if c.is_ascii_digit() => (s, "s"),
211        Some(c) => (&s[..s.len() - c.len_utf8()], &s[s.len() - c.len_utf8()..]),
212        None => return Err(err()),
213    };
214    let n: i64 = num.parse().map_err(|_| err())?;
215    if n <= 0 {
216        return Err(FaucetError::Config(format!(
217            "window slicing: duration '{s}' must be positive"
218        )));
219    }
220    Ok(match unit {
221        "s" => Duration::seconds(n),
222        "m" => Duration::minutes(n),
223        "h" => Duration::hours(n),
224        "d" => Duration::days(n),
225        _ => return Err(err()),
226    })
227}
228
229impl WindowSpec {
230    /// Validate the whole spec at config-load time.
231    pub fn validate(&self) -> Result<(), FaucetError> {
232        parse_step(&self.step)?;
233        if let Some(g) = &self.granularity {
234            parse_step(g)?;
235        }
236        if let Some(l) = &self.lookback {
237            parse_step(l)?;
238        }
239        self.lower.validate("lower")?;
240        if self.has_upper() {
241            self.upper.validate("upper")?;
242        } else if !self.lower.template.contains(WINDOW_END_PLACEHOLDER) {
243            return Err(FaucetError::Config(format!(
244                "window slicing: `upper` may be omitted only when `lower.template` renders the \
245                 window end with `{WINDOW_END_PLACEHOLDER}`; otherwise the window is unbounded above"
246            )));
247        }
248        if self.max_windows == 0 {
249            return Err(FaucetError::Config(
250                "window slicing: `max_windows` must be greater than zero".to_owned(),
251            ));
252        }
253        Ok(())
254    }
255
256    /// The parsed `step` duration.
257    pub fn step_duration(&self) -> Result<Duration, FaucetError> {
258        parse_step(&self.step)
259    }
260
261    /// The parsed `granularity` duration, if any.
262    pub fn granularity_duration(&self) -> Result<Option<Duration>, FaucetError> {
263        self.granularity.as_deref().map(parse_step).transpose()
264    }
265
266    /// The parsed `lookback` duration, if any.
267    pub fn lookback_duration(&self) -> Result<Option<Duration>, FaucetError> {
268        self.lookback.as_deref().map(parse_step).transpose()
269    }
270
271    /// Whether an `upper` bind is configured (it is optional when `lower`
272    /// renders both bounds).
273    pub fn has_upper(&self) -> bool {
274        !self.upper.is_unset()
275    }
276
277    /// The window end as sent to the server: `w.end` minus `granularity`.
278    pub fn rendered_end(&self, w: &Window) -> Result<DateTime<Utc>, FaucetError> {
279        Ok(match self.granularity_duration()? {
280            Some(g) => w.end - g,
281            None => w.end,
282        })
283    }
284
285    /// The rendered lower-bound value for a window (the window **start**). A
286    /// `${window.end}` in the template renders [`Self::rendered_end`], or the
287    /// raw end if `granularity` does not parse (it is validated at load time).
288    pub fn render_lower(&self, w: &Window) -> String {
289        let end = self.rendered_end(w).unwrap_or(w.end);
290        self.lower.render_window(w.start, w.start, end)
291    }
292
293    /// The rendered upper-bound value for a window, applying `granularity` (the
294    /// window **end**, minus `granularity` if set, for inclusive-inclusive APIs).
295    pub fn render_upper(&self, w: &Window) -> Result<String, FaucetError> {
296        let end = self.rendered_end(w)?;
297        Ok(self.upper.render_window(end, w.start, end))
298    }
299
300    /// Every configured bind with its rendered value for a window: `lower`, then
301    /// `upper` when one is configured.
302    pub fn render_binds(&self, w: &Window) -> Result<Vec<(&WindowBind, String)>, FaucetError> {
303        let end = self.rendered_end(w)?;
304        let mut out = vec![(&self.lower, self.lower.render_window(w.start, w.start, end))];
305        if self.has_upper() {
306            out.push((&self.upper, self.upper.render_window(end, w.start, end)));
307        }
308        Ok(out)
309    }
310}
311
312/// Enumerate contiguous half-open `[start, end)` windows from `start` (minus
313/// `lookback`) up to `now`, each `step` wide (the last clamped to `now`).
314///
315/// Returns `(windows, truncated)`: an empty vec when `start >= now` (a no-op
316/// run); `truncated = true` when the sweep hit `max_windows` before reaching
317/// `now` (the caller logs it — the next run resumes from the last window's end,
318/// which is the persisted bookmark).
319pub fn enumerate_windows(
320    start: DateTime<Utc>,
321    now: DateTime<Utc>,
322    step: Duration,
323    lookback: Option<Duration>,
324    max_windows: usize,
325) -> (Vec<Window>, bool) {
326    let mut cur = match lookback {
327        Some(lb) => start - lb,
328        None => start,
329    };
330    let mut out = Vec::new();
331    let mut truncated = false;
332    while cur < now {
333        if out.len() >= max_windows {
334            truncated = true;
335            break;
336        }
337        let end = std::cmp::min(cur + step, now);
338        // `parse_step` guarantees a positive step, so `cur + step > cur`; the
339        // clamp to `now` also keeps `end > cur` because the loop guard is
340        // `cur < now`. This guard is belt-and-braces against a degenerate clock.
341        if end <= cur {
342            break;
343        }
344        out.push(Window { start: cur, end });
345        cur = end;
346    }
347    (out, truncated)
348}
349
350#[cfg(test)]
351mod tests {
352    use super::*;
353    use chrono::TimeZone;
354    use serde_json::json;
355
356    fn dt(s: &str) -> DateTime<Utc> {
357        DateTime::parse_from_rfc3339(s).unwrap().with_timezone(&Utc)
358    }
359
360    #[test]
361    fn parse_step_units() {
362        assert_eq!(parse_step("45s").unwrap(), Duration::seconds(45));
363        assert_eq!(parse_step("30m").unwrap(), Duration::minutes(30));
364        assert_eq!(parse_step("6h").unwrap(), Duration::hours(6));
365        assert_eq!(parse_step("30d").unwrap(), Duration::days(30));
366        assert_eq!(parse_step("3600").unwrap(), Duration::seconds(3600));
367    }
368
369    #[test]
370    fn parse_step_rejects_bad() {
371        assert!(parse_step("0d").is_err());
372        assert!(parse_step("-1h").is_err());
373        assert!(parse_step("").is_err());
374        assert!(parse_step("10y").is_err());
375        assert!(parse_step("abc").is_err());
376    }
377
378    #[test]
379    fn enumerate_contiguous_half_open() {
380        let (ws, trunc) = enumerate_windows(
381            dt("2024-01-01T00:00:00Z"),
382            dt("2024-01-04T00:00:00Z"),
383            Duration::days(1),
384            None,
385            100,
386        );
387        assert!(!trunc);
388        assert_eq!(ws.len(), 3);
389        assert_eq!(ws[0].start, dt("2024-01-01T00:00:00Z"));
390        assert_eq!(ws[0].end, dt("2024-01-02T00:00:00Z"));
391        // Half-open: window N's end equals window N+1's start (no gap, no overlap).
392        assert_eq!(ws[0].end, ws[1].start);
393        assert_eq!(ws[2].end, dt("2024-01-04T00:00:00Z"));
394    }
395
396    #[test]
397    fn last_window_clamps_to_now() {
398        let (ws, _) = enumerate_windows(
399            dt("2024-01-01T00:00:00Z"),
400            dt("2024-01-02T06:00:00Z"),
401            Duration::days(1),
402            None,
403            100,
404        );
405        assert_eq!(ws.len(), 2);
406        assert_eq!(ws[1].start, dt("2024-01-02T00:00:00Z"));
407        assert_eq!(ws[1].end, dt("2024-01-02T06:00:00Z")); // clamped, not +1 day
408    }
409
410    #[test]
411    fn empty_when_start_at_or_after_now() {
412        let (ws, trunc) = enumerate_windows(
413            dt("2024-06-01T00:00:00Z"),
414            dt("2024-06-01T00:00:00Z"),
415            Duration::days(1),
416            None,
417            100,
418        );
419        assert!(ws.is_empty());
420        assert!(!trunc);
421    }
422
423    #[test]
424    fn lookback_extends_the_first_window_backwards() {
425        let (ws, _) = enumerate_windows(
426            dt("2024-01-02T00:00:00Z"),
427            dt("2024-01-03T00:00:00Z"),
428            Duration::days(1),
429            Some(Duration::hours(6)),
430            100,
431        );
432        // First window now starts 6h before the bookmark.
433        assert_eq!(ws[0].start, dt("2024-01-01T18:00:00Z"));
434    }
435
436    #[test]
437    fn max_windows_truncates_and_flags() {
438        let (ws, trunc) = enumerate_windows(
439            dt("2024-01-01T00:00:00Z"),
440            dt("2024-12-31T00:00:00Z"),
441            Duration::days(1),
442            None,
443            5,
444        );
445        assert_eq!(ws.len(), 5);
446        assert!(trunc);
447        // The next run resumes from the last window's end.
448        assert_eq!(ws[4].end, dt("2024-01-06T00:00:00Z"));
449    }
450
451    #[test]
452    fn render_lower_and_upper_with_granularity() {
453        let spec = WindowSpec {
454            step: "1d".into(),
455            lower: WindowBind {
456                into: BindTarget::Query,
457                name: "start".into(),
458                template: "${window}".into(),
459                format: BindFormat::Date,
460                path: None,
461                value_type: Default::default(),
462            },
463            upper: WindowBind {
464                into: BindTarget::Query,
465                name: "end".into(),
466                template: "${window}".into(),
467                format: BindFormat::Date,
468                path: None,
469                value_type: Default::default(),
470            },
471            granularity: Some("1d".into()),
472            lookback: None,
473            max_windows: DEFAULT_MAX_WINDOWS,
474        };
475        let w = Window {
476            start: dt("2024-01-01T00:00:00Z"),
477            end: dt("2024-01-02T00:00:00Z"),
478        };
479        assert_eq!(spec.render_lower(&w), "2024-01-01");
480        // Upper is end - granularity (inclusive-inclusive): 2024-01-01, not -02.
481        assert_eq!(spec.render_upper(&w).unwrap(), "2024-01-01");
482    }
483
484    #[test]
485    fn render_template_and_epoch_format() {
486        let bind = WindowBind {
487            into: BindTarget::Query,
488            name: "since".into(),
489            template: "gte|${window}".into(),
490            format: BindFormat::EpochS,
491            path: None,
492            value_type: Default::default(),
493        };
494        let ts = Utc.timestamp_opt(1_700_000_000, 0).unwrap();
495        assert_eq!(bind.render(ts), "gte|1700000000");
496    }
497
498    #[test]
499    fn validate_catches_misconfig() {
500        let ok = WindowSpec {
501            step: "1d".into(),
502            lower: WindowBind {
503                into: BindTarget::Query,
504                name: "start".into(),
505                template: "${window}".into(),
506                format: BindFormat::Iso8601,
507                path: None,
508                value_type: Default::default(),
509            },
510            upper: WindowBind {
511                into: BindTarget::Query,
512                name: "end".into(),
513                template: "${window}".into(),
514                format: BindFormat::Iso8601,
515                path: None,
516                value_type: Default::default(),
517            },
518            granularity: None,
519            lookback: None,
520            max_windows: DEFAULT_MAX_WINDOWS,
521        };
522        ok.validate().unwrap();
523
524        let mut bad_step = ok.clone();
525        bad_step.step = "0d".into();
526        assert!(bad_step.validate().is_err());
527
528        let mut empty_name = ok.clone();
529        empty_name.lower.name = "  ".into();
530        assert!(empty_name.validate().is_err());
531
532        let mut no_placeholder = ok.clone();
533        no_placeholder.upper.template = "fixed".into();
534        assert!(no_placeholder.validate().is_err());
535
536        let mut zero_windows = ok.clone();
537        zero_windows.max_windows = 0;
538        assert!(zero_windows.validate().is_err());
539    }
540
541    #[test]
542    fn spec_deserializes_from_yaml_shape() {
543        let v = json!({
544            "step": "30d",
545            "lower": {"into": "query", "name": "start_date", "format": "date"},
546            "upper": {"into": "query", "name": "end_date", "format": "date"},
547            "lookback": "1d"
548        });
549        let spec: WindowSpec = serde_json::from_value(v).unwrap();
550        assert_eq!(spec.step, "30d");
551        assert_eq!(spec.lower.template, WINDOW_PLACEHOLDER); // defaulted
552        assert_eq!(spec.max_windows, DEFAULT_MAX_WINDOWS); // defaulted
553        spec.validate().unwrap();
554    }
555
556    fn combined_spec() -> WindowSpec {
557        serde_json::from_value(json!({
558            "step": "7d",
559            "granularity": "1d",
560            "lower": {
561                "into": "body", "path": "/query", "format": "date",
562                "template": "WHERE d BETWEEN '${window.start}' AND '${window.end}'"
563            }
564        }))
565        .unwrap()
566    }
567
568    #[test]
569    fn combined_bind_renders_both_bounds_with_granularity() {
570        let spec = combined_spec();
571        spec.validate().unwrap();
572        assert!(!spec.has_upper());
573        let w = Window {
574            start: dt("2026-09-01T00:00:00Z"),
575            end: dt("2026-09-08T00:00:00Z"),
576        };
577        let want = "WHERE d BETWEEN '2026-09-01' AND '2026-09-07'";
578        assert_eq!(spec.render_lower(&w), want);
579        let binds = spec.render_binds(&w).unwrap();
580        assert_eq!(binds.len(), 1);
581        assert_eq!(binds[0].1, want);
582        let back = serde_json::to_value(&spec).unwrap();
583        assert!(back.get("upper").is_none());
584    }
585
586    #[test]
587    fn omitted_upper_requires_window_end_in_lower() {
588        let mut spec = combined_spec();
589        spec.lower.template = "d >= '${window.start}'".into();
590        let err = spec.validate().unwrap_err().to_string();
591        assert!(err.contains("unbounded above"), "{err}");
592        spec.lower.template = "d >= '${window}'".into();
593        assert!(spec.validate().is_err());
594    }
595
596    #[test]
597    fn two_bind_spec_renders_both_and_accepts_named_placeholders() {
598        let mut spec: WindowSpec = serde_json::from_value(json!({
599            "step": "1d",
600            "granularity": "1s",
601            "lower": {"into": "query", "name": "from"},
602            "upper": {"into": "query", "name": "to", "template": "${window.start}..${window.end}"}
603        }))
604        .unwrap();
605        spec.validate().unwrap();
606        let w = Window {
607            start: dt("2026-01-01T00:00:00Z"),
608            end: dt("2026-01-02T00:00:00Z"),
609        };
610        let binds = spec.render_binds(&w).unwrap();
611        assert_eq!(binds.len(), 2);
612        assert_eq!(binds[0].1, spec.render_lower(&w));
613        assert_eq!(binds[1].1, spec.render_upper(&w).unwrap());
614        assert!(binds[1].1.contains(".."));
615        spec.granularity = Some("bad".into());
616        assert!(spec.render_binds(&w).is_err());
617        assert_eq!(spec.render_lower(&w), spec.lower.render(w.start));
618        assert!(WindowBind::default().is_unset());
619    }
620}