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() {
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,
×tamp_field(),
None,
false,
Some(&test_default_space()),
&MockIncrementalProvider,
);
assert!(result.is_ok(), "Expected success, got error: {:?}", result);
let ranges = result.unwrap();
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() {
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,
×tamp_field(),
None,
false,
Some(&test_default_space()),
&MockIncrementalProvider,
);
assert!(result.is_ok(), "Expected success, got error: {:?}", result);
let ranges = result.unwrap();
assert_eq!(
ranges.query_range,
time_range("2024-01-01 14:00:00", "2024-01-01 17:00:00")
);
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() {
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,
×tamp_field(),
None,
false,
Some(&test_default_space()),
&MockIncrementalProvider,
);
assert!(result.is_ok(), "Expected success, got error: {:?}", result);
let ranges = result.unwrap();
assert_eq!(
ranges.query_range,
time_range("2024-01-01 14:00:00", "2024-01-01 17: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_merges_multiple_sources() {
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,
×tamp_field(),
None,
false,
Some(&test_default_space()),
&MockIncrementalProvider,
);
assert!(result.is_ok(), "Expected success, got error: {:?}", result);
let ranges = result.unwrap();
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() {
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,
×tamp_field(),
None,
false,
Some(&test_default_space()),
&MockIncrementalProvider,
);
assert!(result.is_ok(), "Expected success, got error: {:?}", result);
let ranges = result.unwrap();
assert_eq!(
ranges.query_range,
time_range("2024-01-01 13: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_match_as_data_source_with_asymmetric_stale_ranges() {
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,
×tamp_field(),
None,
false,
Some(&test_default_space()),
&MockIncrementalProvider,
);
assert!(result.is_ok(), "Expected success, got error: {:?}", result);
let ranges = result.unwrap();
assert_eq!(
ranges.replace_range,
time_range("2024-01-01 14:00:00", "2024-01-01 18:00:00")
);
assert_eq!(
ranges.query_range,
time_range("2024-01-01 13:59:55", "2024-01-01 18:00:00")
);
}
#[test]
fn test_match_as_data_source_without_within() {
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,
×tamp_field(),
None,
false,
Some(&test_default_space()),
&MockIncrementalProvider,
);
assert!(result.is_ok(), "Expected success, got error: {:?}", result);
let ranges = result.unwrap();
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")
);
}