hamelin_analysis 0.23.3

Analysis utilities for Hamelin query language
Documentation
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"),
        &timestamp_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"),
        &timestamp_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,
        &timestamp_field(),
    )
    .unwrap();
    assert_eq!(actual, (ts(expected_start)..=ts(expected_end)).into());
    let open_start =
        compute_input_lower_bound_for_query(&statement, output.start(), &timestamp_field())
            .unwrap();
    assert_eq!(open_start, actual.start());
}