hamelin_analysis 0.16.4

Analysis utilities for Hamelin query language
Documentation
//! Tests for boundary conditions of range semantics.
//!
//! TimeRange supports both inclusive end (non-snapped) and exclusive end (snapped).
//! These tests verify correct behavior at boundaries, especially when timestamps
//! fall exactly on partition boundaries.

use rstest::rstest;

use crate::incremental::{compute_incremental_ranges_for_pipeline, TimeRange};
use hamelin_lib::tree::{
    ast::expression::TruncUnit,
    builder::{at_day, at_hour, call, eq, field_ref, hours, pipeline, string},
};

use super::helpers::{
    build_pipeline, stale_ranges, test_default_space, time_range, time_range_exclusive,
    timestamp_field, ts, MockIncrementalProvider,
};

/// Snapping to partition boundaries should produce exclusive-end ranges.
/// End always advances by one unit (via next_truncation_boundary), even if already on a boundary.
#[rstest]
#[case::end_on_hour_boundary(
    "2024-01-01 11:00:00", "2024-01-01 12:00:00",
    (TruncUnit::Hour, 1),
    "2024-01-01 11:00:00", "2024-01-01 13:00:00"
)]
#[case::both_on_hour_boundaries(
    "2024-01-01 10:00:00", "2024-01-01 14:00:00",
    (TruncUnit::Hour, 1),
    "2024-01-01 10:00:00", "2024-01-01 15:00:00"
)]
#[case::end_on_day_boundary(
    "2024-01-01 00:00:00", "2024-01-02 00:00:00",
    (TruncUnit::Day, 1),
    "2024-01-01 00:00:00", "2024-01-03 00:00:00"
)]
fn test_snap_boundary_produces_exclusive_end(
    #[case] stale_start: &str,
    #[case] stale_end: &str,
    #[case] partition_unit: (TruncUnit, u32),
    #[case] expected_start: &str,
    #[case] expected_end: &str,
) {
    let pipeline = build_pipeline(
        pipeline()
            .from(|f| f.table_reference("events"))
            .where_cmd(eq(field_ref("severity"), string("high"))),
    );

    let provider = MockIncrementalProvider;
    let result = compute_incremental_ranges_for_pipeline(
        &pipeline,
        stale_ranges("events", stale_start, stale_end),
        &timestamp_field(),
        Some(partition_unit),
        false,
        &[],
        Some(&test_default_space()),
        &provider,
    )
    .unwrap();

    assert_eq!(
        result.replace_range,
        time_range_exclusive(expected_start, expected_end),
    );
}

/// AGG with end exactly on a truncation boundary should produce inclusive ranges
/// (no snapping) and correct backward-pass query range expansion.
#[rstest]
#[case::hourly_agg(
    "2024-01-01 14:00:00", "2024-01-01 16:00:00",
    at_hour(field_ref("timestamp")),
    "2024-01-01 14:00:00", "2024-01-01 16:00:00",  // replace_range
    "2024-01-01 14:00:00", "2024-01-01 17:00:00"   // query_range
)]
#[case::daily_agg(
    "2024-01-01 00:00:00", "2024-01-03 00:00:00",
    at_day(field_ref("timestamp")),
    "2024-01-01 00:00:00", "2024-01-03 00:00:00",  // replace_range
    "2024-01-01 00:00:00", "2024-01-04 00:00:00"   // query_range
)]
fn test_agg_end_on_boundary(
    #[case] stale_start: &str,
    #[case] stale_end: &str,
    #[case] trunc_expr: impl hamelin_lib::tree::builder::IntoExpressionBuilder,
    #[case] expected_replace_start: &str,
    #[case] expected_replace_end: &str,
    #[case] expected_query_start: &str,
    #[case] expected_query_end: &str,
) {
    let pipeline = build_pipeline(pipeline().from(|f| f.table_reference("events")).agg(|a| {
        a.named_aggregate("count", call("count"))
            .named_group("timestamp", trunc_expr)
    }));

    let provider = MockIncrementalProvider;
    let result = compute_incremental_ranges_for_pipeline(
        &pipeline,
        stale_ranges("events", stale_start, stale_end),
        &timestamp_field(),
        None,
        false,
        &[],
        Some(&test_default_space()),
        &provider,
    )
    .unwrap();

    assert_eq!(
        result.replace_range,
        time_range(expected_replace_start, expected_replace_end),
    );
    assert_eq!(
        result.query_range,
        time_range(expected_query_start, expected_query_end),
    );
}

/// Verify that inclusive ranges use the Inclusive variant
/// and exclusive ranges use the Exclusive variant.
#[test]
fn time_range_inclusive_vs_exclusive() {
    let inclusive = time_range("2024-01-01 10:00:00", "2024-01-01 12:00:00");
    assert_eq!(inclusive.start(), ts("2024-01-01 10:00:00"));
    assert_eq!(inclusive.end(), ts("2024-01-01 12:00:00"));
    assert!(matches!(inclusive, TimeRange::Inclusive(_)));

    let exclusive = time_range_exclusive("2024-01-01 10:00:00", "2024-01-01 12:00:00");
    assert_eq!(exclusive.start(), ts("2024-01-01 10:00:00"));
    assert_eq!(exclusive.end(), ts("2024-01-01 12:00:00"));
    assert!(matches!(exclusive, TimeRange::Exclusive(_)));
}

/// Window lookback with end on hour boundary + hourly snap.
/// Verifies lookback is applied correctly relative to snapped boundaries.
#[test]
fn window_lookback_with_end_on_boundary_and_snap() {
    let pipeline = build_pipeline(
        pipeline()
            .from(|f| f.table_reference("events"))
            .window(|w| {
                w.named_field("sum_amount", call("sum").arg(field_ref("amount")))
                    .within(hours(-1))
            }),
    );

    let provider = MockIncrementalProvider;
    let result = compute_incremental_ranges_for_pipeline(
        &pipeline,
        stale_ranges("events", "2024-01-01 14:00:00", "2024-01-01 16:00:00"),
        &timestamp_field(),
        Some((TruncUnit::Hour, 1)),
        false,
        &[],
        Some(&test_default_space()),
        &provider,
    )
    .unwrap();

    // Snapped replace_range: [14:00, 17:00) — end snaps forward one unit, exclusive
    assert_eq!(
        result.replace_range,
        time_range_exclusive("2024-01-01 14:00:00", "2024-01-01 17:00:00"),
    );
    // Query range: 1h lookback from snapped replace_range start
    assert_eq!(
        result.query_range,
        time_range("2024-01-01 13:00:00", "2024-01-01 17:00:00"),
    );
}