From e9f9cf6f3ca80eb3fa8b17226771ab948e56a8b3 Mon Sep 17 00:00:00 2001 From: Martin Tazl Date: Thu, 1 Oct 2026 11:59:34 +0200 Subject: [PATCH] Add generic ENLYZE power meter pipeline --- PROJECT_KNOWLEDGE.md | 13 +++ README.md | 39 +++++++ config/power-meters.example.yaml | 15 +++ db/power_meter_grafana.sql | 14 +++ db/schema.sql | 24 +++++ src/production_analytics/cli/__main__.py | 97 ++++++++++++++++- .../enlyze/exploration.py | 56 +++++++++- .../power_meter/__init__.py | 0 .../power_meter/config.py | 85 +++++++++++++++ .../power_meter/gateway.py | 85 +++++++++++++++ .../power_meter/report.py | 47 ++++++++ .../power_meter/repository.py | 102 ++++++++++++++++++ .../power_meter/runtime.py | 40 +++++++ tests/test_enlyze_exploration.py | 23 ++++ tests/test_power_meter.py | 44 ++++++++ 15 files changed, 680 insertions(+), 4 deletions(-) create mode 100644 config/power-meters.example.yaml create mode 100644 db/power_meter_grafana.sql create mode 100644 src/production_analytics/power_meter/__init__.py create mode 100644 src/production_analytics/power_meter/config.py create mode 100644 src/production_analytics/power_meter/gateway.py create mode 100644 src/production_analytics/power_meter/report.py create mode 100644 src/production_analytics/power_meter/repository.py create mode 100644 src/production_analytics/power_meter/runtime.py create mode 100644 tests/test_power_meter.py diff --git a/PROJECT_KNOWLEDGE.md b/PROJECT_KNOWLEDGE.md index 9671182..be1e278 100644 --- a/PROJECT_KNOWLEDGE.md +++ b/PROJECT_KNOWLEDGE.md @@ -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 +`///_Powermeter_.pdf`, and persists report +metadata independently of Grafana. diff --git a/README.md b/README.md index 17f6eb2..4d3ef34 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/config/power-meters.example.yaml b/config/power-meters.example.yaml new file mode 100644 index 0000000..af243c9 --- /dev/null +++ b/config/power-meters.example.yaml @@ -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: + # Optional; L2 is not required when it is unavailable. + # phase_active_power: + # L1: + # L3: + poll_interval_seconds: 60 + report_mount_path: /mnt/reports/energy diff --git a/db/power_meter_grafana.sql b/db/power_meter_grafana.sql new file mode 100644 index 0000000..07d4599 --- /dev/null +++ b/db/power_meter_grafana.sql @@ -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); diff --git a/db/schema.sql b/db/schema.sql index 1a69194..3c21596 100644 --- a/db/schema.sql +++ b/db/schema.sql @@ -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) +); diff --git a/src/production_analytics/cli/__main__.py b/src/production_analytics/cli/__main__.py index 8dd5bd5..367ff1a 100644 --- a/src/production_analytics/cli/__main__.py +++ b/src/production_analytics/cli/__main__.py @@ -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 ( diff --git a/src/production_analytics/enlyze/exploration.py b/src/production_analytics/enlyze/exploration.py index a7808cc..a9b0972 100644 --- a/src/production_analytics/enlyze/exploration.py +++ b/src/production_analytics/enlyze/exploration.py @@ -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 diff --git a/src/production_analytics/power_meter/__init__.py b/src/production_analytics/power_meter/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/production_analytics/power_meter/config.py b/src/production_analytics/power_meter/config.py new file mode 100644 index 0000000..07f026c --- /dev/null +++ b/src/production_analytics/power_meter/config.py @@ -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"), + ) diff --git a/src/production_analytics/power_meter/gateway.py b/src/production_analytics/power_meter/gateway.py new file mode 100644 index 0000000..7580dfe --- /dev/null +++ b/src/production_analytics/power_meter/gateway.py @@ -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 diff --git a/src/production_analytics/power_meter/report.py b/src/production_analytics/power_meter/report.py new file mode 100644 index 0000000..48747c0 --- /dev/null +++ b/src/production_analytics/power_meter/report.py @@ -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"{html.escape(str(k))}{html.escape(str(v))}" for k, v in rows) + return f"""

Powermeter report — {html.escape(summary.meter)} — {html.escape(summary.month)}

+ {body}
Source: TimescaleDB power_meter_readings. Source timestamps are preserved exactly.
""" + + +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) diff --git a/src/production_analytics/power_meter/repository.py b/src/production_analytics/power_meter/repository.py new file mode 100644 index 0000000..693198f --- /dev/null +++ b/src/production_analytics/power_meter/repository.py @@ -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)), + ) diff --git a/src/production_analytics/power_meter/runtime.py b/src/production_analytics/power_meter/runtime.py new file mode 100644 index 0000000..6b724ac --- /dev/null +++ b/src/production_analytics/power_meter/runtime.py @@ -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 diff --git a/tests/test_enlyze_exploration.py b/tests/test_enlyze_exploration.py index 3fdd7b1..2fa4fe5 100644 --- a/tests/test_enlyze_exploration.py +++ b/tests/test_enlyze_exploration.py @@ -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"], "") + + @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: diff --git a/tests/test_power_meter.py b/tests/test_power_meter.py new file mode 100644 index 0000000..509f3ab --- /dev/null +++ b/tests/test_power_meter.py @@ -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"),)))