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::trimstrings("FROM events | TRIMSTRINGS 100", "2024-03-15 12:00:00")]
#[case::trimstrings_window(
"FROM events | WINDOW n = count() WITHIN -1h | TRIMSTRINGS 100 | SET timestamp = timestamp + 2h",
"2024-03-15 09:00:00"
)]
#[case::drop_other("FROM events | DROP user", "2024-03-15 12:00:00")]
#[case::explode("FROM events | EXPLODE tag = tags", "2024-03-15 12:00:00")]
#[case::parse("FROM events | PARSE user '*' AS parsed", "2024-03-15 12:00:00")]
#[case::parse_nodrop("FROM events | PARSE user '*' AS parsed NODROP", "2024-03-15 12:00:00")]
#[case::unnest_struct(
"FROM events | SET details = {name: user} | UNNEST details",
"2024-03-15 12:00:00"
)]
#[case::unnest_array("FROM events | UNNEST [{name: user}]", "2024-03-15 12:00:00")]
#[case::row_transforms_between_windows(
"FROM events | WINDOW n = count() WITHIN -1h | EXPLODE tag = tags | PARSE user '*' AS parsed | UNNEST [{name: user}] | TRIMSTRINGS 100 | WINDOW total = sum(n) WITHIN -2h",
"2024-03-15 09:00:00"
)]
#[case::where_sort_within(
"FROM events | WHERE amount > 0 | SORT amount DESC | WITHIN ts('2024-03-01T00:00:00Z') .. ts('2024-04-01T00:00:00Z')",
"2024-03-15 12:00:00"
)]
#[case::distinct_timestamp("FROM events | DISTINCT timestamp, user", "2024-03-15 12:00:00")]
#[case::join_dimension("FROM events | JOIN users ON user == users.id", "2024-03-15 12:00:00")]
#[case::lookup_dimension(
"FROM events | LOOKUP users ON user == users.id",
"2024-03-15 12:00:00"
)]
#[case::lookup_def(
"DEF names = FROM users | SELECT id, name; FROM events | LOOKUP names ON user == names.id",
"2024-03-15 12:00:00"
)]
#[case::lookup_nested_defs("DEF names = FROM users | SELECT id, name; DEF details = FROM names; FROM events | LOOKUP details ON user == details.id", "2024-03-15 12:00:00")]
#[case::lookup_union(
"DEF names = UNION users, users; FROM events | LOOKUP names ON user == names.id",
"2024-03-15 12:00:00"
)]
#[case::nested_timestamp(
"FROM events | NEST record | SELECT timestamp = record.timestamp",
"2024-03-15 12:00:00"
)]
#[case::nested_round_trip("FROM events | NEST record | UNNEST record", "2024-03-15 12:00:00")]
#[case::compound_nest(
"FROM events | NEST outer.inner | SELECT timestamp = outer.inner.timestamp",
"2024-03-15 12:00:00"
)]
#[case::nested_twice(
"FROM events | NEST inner | NEST outer | UNNEST outer | UNNEST inner",
"2024-03-15 12:00:00"
)]
#[case::other_timestamp(
"FROM events | SET other = timestamp | SET timestamp = other",
"2024-03-15 12:00:00"
)]
#[case::nested_parent_projection(
"FROM events | NEST record | SELECT record | SELECT timestamp = record.timestamp",
"2024-03-15 12:00:00"
)]
#[case::nested_parent_alias(
"FROM events | NEST record | SET copied = record | SELECT timestamp = copied.timestamp",
"2024-03-15 12:00:00"
)]
#[case::nest_between_windows(
"FROM events | WINDOW n = count() WITHIN -1h | NEST record | UNNEST record | WINDOW total = sum(n) WITHIN -2h",
"2024-03-15 09:00:00"
)]
#[case::lookup_between_windows(
"FROM events | WINDOW n = count() WITHIN -1h | LOOKUP users ON user == users.id | WINDOW total = sum(n) WITHIN -2h",
"2024-03-15 09: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::explode_timestamp(
"FROM events | EXPLODE timestamp = [ts('2024-03-15T12:00:00Z')]",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::parse_timestamp(
"FROM events | PARSE user '*' AS timestamp",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::unnest_timestamp(
"FROM events | UNNEST {timestamp: ts('2024-03-15T12:00:00Z')}",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::unnest_array_timestamp(
"FROM events | UNNEST [{timestamp: ts('2024-03-15T12:00:00Z')}]",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::unnest_drops_timestamp(
"FROM events | SET timestamp = {name: user} | UNNEST timestamp",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::rows_timestamp(
"ROWS [{timestamp: ts('2024-03-15T12:00:00Z')}]",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::match_pattern(
"MATCH events events WITHIN 1h",
"Command not supported for incremental refresh: MATCH"
)]
#[case::append(
"FROM events | APPEND events",
"DML statements cannot use incremental refresh"
)]
#[case::nest(
"FROM events | NEST record",
"Command not supported for incremental refresh: NEST"
)]
#[case::distinct_without_timestamp(
"FROM events | DISTINCT user",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::nondeterministic_explode(
"FROM events | EXPLODE x = [now()]",
"Query uses non-deterministic function: now()"
)]
#[case::nondeterministic_parse(
"FROM events | PARSE (now() AS string) '*' AS parsed",
"Query uses non-deterministic function: now()"
)]
#[case::nondeterministic_unnest(
"FROM events | UNNEST {x: now()}",
"Query uses non-deterministic function: now()"
)]
#[case::join_timestamped_rhs(
"FROM events | JOIN other_events ON timestamp == other_events.timestamp",
"Command not supported for incremental refresh: JOIN"
)]
#[case::lookup_timestamped_rhs(
"FROM events | LOOKUP other_events ON timestamp == other_events.timestamp",
"Command not supported for incremental refresh: LOOKUP"
)]
#[case::lookup_hidden_timestamp(
"DEF hidden = FROM events | DROP timestamp; FROM events | LOOKUP hidden ON user == hidden.user",
"Command not supported for incremental refresh: LOOKUP"
)]
#[case::lookup_hidden_nested_def(
"DEF hidden = FROM events | SELECT user; DEF names = FROM hidden; FROM events | LOOKUP names ON user == names.user",
"Command not supported for incremental refresh: LOOKUP"
)]
#[case::lookup_hidden_union(
"DEF hidden = FROM events | SELECT user; DEF names = UNION hidden, hidden; FROM events | LOOKUP names ON user == names.user",
"Command not supported for incremental refresh: LOOKUP"
)]
#[case::lookup_nested_join(
"DEF hidden = FROM users | JOIN events ON id == events.user | SELECT id; FROM events | LOOKUP hidden ON user == hidden.id",
"Command not supported for incremental refresh: LOOKUP"
)]
#[case::lookup_clock_alias(
"FROM events | LOOKUP timestamp = users ON user == timestamp.id",
"Command not supported for incremental refresh: LOOKUP"
)]
#[case::nested_overwritten(
"FROM events | NEST record | SET record.timestamp = ts('2024-03-15T12:00:00Z') | SELECT timestamp = record.timestamp",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::nested_parent_overwritten(
"FROM events | NEST record | SET record = {timestamp: ts('2024-03-15T12:00:00Z')} | SELECT timestamp = record.timestamp",
"This query does not compute the timestamp field from the timestamp field"
)]
#[case::other_timestamp_unrelated(
"FROM events | SET other = from_unixtime_micros(amount) | SET timestamp = other",
"This query does not compute the timestamp field from the timestamp field"
)]
#[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::trimstrings_other_timestamp(
"FROM events | SET timestamp = from_unixtime_micros(amount) | TRIMSTRINGS 100",
"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());
}