1use crate::FaucetError;
23use crate::replication::{BindFormat, BindTarget, format_instant};
24use chrono::{DateTime, Duration, Utc};
25use schemars::JsonSchema;
26use serde::{Deserialize, Serialize};
27
28pub const WINDOW_PLACEHOLDER: &str = "${window}";
32
33fn default_window_template() -> String {
34 WINDOW_PLACEHOLDER.to_owned()
35}
36
37pub const DEFAULT_MAX_WINDOWS: usize = 10_000;
41
42fn default_max_windows() -> usize {
43 DEFAULT_MAX_WINDOWS
44}
45
46#[derive(Debug, Clone, Copy, PartialEq, Eq)]
48pub struct Window {
49 pub start: DateTime<Utc>,
51 pub end: DateTime<Utc>,
53}
54
55#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
59#[serde(deny_unknown_fields)]
60pub struct WindowBind {
61 #[serde(default)]
64 pub into: BindTarget,
65 pub name: String,
67 #[serde(default = "default_window_template")]
71 pub template: String,
72 #[serde(default)]
74 pub format: BindFormat,
75}
76
77impl WindowBind {
78 pub fn validate(&self, side: &str) -> Result<(), FaucetError> {
81 if self.name.trim().is_empty() {
82 return Err(FaucetError::Config(format!(
83 "window slicing: `{side}.name` must not be empty"
84 )));
85 }
86 if !self.template.contains(WINDOW_PLACEHOLDER) {
87 return Err(FaucetError::Config(format!(
88 "window slicing: `{side}.template` must contain the `{WINDOW_PLACEHOLDER}` placeholder"
89 )));
90 }
91 Ok(())
92 }
93
94 pub fn render(&self, boundary: DateTime<Utc>) -> String {
96 let formatted = format_instant(boundary, self.format);
97 self.template.replace(WINDOW_PLACEHOLDER, &formatted)
98 }
99}
100
101#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
103#[serde(deny_unknown_fields)]
104pub struct WindowSpec {
105 pub step: String,
109 pub lower: WindowBind,
111 pub upper: WindowBind,
113 #[serde(default, skip_serializing_if = "Option::is_none")]
119 pub granularity: Option<String>,
120 #[serde(default, skip_serializing_if = "Option::is_none")]
123 pub lookback: Option<String>,
124 #[serde(default = "default_max_windows")]
128 pub max_windows: usize,
129}
130
131pub fn parse_step(s: &str) -> Result<Duration, FaucetError> {
135 let s = s.trim();
136 let err = || {
137 FaucetError::Config(format!(
138 "window slicing: '{s}' is not a valid duration — use e.g. 45s, 30m, 6h, 30d"
139 ))
140 };
141 let (num, unit) = match s.chars().last() {
142 Some(c) if c.is_ascii_digit() => (s, "s"),
143 Some(c) => (&s[..s.len() - c.len_utf8()], &s[s.len() - c.len_utf8()..]),
144 None => return Err(err()),
145 };
146 let n: i64 = num.parse().map_err(|_| err())?;
147 if n <= 0 {
148 return Err(FaucetError::Config(format!(
149 "window slicing: duration '{s}' must be positive"
150 )));
151 }
152 Ok(match unit {
153 "s" => Duration::seconds(n),
154 "m" => Duration::minutes(n),
155 "h" => Duration::hours(n),
156 "d" => Duration::days(n),
157 _ => return Err(err()),
158 })
159}
160
161impl WindowSpec {
162 pub fn validate(&self) -> Result<(), FaucetError> {
164 parse_step(&self.step)?;
165 if let Some(g) = &self.granularity {
166 parse_step(g)?;
167 }
168 if let Some(l) = &self.lookback {
169 parse_step(l)?;
170 }
171 self.lower.validate("lower")?;
172 self.upper.validate("upper")?;
173 if self.max_windows == 0 {
174 return Err(FaucetError::Config(
175 "window slicing: `max_windows` must be greater than zero".to_owned(),
176 ));
177 }
178 Ok(())
179 }
180
181 pub fn step_duration(&self) -> Result<Duration, FaucetError> {
183 parse_step(&self.step)
184 }
185
186 pub fn granularity_duration(&self) -> Result<Option<Duration>, FaucetError> {
188 self.granularity.as_deref().map(parse_step).transpose()
189 }
190
191 pub fn lookback_duration(&self) -> Result<Option<Duration>, FaucetError> {
193 self.lookback.as_deref().map(parse_step).transpose()
194 }
195
196 pub fn render_lower(&self, w: &Window) -> String {
198 self.lower.render(w.start)
199 }
200
201 pub fn render_upper(&self, w: &Window) -> Result<String, FaucetError> {
204 let end = match self.granularity_duration()? {
205 Some(g) => w.end - g,
206 None => w.end,
207 };
208 Ok(self.upper.render(end))
209 }
210}
211
212pub fn enumerate_windows(
220 start: DateTime<Utc>,
221 now: DateTime<Utc>,
222 step: Duration,
223 lookback: Option<Duration>,
224 max_windows: usize,
225) -> (Vec<Window>, bool) {
226 let mut cur = match lookback {
227 Some(lb) => start - lb,
228 None => start,
229 };
230 let mut out = Vec::new();
231 let mut truncated = false;
232 while cur < now {
233 if out.len() >= max_windows {
234 truncated = true;
235 break;
236 }
237 let end = std::cmp::min(cur + step, now);
238 if end <= cur {
242 break;
243 }
244 out.push(Window { start: cur, end });
245 cur = end;
246 }
247 (out, truncated)
248}
249
250#[cfg(test)]
251mod tests {
252 use super::*;
253 use chrono::TimeZone;
254 use serde_json::json;
255
256 fn dt(s: &str) -> DateTime<Utc> {
257 DateTime::parse_from_rfc3339(s).unwrap().with_timezone(&Utc)
258 }
259
260 #[test]
261 fn parse_step_units() {
262 assert_eq!(parse_step("45s").unwrap(), Duration::seconds(45));
263 assert_eq!(parse_step("30m").unwrap(), Duration::minutes(30));
264 assert_eq!(parse_step("6h").unwrap(), Duration::hours(6));
265 assert_eq!(parse_step("30d").unwrap(), Duration::days(30));
266 assert_eq!(parse_step("3600").unwrap(), Duration::seconds(3600));
267 }
268
269 #[test]
270 fn parse_step_rejects_bad() {
271 assert!(parse_step("0d").is_err());
272 assert!(parse_step("-1h").is_err());
273 assert!(parse_step("").is_err());
274 assert!(parse_step("10y").is_err());
275 assert!(parse_step("abc").is_err());
276 }
277
278 #[test]
279 fn enumerate_contiguous_half_open() {
280 let (ws, trunc) = enumerate_windows(
281 dt("2024-01-01T00:00:00Z"),
282 dt("2024-01-04T00:00:00Z"),
283 Duration::days(1),
284 None,
285 100,
286 );
287 assert!(!trunc);
288 assert_eq!(ws.len(), 3);
289 assert_eq!(ws[0].start, dt("2024-01-01T00:00:00Z"));
290 assert_eq!(ws[0].end, dt("2024-01-02T00:00:00Z"));
291 assert_eq!(ws[0].end, ws[1].start);
293 assert_eq!(ws[2].end, dt("2024-01-04T00:00:00Z"));
294 }
295
296 #[test]
297 fn last_window_clamps_to_now() {
298 let (ws, _) = enumerate_windows(
299 dt("2024-01-01T00:00:00Z"),
300 dt("2024-01-02T06:00:00Z"),
301 Duration::days(1),
302 None,
303 100,
304 );
305 assert_eq!(ws.len(), 2);
306 assert_eq!(ws[1].start, dt("2024-01-02T00:00:00Z"));
307 assert_eq!(ws[1].end, dt("2024-01-02T06:00:00Z")); }
309
310 #[test]
311 fn empty_when_start_at_or_after_now() {
312 let (ws, trunc) = enumerate_windows(
313 dt("2024-06-01T00:00:00Z"),
314 dt("2024-06-01T00:00:00Z"),
315 Duration::days(1),
316 None,
317 100,
318 );
319 assert!(ws.is_empty());
320 assert!(!trunc);
321 }
322
323 #[test]
324 fn lookback_extends_the_first_window_backwards() {
325 let (ws, _) = enumerate_windows(
326 dt("2024-01-02T00:00:00Z"),
327 dt("2024-01-03T00:00:00Z"),
328 Duration::days(1),
329 Some(Duration::hours(6)),
330 100,
331 );
332 assert_eq!(ws[0].start, dt("2024-01-01T18:00:00Z"));
334 }
335
336 #[test]
337 fn max_windows_truncates_and_flags() {
338 let (ws, trunc) = enumerate_windows(
339 dt("2024-01-01T00:00:00Z"),
340 dt("2024-12-31T00:00:00Z"),
341 Duration::days(1),
342 None,
343 5,
344 );
345 assert_eq!(ws.len(), 5);
346 assert!(trunc);
347 assert_eq!(ws[4].end, dt("2024-01-06T00:00:00Z"));
349 }
350
351 #[test]
352 fn render_lower_and_upper_with_granularity() {
353 let spec = WindowSpec {
354 step: "1d".into(),
355 lower: WindowBind {
356 into: BindTarget::Query,
357 name: "start".into(),
358 template: "${window}".into(),
359 format: BindFormat::Date,
360 },
361 upper: WindowBind {
362 into: BindTarget::Query,
363 name: "end".into(),
364 template: "${window}".into(),
365 format: BindFormat::Date,
366 },
367 granularity: Some("1d".into()),
368 lookback: None,
369 max_windows: DEFAULT_MAX_WINDOWS,
370 };
371 let w = Window {
372 start: dt("2024-01-01T00:00:00Z"),
373 end: dt("2024-01-02T00:00:00Z"),
374 };
375 assert_eq!(spec.render_lower(&w), "2024-01-01");
376 assert_eq!(spec.render_upper(&w).unwrap(), "2024-01-01");
378 }
379
380 #[test]
381 fn render_template_and_epoch_format() {
382 let bind = WindowBind {
383 into: BindTarget::Query,
384 name: "since".into(),
385 template: "gte|${window}".into(),
386 format: BindFormat::EpochS,
387 };
388 let ts = Utc.timestamp_opt(1_700_000_000, 0).unwrap();
389 assert_eq!(bind.render(ts), "gte|1700000000");
390 }
391
392 #[test]
393 fn validate_catches_misconfig() {
394 let ok = WindowSpec {
395 step: "1d".into(),
396 lower: WindowBind {
397 into: BindTarget::Query,
398 name: "start".into(),
399 template: "${window}".into(),
400 format: BindFormat::Iso8601,
401 },
402 upper: WindowBind {
403 into: BindTarget::Query,
404 name: "end".into(),
405 template: "${window}".into(),
406 format: BindFormat::Iso8601,
407 },
408 granularity: None,
409 lookback: None,
410 max_windows: DEFAULT_MAX_WINDOWS,
411 };
412 ok.validate().unwrap();
413
414 let mut bad_step = ok.clone();
415 bad_step.step = "0d".into();
416 assert!(bad_step.validate().is_err());
417
418 let mut empty_name = ok.clone();
419 empty_name.lower.name = " ".into();
420 assert!(empty_name.validate().is_err());
421
422 let mut no_placeholder = ok.clone();
423 no_placeholder.upper.template = "fixed".into();
424 assert!(no_placeholder.validate().is_err());
425
426 let mut zero_windows = ok.clone();
427 zero_windows.max_windows = 0;
428 assert!(zero_windows.validate().is_err());
429 }
430
431 #[test]
432 fn spec_deserializes_from_yaml_shape() {
433 let v = json!({
434 "step": "30d",
435 "lower": {"into": "query", "name": "start_date", "format": "date"},
436 "upper": {"into": "query", "name": "end_date", "format": "date"},
437 "lookback": "1d"
438 });
439 let spec: WindowSpec = serde_json::from_value(v).unwrap();
440 assert_eq!(spec.step, "30d");
441 assert_eq!(spec.lower.template, WINDOW_PLACEHOLDER); assert_eq!(spec.max_windows, DEFAULT_MAX_WINDOWS); spec.validate().unwrap();
444 }
445}