use time::OffsetDateTime;
use crate::encode::schema::col;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Resolution {
Elided,
Required,
}
impl Resolution {
pub const fn is_elided(self) -> bool {
matches!(self, Self::Elided)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct VersionStats {
pub min: i128,
pub max: i128,
}
impl VersionStats {
pub const fn single(version: i128) -> Self {
Self {
min: version,
max: version,
}
}
pub const fn may_contain_corrections(self) -> bool {
self.min != self.max
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct FileStats {
pub version: Option<VersionStats>,
pub interval: Option<(OffsetDateTime, OffsetDateTime)>,
}
impl FileStats {
pub const fn unknown() -> Self {
Self {
version: None,
interval: None,
}
}
pub const fn single(version: i128, from: (OffsetDateTime, OffsetDateTime)) -> Self {
Self {
version: Some(VersionStats::single(version)),
interval: Some(from),
}
}
}
pub fn plan(files: &[FileStats]) -> Resolution {
if files.is_empty() {
return Resolution::Elided;
}
let mut single: Vec<(i128, Option<(OffsetDateTime, OffsetDateTime)>)> =
Vec::with_capacity(files.len());
for file in files {
let Some(stats) = file.version else {
return Resolution::Required;
};
if stats.may_contain_corrections() {
return Resolution::Required;
}
single.push((stats.min, file.interval));
}
if single.iter().any(|(_, span)| span.is_none()) {
let first = single[0].0;
return match single.iter().all(|(v, _)| *v == first) {
true => Resolution::Elided,
false => Resolution::Required,
};
}
let mut known: Vec<(OffsetDateTime, OffsetDateTime, i128)> = single
.into_iter()
.map(|(version, span)| {
let (lo, hi) = span.expect("checked above");
(lo, hi, version)
})
.collect();
known.sort_by_key(|(lo, ..)| *lo);
let mut run_end = known[0].1;
let mut run_version = known[0].2;
for &(lo, hi, version) in &known[1..] {
if lo > run_end {
run_end = hi;
run_version = version;
continue;
}
if version != run_version {
return Resolution::Required;
}
run_end = run_end.max(hi);
}
Resolution::Elided
}
pub fn resolution_sql(table: &str) -> String {
resolution_sql_with_key(table, &default_merge_key(), &[], None)
}
fn recorded_at_ceiling_clause(ceiling: Option<OffsetDateTime>) -> String {
match ceiling {
None => String::new(),
Some(at) => {
let micros = crate::encode::schema::micros(at);
format!(
"\n WHERE \"{rec}\" <= arrow_cast({micros}, 'Timestamp(Microsecond, Some(\"UTC\"))')",
rec = col::RECORDED_AT,
)
}
}
}
fn default_merge_key() -> Vec<String> {
crate::encode::schema::MERGE_KEY
.iter()
.map(|s| (*s).to_string())
.collect()
}
pub fn resolution_sql_with_key(
table: &str,
merge_key: &[String],
extra: &[crate::arrow::datatypes::Field],
recorded_at_ceiling: Option<OffsetDateTime>,
) -> String {
let columns = crate::encode::schema::storage_schema(extra)
.fields()
.iter()
.map(|f| format!("\"{}\"", f.name()))
.collect::<Vec<_>>()
.join(", ");
let partition = merge_key
.iter()
.map(|c| format!("\"{c}\""))
.chain(std::iter::once(format!("\"{}\"", col::VERSION_SCOPE)))
.collect::<Vec<_>>()
.join(", ");
format!(
r#"SELECT {columns} FROM (
SELECT *, ROW_NUMBER() OVER (
PARTITION BY {partition}
ORDER BY "{version}" DESC, "{recorded}" DESC
) AS _meterstore_rank
FROM {table}{ceiling}
) AS _meterstore_resolved WHERE _meterstore_rank = 1"#,
version = col::VERSION,
recorded = col::RECORDED_AT,
ceiling = recorded_at_ceiling_clause(recorded_at_ceiling),
)
}
#[cfg(test)]
mod tests {
use super::*;
const V1: i128 = 20_260_701_000_001;
const V2: i128 = 20_260_715_000_002;
fn day(n: i64) -> (OffsetDateTime, OffsetDateTime) {
let start = time::macros::datetime!(2026-07-01 00:00 UTC) + time::Duration::days(n);
(start, start + time::Duration::minutes(60 * 24 - 15))
}
#[test]
fn an_empty_scan_needs_no_resolution() {
assert_eq!(plan(&[]), Resolution::Elided);
}
#[test]
fn single_version_files_elide_resolution() {
let files = vec![FileStats::single(V1, day(0)); 4];
assert!(plan(&files).is_elided());
}
#[test]
fn a_file_spanning_versions_requires_resolution() {
let files = vec![
FileStats::single(V1, day(0)),
FileStats {
version: Some(VersionStats { min: V1, max: V2 }),
interval: Some(day(1)),
},
];
assert_eq!(plan(&files), Resolution::Required);
}
#[test]
fn files_at_differing_versions_over_disjoint_days_still_elide() {
let files: Vec<FileStats> = (0..365)
.map(|n| FileStats::single(V1 + i128::from(n), day(n)))
.collect();
assert!(plan(&files).is_elided(), "a year of daily windows");
}
#[test]
fn files_at_differing_versions_that_overlap_require_resolution() {
let files = vec![FileStats::single(V1, day(0)), FileStats::single(V2, day(0))];
assert_eq!(plan(&files), Resolution::Required);
let straddling = (day(0).0 + time::Duration::hours(12), day(1).1);
assert_eq!(
plan(&[
FileStats::single(V1, day(0)),
FileStats::single(V2, straddling),
]),
Resolution::Required
);
}
#[test]
fn adjacent_days_do_not_count_as_overlapping() {
assert!(plan(&[FileStats::single(V1, day(0)), FileStats::single(V2, day(1)),]).is_elided());
let touching = (day(0).0, day(1).0);
assert_eq!(
plan(&[
FileStats::single(V1, touching),
FileStats::single(V2, day(1)),
]),
Resolution::Required
);
}
#[test]
fn missing_version_statistics_require_resolution() {
let files = vec![
FileStats::single(V1, day(0)),
FileStats {
version: None,
interval: Some(day(1)),
},
];
assert_eq!(plan(&files), Resolution::Required);
}
#[test]
fn an_unknown_span_falls_back_to_requiring_one_version_everywhere() {
let unknown = FileStats {
version: Some(VersionStats::single(V1)),
interval: None,
};
assert!(plan(&[unknown, FileStats::single(V1, day(9))]).is_elided());
assert_eq!(
plan(&[unknown, FileStats::single(V2, day(9))]),
Resolution::Required
);
assert_eq!(plan(&[FileStats::unknown()]), Resolution::Required);
}
#[test]
fn a_single_file_with_one_version_elides() {
assert!(plan(&[FileStats::single(V1, day(0))]).is_elided());
}
#[test]
fn version_stats_detect_a_mixed_file() {
assert!(!VersionStats::single(V1).may_contain_corrections());
assert!(VersionStats { min: V1, max: V2 }.may_contain_corrections());
}
#[test]
fn resolution_sql_ranks_by_version_descending() {
let sql = resolution_sql("readings_versions");
assert!(sql.contains("ROW_NUMBER()"));
assert!(sql.contains(r#"ORDER BY "version" DESC"#));
assert!(sql.contains("_meterstore_rank = 1"));
}
#[test]
fn resolution_sql_breaks_ties_deterministically() {
let sql = resolution_sql("readings_versions");
assert!(
sql.contains(r#"ORDER BY "version" DESC, "recorded_at" DESC"#),
"{sql}"
);
}
#[test]
fn resolution_sql_aliases_the_derived_table() {
let sql = resolution_sql("readings_versions");
assert!(sql.contains(") AS _meterstore_resolved"), "{sql}");
}
#[test]
fn resolution_sql_does_not_leak_the_ranking_column() {
let sql = resolution_sql("readings_versions");
let outer = sql.split("FROM (").next().unwrap();
assert!(!outer.contains("_meterstore_rank"));
assert!(!outer.contains('*'), "outer projection must be explicit");
}
#[test]
fn resolution_sql_projects_every_storage_column() {
let sql = resolution_sql("readings_versions");
for field in crate::encode::schema::storage_schema(&[]).fields() {
assert!(
sql.contains(field.name().as_str()),
"{} missing from projection",
field.name()
);
}
}
#[test]
fn resolution_sql_partitions_by_the_full_merge_key_and_scope() {
let sql = resolution_sql("readings_versions");
for column in [col::MALO_ID, col::OBIS_CODE, col::FROM, col::VERSION_SCOPE] {
assert!(sql.contains(column), "{column} missing from PARTITION BY");
}
}
#[test]
fn resolution_sql_names_the_requested_table() {
assert!(resolution_sql("custom_table").contains("custom_table"));
}
#[test]
fn a_recorded_at_ceiling_filters_the_inner_scan_before_ranking() {
let at = time::macros::datetime!(2026-07-27 06:00 UTC);
let sql = resolution_sql_with_key("readings_versions", &default_merge_key(), &[], Some(at));
let inner = sql.split("AS _meterstore_resolved").next().unwrap();
assert!(inner.contains(r#""recorded_at" <="#), "{sql}");
assert!(inner.contains("arrow_cast"), "{sql}");
assert!(inner.contains("WHERE"), "ceiling belongs in the inner scan");
}
#[test]
fn no_ceiling_leaves_the_inner_scan_unfiltered() {
let sql = resolution_sql_with_key("readings_versions", &default_merge_key(), &[], None);
let inner = sql.split("AS _meterstore_resolved").next().unwrap();
assert!(
!inner.contains("WHERE"),
"no ceiling → no inner filter: {sql}"
);
}
}
#[cfg(test)]
mod properties {
use super::*;
use proptest::prelude::*;
fn at(hours: i64) -> OffsetDateTime {
OffsetDateTime::UNIX_EPOCH + time::Duration::hours(hours)
}
fn file() -> impl Strategy<Value = FileStats> {
let version = prop_oneof![
1 => Just(None),
6 => (0i128..3, 0i128..3)
.prop_map(|(a, b)| Some(VersionStats { min: a.min(b), max: a.max(b) })),
];
let interval = prop_oneof![
1 => Just(None),
6 => (0i64..8, 0i64..8)
.prop_map(|(a, b)| Some((at(a.min(b)), at(a.max(b))))),
];
(version, interval).prop_map(|(version, interval)| FileStats { version, interval })
}
fn may_share_a_key(a: &FileStats, b: &FileStats) -> bool {
match (a.interval, b.interval) {
(Some((a_lo, a_hi)), Some((b_lo, b_hi))) => a_lo <= b_hi && b_lo <= a_hi,
_ => true,
}
}
fn at_most_one_row_per_key(files: &[FileStats]) -> bool {
if !files
.iter()
.all(|f| f.version.is_some_and(|v| v.min == v.max))
{
return false;
}
for (i, a) in files.iter().enumerate() {
for b in &files[i + 1..] {
let (Some(va), Some(vb)) = (a.version, b.version) else {
return false;
};
if va.min != vb.min && may_share_a_key(a, b) {
return false;
}
}
}
true
}
proptest! {
#[test]
fn eliding_implies_no_key_can_appear_twice(
files in prop::collection::vec(file(), 0..8),
) {
prop_assert!(
!plan(&files).is_elided() || at_most_one_row_per_key(&files),
"elided over {files:?}",
);
}
#[test]
fn disjoint_windows_at_any_versions_always_elide(
versions in prop::collection::vec(0i128..1_000, 1..40),
) {
let files: Vec<FileStats> = versions
.into_iter()
.enumerate()
.map(|(i, v)| {
let day = i as i64 * 24;
FileStats::single(v, (at(day), at(day + 23)))
})
.collect();
prop_assert!(plan(&files).is_elided(), "{files:?}");
}
#[test]
fn widening_a_scan_never_grants_elision(
files in prop::collection::vec(file(), 1..6),
extra in file(),
) {
let before = plan(&files).is_elided();
let mut wider = files.clone();
wider.push(extra);
prop_assert!(before || !plan(&wider).is_elided());
}
#[test]
fn the_decision_is_independent_of_file_order(
files in prop::collection::vec(file(), 0..8),
) {
let forward = plan(&files);
let mut reversed = files.clone();
reversed.reverse();
prop_assert_eq!(forward, plan(&reversed), "{:?}", files);
}
}
}