hamelin_analysis 0.16.4

Analysis utilities for Hamelin query language
Documentation
//! Tests for tabular DEF pipeline (CTE) handling in incremental refresh

use crate::incremental::compute_incremental_ranges_for_query;
use hamelin_lib::tree::builder::{
    at_hour, call, eq, field_ref, hours, pattern_quantified, pipeline, query, seconds, string,
    table_ref,
};

use super::helpers::{
    build_typed_query, multi_stale_ranges, stale_ranges, test_default_space, time_range,
    timestamp_field, MockIncrementalProvider,
};

#[test]
fn test_simple_with_passthrough() {
    // DEF filtered = FROM events | WHERE severity == "high";
    // FROM filtered
    let q = query()
        .main(pipeline().from(|f| f.table_reference("filtered")))
        .def_pipeline(
            "filtered",
            pipeline()
                .from(|f| f.table_reference("events"))
                .where_cmd(eq(field_ref("severity"), string("high"))),
        )
        .build();

    let stale_ranges_map = stale_ranges("events", "2024-01-01 14:00:00", "2024-01-01 16:00:00");

    let typed_query = build_typed_query(q);
    let result = compute_incremental_ranges_for_query(
        &typed_query,
        stale_ranges_map,
        &timestamp_field(),
        None,
        false,
        Some(&test_default_space()),
        &MockIncrementalProvider,
    );

    assert!(result.is_ok(), "Expected success, got error: {:?}", result);
    let ranges = result.unwrap();

    // Both query and replace range should pass through unchanged
    assert_eq!(
        ranges.query_range,
        time_range("2024-01-01 14:00:00", "2024-01-01 16:00:00")
    );
    assert_eq!(
        ranges.replace_range,
        time_range("2024-01-01 14:00:00", "2024-01-01 16:00:00")
    );
}

#[test]
fn test_with_agg_then_filter() {
    // DEF hourly = FROM events | AGG count() BY timestamp@h;
    // FROM hourly | WHERE count > 10
    let q = query()
        .main(
            pipeline()
                .from(|f| f.table_reference("hourly"))
                .where_cmd(eq(field_ref("count"), 10)),
        )
        .def_pipeline(
            "hourly",
            pipeline().from(|f| f.table_reference("events")).agg(|a| {
                a.named_aggregate("count", call("count"))
                    .named_group("timestamp", at_hour(field_ref("timestamp")))
            }),
        )
        .build();

    let stale_ranges_map = stale_ranges("events", "2024-01-01 14:30:00", "2024-01-01 16:45:00");

    let typed_query = build_typed_query(q);
    let result = compute_incremental_ranges_for_query(
        &typed_query,
        stale_ranges_map,
        &timestamp_field(),
        None,
        false,
        Some(&test_default_space()),
        &MockIncrementalProvider,
    );

    assert!(result.is_ok(), "Expected success, got error: {:?}", result);
    let ranges = result.unwrap();

    // Query range should be expanded for hourly aggregation
    assert_eq!(
        ranges.query_range,
        time_range("2024-01-01 14:00:00", "2024-01-01 17:00:00")
    );
    // Replace range should be truncated to hours from the DEF
    assert_eq!(
        ranges.replace_range,
        time_range("2024-01-01 14:00:00", "2024-01-01 16:00:00")
    );
}

