use rstest::rstest;
use crate::incremental::{compute_incremental_ranges_for_pipeline, TimeRange};
use hamelin_lib::tree::{
ast::{expression::TruncUnit, identifier::Identifier},
builder::{at_day, at_hour, call, eq, field_ref, hours, pipeline, string, PipelineBuilder},
};
use std::collections::HashMap;
use super::helpers::{
build_pipeline, stale_ranges, time_range, time_range_exclusive, timestamp_field, ts,
};
#[rstest]
#[case::end_on_hour_boundary(
"2024-01-01 11:00:00", "2024-01-01 12:00:00",
(TruncUnit::Hour, 1),
"2024-01-01 11:00:00", "2024-01-01 13:00:00"
)]
#[case::both_on_hour_boundaries(
"2024-01-01 10:00:00", "2024-01-01 14:00:00",
(TruncUnit::Hour, 1),
"2024-01-01 10:00:00", "2024-01-01 15:00:00"
)]
#[case::end_on_day_boundary(
"2024-01-01 00:00:00", "2024-01-02 00:00:00",
(TruncUnit::Day, 1),
"2024-01-01 00:00:00", "2024-01-03 00:00:00"
)]
fn test_snap_boundary_produces_exclusive_end(
#[case] stale_start: &str,
#[case] stale_end: &str,
#[case] partition_unit: (TruncUnit, u32),
#[case] expected_start: &str,
#[case] expected_end: &str,
) {
let pipeline = build_pipeline(
pipeline()
.from(|f| f.table_reference("events"))
.where_cmd(eq(field_ref("severity"), string("high"))),
);
let result = compute_incremental_ranges_for_pipeline(
&pipeline,
stale_ranges("events", stale_start, stale_end),
×tamp_field(),
Some(partition_unit),
false,
)
.unwrap();
assert_eq!(
result.replace_range,
time_range_exclusive(expected_start, expected_end),
);
}
#[rstest]
#[case::hourly_agg(
"2024-01-01 14:00:00", "2024-01-01 16:00:00",
at_hour(field_ref("timestamp")),
"2024-01-01 14:00:00", "2024-01-01 16:00:00", // replace_range
"2024-01-01 14:00:00", "2024-01-01 17:00:00" // query_range
)]
#[case::daily_agg(
"2024-01-01 00:00:00", "2024-01-03 00:00:00",
at_day(field_ref("timestamp")),
"2024-01-01 00:00:00", "2024-01-03 00:00:00", // replace_range
"2024-01-01 00:00:00", "2024-01-04 00:00:00" // query_range
)]
fn test_agg_end_on_boundary(
#[case] stale_start: &str,
#[case] stale_end: &str,
#[case] trunc_expr: impl hamelin_lib::tree::builder::IntoExpressionBuilder,
#[case] expected_replace_start: &str,
#[case] expected_replace_end: &str,
#[case] expected_query_start: &str,
#[case] expected_query_end: &str,
) {
let pipeline = build_pipeline(pipeline().from(|f| f.table_reference("events")).agg(|a| {
a.named_aggregate("count", call("count"))
.named_group("timestamp", trunc_expr)
}));
let result = compute_incremental_ranges_for_pipeline(
&pipeline,
stale_ranges("events", stale_start, stale_end),
×tamp_field(),
None,
false,
)
.unwrap();
assert_eq!(
result.replace_range,
time_range(expected_replace_start, expected_replace_end),
);
assert_eq!(
result.query_range,
time_range(expected_query_start, expected_query_end),
);
}
#[test]
fn time_range_inclusive_vs_exclusive() {
let inclusive = time_range("2024-01-01 10:00:00", "2024-01-01 12:00:00");
assert_eq!(inclusive.start(), ts("2024-01-01 10:00:00"));
assert_eq!(inclusive.end(), ts("2024-01-01 12:00:00"));
assert!(matches!(inclusive, TimeRange::Inclusive(_)));
let exclusive = time_range_exclusive("2024-01-01 10:00:00", "2024-01-01 12:00:00");
assert_eq!(exclusive.start(), ts("2024-01-01 10:00:00"));
assert_eq!(exclusive.end(), ts("2024-01-01 12:00:00"));
assert!(matches!(exclusive, TimeRange::Exclusive(_)));
}
#[test]
fn window_lookback_with_end_on_boundary_and_snap() {
let pipeline = build_pipeline(
pipeline()
.from(|f| f.table_reference("events"))
.window(|w| {
w.named_field("sum_amount", call("sum").arg(field_ref("amount")))
.within(hours(-1))
}),
);
let result = compute_incremental_ranges_for_pipeline(
&pipeline,
stale_ranges("events", "2024-01-01 14:00:00", "2024-01-01 16:00:00"),
×tamp_field(),
Some((TruncUnit::Hour, 1)),
false,
)
.unwrap();
assert_eq!(
result.replace_range,
time_range_exclusive("2024-01-01 14:00:00", "2024-01-01 17:00:00"),
);
assert_eq!(
result.query_range,
time_range("2024-01-01 13:00:00", "2024-01-01 17:00:00"),
);
}