Expand description
§MeterStore
Hot/cold tiered store for metering time series: PostgreSQL holds the recent interval window, Apache Iceberg holds the history, and a single explicit timestamp — the tiering watermark — separates them.
MeterStore is the persistence layer of the mako
platform. It stores the types defined by the metering crate; it does not
redefine them and it does not compute with them.
§What this crate owns
Exactly three things beyond storage mechanics:
- Correction versioning —
version, implementing MSCONS’s rule that versions are only comparable within a (network operator, month) scope. - The transaction-time axis —
recorded_at, giving bitemporality alongsidemetering’s valid time (from/to). - The tiering boundary —
watermark.
Everything else — DST calendars, unit conversion, quality semantics,
validation, Ersatzwertbildung — belongs to metering. Duplicating any of
it here would create a second implementation to keep correct, and it would
drift.
planner re-exports metering::calendar for convenience, so callers get
DST-correct local days without depending on both crates directly. The
implementation is upstream and there is only one of it. The choice between
the two calendars is upstream too, as
DayBoundary — so the one thing
planner::calendar adds is neither arithmetic nor the choice, but a
storage fact: a stored row carries its sparte, so the store knows which
boundary applies to it. That mapping is planner::day_boundary, written
once.
§The five things a caller usually wants
MeterStore::sql— DataFusion’s ownDataFrameover a table that spans both tiers.MeterStore::queryreturns the same rows plus the boundary they were computed against (P1), andMeterStore::streamreturns that boundary before the rows, for a result too large to hold.MeterStore::series— one measuring point as ameteringMeasurementSeries, version-resolved and tier-split.collectdescribes one channel, because aMeasurementSeriesholds oneobis_code;collect_by_channeldescribes the whole measuring point, splitting one scan into a series per channel rather than reading each in turn.MeterStore::as_of— the same queries against a pinned Iceberg snapshot and an optional version ceiling, for a settlement rerun.MeterStore::completeness— whether a range holds what the DST-aware calendar says it should, because a missing interval is information rather than an empty set.seen_sinceadds the finding a range cannot make about itself: a channel that delivered nothing has no rows to aggregate, so it needs a roster drawn from an earlier window.MeterCatalog— several tables in one session, when a deployment holds more than one stream and needs a statement that mentions both.
Three more that a service exposing SQL will want: MeterStore::scoped
confines a session to one identity value, MeterCatalog::isolated confines
it to one table, and MeterCatalog::scoped confines every table of a
catalog to one identity value — which is what a multi-tenant deployment
putting a whole catalog on a socket needs, since the first two together force
a choice between the cross-table join and the tenant boundary. All three
inject into the plan, so caller-supplied SQL cannot step past them.
MeterStore::append_authoritative is the write path for a value the
operator authors rather than receives.
§What is stored
Two shapes, declared per table by TimeModel. A Lastgang is energy
over [from, to); a Zählerstandsgang is a cumulative register value at an
instant, which BK6-24-174 has made a primary record as voluminous as the
Lastgang derived from it. They are never the same table, because value
would mean two things in one column and no aggregate could tell them apart.
The shape also decides what names a reading. A Marktlokation may be
measured by several Messlokationen, and both meters carry the same OBIS
register at the same instants. A load profile belongs to the market location,
so melo_id labels it; a register belongs to the meter, so melo_id names
it and joins the merge key. A point table does that by default;
TableConfig::identify_by_melo pins it either way.
All four Sparten. value carries the quantity and unit its dimension,
because water is metered and billed in m³ and gas may sit on either side of
the Brennwert conversion — a column named for kilowatt-hours would be wrong
for half of them. A unit the commodity cannot be expressed in is refused at
the write.
Any declared resolution, not only the quarter-hour: completeness asks
metering’s calendar per day, so one-minute data expects 1 440 intervals on
an ordinary day and 1 500 on the 25-hour autumn one.
And two kinds of day. Electricity, heat and water are balanced on the
Berlin calendar day; gas is balanced on the Gastag, 06:00 to 06:00 local.
Grouping a gas Lastgang by the calendar day books its 00:00–06:00 draw into
the neighbouring Bilanzierungstag — six hours a day, every day — so the
bucketing and the expected interval count both follow the row’s Sparte, in
SQL through meter_balancing_day and in Rust through
planner::balancing_day. The month is the same choice one period up —
meter_balancing_month and planner::balancing_month.
The boundary carries up to the month, which is the market’s own rule:
EDI@Energy Allgemeine Festlegungen v6.1c, Kap. 3.1 defines the gas
Bilanzierungsmonat as 01.06 06:00 to 01.07 06:00. So a version scope — which
is (network operator, month) — is cut at 06:00 for gas too, and every
VersionScope constructor takes a Sparte for that reason.
An external engine has neither function, and SQL dialects differ on timestamp
arithmetic — so the answer is stored. Every row carries a balancing_day
column derived once by the encoder, and reading the Iceberg files directly
needs a GROUP BY and no calendar reasoning. It is the single derived value
this crate persists; encode::schema argues the exception.
Identifiers are parsed rather than trusted. malo_id and melo_id are
metering’s MaloId and MeloId on both sides of the encoding: a
MaLo-ID carries a check digit so that a transposition is detectable, and a
store that accepted eleven arbitrary digits would throw that away at the one
point it still mattered.
A deployment’s own columns get the same treatment when they hold an
identifier. checked_column declares one whose values must parse as a
ValueCheck — an Eic, the ENTSO-E code a Bilanzkreis is addressed by;
a MaloId or MeloId, for a reading that references a measuring point
it is not keyed to; or a BdewCode, the Marktpartner-ID of a Lieferant or
Messstellenbetreiber. Values are stored canonicalised, so such a column may
sit in the merge key without two spellings becoming two readings.
Each scheme stops somewhere different, and ValueCheck says where: the EIC
check character and the MaLo check digit are enforced, a MeLo has no check
digit to enforce, and a Marktpartner-ID’s thirteenth digit is deliberately
not checked — BDEW’s Bildungsvorschrift carves out GS1-issued GLNs, which
is the same reason VersionScope does not check it either.
An EIC column may name the object type it holds —
ValueCheck::Eic(Some(EicType::Party)) for a Bilanzkreis, Area for a
Bilanzierungsgebiet. The two share the alphabet, the length and the check
character, so position 3 is the only thing that tells them apart — and unlike
the check character it is expressible as a regular expression, so declaring it
strengthens the hot table’s CHECK as well as the write path.
§Tiers and surfaces
The cold tier accepts any Arc<dyn Catalog>; three are built for you —
IcebergSqlCatalog over the same PostgreSQL as the hot tier,
cold::IcebergRestCatalog over Polaris/Lakekeeper/Nessie/Gravitino
(rest-catalog, on by default), and cold::S3TablesCatalog over an AWS S3
Tables table bucket (s3tables). Settings::connect builds whichever a
configuration file names, along with the pool and every validated table.
Two serving surfaces: a read-only Iceberg REST façade (catalog-facade) for
SQL-catalog deployments, and Flight SQL (flight) for the unified hot + cold
view — the one thing an external client cannot assemble for itself. Flight
serves any SqlSurface, so a whole MeterCatalog goes on a socket as
readily as one table.
And a command line, behind cli: cli is the meterstore binary over the
same public API — check, create, status, archive, maintain,
query, completeness and serve among its verbs — for the questions an
operator asks during an incident and the archival loop a deployment has to
run somewhere.
§Errors carry what to do about them
Error is #[non_exhaustive] and callers match on variants rather than
parsing strings. Error::is_retryable is the split that matters most: a
lost connection and a lock a statement declined to wait for are worth
retrying, and a refused delivery, an invalid configuration and a statement
that will not plan are not — retrying those is a loop on a message that will
never change.
Two are deliberately distinct and are the pair most often conflated.
Error::IntegrityViolation means the store stopped something from
becoming true: an overlapping delivery, two network operators for one
reading, a value restated under an existing version. The producer has to
change. Error::InvariantViolated means something already is true that
should not be — rows below the watermark still in PostgreSQL, two version
scopes for one reading. An operator has to look, and whoever is paged for the
second must not be woken by the first.
§Status
Pre-alpha, and unpublished on purpose: the API is still settling, and
integrating against a real workload is what settles it. What is implemented,
what is measured and what is not yet done are in the
repository README;
testkit is the reference the store is checked against, and is public so a
deployment can run it over its own configuration.
Full documentation: https://hupe1980.github.io/meterstore
Re-exports§
pub use cold::IcebergRestCatalog;pub use cold::S3TablesCatalog;pub use cold::ColdTier;pub use cold::IcebergCold;pub use cold::IcebergSqlCatalog;pub use cold::WarehouseAuth;pub use config::CHECK_VALUES_KEY;pub use config::TableConfig;pub use config::TimeModel;pub use config::VALUE_CHECK_KEY;pub use config::ValidatedTableConfig;pub use config::ValueCheck;pub use config::checked_column;pub use config::coded_column;pub use config::declared_value_check;pub use encode::StoredReadings;pub use encode::canonical_obis;pub use encode::parse_malo;pub use encode::parse_melo;pub use erasure::ErasureRecord;pub use erasure::MIN_ERASURE_SECRET_BYTES;pub use erasure::Retention;pub use erasure::SubjectRef;pub use erasure::SubjectRegistry;pub use error::Error;pub use error::Result;pub use evolution::Compatibility;pub use evolution::SchemaChange;pub use hot::PostgresHot;pub use planner::ReadMode;pub use planner::Resolution;pub use planner::SnapshotSelector;pub use planner::TierSplit;pub use planner::TieredTableProvider;pub use planner::TimeRange;pub use planner::balancing_day;pub use planner::balancing_day_bounds;pub use planner::balancing_day_length;pub use planner::balancing_month;pub use planner::balancing_month_bounds;pub use planner::bilanzierungsmonat;pub use planner::day_boundary;pub use planner::expected_intervals_in_balancing_day;pub use session::AUTHORITATIVE_ATTEMPTS;pub use session::Completeness;pub use session::CompletenessQuery;pub use session::HotWriter;pub use session::Maintenance;pub use session::MaintenanceOutcome;pub use session::MeterCatalog;pub use session::MeterCatalogBuilder;pub use session::MeterStore;pub use session::MeterStoreBuilder;pub use session::QueryDescription;pub use session::QueryResult;pub use session::RETENTION_LABEL;pub use session::ReadingsQuery;pub use session::ResolvedSeries;pub use session::SeriesQuery;pub use session::SqlSurface;pub use session::TableMaintenance;pub use settings::Deployment;pub use settings::PrivacySettings;pub use settings::Settings;pub use tiering::ArchivalOutcome;pub use tiering::Archiver;pub use tiering::ColdStore;pub use tiering::HotStore;pub use tiering::SnapshotInfo;pub use version::ScopedVersion;pub use version::Version;pub use version::VersionScope;pub use watermark::Tier;pub use watermark::TieringWatermark;pub use datafusion::arrow;
Modules§
- cli
- The
meterstorecommand-line tool. - cold
- The cold tier: settled history, in Apache Iceberg.
- config
- Table and archival configuration.
- encode
- Encoding
meteringtypes into Arrow, and back. - erasure
- Making Article 17 erasure possible over an append-only lake.
- error
- Error types for MeterStore.
- evolution
- Schema evolution, and the quarantine that stops an unsafe one.
- hot
- The hot tier: recent intervals, in PostgreSQL.
- observe
- Metrics.
- planner
- Query planning across the two tiers.
- prelude
- Common imports for working with MeterStore.
- serve
- Serving surfaces for engines that are not this process.
- session
- The session layer: a handle that wires both tiers into a query engine.
- settings
- The TOML front end over the same validated configuration.
- testkit
- A harness, a workload generator, and the assertions that matter.
- tiering
- Moving settled intervals from the hot store to the cold store.
- version
- MSCONS correction versions and the scopes they are comparable within.
- watermark
- The tiering watermark — the boundary between hot and cold storage.
Enums§
- EicType
- The EIC object type, re-exported from
metering.