hamelin_analysis 0.11.1

Analysis utilities for Hamelin query language
Documentation
//! Tests for strategy detection (which incremental strategies are supported)

use rstest::rstest;
use std::{collections::HashMap, sync::Arc};

use crate::incremental::{
    detect_supported_strategies_for_pipeline, detect_supported_strategies_for_query,
    IncrementalStrategyKind,
};
use hamelin_lib::tree::builder::{
    at_hour, call, eq, field_ref, hours, pipeline, string, PipelineBuilder,
};
use hamelin_lib::{
    parse,
    tree::{
        ast::query::Query,
        options::{TemplateParameterKind, TypeCheckOptions},
    },
    type_check_with_options,
};

use super::helpers::{build_pipeline, MockIncrementalProvider};

#[rstest]
#[case::simple_filter(
    pipeline()
        .from(|f| f.table_reference("events"))
        .where_cmd(eq(field_ref("severity"), string("high"))),
    vec![IncrementalStrategyKind::CascadedAppend, IncrementalStrategyKind::TimeRangeRefresh]
)]
#[case::set_and_select(
    pipeline()
        .from(|f| f.table_reference("events"))
        .set_cmd(|l| l.named_field("risk_score", field_ref("severity")))
        .select(|s| s.field("timestamp").field("user").field("risk_score")),
    vec![IncrementalStrategyKind::CascadedAppend, IncrementalStrategyKind::TimeRangeRefresh]
)]
#[case::drop_non_timestamp(
    pipeline()
        .from(|f| f.table_reference("events"))
        .drop(|d| d.field("extra_field")),
    vec![IncrementalStrategyKind::CascadedAppend, IncrementalStrategyKind::TimeRangeRefresh]
)]
#[case::window_disqualifies_cascaded_append(
    pipeline()
        .from(|f| f.table_reference("events"))
        .window(|w| w
            .named_field("count", call("count"))
            .within(hours(-1))
        ),
    vec![IncrementalStrategyKind::TimeRangeRefresh]
)]
#[case::agg_disqualifies_cascaded_append(
    pipeline()
        .from(|f| f.table_reference("events"))
        .agg(|a| a
            .named_aggregate("count", call("count"))
            .named_group("timestamp", at_hour(field_ref("timestamp")))
        ),
    vec![IncrementalStrategyKind::TimeRangeRefresh]
)]
#[case::parse_disqualifies_cascaded_append(
    pipeline()
        .from(|f| f.table_reference("events"))
        .parse(|p| p
            .pattern("*")
            .identifier("entire_event")
        ),
    vec![IncrementalStrategyKind::TimeRangeRefresh]
)]
#[case::explode_allows_both_strategies(
    pipeline()
        .from(|f| f.table_reference("events"))
        .explode(|e| e.named_field("tag", field_ref("tags"))),
    vec![IncrementalStrategyKind::CascadedAppend, IncrementalStrategyKind::TimeRangeRefresh]
)]
// TODO: Fix UNNEST - may need schema or builder adjustments
// #[case::unnest_allows_both_strategies(
//     pipeline()
//         .from(|f| f.table_reference("events"))
//         .unnest(field_ref("tags")),
//     vec![IncrementalStrategyKind::CascadedAppend, IncrementalStrategyKind::TimeRangeRefresh]
// )]
#[case::nest_disqualifies_cascaded_append(
    pipeline()
        .from(|f| f.table_reference("events"))
        .nest("user_info"),
    vec![IncrementalStrategyKind::TimeRangeRefresh]
)]
#[case::suppress_disqualifies_cascaded_append(
    pipeline()
        .from(|f| f.table_reference("events"))
        .suppress(hours(1), |s| s),
    vec![IncrementalStrategyKind::TimeRangeRefresh]
)]
#[case::non_deterministic_disqualifies_all_strategies_in_where(
    pipeline()
        .from(|f| f.table_reference("events"))
        .where_cmd(eq(field_ref("timestamp"), call("now"))),
    vec![]
)]
#[case::non_deterministic_disqualifies_all_strategies_in_set(
    pipeline()
        .from(|f| f.table_reference("events"))
        .set_cmd(|l| l.named_field("current_time", call("today"))),
    vec![]
)]
#[case::non_deterministic_disqualifies_all_strategies_in_window(
    pipeline()
        .from(|f| f.table_reference("events"))
        .window(|w| w.named_field("current", call("yesterday"))),
    vec![]
)]
#[case::non_deterministic_disqualifies_all_strategies_in_agg(
    pipeline()
        .from(|f| f.table_reference("events"))
        .agg(|a| a
            .named_aggregate("count", call("count"))
            .named_aggregate("as_of", call("tomorrow"))
            .group_by("timestamp")
        ),
    vec![]
)]
#[case::non_deterministic_disqualifies_all_strategies_in_within(
    pipeline()
        .from(|f| f.table_reference("events"))
        .within(hours(1)),
    vec![]
)]
#[case::join_disqualifies_all_strategies_without_allow_lookups(
    pipeline()
        .from(|f| f.table_reference("events"))
        .join("users", eq(field_ref("events.user_id"), field_ref("users.id"))),
    vec![]
)]
fn test_strategy_detection(
    #[case] pipeline_builder: PipelineBuilder,
    #[case] expected_strategies: Vec<IncrementalStrategyKind>,
) {
    let pipeline = build_pipeline(pipeline_builder);
    let result = detect_supported_strategies_for_pipeline(pipeline.into(), false);
    assert_eq!(result.supported, expected_strategies);
}

