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
}
}
pub fn plan(files: &[Option<VersionStats>]) -> Resolution {
if files.is_empty() {
return Resolution::Elided;
}
let mut seen: Option<i128> = None;
for file in files {
let Some(stats) = file else {
return Resolution::Required;
};
if stats.may_contain_corrections() {
return Resolution::Required;
}
match seen {
Some(v) if v != stats.min => return Resolution::Required,
_ => seen = Some(stats.min),
}
}
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 = at.unix_timestamp_nanos() / 1_000;
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
) AS _meterstore_rank
FROM {table}{ceiling}
) AS _meterstore_resolved WHERE _meterstore_rank = 1"#,
version = col::VERSION,
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;
#[test]
fn an_empty_scan_needs_no_resolution() {
assert_eq!(plan(&[]), Resolution::Elided);
}
#[test]
fn single_version_files_elide_resolution() {
let files = vec![Some(VersionStats::single(V1)); 4];
assert!(plan(&files).is_elided());
}
#[test]
fn a_file_spanning_versions_requires_resolution() {
let files = vec![
Some(VersionStats::single(V1)),
Some(VersionStats { min: V1, max: V2 }),
];
assert_eq!(plan(&files), Resolution::Required);
}
#[test]
fn files_at_differing_versions_require_resolution() {
let files = vec![
Some(VersionStats::single(V1)),
Some(VersionStats::single(V2)),
];
assert_eq!(plan(&files), Resolution::Required);
}
#[test]
fn missing_statistics_require_resolution() {
let files = vec![Some(VersionStats::single(V1)), None];
assert_eq!(plan(&files), Resolution::Required);
}
#[test]
fn a_single_file_with_one_version_elides() {
assert!(plan(&[Some(VersionStats::single(V1))]).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_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 stats() -> impl Strategy<Value = Option<VersionStats>> {
prop_oneof![
1 => Just(None),
6 => (0i128..4, 0i128..4)
.prop_map(|(a, b)| Some(VersionStats { min: a.min(b), max: a.max(b) })),
]
}
proptest! {
#[test]
fn resolution_is_skipped_only_when_one_version_is_provable(
files in prop::collection::vec(stats(), 0..8),
) {
let provable = files.iter().all(|f| f.is_some_and(|s| s.min == s.max))
&& files
.iter()
.filter_map(|f| f.map(|s| s.min))
.collect::<std::collections::BTreeSet<_>>()
.len()
<= 1;
prop_assert_eq!(plan(&files).is_elided(), provable, "{:?}", files);
}
#[test]
fn widening_a_scan_never_grants_elision(
files in prop::collection::vec(stats(), 1..6),
extra in stats(),
) {
let before = plan(&files).is_elided();
let mut wider = files.clone();
wider.push(extra);
prop_assert!(before || !plan(&wider).is_elided());
}
}
}