use std::sync::Arc;
use crate::arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit};
use crate::error::{Error, Result};
pub const VALUE_PRECISION: u8 = 18;
pub const VALUE_SCALE: i8 = 6;
pub const VERSION_PRECISION: u8 = 20;
pub const VERSION_SCALE: i8 = 0;
pub const TS_UNIT: TimeUnit = TimeUnit::Microsecond;
pub mod col {
pub const MALO_ID: &str = "malo_id";
pub const MELO_ID: &str = "melo_id";
pub const OBIS_CODE: &str = "obis_code";
pub const SPARTE: &str = "sparte";
pub const FROM: &str = "from";
pub const TO: &str = "to";
pub const VALUE: &str = "value";
pub const UNIT: &str = "unit";
pub const QUALITY: &str = "quality";
pub const RESOLUTION: &str = "resolution";
pub const SOURCE_KIND: &str = "source_kind";
pub const SOURCE_DETAIL: &str = "source_detail";
pub const PROVENANCE: &str = "provenance";
pub const VERSION: &str = "version";
pub const VERSION_SCOPE: &str = "version_scope";
pub const RECORDED_AT: &str = "recorded_at";
pub const BALANCING_DAY: &str = "balancing_day";
}
pub fn timestamp_type() -> DataType {
DataType::Timestamp(TS_UNIT, Some("UTC".into()))
}
#[must_use]
pub fn micros(instant: time::OffsetDateTime) -> i64 {
i64::try_from(instant.unix_timestamp_nanos() / 1_000).unwrap_or(i64::MAX)
}
pub fn instant(micros: i64) -> Result<time::OffsetDateTime> {
time::OffsetDateTime::from_unix_timestamp_nanos(i128::from(micros) * 1_000)
.map_err(|e| Error::decode("timestamp", format!("{micros}: {e}")))
}
#[must_use]
pub fn date32(date: time::Date) -> i32 {
(date - EPOCH).whole_days() as i32
}
pub fn date_of(days: i32) -> Result<time::Date> {
EPOCH
.checked_add(time::Duration::days(i64::from(days)))
.ok_or_else(|| Error::decode(col::BALANCING_DAY, format!("{days} is out of range")))
}
#[must_use]
pub fn timestamp_scalar(instant: time::OffsetDateTime) -> datafusion::common::ScalarValue {
datafusion::common::ScalarValue::TimestampMicrosecond(Some(micros(instant)), Some("UTC".into()))
}
const EPOCH: time::Date = time::macros::date!(1970 - 01 - 01);
pub fn storage_schema(extra: &[Field]) -> SchemaRef {
let mut fields = vec![
Field::new(col::MALO_ID, DataType::Utf8, false),
Field::new(col::MELO_ID, DataType::Utf8, true),
Field::new(col::OBIS_CODE, DataType::Utf8, false),
Field::new(col::SPARTE, DataType::Utf8, false),
Field::new(col::FROM, timestamp_type(), false),
Field::new(col::TO, timestamp_type(), true),
Field::new(
col::VALUE,
DataType::Decimal128(VALUE_PRECISION, VALUE_SCALE),
false,
),
Field::new(col::UNIT, DataType::Utf8, false),
Field::new(col::QUALITY, DataType::Utf8, false),
Field::new(col::RESOLUTION, DataType::Utf8, true),
Field::new(col::SOURCE_KIND, DataType::Utf8, false),
Field::new(col::SOURCE_DETAIL, DataType::Utf8, true),
Field::new(col::PROVENANCE, DataType::Utf8, true),
Field::new(
col::VERSION,
DataType::Decimal128(VERSION_PRECISION, VERSION_SCALE),
false,
),
Field::new(col::VERSION_SCOPE, DataType::Utf8, false),
Field::new(col::RECORDED_AT, timestamp_type(), false),
Field::new(col::BALANCING_DAY, DataType::Date32, false),
];
fields.extend_from_slice(extra);
Arc::new(Schema::new(fields))
}
const _: () = assert!(
VALUE_PRECISION <= 18,
"value precision above 18 loses DELTA_BINARY_PACKED encoding"
);
pub const MERGE_KEY: [&str; 3] = [col::MALO_ID, col::OBIS_CODE, col::FROM];
pub const BLOOM_FILTER_COLUMNS: [&str; 2] = [col::MALO_ID, col::OBIS_CODE];
pub const DELTA_ENCODED_COLUMNS: [&str; 3] = [col::FROM, col::TO, col::VALUE];
pub const SORT_COLUMNS: [&str; 2] = [col::MALO_ID, col::FROM];
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_storage_encoding_holds_every_representable_instant() {
use time::{Date, OffsetDateTime, PrimitiveDateTime, Time};
for date in [Date::MIN, Date::MAX] {
for at in [Time::MIDNIGHT, Time::MAX] {
let instant = PrimitiveDateTime::new(date, at).assume_utc();
let nanos = instant.unix_timestamp_nanos() / 1_000;
assert!(
i64::try_from(nanos).is_ok(),
"{instant} does not fit the storage encoding — has `large-dates` \
been enabled somewhere in the graph?"
);
assert_eq!(micros(instant), nanos as i64);
}
}
let mut at = time::macros::datetime!(1970-01-01 00:00 UTC);
while at < time::macros::datetime!(2100-01-01 00:00 UTC) {
assert_eq!(instant(micros(at)).unwrap(), at);
at += time::Duration::days(97);
}
assert_eq!(
instant(micros(OffsetDateTime::UNIX_EPOCH)).unwrap(),
OffsetDateTime::UNIX_EPOCH
);
}
#[test]
fn a_date_round_trips_through_the_date32_encoding() {
use time::macros::date;
for d in [
date!(1970 - 01 - 01),
date!(2026 - 03 - 29),
date!(2026 - 10 - 25),
date!(2100 - 12 - 31),
] {
assert_eq!(date_of(date32(d)).unwrap(), d);
}
assert_eq!(date32(date!(1970 - 01 - 01)), 0);
assert!(date_of(i32::MAX).is_err());
assert!(date_of(i32::MIN).is_err());
assert!(instant(i64::MAX).is_err());
}
#[test]
fn a_timestamp_literal_carries_the_columns_exact_type() {
let scalar = timestamp_scalar(time::macros::datetime!(2026-07-20 00:00 UTC));
assert_eq!(
scalar.data_type(),
*storage_schema(&[])
.field_with_name(col::FROM)
.unwrap()
.data_type()
);
}
#[test]
fn schema_has_expected_core_columns() {
let s = storage_schema(&[]);
assert_eq!(s.fields().len(), 17);
assert_eq!(s.field(0).name(), col::MALO_ID);
assert_eq!(s.field(15).name(), col::RECORDED_AT);
assert_eq!(s.field(16).name(), col::BALANCING_DAY);
}
#[test]
fn the_value_column_is_never_dimensionless() {
let s = storage_schema(&[]);
assert!(
s.field_with_name("value_kwh").is_err(),
"the old name is gone"
);
for name in [col::VALUE, col::UNIT, col::SPARTE] {
let f = s.field_with_name(name).unwrap();
assert!(!f.is_nullable(), "{name} must be present on every row");
}
}
#[test]
fn sparte_stays_out_of_the_merge_key() {
assert!(!MERGE_KEY.contains(&col::SPARTE));
assert!(!MERGE_KEY.contains(&col::UNIT));
}
#[test]
fn extra_columns_append_without_shifting_core_positions() {
let base = storage_schema(&[]);
let extended = storage_schema(&[
Field::new("bilanzkreis", DataType::Utf8, true),
Field::new("netzgebiet", DataType::Utf8, true),
]);
for (i, f) in base.fields().iter().enumerate() {
assert_eq!(extended.field(i).name(), f.name());
assert_eq!(extended.field(i).data_type(), f.data_type());
}
assert_eq!(extended.fields().len(), 19);
}
#[test]
fn the_interval_end_is_stored_and_never_derived() {
let s = storage_schema(&[]);
let to = s.field_with_name(col::TO).unwrap();
assert_eq!(to.data_type(), ×tamp_type());
assert!(to.is_nullable(), "a point row has no end");
assert!(!s.field_with_name(col::FROM).unwrap().is_nullable());
}
#[test]
fn timestamps_carry_utc_zone() {
let s = storage_schema(&[]);
for name in [col::FROM, col::TO, col::RECORDED_AT] {
match s.field_with_name(name).unwrap().data_type() {
DataType::Timestamp(TimeUnit::Microsecond, Some(tz)) => {
assert_eq!(tz.as_ref(), "UTC", "{name} must be UTC-tagged");
}
other => panic!("{name} has unexpected type {other:?}"),
}
}
}
#[test]
fn merge_key_columns_all_exist_and_are_non_nullable() {
let s = storage_schema(&[]);
for name in MERGE_KEY {
let f = s.field_with_name(name).unwrap();
assert!(!f.is_nullable(), "{name} is part of the merge key");
}
}
#[test]
fn tuning_column_lists_reference_real_columns() {
let s = storage_schema(&[]);
for name in BLOOM_FILTER_COLUMNS
.iter()
.chain(&DELTA_ENCODED_COLUMNS)
.chain(&SORT_COLUMNS)
{
assert!(s.field_with_name(name).is_ok(), "{name} not in schema");
}
}
}