#[rstest]
#[case::join_allows_both_strategies_with_allow_lookups(
    pipeline()
        .from(|f| f.table_reference("events"))
        .join("users", eq(field_ref("events.user_id"), field_ref("users.id"))),
    vec![IncrementalStrategyKind::CascadedAppend, IncrementalStrategyKind::TimeRangeRefresh]
)]
#[case::lookup_allows_both_strategies_with_allow_lookups(
    pipeline()
        .from(|f| f.table_reference("events"))
        .lookup("users", |l| l.on(eq(field_ref("events.user_id"), field_ref("users.id")))),
    vec![IncrementalStrategyKind::CascadedAppend, IncrementalStrategyKind::TimeRangeRefresh]
)]
fn test_strategy_detection_allow_lookups(
    #[case] pipeline_builder: PipelineBuilder,
    #[case] expected_strategies: Vec<IncrementalStrategyKind>,
) {
    let pipeline = build_pipeline(pipeline_builder);
    let result = detect_supported_strategies_for_pipeline(pipeline.into(), true);
    assert_eq!(result.supported, expected_strategies);
}

#[test]
fn templated_table_reference_disqualifies_strategy_detection() {
    let mut template_parameters = HashMap::new();
    template_parameters.insert(
        "source".to_string(),
        TemplateParameterKind::IdentifierFragment(vec!["events".to_string()]),
    );
    let query = parse("FROM ${source}")
        .into_result()
        .expect("query should parse");
    let typed = type_check_with_options::<Query>(
        query,
        TypeCheckOptions::builder()
            .provider(Arc::new(MockIncrementalProvider))
            .template_parameters(Arc::new(template_parameters))
            .build(),
    )
    .into_result()
    .expect("query should typecheck");

    let result = detect_supported_strategies_for_query(&typed, false);

    assert!(result.supported.is_empty());
    assert!(result
        .rejections
        .iter()
        .any(|r| r.contains("substitute template parameters")));
}

#[test]
fn expression_template_parameter_disqualifies_query_strategy_detection() {
    let mut template_parameters = HashMap::new();
    template_parameters.insert(
        "sev".to_string(),
        TemplateParameterKind::Primitive(hamelin_lib::types::STRING),
    );
    let query = parse("FROM events | WHERE severity == ${sev}")
        .into_result()
        .expect("query should parse");
    let typed = type_check_with_options::<Query>(
        query,
        TypeCheckOptions::builder()
            .provider(Arc::new(MockIncrementalProvider))
            .template_parameters(Arc::new(template_parameters))
            .build(),
    )
    .into_result()
    .expect("query should typecheck");

    let result = detect_supported_strategies_for_query(&typed, false);

    assert!(result.supported.is_empty());
    assert!(result
        .rejections
        .iter()
        .any(|r| r.contains("substitute template parameters")));
}