#[test]
fn test_multiple_defs_in_sequence() {
    // DEF filtered = FROM events | WHERE severity == "high";
    // DEF hourly = FROM filtered | AGG count() BY timestamp@h;
    // FROM hourly
    let q = query()
        .main(pipeline().from(|f| f.table_reference("hourly")))
        .def_pipeline(
            "filtered",
            pipeline()
                .from(|f| f.table_reference("events"))
                .where_cmd(eq(field_ref("severity"), string("high"))),
        )
        .def_pipeline(
            "hourly",
            pipeline().from(|f| f.table_reference("filtered")).agg(|a| {
                a.named_aggregate("count", call("count"))
                    .named_group("timestamp", at_hour(field_ref("timestamp")))
            }),
        )
        .build();

    let stale_ranges_map = stale_ranges("events", "2024-01-01 14:30:00", "2024-01-01 16:45:00");

    let typed_query = build_typed_query(q);
    let result = compute_incremental_ranges_for_query(
        &typed_query,
        stale_ranges_map,
        &timestamp_field(),
        None,
        false,
        Some(&test_default_space()),
        &MockIncrementalProvider,
    );

    assert!(result.is_ok(), "Expected success, got error: {:?}", result);
    let ranges = result.unwrap();

    // Query range should include the hourly aggregation expansion
    assert_eq!(
        ranges.query_range,
        time_range("2024-01-01 14:00:00", "2024-01-01 17:00:00")
    );
    // Replace range should be truncated to hours (from the last DEF)
    assert_eq!(
        ranges.replace_range,
        time_range("2024-01-01 14:00:00", "2024-01-01 16:00:00")
    );
}

#[test]
fn test_with_merges_multiple_sources() {
    // DEF combined = FROM events, other_events;
    // FROM combined
    let q = query()
        .main(pipeline().from(|f| f.table_reference("combined")))
        .def_pipeline(
            "combined",
            pipeline().from(|f| f.table_reference("events").table_reference("other_events")),
        )
        .build();

    let stale_ranges_map = multi_stale_ranges(&[
        ("events", "2024-01-01 10:00:00", "2024-01-01 12:00:00"),
        ("other_events", "2024-01-01 14:00:00", "2024-01-01 16:00:00"),
    ]);

    let typed_query = build_typed_query(q);
    let result = compute_incremental_ranges_for_query(
        &typed_query,
        stale_ranges_map,
        &timestamp_field(),
        None,
        false,
        Some(&test_default_space()),
        &MockIncrementalProvider,
    );

    assert!(result.is_ok(), "Expected success, got error: {:?}", result);
    let ranges = result.unwrap();

    // Both query and replace range should merge the two source ranges
    assert_eq!(
        ranges.query_range,
        time_range("2024-01-01 10:00:00", "2024-01-01 16:00:00")
    );
    assert_eq!(
        ranges.replace_range,
        time_range("2024-01-01 10:00:00", "2024-01-01 16:00:00")
    );
}

#[test]
fn test_with_window_lookback() {
    // DEF windowed = FROM events | WINDOW sum(amount) WITHIN -1h;
    // FROM windowed
    let q = query()
        .main(pipeline().from(|f| f.table_reference("windowed")))
        .def_pipeline(
            "windowed",
            pipeline()
                .from(|f| f.table_reference("events"))
                .window(|w| {
                    w.named_field("sum_amount", call("sum").arg(field_ref("amount")))
                        .within(hours(-1))
                }),
        )
        .build();

    let stale_ranges_map = stale_ranges("events", "2024-01-01 14:00:00", "2024-01-01 16:00:00");

    let typed_query = build_typed_query(q);
    let result = compute_incremental_ranges_for_query(
        &typed_query,
        stale_ranges_map,
        &timestamp_field(),
        None,
        false,
        Some(&test_default_space()),
        &MockIncrementalProvider,
    );

    assert!(result.is_ok(), "Expected success, got error: {:?}", result);
    let ranges = result.unwrap();

    // Query range should include 1h lookback from the DEF
    assert_eq!(
        ranges.query_range,
        time_range("2024-01-01 13:00:00", "2024-01-01 16:00:00")
    );
    // Replace range unchanged
    assert_eq!(
        ranges.replace_range,
        time_range("2024-01-01 14:00:00", "2024-01-01 16:00:00")
    );
}

