use std::sync::Arc;
use hamelin_lib::{parse_and_typecheck_with_options, tree::options::TypeCheckOptions};
use rstest::rstest;
use super::helpers::{test_default_space, timestamp_field, ts, MockIncrementalProvider};
use crate::incremental::compute_input_lower_bound_for_query;
#[rstest]
#[case::select_timestamp("FROM events | SELECT timestamp", "2024-03-15 12:00:00")]
#[case::drop_other("FROM events | DROP user", "2024-03-15 12:00:00")]
#[case::grouping_shift(
"FROM events | WINDOW n = count() BY timestamp = timestamp + 2h WITHIN -1h",
"2024-03-15 09:00:00"
)]
#[case::explicit_clock(
"FROM events | WINDOW n = count() SORT timestamp ASC WITHIN -1h",
"2024-03-15 11:00:00"
)]
#[case::passthrough("FROM events", "2024-03-15 12:00:00")]
#[case::window("FROM events | WINDOW n = count() WITHIN -1h", "2024-03-15 11:00:00")]
#[case::composed_windows(
"FROM events | WINDOW n = count() WITHIN -1h | WINDOW total = sum(n) WITHIN -2h",
"2024-03-15 09:00:00"
)]
#[case::shift_then_window(
"FROM events | SET timestamp = timestamp + 2h | WINDOW n = count() WITHIN -1h",
"2024-03-15 09:00:00"
)]
#[case::window_then_shift(
"FROM events | WINDOW n = count() WITHIN -1h | SET timestamp = timestamp + 2h",
"2024-03-15 09:00:00"
)]
#[case::negative_shift("FROM events | SET timestamp = timestamp - 2h", "2024-03-15 12:00:00")]
#[case::bucket("FROM events | AGG n = count() BY timestamp@d", "2024-03-15 12:00:00")]
#[case::shifted_bucket(
"FROM events | AGG n = count() BY timestamp = timestamp@d + 12h",
"2024-03-15 00:00:00"
)]
#[case::calendar("FROM events | WINDOW n = count() WITHIN -1mon", "2024-02-15 12:00:00")]
#[case::future_window("FROM events | WINDOW n = count() WITHIN 1h", "2024-03-15 12:00:00")]
#[case::range_window(
"FROM events | WINDOW n = count() WITHIN -2h .. 1h",
"2024-03-15 10:00:00"
)]
#[case::definition(
"DEF shifted = FROM events | SET timestamp = timestamp + 2h; FROM shifted | WINDOW n = count() WITHIN -1h",
"2024-03-15 09:00:00"
)]
#[case::nested_definitions(
"DEF a = FROM events | WINDOW n = count() WITHIN -1h; DEF b = FROM a | SET timestamp = timestamp + 2h; FROM b | WINDOW total = sum(n) WITHIN -2h",
"2024-03-15 07:00:00"
)]
#[case::union(
"DEF shifted = FROM events | SET timestamp = timestamp + 2h; UNION shifted, other_events",
"2024-03-15 10:00:00"
)]
#[case::suppression("FROM events | SUPPRESS 1d BY user", "2024-03-15 00:00:00")]
#[case::unused_definition(
"DEF unused = FROM events | WINDOW n = count(); FROM events",
"2024-03-15 12:00:00"
)]
#[case::shared_definition("DEF a = FROM events | WINDOW n = count() WITHIN -1h; DEF b = FROM a | SET timestamp = timestamp + 2h; DEF c = FROM a | SET timestamp = timestamp + 4h; UNION b, c | WINDOW total = sum(n) WITHIN -2h", "2024-03-15 05:00:00")]
fn input_lower_bound(#[case] source: &str, #[case] expected: &str) {
let statement = parse_and_typecheck_with_options(
source,
TypeCheckOptions::builder()
.provider(Arc::new(MockIncrementalProvider))
.maybe_default_space(Some(test_default_space()))
.build(),
)
.into_result()
.unwrap();
let actual = compute_input_lower_bound_for_query(
&statement,
ts("2024-03-15 12:00:00"),
×tamp_field(),
)
.unwrap();
assert_eq!(actual, ts(expected));
}
#[rstest]
#[case::select_removed(
"FROM events | SELECT user",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::drop_removed(
"FROM events | DROP timestamp",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::alternate_clock(
"FROM events | SET other = timestamp | WINDOW n = count() SORT other WITHIN -1h",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::shifted_clock(
"FROM events | WINDOW n = count() SORT timestamp + 2h WITHIN -1h",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::descending_clock(
"FROM events | WINDOW n = count() SORT timestamp DESC WITHIN -1h",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::untraceable_grouping(
"FROM events | SET other = timestamp | WINDOW n = count() BY timestamp = other WITHIN -1h",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::unbounded_window(
"FROM events | WINDOW n = count()",
"WINDOW .. WITHIN has unbound range"
)]
#[case::other_timestamp(
"FROM events | SET other = timestamp | SET timestamp = other",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::constant_timestamp(
"FROM events | SET timestamp = ts('2024-03-15T12:00:00Z')",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::limit(
"FROM events | LIMIT 10",
"Command not supported for incremental refresh: LIMIT"
)]
#[case::nondeterministic_grouping(
"FROM events | WINDOW n = count() BY timestamp = timestamp + (now() - ts('2024-03-15T00:00:00Z')) WITHIN -1h",
"Query uses non-deterministic function: now()"
)]
#[case::nondeterministic_timestamp(
"FROM events | SET timestamp = timestamp + (now() - ts('2024-03-15T00:00:00Z'))",
"Query uses non-deterministic function: now()"
)]
fn input_lower_bound_rejects_unsafe_bounds(#[case] source: &str, #[case] expected: &str) {
let statement = parse_and_typecheck_with_options(
source,
TypeCheckOptions::builder()
.provider(Arc::new(MockIncrementalProvider))
.maybe_default_space(Some(test_default_space()))
.build(),
)
.into_result()
.unwrap();
let error = compute_input_lower_bound_for_query(
&statement,
ts("2024-03-15 12:00:00"),
×tamp_field(),
)
.unwrap_err();
assert_eq!(error.to_string(), expected);
}
#[rstest]
#[case::composed_windows(
"FROM events | WINDOW n = count() WITHIN -1h | WINDOW total = sum(n) WITHIN -2h",
"2024-03-15 09:00:00",
"2024-03-15 13:00:00"
)]
#[case::shift_before_window(
"FROM events | SET timestamp = timestamp + 2h | WINDOW n = count() WITHIN -1h",
"2024-03-15 09:00:00",
"2024-03-15 13:00:00"
)]
#[case::shift_after_window(
"FROM events | WINDOW n = count() WITHIN -1h | SET timestamp = timestamp + 2h",
"2024-03-15 09:00:00",
"2024-03-15 11:00:00"
)]
#[case::composed_forward_windows(
"FROM events | WINDOW n = count() WITHIN 1h | WINDOW total = sum(n) WITHIN 2h",
"2024-03-15 12:00:00",
"2024-03-15 16:00:00"
)]
#[case::grouping_shift(
"FROM events | WINDOW n = count() BY timestamp = timestamp + 2h WITHIN -1h",
"2024-03-15 09:00:00",
"2024-03-15 13:00:00"
)]
fn bounded_backward_requirements(
#[case] source: &str,
#[case] expected_start: &str,
#[case] expected_end: &str,
) {
let statement = parse_and_typecheck_with_options(
source,
TypeCheckOptions::builder()
.provider(Arc::new(MockIncrementalProvider))
.maybe_default_space(Some(test_default_space()))
.build(),
)
.into_result()
.unwrap();
let output = (ts("2024-03-15 12:00:00")..=ts("2024-03-15 13:00:00")).into();
let actual = crate::incremental::compute_query_range_backward_pass(
&statement.pipeline,
&output,
×tamp_field(),
)
.unwrap();
assert_eq!(actual, (ts(expected_start)..=ts(expected_end)).into());
let open_start =
compute_input_lower_bound_for_query(&statement, output.start(), ×tamp_field())
.unwrap();
assert_eq!(open_start, actual.start());
}