use rstest::rstest;
use crate::incremental::{detect_supported_strategies_for_pipeline, IncrementalStrategyKind};
use hamelin_lib::tree::builder::{
at_hour, call, eq, field_ref, hours, pipeline, string, PipelineBuilder,
};
use super::helpers::build_pipeline;
#[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]
)]
#[case::nest_disqualifies_cascaded_append(
pipeline()
.from(|f| f.table_reference("events"))
.nest("user_info"),
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);
}