/// MATCH as data source with asymmetric CTE stale ranges.
///
/// The two CTEs have different stale ranges, so the MATCH replace_range should be
/// the bounding range (min start, max end) across all pattern references.
/// Query range expands backward by the WITHIN interval.
#[test]
fn test_match_as_data_source_with_asymmetric_stale_ranges() {
    // DEF source_a = FROM events | SELECT timestamp, user, severity;
    // DEF source_b = FROM other_events | SELECT timestamp, data;
    // MATCH source_a source_b=source_b+ AGG event_end = max(timestamp) BY user WITHIN 5s
    // | SELECT timestamp, user
    let q = query()
        .main(
            pipeline()
                .match_cmd(|m| {
                    m.pattern(pattern_quantified(table_ref("source_a"), None))
                        .pattern(pattern_quantified(
                            table_ref("source_b"),
                            Some(hamelin_lib::tree::builder::quantifier_at_least_one()),
                        ))
                        .agg("event_end", call("max").arg(field_ref("timestamp")))
                        .group_by("user")
                        .within(seconds(5))
                })
                .select(|s| s.field("timestamp").field("user")),
        )
        .def_pipeline(
            "source_a",
            pipeline()
                .from(|f| f.table_reference("events"))
                .select(|s| s.field("timestamp").field("user").field("severity")),
        )
        .def_pipeline(
            "source_b",
            pipeline()
                .from(|f| f.table_reference("other_events"))
                .select(|s| s.field("timestamp").field("data")),
        )
        .build();

    let stale_ranges_map = multi_stale_ranges(&[
        ("events", "2024-01-01 14:00:00", "2024-01-01 16:00:00"),
        ("other_events", "2024-01-01 15:00:00", "2024-01-01 18:00:00"),
    ]);

    let typed_query = build_typed_query(q);
    let result = compute_incremental_ranges_for_query(
        &typed_query,
        stale_ranges_map,
        &timestamp_field(),
        None,
        false,
        Some(&test_default_space()),
        &MockIncrementalProvider,
    );

    assert!(result.is_ok(), "Expected success, got error: {:?}", result);
    let ranges = result.unwrap();

    // Replace range is the bounding range across both CTEs
    assert_eq!(
        ranges.replace_range,
        time_range("2024-01-01 14:00:00", "2024-01-01 18:00:00")
    );

    // Query range expands backward by 5s from the bounding range
    assert_eq!(
        ranges.query_range,
        time_range("2024-01-01 13:59:55", "2024-01-01 18:00:00")
    );
}

/// MATCH as data source without WITHIN — no query range expansion.
///
/// When MATCH has no WITHIN clause, the backward pass hits MATCH and breaks
/// with no expansion. Query range equals replace range (plus any CTE expansion).
#[test]
fn test_match_as_data_source_without_within() {
    // DEF source_a = FROM events | SELECT timestamp, user;
    // MATCH source_a AGG event_end = max(timestamp) BY user
    // | SELECT timestamp, user
    let q = query()
        .main(
            pipeline()
                .match_cmd(|m| {
                    m.pattern(pattern_quantified(table_ref("source_a"), None))
                        .agg("event_end", call("max").arg(field_ref("timestamp")))
                        .group_by("user")
                })
                .select(|s| s.field("timestamp").field("user")),
        )
        .def_pipeline(
            "source_a",
            pipeline()
                .from(|f| f.table_reference("events"))
                .select(|s| s.field("timestamp").field("user")),
        )
        .build();

    let stale_ranges_map = stale_ranges("events", "2024-01-01 14:00:00", "2024-01-01 16:00:00");

    let typed_query = build_typed_query(q);
    let result = compute_incremental_ranges_for_query(
        &typed_query,
        stale_ranges_map,
        &timestamp_field(),
        None,
        false,
        Some(&test_default_space()),
        &MockIncrementalProvider,
    );

    assert!(result.is_ok(), "Expected success, got error: {:?}", result);
    let ranges = result.unwrap();

    // Both ranges should be identical — no WITHIN expansion
    assert_eq!(
        ranges.replace_range,
        time_range("2024-01-01 14:00:00", "2024-01-01 16:00:00")
    );
    assert_eq!(
        ranges.query_range,
        time_range("2024-01-01 14:00:00", "2024-01-01 16:00:00")
    );
}