Add generic ENLYZE power meter pipeline

This commit is contained in:
2026-10-01 11:59:34 +02:00
parent 4a3ed625ce
commit e9f9cf6f3c
15 changed files with 680 additions and 4 deletions
+13
View File
@@ -158,3 +158,16 @@ The exploration deliverable is verified API observations and sanitized
fixtures. The next product deliverable is live incremental material
integration, using a generic material-rate integrator rather than a separate
historical-only calculation path.
## Power-meter pipeline
The generic `power_meter` package collects configured ENLYZE meters into
TimescaleDB/PostgreSQL. The B2 example uses only verified UUIDs for machine,
energy total, active power total and grid frequency. Power-outage counting is
supported by schema/repository but remains unset until the real UUID is present
in local raw exploration data. Counter decreases start a new epoch and record
the new counter value as the increment; the first observation establishes a
baseline with zero increment. Optional phase power is stored as JSONB and does
not require L2. `report monthly` reads the database directly, writes
`<mount>/<meter>/<year>/<meter>_Powermeter_<YYYY-MM>.pdf`, and persists report
metadata independently of Grafana.
+39
View File
@@ -56,6 +56,10 @@ The second command saves a sanitized, reviewable fixture under
`fixtures/enlyze/` by default. Raw captures belong in the ignored
`data/raw/enlyze/` directory and must not be committed.
The `/v2/variables` exploration path automatically follows cursor pagination,
so the emitted response contains all variable pages. Non-2xx responses retain a
bounded, sanitized response body in the error for troubleshooting.
The OpenAPI server URL is `https://app.enlyze.com/api/`; its operation paths
begin with `/v2/`. Set `ENLYZE_BASE_URL` to that server URL, not to an
operation path. The documented read-only time-series operation is exposed for
@@ -217,6 +221,41 @@ start times, not source-sample timestamps. Schema application is manual.
Bento 1 bentonite source and gate selection remain intentionally undefined pending
process validation; the generic integration core is unchanged and reusable.
## Generic ENLYZE power meters
`config/power-meters.example.yaml` is the configuration-only B2 starting point.
It contains the verified machine, energy, active-power and frequency UUIDs; K1
and K6 are added as more `meters` entries without new code. Optional phase
channels are a mapping and may omit L2. `power_outages` is intentionally absent
until its real UUID is verified from the raw second `/v2/variables` page/register
1828; sanitized placeholders are never accepted.
Apply `db/schema.sql`, then run a bounded collection/backfill with ENLYZE and
PostgreSQL values from the environment or secret files:
```bash
production-analytics run power-meter --config config/power-meters.example.yaml \
--meter B2 --start 2026-09-01T00:00:00Z --end 2026-10-01T00:00:00Z
```
Generate the independent monthly PDF (under the configured mounted SMB path):
```bash
production-analytics report monthly --config config/power-meters.example.yaml \
--meter B2 --month 2026-09
```
The collector stores raw counters plus reset-safe increments in
`power_meter_readings`. Monthly reports preserve exact source timestamps for
the first/last reading, consumption, outage delta, frequency extrema and their
timestamps. PDF rendering needs WeasyPrint or `wkhtmltopdf`; report metadata is
stored in `power_meter_monthly_reports`. Grafana query templates are in
`db/power_meter_grafana.sql`.
For continuous collection, add `--follow`; without it the same command is a
bounded collector/backfill over `--start` and `--end` (or the latest polling
window when omitted).
## ERP current workplace status
The read-only ERP adapter reads production-order/article context from MSSQL
+15
View File
@@ -0,0 +1,15 @@
meters:
B2:
machine_uuid: 0d7955e9-5cde-4e14-9af7-5b041dc796f0
signals:
energy_total: 0f5c0853-1aa7-4fa3-b267-ea2669024e46
active_power_total: 41e5729e-b724-49d4-ad24-666d959e121d
grid_frequency: 127ddf56-9358-43ce-84fe-f7631486cfc9
# Add only after verifying the real UUID in the raw /v2/variables page.
# power_outages: <verified-uuid>
# Optional; L2 is not required when it is unavailable.
# phase_active_power:
# L1: <verified-uuid>
# L3: <verified-uuid>
poll_interval_seconds: 60
report_mount_path: /mnt/reports/energy
+14
View File
@@ -0,0 +1,14 @@
-- Set :meter, :from and :to from Grafana variables.
-- Current power:
SELECT timestamp AS "time", active_power_total_kw AS "value"
FROM power_meter_readings WHERE meter = '${meter}' AND $__timeFilter(timestamp);
-- Energy today / rolling 24h (counter increments are reset-safe).
SELECT COALESCE(sum(energy_delta_kwh), 0) AS kwh
FROM power_meter_readings
WHERE meter = '${meter}' AND timestamp >= now() - interval '24 hours';
-- Frequency and outages in selected period.
SELECT min(grid_frequency_hz) AS frequency_min_hz, max(grid_frequency_hz) AS frequency_max_hz,
COALESCE(sum(outage_delta), 0) AS outages
FROM power_meter_readings WHERE meter = '${meter}' AND $__timeFilter(timestamp);
+24
View File
@@ -46,3 +46,27 @@ CREATE TABLE IF NOT EXISTS material_application_snapshots (
CREATE INDEX IF NOT EXISTS material_application_snapshots_machine_time_idx
ON material_application_snapshots (calculation_id, machine_id, timestamp);
CREATE TABLE IF NOT EXISTS power_meter_readings (
meter text NOT NULL,
timestamp timestamptz NOT NULL,
energy_total_kwh double precision NOT NULL,
energy_delta_kwh double precision NOT NULL,
active_power_total_kw double precision NOT NULL,
grid_frequency_hz double precision NOT NULL,
power_outages double precision NULL,
outage_delta double precision NULL,
phase_active_power_kw jsonb NOT NULL DEFAULT '{}'::jsonb,
PRIMARY KEY (meter, timestamp)
);
CREATE INDEX IF NOT EXISTS power_meter_readings_time_idx
ON power_meter_readings (meter, timestamp DESC);
CREATE TABLE IF NOT EXISTS power_meter_monthly_reports (
meter text NOT NULL,
month text NOT NULL,
report_path text NOT NULL,
summary jsonb NOT NULL,
generated_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (meter, month)
);
+96 -1
View File
@@ -113,6 +113,21 @@ def _parser() -> argparse.ArgumentParser:
material.add_argument("--state-directory", type=Path, default=Path("data/state/material"))
material.add_argument("--secrets-file", type=Path, default=Path("secrets/enlyze.env"))
material.add_argument("--erp-secrets-file", type=Path, default=Path("secrets/erp.env"))
power = runners.add_parser("power-meter", help="Collect configured ENLYZE power-meter telemetry")
power.add_argument("--config", type=Path, required=True)
power.add_argument("--meter")
power.add_argument("--start", help="ISO timestamp; defaults to one polling window ago")
power.add_argument("--end", help="ISO timestamp; defaults to now")
power.add_argument("--follow", action="store_true", help="Keep polling until interrupted")
power.add_argument("--secrets-file", type=Path, default=Path("secrets/enlyze.env"))
power.add_argument("--db-secrets-file", type=Path, default=Path("secrets/postgres.env"))
report = namespaces.add_parser("report", help="Generate independent reports")
report_commands = report.add_subparsers(dest="command", required=True)
monthly = report_commands.add_parser("monthly", help="Render one configured monthly PDF")
monthly.add_argument("--config", type=Path, required=True)
monthly.add_argument("--meter")
monthly.add_argument("--month", required=True, help="YYYY-MM")
monthly.add_argument("--db-secrets-file", type=Path, default=Path("secrets/postgres.env"))
efficiency = runners.add_parser("material-efficiency", help="Persist ERP feedback KPI points")
efficiency.add_argument("--config", type=Path, required=True)
efficiency.add_argument("--secrets-file", type=Path, default=Path("secrets/enlyze.env"))
@@ -183,12 +198,81 @@ def _run_material_efficiency(args: argparse.Namespace) -> int:
return 0
def _runtime_environment(*paths: Path) -> dict[str, str]:
environment = dict(os.environ)
for path in paths:
environment.update(load_secret_file(path))
return environment
def _parse_cli_timestamp(value: str) -> object:
from datetime import datetime
parsed = datetime.fromisoformat(value.replace("Z", "+00:00"))
if parsed.tzinfo is None:
raise ExplorationError("timestamps must include a timezone")
return parsed
def _run_power_meter(args: argparse.Namespace) -> int:
from datetime import UTC, datetime, timedelta
import time
from production_analytics.enlyze.exploration import ConfigurationError
from production_analytics.power_meter.config import load_power_meter_config
from production_analytics.power_meter.gateway import PowerMeterGateway
from production_analytics.power_meter.repository import PostgresPowerMeterRepository
from production_analytics.power_meter.runtime import PowerMeterCollector
try:
config = load_power_meter_config(args.config, args.meter)
environment = _runtime_environment(args.secrets_file, args.db_secrets_file)
settings = ExplorationSettings.from_environment(environment)
from production_analytics.service.postgres_material import PostgresSettings
db = PostgresSettings.from_environment(environment)
collector = PowerMeterCollector(config, PowerMeterGateway(ExplorationClient(settings)),
PostgresPowerMeterRepository(db))
if args.follow:
while True:
count = collector.poll_once()
print(f"Collected {count} {config.name} power-meter readings.", flush=True)
time.sleep(config.poll_interval_seconds)
now = datetime.now(UTC)
end = _parse_cli_timestamp(args.end) if args.end else now
start = (_parse_cli_timestamp(args.start) if args.start else
end - timedelta(seconds=config.poll_interval_seconds * 2))
count = collector.backfill(start, end)
print(f"Collected {count} {config.name} power-meter readings.")
return 0
except (CalculationConfigError, ConfigurationError, ExplorationError, ValueError) as exc:
print(f"Config/startup error: {exc}", file=sys.stderr)
return 2
def _run_monthly_report(args: argparse.Namespace) -> int:
from production_analytics.enlyze.exploration import ConfigurationError
from production_analytics.power_meter.config import load_power_meter_config
from production_analytics.power_meter.repository import PostgresPowerMeterRepository
from production_analytics.power_meter.runtime import monthly_report
try:
config = load_power_meter_config(args.config, args.meter)
environment = _runtime_environment(args.db_secrets_file)
from production_analytics.service.postgres_material import PostgresSettings
path = monthly_report(config, PostgresPowerMeterRepository(PostgresSettings.from_environment(environment)), args.month)
print(path)
return 0
except (CalculationConfigError, ConfigurationError, ValueError, RuntimeError) as exc:
print(f"Report error: {exc}", file=sys.stderr)
return 2
def main(argv: Sequence[str] | None = None) -> int:
args = _parser().parse_args(argv)
if args.namespace == "run":
if args.command == "material-efficiency":
return _run_material_efficiency(args)
if args.command == "power-meter":
return _run_power_meter(args)
return _run_material(args)
if args.namespace == "report":
return _run_monthly_report(args)
if args.namespace != "enlyze" or args.command not in {"raw", "timeseries"}:
return 2
@@ -207,7 +291,18 @@ def main(argv: Sequence[str] | None = None) -> int:
)
client = ExplorationClient(settings)
if args.command == "raw":
response = client.get(args.path, dict(args.query))
query = dict(args.query)
if args.path == "/v2/variables":
pages = client.get_all_pages(args.path, query)
response = pages[0]
if len(pages) > 1:
merged = dict(response.body)
merged["data"] = [item for page in pages for item in page.body.get("data", [])]
merged["metadata"] = {"next_cursor": None, "pages": len(pages)}
response = type(response)(response.status_code, response.path,
response.headers, merged)
else:
response = client.get(args.path, query)
request = {"method": "GET", "path": response.path}
else:
if (
+53 -3
View File
@@ -21,12 +21,20 @@ class ConfigurationError(ExplorationError):
class AuthenticationError(ExplorationError):
"""Raised for HTTP 401/403 without exposing credentials or response bodies."""
"""Raised for HTTP 401/403 with a bounded, sanitized response body."""
def __init__(self, message: str, *, response_body: object | None = None) -> None:
super().__init__(message)
self.response_body = response_body
class HttpResponseError(ExplorationError):
"""Raised for a non-successful HTTP response."""
def __init__(self, message: str, *, response_body: object | None = None) -> None:
super().__init__(message)
self.response_body = response_body
class NonJsonResponseError(ExplorationError):
"""Raised when a successful response cannot be decoded as JSON."""
@@ -82,6 +90,30 @@ class ExplorationClient:
"""Issue a documented read-only POST operation with a JSON body."""
return self._request_json("POST", path, payload=payload)
def get_all_pages(
self, path: str, query: Mapping[str, str] | None = None,
) -> list[ExplorationResponse]:
"""Follow the ENLYZE cursor contract and return every response page."""
request_query = dict(query or {})
pages: list[ExplorationResponse] = []
followed: set[str] = set()
while True:
response = self.get(path, request_query)
pages.append(response)
body = response.body
if not isinstance(body, dict) or "metadata" not in body:
return pages
metadata = body["metadata"]
if not isinstance(metadata, dict) or "next_cursor" not in metadata:
raise ExplorationError("Invalid paginated response: metadata.next_cursor is missing.")
cursor = metadata["next_cursor"]
if cursor is None:
return pages
if not isinstance(cursor, str) or not cursor or cursor in followed:
raise ExplorationError("Invalid paginated response: repeated or invalid cursor.")
followed.add(cursor)
request_query["cursor"] = cursor
def _request_json(
self,
method: str,
@@ -118,9 +150,27 @@ class ExplorationClient:
status_code = response.status
response_headers = dict(response.headers.items())
except HTTPError as error:
response_body: object | None = None
try:
raw_error_body = error.read(8192)
try:
from production_analytics.enlyze.sanitize import sanitize
response_body = sanitize(json.loads(raw_error_body))
except (UnicodeDecodeError, json.JSONDecodeError):
response_body = raw_error_body.decode("utf-8", errors="replace")[:2000]
except OSError:
response_body = None
if error.code in {401, 403}:
raise AuthenticationError(f"Authentication or authorization failed (HTTP {error.code}).") from error
raise HttpResponseError(f"HTTP request failed with status {error.code}.") from error
suffix = f" Response: {response_body!r}" if response_body is not None else ""
raise AuthenticationError(
f"Authentication or authorization failed (HTTP {error.code}).{suffix}",
response_body=response_body,
) from error
suffix = f" Response: {response_body!r}" if response_body is not None else ""
raise HttpResponseError(
f"HTTP request failed with status {error.code}.{suffix}",
response_body=response_body,
) from error
except URLError as error:
raise HttpResponseError("HTTP request could not be completed.") from error
@@ -0,0 +1,85 @@
"""Strict, meter-generic configuration for ENLYZE power meters."""
from dataclasses import dataclass
from pathlib import Path
from typing import Any
import yaml
from production_analytics.calculations.config import CalculationConfigError, _UniqueLoader
@dataclass(frozen=True, slots=True)
class PowerMeterSignals:
energy_total: str
active_power_total: str
grid_frequency: str
power_outages: str | None = None
phase_active_power: tuple[tuple[str, str], ...] = ()
@dataclass(frozen=True, slots=True)
class PowerMeterConfig:
name: str
machine_uuid: str
signals: PowerMeterSignals
poll_interval_seconds: float = 60.0
report_mount_path: str = "/mnt/reports/energy"
def _text(value: Any, label: str) -> str:
if not isinstance(value, str) or not value.strip() or "\x00" in value:
raise CalculationConfigError(f"{label} must be a non-empty string")
if value.startswith("redacted-") or value.startswith("<"):
raise CalculationConfigError(f"{label} must contain a real local identifier")
return value.strip()
def load_power_meter_config(path: str | Path, meter_name: str | None = None) -> PowerMeterConfig:
try:
with Path(path).open(encoding="utf-8") as stream:
document = yaml.load(stream, Loader=_UniqueLoader)
except (OSError, UnicodeError, yaml.YAMLError) as exc:
raise CalculationConfigError("Cannot read power-meter configuration") from exc
if not isinstance(document, dict) or set(document) != {"meters"}:
raise CalculationConfigError("Configuration must contain only 'meters'")
meters = document["meters"]
if not isinstance(meters, dict) or not meters:
raise CalculationConfigError("meters must be a non-empty mapping")
if meter_name is None:
if len(meters) != 1:
raise CalculationConfigError("Multiple meters: specify --meter")
meter_name = next(iter(meters))
if meter_name not in meters or not isinstance(meters[meter_name], dict):
raise CalculationConfigError("Requested meter was not found")
entry = meters[meter_name]
if set(entry) - {"machine_uuid", "signals", "poll_interval_seconds", "report_mount_path"}:
raise CalculationConfigError("Unknown power-meter configuration field")
machine = _text(entry.get("machine_uuid"), f"meters.{meter_name}.machine_uuid")
raw = entry.get("signals")
if not isinstance(raw, dict):
raise CalculationConfigError("signals must be a mapping")
required = {"energy_total", "active_power_total", "grid_frequency"}
if not required <= raw.keys():
raise CalculationConfigError("signals requires energy_total, active_power_total, grid_frequency")
phase = raw.get("phase_active_power", {})
if not isinstance(phase, dict) or any(not isinstance(k, str) for k in phase):
raise CalculationConfigError("phase_active_power must be a mapping")
phase_items = tuple((_text(k, "phase name"), _text(v, f"phase_active_power.{k}"))
for k, v in phase.items())
if len({k for k, _ in phase_items}) != len(phase_items):
raise CalculationConfigError("phase names must be unique")
interval = entry.get("poll_interval_seconds", 60.0)
if isinstance(interval, bool) or not isinstance(interval, (int, float)) or interval <= 0:
raise CalculationConfigError("poll_interval_seconds must be greater than zero")
return PowerMeterConfig(
name=_text(meter_name, "meter name"), machine_uuid=machine,
signals=PowerMeterSignals(
energy_total=_text(raw["energy_total"], "energy_total"),
active_power_total=_text(raw["active_power_total"], "active_power_total"),
grid_frequency=_text(raw["grid_frequency"], "grid_frequency"),
power_outages=_text(raw["power_outages"], "power_outages") if raw.get("power_outages") else None,
phase_active_power=phase_items,
),
poll_interval_seconds=float(interval),
report_mount_path=_text(entry.get("report_mount_path", "/mnt/reports/energy"), "report_mount_path"),
)
@@ -0,0 +1,85 @@
"""ENLYZE adapter for generic power-meter telemetry."""
from dataclasses import dataclass
from datetime import UTC, datetime
from math import isfinite
from production_analytics.enlyze.exploration import ExplorationClient
from .config import PowerMeterConfig
@dataclass(frozen=True, slots=True)
class PowerMeterReading:
timestamp: datetime
energy_total_kwh: float
active_power_total_kw: float
grid_frequency_hz: float
power_outages: float | None
phase_active_power_kw: tuple[tuple[str, float], ...]
def _number(value: object, label: str) -> float:
if isinstance(value, bool) or not isinstance(value, (int, float)) or not isfinite(float(value)):
raise ValueError(f"{label} must be a finite number")
return float(value)
class PowerMeterGateway:
def __init__(self, client: ExplorationClient) -> None:
self.client = client
def read(self, config: PowerMeterConfig, start: datetime, end: datetime) -> list[PowerMeterReading]:
if start.tzinfo is None or end.tzinfo is None:
raise ValueError("power-meter window must be timezone-aware")
if end < start:
raise ValueError("power-meter end must not precede start")
ids = [config.signals.energy_total, config.signals.active_power_total,
config.signals.grid_frequency]
ids += [value for _, value in config.signals.phase_active_power]
if config.signals.power_outages:
ids.append(config.signals.power_outages)
body = {"machine": config.machine_uuid, "start": start.astimezone(UTC).isoformat(),
"end": end.astimezone(UTC).isoformat(),
"variables": [{"uuid": value} for value in dict.fromkeys(ids)]}
response = self.client.post_json("/v2/timeseries", body)
result: list[PowerMeterReading] = []
page_body = response.body
while True:
if not isinstance(page_body, dict) or not isinstance(page_body.get("data"), dict):
raise ValueError("Invalid power-meter timeseries response")
data = page_body["data"]
columns, records = data.get("columns"), data.get("records")
if not isinstance(columns, list) or not isinstance(records, list):
raise ValueError("Invalid power-meter columns or records")
index = {column: i for i, column in enumerate(columns)}
required = ["time", config.signals.energy_total, config.signals.active_power_total,
config.signals.grid_frequency]
if any(name not in index for name in required):
raise ValueError("Power-meter response is missing a required column")
for row_number, row in enumerate(records):
if not isinstance(row, list) or len(row) < len(columns):
raise ValueError(f"Malformed power-meter record {row_number}")
try:
timestamp = datetime.fromisoformat(str(row[index["time"]]).replace("Z", "+00:00"))
if timestamp.tzinfo is None:
raise ValueError("timestamp is not timezone-aware")
phase = tuple((name, _number(row[index[signal]], name))
for name, signal in config.signals.phase_active_power
if signal in index and row[index[signal]] is not None)
result.append(PowerMeterReading(
timestamp=timestamp.astimezone(UTC),
energy_total_kwh=_number(row[index[config.signals.energy_total]], "energy_total"),
active_power_total_kw=_number(row[index[config.signals.active_power_total]], "active_power_total"),
grid_frequency_hz=_number(row[index[config.signals.grid_frequency]], "grid_frequency"),
power_outages=(None if not config.signals.power_outages or
row[index[config.signals.power_outages]] is None else
_number(row[index[config.signals.power_outages]], "power_outages")),
phase_active_power_kw=phase,
))
except (IndexError, KeyError, TypeError, ValueError) as exc:
raise ValueError(f"Malformed power-meter record {row_number}") from exc
metadata = page_body.get("metadata", {})
cursor = metadata.get("next_cursor") if isinstance(metadata, dict) else None
if cursor is None:
return sorted(result, key=lambda item: item.timestamp)
body["cursor"] = cursor
page_body = self.client.post_json("/v2/timeseries", body).body
@@ -0,0 +1,47 @@
"""Independent HTML-to-PDF monthly power-meter reporting."""
from dataclasses import asdict
from pathlib import Path
import html
import subprocess
from .repository import MonthlySummary
def render_html(summary: MonthlySummary) -> str:
def value(item, suffix=""):
return "—" if item is None else f"{item}{suffix}"
rows = [
("Source start", value(summary.first_timestamp)),
("Source end", value(summary.last_timestamp)),
("Start meter reading", value(summary.first_energy_kwh, " kWh")),
("End meter reading", value(summary.last_energy_kwh, " kWh")),
("Monthly consumption", value(round(summary.consumption_kwh, 3), " kWh")),
("Outage counter delta", value(summary.outages)),
("Frequency minimum", value(summary.frequency_min_hz, " Hz") + (f" ({summary.frequency_min_timestamp})" if summary.frequency_min_timestamp else "")),
("Frequency maximum", value(summary.frequency_max_hz, " Hz") + (f" ({summary.frequency_max_timestamp})" if summary.frequency_max_timestamp else "")),
("Deviation from 50 Hz", value(None if summary.frequency_min_hz is None else round(summary.frequency_min_hz - 50, 3), " Hz") + " min; " + value(None if summary.frequency_max_hz is None else round(summary.frequency_max_hz - 50, 3), " Hz") + " max"),
]
body = "".join(f"<tr><th>{html.escape(str(k))}</th><td>{html.escape(str(v))}</td></tr>" for k, v in rows)
return f"""<!doctype html><html><head><meta charset='utf-8'><style>
body {{ font-family: sans-serif; margin: 28mm 20mm; color: #17202a; }} h1 {{ font-size: 20pt; }}
table {{ border-collapse: collapse; width: 100%; font-size: 10pt; }} th,td {{ border-bottom: 1px solid #d5d8dc; padding: 7px; text-align:left; }} th {{ width: 42%; background:#f3f5f7; }}
footer {{ margin-top:25px; color:#657; font-size:8pt; }}
</style></head><body><h1>Powermeter report — {html.escape(summary.meter)} — {html.escape(summary.month)}</h1>
<table>{body}</table><footer>Source: TimescaleDB power_meter_readings. Source timestamps are preserved exactly.</footer></body></html>"""
def write_pdf(summary: MonthlySummary, output: Path) -> None:
output.parent.mkdir(parents=True, exist_ok=True)
html_path = output.with_suffix(".html")
html_path.write_text(render_html(summary), encoding="utf-8")
try:
from weasyprint import HTML
HTML(filename=str(html_path)).write_pdf(str(output))
except ImportError:
try:
subprocess.run(["wkhtmltopdf", str(html_path), str(output)], check=True,
capture_output=True, text=True)
except (FileNotFoundError, subprocess.CalledProcessError) as exc:
raise RuntimeError("PDF requires WeasyPrint or wkhtmltopdf") from exc
finally:
html_path.unlink(missing_ok=True)
@@ -0,0 +1,102 @@
"""TimescaleDB persistence and monthly aggregation for power meters."""
from dataclasses import asdict, dataclass
from datetime import datetime
import json
from .gateway import PowerMeterReading
@dataclass(frozen=True, slots=True)
class MonthlySummary:
meter: str
month: str
first_timestamp: datetime | None
last_timestamp: datetime | None
first_energy_kwh: float | None
last_energy_kwh: float | None
consumption_kwh: float
outages: float
frequency_min_hz: float | None
frequency_min_timestamp: datetime | None
frequency_max_hz: float | None
frequency_max_timestamp: datetime | None
def counter_delta(previous: float | None, current: float) -> float:
"""Return a non-negative increment; a lower value starts a new counter epoch."""
if previous is None:
return 0.0
return current - previous if current >= previous else current
class PostgresPowerMeterRepository:
def __init__(self, settings) -> None:
self.settings = settings
def _connect(self):
import psycopg
return psycopg.connect(host=self.settings.host, port=self.settings.port,
dbname=self.settings.dbname, user=self.settings.user,
password=self.settings.password, connect_timeout=10,
options="-c statement_timeout=10000")
def latest(self, meter: str, before: datetime):
with self._connect() as connection:
return connection.execute(
"SELECT energy_total_kwh, power_outages FROM power_meter_readings "
"WHERE meter=%s AND timestamp < %s ORDER BY timestamp DESC LIMIT 1",
(meter, before)).fetchone()
def write(self, meter: str, reading: PowerMeterReading) -> None:
previous = self.latest(meter, reading.timestamp)
energy_delta = counter_delta(None if previous is None else previous[0], reading.energy_total_kwh)
outage_delta = None if reading.power_outages is None else counter_delta(
None if previous is None or previous[1] is None else previous[1], reading.power_outages)
phases = dict(reading.phase_active_power_kw)
with self._connect() as connection:
connection.execute(
"""INSERT INTO power_meter_readings
(meter,timestamp,energy_total_kwh,energy_delta_kwh,active_power_total_kw,
grid_frequency_hz,power_outages,outage_delta,phase_active_power_kw)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s)
ON CONFLICT (meter,timestamp) DO UPDATE SET
energy_total_kwh=EXCLUDED.energy_total_kwh,
energy_delta_kwh=EXCLUDED.energy_delta_kwh,
active_power_total_kw=EXCLUDED.active_power_total_kw,
grid_frequency_hz=EXCLUDED.grid_frequency_hz,
power_outages=EXCLUDED.power_outages,
outage_delta=EXCLUDED.outage_delta,
phase_active_power_kw=EXCLUDED.phase_active_power_kw""",
(meter, reading.timestamp, reading.energy_total_kwh, energy_delta,
reading.active_power_total_kw, reading.grid_frequency_hz,
reading.power_outages, outage_delta, json.dumps(phases)),
)
def monthly(self, meter: str, start: datetime, end: datetime, month: str) -> MonthlySummary:
query = """WITH points AS (
SELECT * FROM power_meter_readings WHERE meter=%s AND timestamp >= %s AND timestamp < %s
), extrema AS (
SELECT min(grid_frequency_hz) AS min_hz, max(grid_frequency_hz) AS max_hz,
(array_agg(timestamp ORDER BY grid_frequency_hz ASC, timestamp ASC))[1] AS min_at,
(array_agg(timestamp ORDER BY grid_frequency_hz DESC, timestamp ASC))[1] AS max_at
FROM points
) SELECT (SELECT timestamp FROM points ORDER BY timestamp LIMIT 1),
(SELECT timestamp FROM points ORDER BY timestamp DESC LIMIT 1),
(SELECT energy_total_kwh FROM points ORDER BY timestamp LIMIT 1),
(SELECT energy_total_kwh FROM points ORDER BY timestamp DESC LIMIT 1),
COALESCE((SELECT sum(energy_delta_kwh) FROM points),0),
COALESCE((SELECT sum(outage_delta) FROM points),0),
min_hz, min_at, max_hz, max_at FROM extrema"""
with self._connect() as connection:
row = connection.execute(query, (meter, start, end)).fetchone()
return MonthlySummary(meter, month, *row)
def save_report_metadata(self, summary: MonthlySummary, path: str) -> None:
with self._connect() as connection:
connection.execute(
"""INSERT INTO power_meter_monthly_reports
(meter,month,report_path,summary) VALUES (%s,%s,%s,%s::jsonb)
ON CONFLICT (meter,month) DO UPDATE SET report_path=EXCLUDED.report_path,
summary=EXCLUDED.summary, generated_at=now()""",
(summary.meter, summary.month, path, json.dumps(asdict(summary), default=str)),
)
@@ -0,0 +1,40 @@
"""One-shot collection/backfill and monthly report orchestration."""
from datetime import UTC, datetime, timedelta
from pathlib import Path
from .config import PowerMeterConfig
from .gateway import PowerMeterGateway
from .repository import PostgresPowerMeterRepository
from .report import write_pdf
class PowerMeterCollector:
def __init__(self, config: PowerMeterConfig, gateway: PowerMeterGateway,
repository: PostgresPowerMeterRepository) -> None:
self.config, self.gateway, self.repository = config, gateway, repository
def backfill(self, start: datetime, end: datetime) -> int:
readings = self.gateway.read(self.config, start, end)
for reading in readings:
self.repository.write(self.config.name, reading)
return len(readings)
def poll_once(self, now: datetime | None = None) -> int:
now = now or datetime.now(UTC)
start = now - timedelta(seconds=self.config.poll_interval_seconds * 2)
return self.backfill(start, now)
def monthly_report(config: PowerMeterConfig, repository: PostgresPowerMeterRepository,
month: str) -> Path:
start = datetime.fromisoformat(month + "-01").replace(tzinfo=UTC)
end = datetime.fromisoformat(month + "-01").replace(tzinfo=UTC)
if start.month == 12:
end = end.replace(year=start.year + 1, month=1)
else:
end = end.replace(month=start.month + 1)
summary = repository.monthly(config.name, start, end, month)
output = Path(config.report_mount_path) / config.name / month[:4] / f"{config.name}_Powermeter_{month}.pdf"
write_pdf(summary, output)
repository.save_report_metadata(summary, str(output))
return output
+23
View File
@@ -9,6 +9,7 @@ from production_analytics.enlyze.exploration import (
AuthenticationError,
ExplorationClient,
ExplorationSettings,
HttpResponseError,
NonJsonResponseError,
load_secret_file,
)
@@ -78,6 +79,28 @@ class ExplorationClientTests(unittest.TestCase):
with self.assertRaisesRegex(AuthenticationError, "HTTP 401"):
ExplorationClient(self.settings).get("/observed/path")
@patch("production_analytics.enlyze.exploration.urlopen")
def test_http_error_keeps_sanitized_bounded_body(self, urlopen_mock: object) -> None:
urlopen_mock.side_effect = HTTPError(
"https://example.invalid", 422, "Unprocessable", {},
io.BytesIO(b'{"message":"bad variable","token":"must not echo"}'),
) # type: ignore[attr-defined]
with self.assertRaises(HttpResponseError) as caught:
ExplorationClient(self.settings).get("/observed/path")
self.assertEqual(caught.exception.response_body["message"], "bad variable")
self.assertEqual(caught.exception.response_body["token"], "<redacted-secret>")
@patch("production_analytics.enlyze.exploration.urlopen")
def test_get_all_pages_follows_cursor(self, urlopen_mock: object) -> None:
urlopen_mock.side_effect = [
_Response(b'{"data":[1],"metadata":{"next_cursor":"next"}}'),
_Response(b'{"data":[2],"metadata":{"next_cursor":null}}'),
]
pages = ExplorationClient(self.settings).get_all_pages("/v2/variables", {"machine": "m"})
self.assertEqual([page.body["data"] for page in pages], [[1], [2]])
class SanitizationTests(unittest.TestCase):
def test_sanitization_preserves_shape_and_redacts_sensitive_values(self) -> None:
+44
View File
@@ -0,0 +1,44 @@
from datetime import UTC, datetime
from unittest.mock import Mock
import pytest
from production_analytics.power_meter.config import load_power_meter_config
from production_analytics.power_meter.gateway import PowerMeterGateway
from production_analytics.power_meter.repository import counter_delta
def test_b2_config_has_only_verified_signals() -> None:
config = load_power_meter_config("config/power-meters.example.yaml", "B2")
assert config.machine_uuid == "0d7955e9-5cde-4e14-9af7-5b041dc796f0"
assert config.signals.energy_total == "0f5c0853-1aa7-4fa3-b267-ea2669024e46"
assert config.signals.power_outages is None
def test_counter_reset_starts_new_epoch() -> None:
assert counter_delta(None, 7) == 0
assert counter_delta(100, 105) == 5
assert counter_delta(100, 4) == 4
def test_gateway_keeps_optional_phases_and_handles_pagination() -> None:
config = load_meter()
client = Mock()
client.post_json.side_effect = [
Mock(body={"data": {"columns": ["time", "e", "p", "f", "l1"],
"records": [["2026-09-01T00:00:00Z", 10, 4, 50, 2]]},
"metadata": {"next_cursor": "next"}}),
Mock(body={"data": {"columns": ["time", "e", "p", "f", "l1"],
"records": [["2026-09-01T00:01:00Z", 11, 5, 49.9, 3]]},
"metadata": {"next_cursor": None}}),
]
readings = PowerMeterGateway(client).read(config, datetime(2026, 9, 1, tzinfo=UTC),
datetime(2026, 9, 2, tzinfo=UTC))
assert [item.energy_total_kwh for item in readings] == [10, 11]
assert readings[0].phase_active_power_kw == (("L1", 2.0),)
assert client.post_json.call_args_list[1].args[1]["cursor"] == "next"
def load_meter():
from production_analytics.power_meter.config import PowerMeterConfig, PowerMeterSignals
return PowerMeterConfig("B2", "machine", PowerMeterSignals("e", "p", "f", phase_active_power=(("L1", "l1"),)))