From f47f1f01297f02d1a221304eba4722b6e67d2542 Mon Sep 17 00:00:00 2001 From: glenn schrooyen Date: Mon, 24 Aug 2026 21:52:51 +0200 Subject: [PATCH] TEL-01: P1 ingestion, with the derivation and the age the EMS owns A Belgian P1 meter publishes two UNSIGNED registers, not one signed figure. Until now the add-on asked the installer to bridge that gap with a template sensor, which put the sign convention of the whole control loop in a text box. This moves it into the EMS: net = import - export, derived once, in one place, with a test that fails if anyone inverts it. Two transports behind one contract, chosen by `meter_source`: the HA WebSocket subscribing to the DSMR integration's entities, and MQTT on a configurable topic. Everything downstream reads P1Ingest, so switching is a config edit. `meter_source: off` is the default and keeps the existing meter_entity path, so no installed system changes until it opts in. The other half is the timestamp. Every accepted sample is stamped at ingest with a monotonic clock, `meter_max_age_s` is applied to it, and the age is published as sensor.p1_sample_age_s for the ESP32's stale-input watchdog. That entity is recomputed against the clock every second rather than only when a telegram lands, because HA pushes state only on change: a meter frozen at a constant reading emits nothing and looks, to anything watching the value, exactly like a meter that has died. The age tells them apart. Deliberately absent: any fallback to an inverter-side power figure. The inverter's own AC power correlates 0.998 with battery power and 0.09 with the real meter, so failing over to it means regulating against your own output. A gap stays a gap - a reconnect emits no synthetic sample, and a rejected telegram never resolves to 0 W or refreshes the timestamp. Quarter-hour averages are time-weighted over clock-aligned blocks rather than a mean of samples, so a cadence change cannot bias the capacity-tariff figure, and only offtake is accumulated so a quarter of pure export averages to 0 kW. Per-phase import is kept separately: on an unbalanced three-phase load the phase sum and the connection net are different numbers, and only one of them is billed. test_p1.py: 99 checks, runnable with a bare interpreter and no meter. Includes an end-to-end run of the HA transport against a fake Home Assistant websocket. Stacked on SAFETY-04; nothing here touches control.py. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01Du77usMj8XNKNFZGmUiWDa --- goodwe_controller/DOCS.md | 54 +++ goodwe_controller/app/main.py | 70 +++- goodwe_controller/app/mqtt.py | 8 + goodwe_controller/app/p1.py | 674 ++++++++++++++++++++++++++++++++++ goodwe_controller/config.yaml | 32 ++ goodwe_controller/test_p1.py | 504 +++++++++++++++++++++++++ 6 files changed, 1331 insertions(+), 11 deletions(-) create mode 100644 goodwe_controller/app/p1.py create mode 100644 goodwe_controller/test_p1.py diff --git a/goodwe_controller/DOCS.md b/goodwe_controller/DOCS.md index 9e0de76..e7e2b7c 100644 --- a/goodwe_controller/DOCS.md +++ b/goodwe_controller/DOCS.md @@ -48,6 +48,60 @@ Use the ESP32's readings rather than the inverter's cloud or dongle sensors: those serve cached values, and a stale reading here ends the maintenance charge phase having charged nothing. +### P1 meter ingestion + +`meter_entity` above expects one signed sensor, which usually means a template +someone wrote by hand. A Belgian P1 meter does not publish one: it publishes two +**unsigned** registers, consumption and injection. Setting `meter_source` moves +that subtraction into the add-on, where it is done once and tested, and replaces +`meter_entity` entirely. + +| option | default | meaning | +|---|---|---| +| `meter_source` | `off` | `off` keeps `meter_entity`. `ha_dsmr` subscribes to the DSMR integration over the HA WebSocket; `mqtt_p1` reads a topic | +| `meter_phases` | 1 | 1 or 3. Must match the telegram, or every telegram is rejected and logged | +| `meter_max_age_s` | 30 | Beyond this the reading is stale: grid power reads as *missing*, and the existing failsafe commands 0 W | +| `meter_mqtt_topic` | | `mqtt_p1` only | +| `p1_import_entity` | | The **unsigned** consumption sensor. Do not point this at a signed template | +| `p1_export_entity` | | The **unsigned** injection sensor | +| `p1_phase_import_entities` | `[]` | L1..L3, in order. Needed for the capacity-tariff peak on a three-phase connection | +| `p1_phase_export_entities` | `[]` | L1..L3, in order | + +There is **no fallback to an inverter-side power figure**, deliberately. The +inverter's own AC power tracks its battery almost perfectly and the real meter +hardly at all, so a controller that failed over to it would be regulating +against its own output while looking healthy. + +The `mqtt_p1` payload is one JSON object per telegram, and the schema is strict — +a key it does not recognise is a telegram from something other than what was +tested, and guessing a key here means guessing a kilowatt: + +```json +{"import_w": 1234.0, + "export_w": 0.0, + "phases": [{"import_w": 500, "export_w": 0}, + {"import_w": 400, "export_w": 0}, + {"import_w": 334, "export_w": 0}], + "timestamp": "2026-08-24T18:00:05+02:00"} +``` + +`phases` and `timestamp` are optional; `timestamp` must carry a UTC offset. Where +it is present it is used for the age, which is what stops a retained message +replayed on reconnect from presenting a ten-minute-old reading as current. + +#### `sensor.p1_sample_age_s` + +Published over MQTT discovery whenever a broker is available: **seconds since the +newest accepted telegram**, refreshed every second rather than only when a +telegram lands. The ESP32's stale-input watchdog subscribes to this exact entity +id, so do not rename it. + +The reason it is recomputed against the clock is that Home Assistant only pushes +a state when the state *changes*. A meter sitting at a genuinely constant reading +emits nothing, which is indistinguishable — to anything watching the value — from +a meter that has died. Watching the age instead separates the two: it climbs when +telegrams stop and resets when they arrive, whatever the reading says. + ### Control | option | default | meaning | diff --git a/goodwe_controller/app/main.py b/goodwe_controller/app/main.py index f42ec5c..b9de864 100644 --- a/goodwe_controller/app/main.py +++ b/goodwe_controller/app/main.py @@ -39,6 +39,7 @@ from .control import Tuning, compute, maintenance_charge_floor, peak_at_risk from .hass import HomeAssistant from .maintenance import IDLE, MaintConfig, Maintenance from .mqtt import MqttPublisher +from .p1 import P1Ingest, build_source from . import web OPTIONS_PATH = "/data/options.json" @@ -90,6 +91,13 @@ class Controller: store, ) + # P1 ingestion (TEL-01). `meter_source: off` keeps the original + # single-entity meter_entity path, so an existing install is unchanged + # until it opts in. + self.p1 = P1Ingest(phases=int(opts.get("meter_phases", 1)), + max_age_s=float(opts.get("meter_max_age_s", 30))) + self.p1_enabled = str(opts.get("meter_source", "off")) not in ("off", "") + # live state self.auto = bool(store.data.get("auto", opts.get("auto_start", False))) self.target = 0.0 @@ -121,8 +129,18 @@ class Controller: # -- io ------------------------------------------------------------------ async def read_inputs(self) -> None: o = self.o - self.grid = await self.hass.number(o.get("meter_entity", ""), - bool(o.get("meter_invert"))) + if self.p1_enabled: + # ⚠️ P1 is the only authoritative measurement of what the utility + # sees (§5.1). When it is stale this is None, which falls into the + # existing "inputs missing -> command 0 W" path below. There is + # deliberately NO fallback to an inverter-side figure: the + # inverter's own AC power correlates 0.998 with battery power and + # 0.09 with the real meter, so a controller that failed over to it + # would be regulating against its own output. + self.grid = self.p1.net_w + else: + self.grid = await self.hass.number(o.get("meter_entity", ""), + bool(o.get("meter_invert"))) self.soc = await self.hass.number(o.get("soc_entity", "")) self.batt = await self.hass.number(o.get("batt_entity", ""), bool(o.get("batt_invert"))) @@ -313,6 +331,8 @@ class Controller: "soc": self.soc, "phase": self.maint.phase, "status": "running" if self.auto else "stopped", + # Recomputed here, once a second, on purpose - see P1Ingest. + "p1_age": round(self.p1.published_age_s, 1), }) async def shutdown(self) -> None: @@ -326,11 +346,27 @@ class Controller: def checks(self) -> list: o = self.o out = [] - for label, value, entity in ( - ("grid power", self.grid, o.get("meter_entity")), - ("battery SoC", self.soc, o.get("soc_entity")), - ("battery power", self.batt, o.get("batt_entity")), - ): + if self.p1_enabled: + age = self.p1.published_age_s + if self.p1.stale: + out.append({"ok": False, "warn": False, + "text": f"P1 meter ({o.get('meter_source')}): no reading for " + f"{age:.0f} s (limit {self.p1.max_age_s:.0f} s)" + + (f" - last error: {self.p1.last_error}" + if self.p1.last_error else "")}) + else: + out.append({"ok": True, "warn": False, + "text": f"P1 meter ({o.get('meter_source')}): {self.p1.net_w:g} W, " + f"{age:.0f} s old, {self.p1.samples} telegrams, " + f"{self.p1.parse_errors} rejected"}) + rows = [("battery SoC", self.soc, o.get("soc_entity")), + ("battery power", self.batt, o.get("batt_entity"))] + if not self.p1_enabled: + # In P1 mode the check above replaces this one; leaving both in + # would report "no entity configured" for a meter_entity that is + # correctly unused, i.e. a permanent false NOT READY. + rows.insert(0, ("grid power", self.grid, o.get("meter_entity"))) + for label, value, entity in rows: if not entity: out.append({"ok": False, "warn": False, "text": f"{label}: no entity configured"}) elif value is None: @@ -457,6 +493,7 @@ async def amain() -> None: # observability, and the battery does not care. Caught broadly and on # purpose: this crashed the add-on once already (paho 1.x vs 2.x) and # took the control loop down with it. + broker = None try: broker = await hass.mqtt_service() pub = MqttPublisher( @@ -484,13 +521,24 @@ async def amain() -> None: with contextlib.suppress(NotImplementedError): loop.add_signal_handler(sig, stop.set) - task = asyncio.create_task(controller.run_control()) + tasks = [asyncio.create_task(controller.run_control())] + + # P1 ingestion runs as its own long-lived task. ⚠️ It must not be driven + # off the control loop: telegrams arrive every ~5 s and the loop would + # decimate them, so the 15-minute average - the capacity-tariff billing + # unit - would be computed from a fraction of the data. + p1_source = build_source(opts, controller.p1, session, broker) + if p1_source is not None: + tasks.append(asyncio.create_task(p1_source.run())) + await stop.wait() await controller.shutdown() - task.cancel() - with contextlib.suppress(asyncio.CancelledError): - await task + for task in tasks: + task.cancel() + for task in tasks: + with contextlib.suppress(asyncio.CancelledError): + await task await runner.cleanup() _LOG.info("stopped") diff --git a/goodwe_controller/app/mqtt.py b/goodwe_controller/app/mqtt.py index 7c8fadf..9834cc5 100644 --- a/goodwe_controller/app/mqtt.py +++ b/goodwe_controller/app/mqtt.py @@ -45,6 +45,14 @@ SENSORS = [ ("soc", "goodwe_battery_soc", "Battery SoC", "%", "battery", "measurement", None), ("phase", "goodwe_maintenance_phase", "Maintenance phase", None, None, None, "mdi:battery-sync"), ("status", "goodwe_controller_status", "Controller status", None, None, None, "mdi:heart-pulse"), + # ⚠️ This one deliberately breaks the goodwe_ prefix above: the entity id + # must be exactly `sensor.p1_sample_age_s`, because SAFETY-01's firmware + # watchdog subscribes to that literal id and the ENV-01 simulation rig + # asserts on it. Renaming it silently disarms a safety layer. It is seconds + # since the newest accepted P1 telegram, republished every second so that a + # meter frozen at a constant value still shows a climbing age - which is the + # false-trip that this entity exists to remove. + ("p1_age", "p1_sample_age_s", "P1 sample age", "s", "duration", "measurement", None), ] BASE = "goodwe_ctl" diff --git a/goodwe_controller/app/p1.py b/goodwe_controller/app/p1.py new file mode 100644 index 0000000..5f4b2a1 --- /dev/null +++ b/goodwe_controller/app/p1.py @@ -0,0 +1,674 @@ +"""P1 meter ingestion - the only authoritative measurement of real grid exchange. + +Everything downstream trusts this module: the safety checks, the capacity-tariff +peak, the optimizer, the control loop's sign. So three things happen here and +nowhere else. + + 1. The IMPORT/EXPORT DERIVATION. A Belgian P1 meter exposes two UNSIGNED + registers - consumption and injection - never one signed figure. Net power + is `import_w - export_w`, positive = import, and that subtraction is done + exactly once, here (spec §5.2: "the derivation is the EMS's job, not a + template the user has to write"). A second copy of it somewhere else is a + second chance to invert the control loop. + + 2. THE INGEST TIMESTAMP. Every accepted sample is stamped on arrival. A value + with no age is a value that cannot be trusted (§5.2), and staleness is the + failsafe trigger (§11.2). + + 3. VALIDATION. This is untrusted external data at the edge of a safety chain. + A malformed telegram must not become a plausible-looking number, and it + must never resolve to 0 W - a fabricated zero is indistinguishable from a + balanced house and defeats the very staleness trigger this module feeds. + +⚠️ There is deliberately NO fallback to an inverter-side power figure. The +inverter's own AC power correlates 0.998 with battery power and 0.09 with the +real meter (§5.1) - regulating on it means regulating against your own output. +When the transport dies the correct behaviour is a gap: no sample, a growing +age, and the existing "inputs missing -> command 0 W" path in main.py. +""" + +import asyncio +import json +import logging +import math +import os +import time +from dataclasses import dataclass +from datetime import datetime, timezone + +import aiohttp + +try: + import paho.mqtt.client as mqtt +except ImportError: # pragma: no cover - container always has it + mqtt = None + +_LOG = logging.getLogger("goodwe.p1") + +SOURCE_HA = "ha_dsmr" +SOURCE_MQTT = "mqtt_p1" + +QUARTER_S = 900 + +# ⚠️ Plausibility ceiling, not a clamp - anything above it is rejected as an +# anomaly rather than averaged in. Chosen to sit above the largest Belgian +# residential connection (3x63 A ~ 43 kW) and BELOW 65535: §20 open question 5 +# records an HA sensor reporting 64954 for -582 W, i.e. an unsigned 16-bit +# register decoded without its sign. That corruption reads as a perfectly +# plausible 65 kW if you only bound it at "some big number". +PLAUSIBLE_MAX_W = 50_000.0 + + +# --------------------------------------------------------------------------- # +# the sample +# --------------------------------------------------------------------------- # +class P1Error(ValueError): + """A telegram that must be rejected rather than believed.""" + + +@dataclass(frozen=True) +class P1Sample: + """One telegram, validated, derived and stamped. + + Frozen on purpose: this object is handed to readers on other tasks (and, + for the MQTT transport, produced on paho's network thread). Immutability is + what makes "read the latest sample" safe without a lock. + """ + + ingest_ts: datetime # tz-aware UTC, set at ingest + ingest_mono: float # time.monotonic() at ingest - see age_s() + telegram_ts: datetime | None # from the telegram, where the source has one + source: str # SOURCE_HA | SOURCE_MQTT + import_w: float # unsigned magnitude, as the meter reports it + export_w: float # unsigned magnitude + net_w: float # import_w - export_w (+ import, - export) + per_phase_w: tuple[float, ...] | None # signed net, len == phases + per_phase_import_w: tuple[float, ...] | None # offtake only, for the tariff + + def age_s(self, now_mono: float | None = None, + now_utc: datetime | None = None) -> float: + """Seconds since this sample was ingested, never negative. + + ⚠️ Measured with time.monotonic(), not the wall clock. An NTP step on a + Pi that just booted moves the wall clock by minutes; using it here would + either fake a stale meter or, worse, hide a real one. + + Where the telegram carries its own timestamp we take the WORSE of the + two ages. That is what stops an MQTT retained message - replayed on + reconnect with a fresh receive time - from presenting a ten-minute-old + reading as brand new. + """ + now_mono = time.monotonic() if now_mono is None else now_mono + age = max(0.0, now_mono - self.ingest_mono) + if self.telegram_ts is not None: + now_utc = datetime.now(timezone.utc) if now_utc is None else now_utc + age = max(age, (now_utc - self.telegram_ts).total_seconds()) + return max(0.0, age) + + +def _watts(value, what: str) -> float: + """Parse one power figure, or raise. Never returns a substituted default. + + ⚠️ Strings are refused even when float() would happily take them. A JSON + telegram carrying "1200" where a number belongs is a payload from a source + that is not the one we validated against, and the next surprise it has may + not be a benign one. Transports that legitimately deal in text (HA entity + states are always strings) convert before they get here, so this stays the + strict edge for structured payloads. + """ + if isinstance(value, (bool, str, bytes)) or value is None: + raise P1Error(f"{what}: not a number ({value!r})") + try: + out = float(value) + except (TypeError, ValueError): + raise P1Error(f"{what}: not a number ({value!r})") from None + if not math.isfinite(out): + raise P1Error(f"{what}: not finite ({value!r})") + if abs(out) > PLAUSIBLE_MAX_W: + raise P1Error(f"{what}: {out:g} W is outside plausible meter range") + return out + + +def make_sample(source: str, import_w, export_w, *, phases: int, + phase_import_w=None, phase_export_w=None, + telegram_ts: datetime | None = None, + ingest_ts: datetime | None = None, + ingest_mono: float | None = None) -> P1Sample: + """Validate, derive net power, stamp. Raises P1Error on anything doubtful. + + `import_w`/`export_w` are the two unsigned Belgian registers. Per-phase + figures are equally unsigned and equally split, so each phase gets the same + derivation. + """ + imp = _watts(import_w, "import") + exp = _watts(export_w, "export") + # ⚠️ Both registers are magnitudes. A negative one means the upstream + # already applied a sign we are about to apply again - reject it rather + # than silently double-signing the control loop. + if imp < 0 or exp < 0: + raise P1Error(f"unsigned registers cannot be negative (import={imp:g} export={exp:g})") + + per_phase = per_phase_import = None + if phase_import_w is not None or phase_export_w is not None: + pi = list(phase_import_w or []) + pe = list(phase_export_w or [0.0] * len(pi)) + if len(pi) != phases or len(pe) != phases: + raise P1Error( + f"phase count mismatch: telegram has {len(pi)} import / {len(pe)} export " + f"phases, meter_phases is {phases}") + vals = [_watts(a, f"L{i + 1} import") - _watts(b, f"L{i + 1} export") + for i, (a, b) in enumerate(zip(pi, pe))] + per_phase = tuple(vals) + per_phase_import = tuple(max(v, 0.0) for v in vals) + + if telegram_ts is not None and telegram_ts.tzinfo is None: + raise P1Error("telegram timestamp has no timezone") + + return P1Sample( + ingest_ts=ingest_ts or datetime.now(timezone.utc), + ingest_mono=time.monotonic() if ingest_mono is None else ingest_mono, + telegram_ts=telegram_ts, + source=source, + import_w=imp, + export_w=exp, + net_w=imp - exp, + per_phase_w=per_phase, + per_phase_import_w=per_phase_import, + ) + + +# --------------------------------------------------------------------------- # +# the 15-minute average +# --------------------------------------------------------------------------- # +@dataclass(frozen=True) +class QuarterBlock: + start: datetime # UTC, aligned to :00/:15/:30/:45 + offtake_avg_w: float # billed figure: net offtake only + per_phase_offtake_avg_w: tuple[float, ...] | None + + +class QuarterAverager: + """Time-weighted average of net offtake over clock-aligned 15-min blocks. + + Samples arrive irregularly (~1-10 s), so a plain mean over samples would + weight a burst of fast telegrams the same as a slow one and produce a figure + that is not the billed quantity. Each sample's value is therefore HELD until + the next arrives and integrated over that interval: sum(value * dt) / dt. + + ⚠️ Only OFFTAKE is accumulated (§9.1) - the capacity tariff bills the highest + quarter-hour average offtake, and a quarter of pure export averages to 0 kW, + not to a negative one. The signed series stays available for control; this + accumulator is for the meter's bill. + + ⚠️ Blocks are found by flooring epoch seconds to 900. That IS clock-aligned + and DST-proof for Belgium, because every offset in that tz is a whole number + of hours, so a 900 s grid in UTC lands on :00/:15/:30/:45 local before and + after a transition - no tz database, no DST special case. + ponytail: the ceiling is a timezone with a sub-hour offset (India +05:30, + Nepal, Chatham). Those need real tz-aware boundary maths; upgrade path is to + compute the boundary with zoneinfo instead of the modulo, everything else + here is unchanged. + + A sample that straddles a boundary is split at the boundary and its two + halves credited to the two blocks, never attributed wholly to either. + """ + + def __init__(self, phases: int = 1): + self.phases = phases + self._block: int | None = None # epoch seconds of the block start + self._acc = 0.0 # W*s of offtake in the open block + self._pp_acc = [0.0] * phases + self._elapsed = 0.0 # seconds integrated in the open block + self._last_t: float | None = None # epoch seconds of the held sample + self._last_net = 0.0 + self._last_pp: tuple[float, ...] | None = None + + # -- reading ------------------------------------------------------------ + @property + def block_start(self) -> datetime | None: + if self._block is None: + return None + return datetime.fromtimestamp(self._block, timezone.utc) + + @property + def elapsed_s(self) -> float: + """Seconds already integrated into the open block. + + Exposed alongside the partial accumulator because SAFETY-07 projects the + end-of-quarter average and cannot do that from a finished average. + """ + return self._elapsed + + @property + def partial_ws(self) -> float: + """Offtake watt-seconds accumulated in the open block so far.""" + return self._acc + + @property + def offtake_avg_w(self) -> float: + """Average offtake over the part of the open block seen so far.""" + return self._acc / self._elapsed if self._elapsed > 0 else 0.0 + + @property + def per_phase_offtake_avg_w(self) -> tuple[float, ...] | None: + if self._last_pp is None or self._elapsed <= 0: + return None + return tuple(a / self._elapsed for a in self._pp_acc) + + # -- writing ------------------------------------------------------------ + def add(self, sample: P1Sample) -> list[QuarterBlock]: + """Integrate up to this sample, then hold its value. Returns any blocks + that closed in the process (usually none, occasionally one).""" + t = sample.ingest_ts.timestamp() + closed: list[QuarterBlock] = [] + + if self._last_t is None: + self._block = int(t // QUARTER_S) * QUARTER_S + self._last_t, self._last_net = t, sample.net_w + self._last_pp = sample.per_phase_w + return closed + if t <= self._last_t: + # Out-of-order or duplicate arrival: integrating a negative dt would + # subtract energy that really happened. Drop it, keep the held value. + return closed + + cursor = self._last_t + while True: + end = self._block + QUARTER_S + stop = min(t, end) + dt = stop - cursor + if dt > 0: + self._acc += max(self._last_net, 0.0) * dt + if self._last_pp is not None: + for i, v in enumerate(self._last_pp[: self.phases]): + self._pp_acc[i] += max(v, 0.0) * dt + self._elapsed += dt + cursor = stop + if stop < end: + break + closed.append(QuarterBlock( + start=datetime.fromtimestamp(self._block, timezone.utc), + # A closed block is always divided by the full 900 s, never by + # the seconds we happened to observe - a gap in coverage must + # drag the billed average down, not be averaged away. + offtake_avg_w=self._acc / QUARTER_S, + per_phase_offtake_avg_w=( + tuple(a / QUARTER_S for a in self._pp_acc) + if self._last_pp is not None else None), + )) + self._block = end + self._acc = 0.0 + self._pp_acc = [0.0] * self.phases + self._elapsed = 0.0 + + self._last_t, self._last_net = t, sample.net_w + self._last_pp = sample.per_phase_w + return closed + + +# --------------------------------------------------------------------------- # +# what the rest of the add-on talks to +# --------------------------------------------------------------------------- # +class P1Ingest: + """Holds the latest sample and the rolling quarter-hour average. + + ⚠️ Staleness is DERIVED from the stored sample, not carried as a separate + flag. That is what makes "flag before the value is visible" free: there is + one immutable object and a single attribute rebind to publish it, so a + reader can never see a fresh value with a stale flag or the reverse. + """ + + def __init__(self, phases: int = 1, max_age_s: float = 30.0): + self.phases = phases + self.max_age_s = float(max_age_s) + self.averager = QuarterAverager(phases) + self.blocks: list[QuarterBlock] = [] + self.samples = 0 + self.parse_errors = 0 + self.last_error: str | None = None + self.started_mono = time.monotonic() + self._last: P1Sample | None = None + + @property + def last(self) -> P1Sample | None: + return self._last + + def submit(self, sample: P1Sample) -> None: + self._last = sample + self.samples += 1 + for block in self.averager.add(sample): + self.blocks.append(block) + del self.blocks[:-96] # a day of quarters; STATE-01 owns real retention + + def reject(self, err: Exception | str) -> None: + """A malformed telegram or an unavailable entity. + + ⚠️ The last good sample and ITS timestamp are left untouched. The reading + does not become 0 W and it does not become fresh - the age keeps growing, + which is precisely the signal a rejected telegram should produce. + """ + self.parse_errors += 1 + self.last_error = str(err) + _LOG.warning("P1 telegram rejected: %s", err) + + # -- what consumers read ------------------------------------------------- + def age_s(self) -> float | None: + """Age of the newest accepted sample, or None if there has never been one.""" + return None if self._last is None else self._last.age_s() + + @property + def published_age_s(self) -> float: + """The figure behind `sensor.p1_sample_age_s`. + + Seconds since the newest accepted telegram, or since this ingester + started when none has ever arrived. + + ⚠️ Always a number and never `unknown`, because SAFETY-01's firmware + watchdog subscribes to it: an entity that simply stops existing is + indistinguishable, from the firmware's side, from a meter that is fine. + And ⚠️ it is recomputed against the clock on every publish rather than + stamped once per telegram, so a meter that freezes at a constant reading + still produces a visibly climbing age. That is the whole point of this + entity - HA pushes state changes, so a genuinely constant P1 value emits + nothing at all, and a watchdog watching the value would sit there + believing the last update was recent. + """ + age = self.age_s() + return max(0.0, time.monotonic() - self.started_mono) if age is None else age + + @property + def stale(self) -> bool: + """True when there is no sample, or the newest one is past max_age_s.""" + age = self.age_s() + return age is None or age > self.max_age_s + + @property + def net_w(self) -> float | None: + """Signed net grid power, or None when stale. Never a substituted zero.""" + return None if self.stale else self._last.net_w + + @property + def per_phase_import_w(self) -> tuple[float, ...] | None: + if self.stale or self._last is None: + return None + return self._last.per_phase_import_w + + +# --------------------------------------------------------------------------- # +# transport 1: Home Assistant WebSocket (the DSMR integration's entities) +# --------------------------------------------------------------------------- # +WS_URL = "ws://supervisor/core/websocket" +BAD_STATES = ("unknown", "unavailable", "none", "") + + +class HaDsmrSource: + """Subscribes to state_changed for the configured DSMR entities. + + ⚠️ WebSocket, not REST polling. REST returns states, but polling at the + 30 s planning tick decimates a 5 s telegram stream and the quarter-hour + average would then be computed from a sixth of the data (§5.3). "Consume + every telegram" means event-driven. + + ⚠️ One telegram updates several entities, and HA emits one state_changed per + entity. Building a sample on each event would mix a new import reading with + a stale export one for a few milliseconds every 5 s. A short debounce + coalesces the burst back into the single telegram it came from. + """ + + DEBOUNCE_S = 0.35 + + def __init__(self, session: aiohttp.ClientSession, ingest: P1Ingest, + entities: dict, token: str | None = None): + self.session = session + self.ingest = ingest + self.entities = entities # {"import": id, "export": id, "phase_import": [...], ...} + self.token = token or os.environ.get("SUPERVISOR_TOKEN", "") + self.ids = self._wanted() + self.cache: dict[str, float] = {} + self.connected = False + self._pending: asyncio.Task | None = None + + def _wanted(self) -> set[str]: + out = set() + for key in ("import", "export"): + if self.entities.get(key): + out.add(self.entities[key]) + for key in ("phase_import", "phase_export"): + out.update(e for e in self.entities.get(key) or [] if e) + return out + + async def run(self) -> None: + """Long-lived task: connect, subscribe, reconnect with backoff, forever. + + ⚠️ A reconnect emits nothing. A gap must stay a gap - a synthetic sample + on reconnect would reset the age and hide the outage from the very + watchdog that exists to catch it. + """ + backoff = 1.0 + while True: + try: + await self._session_once() + backoff = 1.0 + except asyncio.CancelledError: + raise + except Exception as err: # noqa: BLE001 - any transport fault retries + _LOG.warning("P1 HA websocket: %s - reconnecting in %.0fs", err, backoff) + finally: + self.connected = False + await asyncio.sleep(backoff) + backoff = min(backoff * 2, 30.0) + + async def _session_once(self) -> None: + async with self.session.ws_connect(WS_URL, heartbeat=30) as ws: + hello = await ws.receive_json() + if hello.get("type") == "auth_required": + await ws.send_json({"type": "auth", "access_token": self.token}) + reply = await ws.receive_json() + if reply.get("type") != "auth_ok": + raise RuntimeError(f"auth rejected: {reply.get('message', reply)}") + await ws.send_json({"id": 1, "type": "subscribe_events", + "event_type": "state_changed"}) + await ws.send_json({"id": 2, "type": "get_states"}) + self.connected = True + _LOG.info("P1 ingest: subscribed to %s", ", ".join(sorted(self.ids))) + + async for msg in ws: + if msg.type is not aiohttp.WSMsgType.TEXT: + continue + payload = json.loads(msg.data) + if payload.get("id") == 2 and payload.get("type") == "result": + for obj in payload.get("result") or []: + self._absorb(obj.get("entity_id"), obj.get("state")) + self._schedule() + elif payload.get("type") == "event": + data = (payload.get("event") or {}).get("data") or {} + if data.get("entity_id") not in self.ids: + continue + new = data.get("new_state") or {} + self._absorb(data.get("entity_id"), new.get("state")) + self._schedule() + raise RuntimeError("websocket closed") + + def _absorb(self, entity_id: str | None, state) -> None: + if not entity_id or entity_id not in self.ids: + return + raw = str(state).strip().lower() + if raw in BAD_STATES: + # ⚠️ An `unavailable` DSMR entity is a missing reading, not 0 W. + # Forget the cached value so no sample can be built from a mixture + # of a live register and one that stopped reporting. + self.cache.pop(entity_id, None) + self.ingest.reject(f"{entity_id} is {raw}") + return + try: + self.cache[entity_id] = float(raw) + except ValueError: + self.cache.pop(entity_id, None) + self.ingest.reject(f"{entity_id} is not numeric: {raw!r}") + + def _schedule(self) -> None: + if self._pending and not self._pending.done(): + return + self._pending = asyncio.get_running_loop().create_task(self._after_debounce()) + + async def _after_debounce(self) -> None: + await asyncio.sleep(self.DEBOUNCE_S) + self.build() + + def build(self) -> bool: + """Assemble one sample from the cache. Returns True if one was accepted.""" + imp_id, exp_id = self.entities.get("import"), self.entities.get("export") + if imp_id not in self.cache or exp_id not in self.cache: + return False + pi = [self.cache.get(e) for e in self.entities.get("phase_import") or []] + pe = [self.cache.get(e) for e in self.entities.get("phase_export") or []] + if pi and (None in pi or (pe and None in pe)): + return False # incomplete phase set: wait, do not guess + try: + self.ingest.submit(make_sample( + SOURCE_HA, self.cache[imp_id], self.cache[exp_id], + phases=self.ingest.phases, + phase_import_w=pi or None, + phase_export_w=pe or None, + # ⚠️ No telegram_ts: HA's last_changed is when the STATE changed, + # which for a constant reading is minutes ago even though the + # telegram is current. Using it as a telegram time would fake + # staleness on a genuinely steady meter. + )) + return True + except P1Error as err: + self.ingest.reject(err) + return False + + +# --------------------------------------------------------------------------- # +# transport 2: MQTT +# --------------------------------------------------------------------------- # +def parse_mqtt_payload(raw: bytes | str, phases: int, *, + now: datetime | None = None) -> P1Sample: + """One JSON telegram from the configured topic. Raises P1Error. + + The accepted document, documented in DOCS.md: + + {"import_w": 1234.0, "export_w": 0.0, + "phases": [{"import_w": 500, "export_w": 0}, ...], # optional + "timestamp": "2026-08-24T18:00:05+02:00"} # optional + + ponytail: one strict schema rather than sniffing the half-dozen P1-bridge + dialects in the wild. The upgrade path is a `meter_mqtt_format` option + selecting a parser; a lenient parser is the wrong default at a safety + boundary, where guessing a key means guessing a kilowatt. + """ + try: + doc = json.loads(raw) + except (ValueError, TypeError) as err: + raise P1Error(f"payload is not JSON: {err}") from None + if not isinstance(doc, dict): + raise P1Error(f"payload is not a JSON object ({type(doc).__name__})") + + ts = None + if doc.get("timestamp"): + try: + ts = datetime.fromisoformat(str(doc["timestamp"])) + except ValueError: + raise P1Error(f"unparseable timestamp {doc['timestamp']!r}") from None + if ts.tzinfo is None: + raise P1Error("timestamp has no UTC offset") + + pi = pe = None + if "phases" in doc: + rows = doc["phases"] + if not isinstance(rows, list) or not all(isinstance(r, dict) for r in rows): + raise P1Error("'phases' must be a list of objects") + pi = [r.get("import_w") for r in rows] + pe = [r.get("export_w", 0.0) for r in rows] + + return make_sample(SOURCE_MQTT, doc.get("import_w"), doc.get("export_w"), + phases=phases, phase_import_w=pi, phase_export_w=pe, + telegram_ts=ts, ingest_ts=now) + + +class MqttP1Source: + """Subscribes to one topic and submits every message that parses. + + Uses paho's own reconnect loop on its own thread, then hops back onto the + event loop with call_soon_threadsafe so the ingest state is only ever + mutated from one thread. + """ + + def __init__(self, ingest: P1Ingest, topic: str, host, port=1883, + username=None, password=None): + self.ingest = ingest + self.topic = topic + self.host, self.port = host, int(port or 1883) + self.username, self.password = username, password + self.client = None + self.loop = None + + async def run(self) -> None: + if mqtt is None or not self.host or not self.topic: + _LOG.error("P1 MQTT source not usable (broker=%s topic=%r) - no meter data", + self.host, self.topic) + return + self.loop = asyncio.get_running_loop() + try: + self.client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, + client_id="goodwe_p1_ingest") + except AttributeError: # paho 1.x, which is what Alpine ships + self.client = mqtt.Client(client_id="goodwe_p1_ingest") + if self.username: + self.client.username_pw_set(self.username, self.password or "") + self.client.on_connect = lambda *_a, **_k: self.client.subscribe(self.topic, qos=0) + self.client.on_message = self._on_message + self.client.reconnect_delay_set(min_delay=1, max_delay=30) + self.client.connect_async(self.host, self.port, keepalive=60) + self.client.loop_start() + _LOG.info("P1 ingest: MQTT %s:%s topic %s", self.host, self.port, self.topic) + try: + while True: + await asyncio.sleep(3600) + finally: + self.client.loop_stop() + self.client.disconnect() + + def _on_message(self, _client, _userdata, msg) -> None: + # Runs on paho's network thread. + if self.loop is None: + return + self.loop.call_soon_threadsafe(self._handle, msg.payload) + + def _handle(self, payload) -> None: + try: + self.ingest.submit(parse_mqtt_payload(payload, self.ingest.phases)) + except P1Error as err: + self.ingest.reject(err) + + +# --------------------------------------------------------------------------- # +# selection +# --------------------------------------------------------------------------- # +def build_source(opts: dict, ingest: P1Ingest, session, broker: dict | None): + """Return the transport named by `meter_source`, or None if disabled. + + This is the whole of AC 1's "switchable": every consumer reads P1Ingest, so + changing transport is a config edit, never a code path. + """ + source = str(opts.get("meter_source", "off") or "off").strip() + if source in ("off", ""): + return None + if source == SOURCE_HA: + return HaDsmrSource(session, ingest, { + "import": opts.get("p1_import_entity", ""), + "export": opts.get("p1_export_entity", ""), + "phase_import": opts.get("p1_phase_import_entities") or [], + "phase_export": opts.get("p1_phase_export_entities") or [], + }) + if source == SOURCE_MQTT: + broker = broker or {} + return MqttP1Source(ingest, str(opts.get("meter_mqtt_topic", "")), + broker.get("host"), broker.get("port", 1883), + broker.get("username"), broker.get("password")) + if source: + _LOG.error("meter_source %r is not %s or %s - P1 ingestion disabled", + source, SOURCE_HA, SOURCE_MQTT) + return None diff --git a/goodwe_controller/config.yaml b/goodwe_controller/config.yaml index f27b673..82c3499 100644 --- a/goodwe_controller/config.yaml +++ b/goodwe_controller/config.yaml @@ -39,6 +39,23 @@ options: batt_invert: false setpoint_entity: "" + # --- P1 meter ingestion (specs §5.2 / §14 `meter:`) ------------------------- + # `off` keeps the original single meter_entity path above, so an existing + # install is untouched until it opts in. ha_dsmr subscribes to the DSMR + # integration's entities over the HA WebSocket; mqtt_p1 reads the topic below. + meter_source: "off" + meter_phases: 1 + meter_max_age_s: 30 + meter_mqtt_topic: "" + # The two UNSIGNED Belgian registers. The EMS derives net power from them + # (import - export); do NOT point these at a signed template sensor. + p1_import_entity: "" + p1_export_entity: "" + # Optional, in L1..L3 order. Required for the capacity-tariff peak on a + # three-phase connection; the list length must equal meter_phases. + p1_phase_import_entities: [] + p1_phase_export_entities: [] + # --- control --------------------------------------------------------------- max_w: 2000 gain: 0.6 @@ -81,6 +98,21 @@ schema: batt_invert: bool setpoint_entity: str + meter_source: list(off|ha_dsmr|mqtt_p1) + # ⚠️ 2 is accepted by this range but is not a real Belgian connection. A + # telegram whose phase count disagrees is rejected at ingest and logged, so a + # mis-set 2 shows up immediately as "0 telegrams accepted" rather than as a + # quietly wrong number. + meter_phases: int(1,3) + meter_max_age_s: int(5,300) + meter_mqtt_topic: str? + p1_import_entity: str? + p1_export_entity: str? + p1_phase_import_entities: + - str + p1_phase_export_entities: + - str + max_w: int(100,5000) gain: float(0.05,1.0) slew_w: int(50,5000) diff --git a/goodwe_controller/test_p1.py b/goodwe_controller/test_p1.py new file mode 100644 index 0000000..07f706c --- /dev/null +++ b/goodwe_controller/test_p1.py @@ -0,0 +1,504 @@ +"""Runnable check for P1 ingestion. `python3 test_p1.py` + +No framework, no fixtures, no meter - it has to run on a tech's laptop and in CI +with nothing installed and nothing plugged in (§17: "usable without real +hardware"). + +Every assert here is a rule whose absence poisons something downstream: an +inverted sign inverts the control loop, a fabricated zero hides a dead meter, a +naive per-phase sum misbills the capacity tariff, and a plain mean of samples +computes the wrong quarter-hour figure whenever the telegram cadence changes. +""" + +import asyncio +import sys +import time +from datetime import datetime, timedelta, timezone + +import aiohttp # already required by app.p1, so this adds no new dependency + +from app.p1 import ( + P1Error, P1Ingest, HaDsmrSource, QuarterAverager, SOURCE_HA, SOURCE_MQTT, + make_sample, parse_mqtt_payload, +) + +fails = [] +total = 0 + +TZ_BE_SUMMER = timezone(timedelta(hours=2)) +TZ_BE_WINTER = timezone(timedelta(hours=1)) +BASE = datetime(2026, 8, 24, 10, 0, 0, tzinfo=timezone.utc) # a quarter boundary + + +def check(name, cond): + global total + total += 1 + if cond: + print(f" ok {name}") + else: + print(f" FAIL {name}") + fails.append(name) + + +def raises(name, fn): + global total + total += 1 + try: + fn() + except P1Error: + print(f" ok {name}") + return + except Exception as err: # noqa: BLE001 + print(f" FAIL {name} (raised {type(err).__name__}, wanted P1Error)") + fails.append(name) + return + print(f" FAIL {name} (no error raised)") + fails.append(name) + + +def sample(net_import, net_export=0.0, at=BASE, phases=1, pi=None, pe=None): + return make_sample(SOURCE_HA, net_import, net_export, phases=phases, + phase_import_w=pi, phase_export_w=pe, + ingest_ts=at, ingest_mono=at.timestamp()) + + +# --------------------------------------------------------------------------- # +print("import/export -> signed net (the derivation the EMS owns)") + +s = sample(1500.0, 0.0) +check("pure import is positive", s.net_w == 1500.0) + +s = sample(0.0, 900.0) +check("pure export is negative", s.net_w == -900.0) + +# The single test that catches an inverted control loop. +s = sample(120.0, 2000.0) +check("export-dominant telegram yields negative net", s.net_w == -1880.0) + +s = sample(2000.0, 120.0) +check("import-dominant telegram yields positive net", s.net_w == 1880.0) + +# Both registers non-zero at once is real: a three-phase house can import on one +# phase and export on another in the same telegram. +s = sample(400.0, 400.0) +check("both registers equal nets to exactly zero", s.net_w == 0.0) + +s = sample(0.0, 0.0) +check("both registers zero is a valid balanced reading", s.net_w == 0.0 + and s.import_w == 0.0 and s.export_w == 0.0) + +check("the magnitudes survive the derivation", + sample(300.0, 50.0).import_w == 300.0 and sample(300.0, 50.0).export_w == 50.0) + +# --------------------------------------------------------------------------- # +print("rejecting a telegram instead of believing it") + +raises("negative 'unsigned' import is rejected", lambda: sample(-100.0, 0.0)) +raises("negative 'unsigned' export is rejected", lambda: sample(0.0, -100.0)) +raises("a non-numeric register is rejected", lambda: sample("n/a", 0.0)) +raises("None is rejected, not read as zero", lambda: sample(None, 0.0)) +raises("NaN is rejected", lambda: sample(float("nan"), 0.0)) +raises("infinity is rejected", lambda: sample(float("inf"), 0.0)) +# §20 open question 5: an HA sensor reporting 64954 for -582 W, i.e. an unsigned +# 16-bit register decoded without its sign. Must not average in as 65 kW. +raises("the 64954 signed-decode contamination is rejected", + lambda: sample(64954.0, 0.0)) +raises("a naive timestamp is rejected", + lambda: make_sample(SOURCE_MQTT, 100, 0, phases=1, + telegram_ts=datetime(2026, 8, 24, 10, 0, 0))) + +# --------------------------------------------------------------------------- # +print("single- and three-phase") + +s = sample(800.0, 0.0, phases=1, pi=[800.0], pe=[0.0]) +check("single phase accepts one phase", s.per_phase_w == (800.0,)) +check("per-phase tuple, not list", isinstance(s.per_phase_w, tuple)) + +raises("three phases configured, one delivered -> rejected", + lambda: sample(800.0, 0.0, phases=3, pi=[800.0], pe=[0.0])) +raises("one phase configured, three delivered -> rejected", + lambda: sample(800.0, 0.0, phases=1, pi=[300.0, 300.0, 200.0], + pe=[0.0, 0.0, 0.0])) + +s = sample(600.0, 0.0) +check("no phase data leaves per-phase None, not a fabricated tuple", + s.per_phase_w is None and s.per_phase_import_w is None) + +# §17's three-phase unbalanced-load regression case. L1 imports hard, L2 exports, +# L3 idles: the connection nets to 600 W of offtake while 2100 W is drawn across +# the phases. This is the case a naive per-phase sum gets wrong. +s = sample(600.0, 0.0, phases=3, pi=[2000.0, 0.0, 100.0], pe=[0.0, 1500.0, 0.0]) +check("unbalanced: per-phase net keeps the export phase negative", + s.per_phase_w == (2000.0, -1500.0, 100.0)) +check("unbalanced: per-phase IMPORT clamps the exporting phase to zero", + s.per_phase_import_w == (2000.0, 0.0, 100.0)) +check("unbalanced: per-phase import sums to 2100 W, the connection nets 600 W", + sum(s.per_phase_import_w) == 2100.0 and s.net_w == 600.0) +check("unbalanced: the naive sum is NOT the billed figure", + sum(s.per_phase_import_w) != s.net_w) + +# --------------------------------------------------------------------------- # +print("rolling 15-minute average (time-weighted, clock-aligned, offtake only)") + +# Irregular spacing, hand-computed: +# 0->3 s held at 1000 W -> 3000 Ws +# 3->13 s held at 0 W -> 0 Ws +# 13->20 s held at 2000 W -> 14000 Ws +# total 17000 Ws over 20 s -> 850 W +a = QuarterAverager(1) +a.add(sample(1000.0, at=BASE)) +a.add(sample(0.0, at=BASE + timedelta(seconds=3))) +a.add(sample(2000.0, at=BASE + timedelta(seconds=13))) +a.add(sample(2000.0, at=BASE + timedelta(seconds=20))) +check("irregular spacing integrates to the hand-computed 850 W", + abs(a.offtake_avg_w - 850.0) < 1e-9) +check("elapsed-seconds-in-block is exposed for SAFETY-07", a.elapsed_s == 20.0) +check("partial accumulator is exposed for SAFETY-07", a.partial_ws == 17000.0) +# A plain mean over the four samples would be 1250 W. The whole point of the +# time weighting is that these two numbers differ. +check("a plain mean of those samples would have said 1250 W, not 850", + abs((1000 + 0 + 2000 + 2000) / 4 - 1250.0) < 1e-9 and a.offtake_avg_w != 1250.0) + +# Cadence change mid-block: 1 s telegrams for a minute, then 10 s telegrams. +# 0->60 s held at 1000 W -> 60000 Ws +# 60->600 s held at 100 W -> 54000 Ws +# 114000 Ws over 600 s -> 190 W +a = QuarterAverager(1) +for i in range(0, 61): + a.add(sample(1000.0 if i < 60 else 100.0, at=BASE + timedelta(seconds=i))) +for i in range(70, 601, 10): + a.add(sample(100.0, at=BASE + timedelta(seconds=i))) +check("a cadence change does not bias the average (190 W)", + abs(a.offtake_avg_w - 190.0) < 1e-9) +check("elapsed tracks the whole 600 s despite the cadence change", a.elapsed_s == 600.0) +# The decimation trap: the fast minute contributes 60 of 114 samples but only +# 10 % of the time, so a per-sample mean lands near 574 W - three times high. +naive = (60 * 1000 + 54 * 100) / 114 +check("a per-sample mean would have said ~574 W", 570 < naive < 578) + +# Straddling the boundary: 890 s into a block, next telegram 20 s later. Ten +# seconds belong to each block and must be split, not attributed to one. +a = QuarterAverager(1) +a.add(sample(1000.0, at=BASE + timedelta(seconds=890))) +closed = a.add(sample(1000.0, at=BASE + timedelta(seconds=910))) +check("crossing a boundary closes exactly one block", len(closed) == 1) +check("the closed block keeps only its own 10 s (10000/900 W)", + abs(closed[0].offtake_avg_w - 10000.0 / 900.0) < 1e-9) +check("the closed block divides by the full 900 s, so a gap drags it down", + closed[0].offtake_avg_w < 1000.0) +check("the new block carries the other 10 s", a.elapsed_s == 10.0 + and abs(a.offtake_avg_w - 1000.0) < 1e-9) +check("the closed block starts on a quarter boundary", + closed[0].start == BASE and closed[0].start.minute % 15 == 0) + +# A whole block of pure export: the capacity tariff bills offtake, so this is +# 0 kW, never a negative peak. +a = QuarterAverager(1) +a.add(sample(0.0, 3000.0, at=BASE)) +closed = a.add(sample(0.0, 3000.0, at=BASE + timedelta(seconds=900))) +check("a block of pure export averages to 0 W of offtake", + len(closed) == 1 and closed[0].offtake_avg_w == 0.0) +check("...while the signed sample itself stays negative", + sample(0.0, 3000.0).net_w == -3000.0) + +# Clock alignment holds either side of a DST change, because every Belgian UTC +# offset is a whole number of hours and the block grid is 900 s of UTC. +for label, tz in (("summer (+02:00)", TZ_BE_SUMMER), ("winter (+01:00)", TZ_BE_WINTER)): + a = QuarterAverager(1) + odd = datetime(2026, 8, 24, 13, 7, 23, tzinfo=tz) + a.add(sample(500.0, at=odd)) + local = a.block_start.astimezone(tz) + check(f"block boundary is local :00/:15/:30/:45 in {label}", + local.minute in (0, 15, 30, 45) and local.second == 0 and local.microsecond == 0) + +# Three-phase averaging keeps the phases apart. +a = QuarterAverager(3) +a.add(sample(600.0, 0.0, at=BASE, phases=3, pi=[2000.0, 0.0, 100.0], pe=[0.0, 1500.0, 0.0])) +a.add(sample(600.0, 0.0, at=BASE + timedelta(seconds=100), phases=3, + pi=[2000.0, 0.0, 100.0], pe=[0.0, 1500.0, 0.0])) +check("per-phase offtake averages are held separately", + a.per_phase_offtake_avg_w == (2000.0, 0.0, 100.0)) +check("the block's own average is the connection net, not the phase sum", + abs(a.offtake_avg_w - 600.0) < 1e-9) + +# An out-of-order arrival must not subtract energy that really happened. +a = QuarterAverager(1) +a.add(sample(1000.0, at=BASE)) +a.add(sample(1000.0, at=BASE + timedelta(seconds=10))) +before = a.partial_ws +a.add(sample(1000.0, at=BASE + timedelta(seconds=5))) +check("an out-of-order telegram is dropped, not integrated backwards", + a.partial_ws == before and a.elapsed_s == 10.0) + +# --------------------------------------------------------------------------- # +print("ingest timestamp, age and staleness") + +now = time.monotonic() +s = make_sample(SOURCE_HA, 1000, 0, phases=1, ingest_mono=now) +check("a sample carries a tz-aware ingest timestamp", + s.ingest_ts.tzinfo is not None) +check("age is ~0 immediately after ingest", s.age_s(now_mono=now) == 0.0) +check("age grows with elapsed time", s.age_s(now_mono=now + 12.5) == 12.5) +check("age never goes negative on a clock step", + s.age_s(now_mono=now - 100.0) == 0.0) + +ing = P1Ingest(phases=1, max_age_s=30.0) +check("no sample yet is stale, not zero", ing.stale is True and ing.net_w is None) +check("no sample yet has no age at all", ing.age_s() is None) + +ing.submit(make_sample(SOURCE_HA, 1234, 0, phases=1, ingest_mono=time.monotonic())) +check("a fresh sample is not stale", ing.stale is False) +check("a fresh sample exposes signed net power", ing.net_w == 1234.0) +check("a fresh sample's age is small", 0 <= ing.age_s() < 1.0) + +ing.submit(make_sample(SOURCE_HA, 1234, 0, phases=1, + ingest_mono=time.monotonic() - 29.0)) +check("29 s old with max_age_s 30 is still usable", ing.stale is False) +ing.submit(make_sample(SOURCE_HA, 1234, 0, phases=1, + ingest_mono=time.monotonic() - 31.0)) +check("31 s old with max_age_s 30 is stale", ing.stale is True) +check("a stale sample reads as None, never as 0 W", ing.net_w is None) +check("...and its per-phase import is None too", ing.per_phase_import_w is None) + +# The MQTT retained-message trap: replayed on reconnect with a fresh receive +# time but a ten-minute-old telegram time. Fresh ingest must not launder it. +old = datetime.now(timezone.utc) - timedelta(minutes=10) +ing = P1Ingest(phases=1, max_age_s=30.0) +ing.submit(make_sample(SOURCE_MQTT, 1000, 0, phases=1, telegram_ts=old, + ingest_mono=time.monotonic())) +check("a retained telegram is stale on arrival despite a fresh receive time", + ing.stale is True and ing.age_s() > 590) + +# A rejected telegram must not refresh anything and must not become 0 W. +ing = P1Ingest(phases=1, max_age_s=30.0) +ing.submit(make_sample(SOURCE_HA, 1500, 0, phases=1, + ingest_mono=time.monotonic() - 10.0)) +stamp = ing.last.ingest_mono +ing.reject("malformed telegram") +check("a rejected telegram counts as a parse error", ing.parse_errors == 1) +check("a rejected telegram leaves the last good value in place", + ing.last.net_w == 1500.0) +check("a rejected telegram does not refresh the timestamp", + ing.last.ingest_mono == stamp) +check("a rejected telegram does not resolve to 0 W", ing.net_w == 1500.0) + +# Meter goes stale but stays connected (§17 failure injection): no new sample +# ever arrives, and the age keeps climbing past max_age_s on its own. +ing = P1Ingest(phases=1, max_age_s=30.0) +ing.submit(make_sample(SOURCE_HA, 800, 0, phases=1, + ingest_mono=time.monotonic() - 120.0)) +check("connected-but-silent meter trips staleness with no new telegram", + ing.stale is True and ing.age_s() > 100) + +# sensor.p1_sample_age_s: SAFETY-01's firmware subscribes to this, so it must +# always be a number and must climb while nothing arrives. +ing = P1Ingest(phases=1, max_age_s=30.0) +ing.started_mono = time.monotonic() - 45.0 +check("the published age is a number before the first telegram ever arrives", + isinstance(ing.published_age_s, float) and ing.published_age_s > 44) +check("...and it is never None, unlike the raw age", + ing.age_s() is None and ing.published_age_s is not None) +ing.submit(make_sample(SOURCE_HA, 500, 0, phases=1, ingest_mono=time.monotonic())) +check("a telegram resets the published age", ing.published_age_s < 1.0) +# The false-trip this entity exists to remove: the meter keeps sending, the +# VALUE never changes, and the age must still reflect that it is being sent. +for _ in range(3): + ing.submit(make_sample(SOURCE_HA, 500, 0, phases=1, ingest_mono=time.monotonic())) +check("an unchanging meter value still reads as fresh while telegrams arrive", + ing.stale is False and ing.published_age_s < 1.0) +ing.submit(make_sample(SOURCE_HA, 500, 0, phases=1, + ingest_mono=time.monotonic() - 90.0)) +check("the same unchanging value reads as stale once the telegrams stop", + ing.stale is True and ing.published_age_s > 89) + +# --------------------------------------------------------------------------- # +print("MQTT transport: parsing a telegram off the topic") + +good = '{"import_w": 1200.5, "export_w": 0}' +s = parse_mqtt_payload(good, 1) +check("a well-formed payload parses", s.net_w == 1200.5 and s.source == SOURCE_MQTT) + +s = parse_mqtt_payload('{"import_w": 0, "export_w": 2500}', 1) +check("an export payload parses to a negative net", s.net_w == -2500.0) + +s = parse_mqtt_payload( + '{"import_w": 600, "export_w": 0, "phases":' + ' [{"import_w":2000,"export_w":0},{"import_w":0,"export_w":1500},' + ' {"import_w":100,"export_w":0}]}', 3) +check("a three-phase payload parses per-phase", + s.per_phase_w == (2000.0, -1500.0, 100.0)) + +s = parse_mqtt_payload( + '{"import_w": 100, "export_w": 0, "timestamp": "2026-08-24T12:00:00+02:00"}', 1) +check("a telegram timestamp is kept when the payload has one", + s.telegram_ts == datetime(2026, 8, 24, 12, 0, 0, tzinfo=TZ_BE_SUMMER)) + +raises("a non-JSON payload is rejected", lambda: parse_mqtt_payload("not json", 1)) +raises("a JSON array is rejected", lambda: parse_mqtt_payload("[1,2,3]", 1)) +raises("a payload missing import_w is rejected", + lambda: parse_mqtt_payload('{"export_w": 0}', 1)) +raises("a payload with a null register is rejected", + lambda: parse_mqtt_payload('{"import_w": null, "export_w": 0}', 1)) +raises("a payload with a string register is rejected", + lambda: parse_mqtt_payload('{"import_w": "1200", "export_w": 0}', 1)) +raises("a timestamp with no UTC offset is rejected", + lambda: parse_mqtt_payload( + '{"import_w":1,"export_w":0,"timestamp":"2026-08-24T12:00:00"}', 1)) +raises("a phase count that disagrees with meter_phases is rejected", + lambda: parse_mqtt_payload( + '{"import_w":1,"export_w":0,"phases":[{"import_w":1,"export_w":0}]}', 3)) + +# --------------------------------------------------------------------------- # +print("HA DSMR transport: building a sample out of entity states") + +ENT = {"import": "sensor.p1_import", "export": "sensor.p1_export", + "phase_import": [], "phase_export": []} +ing = P1Ingest(phases=1, max_age_s=30.0) +src = HaDsmrSource(None, ing, ENT, token="x") + +src._absorb("sensor.p1_import", "1000") +check("one entity alone does not build a sample", src.build() is False + and ing.last is None) +src._absorb("sensor.p1_export", "0") +check("both entities present builds one sample", src.build() is True + and ing.net_w == 1000.0) + +# ⚠️ The reason the debounce exists: HA emits one state_changed per entity, so +# mid-telegram the cache briefly holds a new import with the old export. +src._absorb("sensor.p1_import", "0") +src._absorb("sensor.p1_export", "2500") +src.build() +check("a coalesced telegram lands as one consistent sample", ing.net_w == -2500.0) + +before = ing.last +src._absorb("sensor.p1_export", "unavailable") +check("an unavailable entity is a parse error", ing.parse_errors == 1) +check("an unavailable entity does not build a sample from a stale half", + src.build() is False) +check("an unavailable entity leaves the last good sample untouched", + ing.last is before and ing.net_w == -2500.0) + +src._absorb("sensor.p1_export", "unknown") +check("an unknown entity is treated the same way", ing.parse_errors == 2) +src._absorb("sensor.p1_export", "banana") +check("a non-numeric entity state is a parse error, not 0 W", + ing.parse_errors == 3 and ing.net_w == -2500.0) + +src._absorb("sensor.p1_export", "64954") +src._absorb("sensor.p1_import", "0") +check("the contaminated signed decode is refused at the transport too", + src.build() is False and ing.parse_errors == 4) + +ing3 = P1Ingest(phases=3, max_age_s=30.0) +ENT3 = {"import": "sensor.p1_import", "export": "sensor.p1_export", + "phase_import": ["sensor.l1_i", "sensor.l2_i", "sensor.l3_i"], + "phase_export": ["sensor.l1_e", "sensor.l2_e", "sensor.l3_e"]} +src3 = HaDsmrSource(None, ing3, ENT3, token="x") +for eid, val in (("sensor.p1_import", "600"), ("sensor.p1_export", "0"), + ("sensor.l1_i", "2000"), ("sensor.l2_i", "0")): + src3._absorb(eid, val) +check("an incomplete phase set waits instead of guessing", src3.build() is False) +for eid, val in (("sensor.l3_i", "100"), ("sensor.l1_e", "0"), + ("sensor.l2_e", "1500"), ("sensor.l3_e", "0")): + src3._absorb(eid, val) +check("a complete three-phase set builds", src3.build() is True) +check("three-phase entities produce the unbalanced per-phase tuple", + ing3.last.per_phase_w == (2000.0, -1500.0, 100.0)) +check("the sample's per-phase tuple length matches meter_phases", + len(ing3.last.per_phase_w) == ing3.phases == 3) + +# --------------------------------------------------------------------------- # +print("HA DSMR transport: end to end against a fake Home Assistant") +# Everything above pokes at build()/_absorb() directly. This one drives the real +# thing over a real websocket - auth handshake, subscribe_events, get_states, +# per-entity events - because "the transport works" is otherwise an untested +# claim about a protocol nobody re-reads. + + +async def _e2e(): + from aiohttp import web + import app.p1 as p1mod + + done = asyncio.Event() + + async def fake_ha(request): + ws = web.WebSocketResponse() + await ws.prepare(request) + await ws.send_json({"type": "auth_required", "ha_version": "2026.8"}) + auth = await ws.receive_json() + assert auth["type"] == "auth" and auth["access_token"] == "tok" + await ws.send_json({"type": "auth_ok"}) + sub = await ws.receive_json() + assert sub["type"] == "subscribe_events" + assert sub["event_type"] == "state_changed" + await ws.send_json({"id": sub["id"], "type": "result", "success": True}) + get = await ws.receive_json() + assert get["type"] == "get_states" + await ws.send_json({"id": get["id"], "type": "result", "success": True, "result": [ + {"entity_id": "sensor.p1_import", "state": "1000"}, + {"entity_id": "sensor.p1_export", "state": "0"}, + {"entity_id": "sensor.something_else", "state": "hello"}, + ]}) + # One telegram, two state_changed events - exactly how HA emits it. + for imp, exp in ((1500, 0), (0, 800)): + await asyncio.sleep(0.5) + for eid, val in (("sensor.p1_import", imp), ("sensor.p1_export", exp)): + await ws.send_json({"type": "event", "event": {"data": { + "entity_id": eid, + "new_state": {"entity_id": eid, "state": str(val)}}}}) + await asyncio.sleep(0.5) + await ws.send_json({"type": "event", "event": {"data": { + "entity_id": "sensor.p1_export", + "new_state": {"entity_id": "sensor.p1_export", "state": "unavailable"}}}}) + await asyncio.sleep(0.5) + done.set() + return ws + + srv = web.Application() + srv.router.add_get("/ws", fake_ha) + runner = web.AppRunner(srv) + await runner.setup() + site = web.TCPSite(runner, "127.0.0.1", 0) + await site.start() + port = site._server.sockets[0].getsockname()[1] + p1mod.WS_URL = f"http://127.0.0.1:{port}/ws" + + ing = P1Ingest(phases=1, max_age_s=30.0) + async with aiohttp.ClientSession() as sess: + src = HaDsmrSource(sess, ing, dict(ENT), token="tok") + task = asyncio.get_running_loop().create_task(src.run()) + try: + await asyncio.wait_for(done.wait(), 20) + await asyncio.sleep(0.5) + finally: + task.cancel() + try: + await task + except asyncio.CancelledError: + pass + await runner.cleanup() + return ing, src + + +live, wire = asyncio.run(_e2e()) +check("the websocket handshake and subscription complete", live.samples >= 1) +# Six state_changed events arrived (two per telegram). The debounce is what +# makes that three consistent samples instead of six half-updated ones. +check("three telegrams produce three samples, not six", live.samples == 3) +check("the final export-dominant telegram nets negative", + live.last.net_w == -800.0) +check("the sample was built over the wire, tagged with its transport", + live.last.source == SOURCE_HA) +check("an entity we did not subscribe to is never cached", + "sensor.something_else" not in wire.cache and len(wire.cache) == 1) +check("a mid-stream unavailable is a parse error, not a sample", + live.parse_errors == 1 and live.samples == 3) +check("the last good reading survives the unavailable", live.net_w == -800.0) +check("the averager integrated the live stream", live.averager.elapsed_s > 0.5) + +print() +if fails: + print(f"{len(fails)} of {total} FAILED: {', '.join(fails)}") + sys.exit(1) +print(f"{total} checks") +print("all checks passed")