T-1, and it was a fleet-wide trip to zero. publish() emitted p1_age unconditionally, and published_age_s counts from P1Ingest.__init__ when no sample has ever arrived. With meter_source defaulting to off, every existing install would have published sensor.p1_sample_age_s climbing without bound; the ESP32 does `has_state() && state >= max_age_s` and forces the layer-1 failsafe, so each of them would have pinned its inverter at 0 W within 30 s. Exactly the opposite of the zero-regression the off default was for. The key is now omitted from the payload AND from MQTT discovery when P1 is off, so the entity does not exist at all - which is the status quo, and what has_state() is testing for. The predicate is one function, is_enabled(), because the grid reading, the task start and the discovery announcement have to agree or this comes back. T-2, connect no longer manufactures a sample. get_states returns whatever HA currently holds, which after a Core restart is a RestoreEntity value of unknown age; stamping it with ingest_ts=now reset the age and reported a fresh meter that could have been dead for an hour. run()'s own docstring already said a reconnect must emit nothing - the code disagreed with it, and a test asserted the violation. The cache is still primed, so the first real state_changed builds a complete sample; the age just stays honest until one arrives. T-3, gaps are no longer filled with the last held value. The averager held a sample forward across any interval, so a meter dying at 5 kW and returning ten minutes later credited 5 kW x 600 s to the capacity-tariff accumulator - a fabricated peak on a permanent record. The hold is capped at max_age_s: past that the stretch is walked so block boundaries still land correctly, but nothing accumulates and elapsed does not grow, which is what finally makes the comment about a gap dragging the billed average down true. Same threshold for control and billing: a reading too old to steer by is too old to bill by. T-4, the out-of-order/duplicate guard is covered. It was untested, and the reason is worth recording: the obvious assertion passes without the guard, because the negative interval is separately refused by the covered > 0 test. What the guard prevents is the timestamp REWIND, which only shows up one sample later as a re-integrated window. The test now goes one sample later. T-6, DOCS was wrong about latency. meter_max_age_s and stale_input_s stack, so meter death to 0 W is 45 s and not 30. Documented as a table with both clocks. Also documented the T-5 asymmetry rather than papering over it: the age measures arrival, not change, so a stuck MQTT bridge republishing its last telegram still looks fresh. Correct on ha_dsmr, not detectable on mqtt_p1 without a change-detector. Written up as a known limit. Writing the T-1 test caught a second defect in the test itself: it recorded only MQTT topics, and object_id lives in the payload, so "the age sensor is not announced" had been passing for the wrong reason. test_p1.py: 99 -> 122 checks. 14 mutations run, all 14 red, files restored byte-identical - including one per fix above. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Du77usMj8XNKNFZGmUiWDa
645 lines
29 KiB
Python
645 lines
29 KiB
Python
"""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)
|
|
# ⚠️ The assertion above is NOT sufficient on its own, and that is the whole
|
|
# lesson: deleting the guard still passes it, because the negative interval is
|
|
# separately refused by the `covered > 0` test. What the guard actually prevents
|
|
# is the REWIND - without it the held timestamp moves back to +5 s and the next
|
|
# telegram re-integrates the 5..10 s window that was already counted. The damage
|
|
# only becomes visible one sample later, so the test has to go one sample later.
|
|
a.add(sample(1000.0, at=BASE + timedelta(seconds=20)))
|
|
check("...and the held timestamp is not rewound, so the next telegram "
|
|
"cannot double-count", a.elapsed_s == 20.0 and a.partial_ws == 20000.0)
|
|
|
|
# A duplicate telegram (identical timestamp) is the same rule.
|
|
a = QuarterAverager(1)
|
|
a.add(sample(1000.0, at=BASE))
|
|
a.add(sample(1000.0, at=BASE + timedelta(seconds=10)))
|
|
a.add(sample(4000.0, at=BASE + timedelta(seconds=10)))
|
|
a.add(sample(1000.0, at=BASE + timedelta(seconds=20)))
|
|
check("a duplicate timestamp neither re-integrates nor replaces the held value",
|
|
a.elapsed_s == 20.0 and a.partial_ws == 20000.0)
|
|
|
|
# A gap must not be filled with the last held value. The meter dies at 5 kW and
|
|
# returns ten minutes later; hold-forward would credit 5 kW x 600 s to the
|
|
# capacity-tariff accumulator - a fabricated peak, on a permanent record, from
|
|
# data nobody measured.
|
|
a = QuarterAverager(1, max_hold_s=30.0)
|
|
a.add(sample(5000.0, at=BASE))
|
|
a.add(sample(5000.0, at=BASE + timedelta(seconds=600)))
|
|
check("a 600 s gap is held for at most max_hold_s, not for the whole gap",
|
|
a.partial_ws == 5000.0 * 30.0)
|
|
check("the unobserved stretch does not count as elapsed time", a.elapsed_s == 30.0)
|
|
closed = a.add(sample(5000.0, at=BASE + timedelta(seconds=900)))
|
|
check("the outage drags the billed quarter down instead of inventing a peak",
|
|
len(closed) == 1 and abs(closed[0].offtake_avg_w - 300000.0 / 900.0) < 1e-9)
|
|
check("...nowhere near the 5000 W a hold-forward would have billed",
|
|
closed[0].offtake_avg_w < 400.0)
|
|
|
|
# The cap must not disturb a normally-spaced stream.
|
|
a = QuarterAverager(1, max_hold_s=30.0)
|
|
for i in range(0, 121, 5): # a healthy 5 s telegram cadence
|
|
a.add(sample(2000.0, at=BASE + timedelta(seconds=i)))
|
|
check("a healthy 5 s cadence is untouched by the hold cap",
|
|
a.elapsed_s == 120.0 and abs(a.offtake_avg_w - 2000.0) < 1e-9)
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
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)
|
|
# ⚠️ TWO, not three. get_states primes the cache but must NOT build a sample:
|
|
# HA returns whatever it currently holds, which after a Core restart is a
|
|
# RestoreEntity value of unknown age, and stamping that with ingest_ts=now
|
|
# resets the age and reports a fresh meter that may have been dead for an hour.
|
|
# Only the two real state_changed telegrams become samples. Four state_changed
|
|
# events arrived (two per telegram); the debounce is what makes those two
|
|
# consistent samples rather than four half-updated ones.
|
|
check("connecting does not manufacture a sample from cached HA state",
|
|
live.samples == 2)
|
|
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 == 2)
|
|
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.2)
|
|
|
|
# The reason get_states still matters: it is what lets the FIRST real telegram
|
|
# build a complete sample instead of waiting for every entity to change once.
|
|
check("the primed cache let the first telegram build immediately",
|
|
live.samples == 2 and live.last.import_w == 0.0)
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
print("the age sensor must not exist when P1 is off")
|
|
# ⚠️ This is a fleet-wide regression guard, not a nicety. The ESP32 watchdog
|
|
# does `id(p1_age_s).has_state() && id(p1_age_s).state >= max_age_s` and forces
|
|
# the layer-1 failsafe. published_age_s counts from P1Ingest.__init__, so if the
|
|
# age were published with meter_source off it would climb past 30 s on every
|
|
# existing install within half a minute and pin the inverter at 0 W forever.
|
|
|
|
from app.p1 import is_enabled # noqa: E402
|
|
from app.mqtt import SENSORS, MqttPublisher # noqa: E402
|
|
|
|
check("meter_source off is disabled", is_enabled({"meter_source": "off"}) is False)
|
|
check("a missing meter_source is disabled", is_enabled({}) is False)
|
|
check("an empty meter_source is disabled", is_enabled({"meter_source": ""}) is False)
|
|
check("ha_dsmr is enabled", is_enabled({"meter_source": SOURCE_HA}) is True)
|
|
check("mqtt_p1 is enabled", is_enabled({"meter_source": SOURCE_MQTT}) is True)
|
|
|
|
# The entity id SAFETY-01's firmware subscribes to, pinned by object_id.
|
|
row = [s for s in SENSORS if s[0] == "p1_age"]
|
|
check("the age sensor is declared exactly once", len(row) == 1)
|
|
check("its object_id pins entity_id to sensor.p1_sample_age_s",
|
|
row[0][1] == "p1_sample_age_s")
|
|
check("it is published in seconds", row[0][3] == "s")
|
|
|
|
|
|
class _RecordingClient:
|
|
def __init__(self):
|
|
self.sent = []
|
|
|
|
def publish(self, topic, payload=None, retain=False):
|
|
# Topic AND payload: object_id, the thing that actually pins the entity
|
|
# id, only appears in the discovery payload. Recording topics alone made
|
|
# the "is not announced" check pass for the wrong reason.
|
|
self.sent.append(f"{topic} {payload}")
|
|
|
|
|
|
def _announced(omit):
|
|
pub = MqttPublisher(None, 1883, omit=omit) # host None -> never connects
|
|
pub.client = _RecordingClient()
|
|
pub._announce()
|
|
return " ".join(pub.client.sent)
|
|
|
|
|
|
check("with P1 off the age sensor is never announced",
|
|
"p1_sample_age_s" not in _announced(("p1_age",)))
|
|
check("the other status entities are still announced with P1 off",
|
|
"goodwe_grid_power" in _announced(("p1_age",)))
|
|
check("with P1 on the age sensor IS announced",
|
|
"p1_sample_age_s" in _announced(()))
|
|
|
|
# And the publish dict itself, through the real Controller.
|
|
from app.main import Controller # noqa: E402
|
|
|
|
|
|
class _Store:
|
|
data = {}
|
|
|
|
def set(self, *a):
|
|
pass
|
|
|
|
def get_time(self, *a):
|
|
return None
|
|
|
|
|
|
class _Pub:
|
|
def __init__(self):
|
|
self.last = {}
|
|
|
|
def publish(self, values):
|
|
self.last = values
|
|
|
|
def close(self):
|
|
pass
|
|
|
|
|
|
pub_off = _Pub()
|
|
Controller({"meter_source": "off"}, None, _Store(), pub_off).publish()
|
|
check("with P1 off, p1_age is absent from the published payload",
|
|
"p1_age" not in pub_off.last)
|
|
check("...while the normal status keys are still published",
|
|
"setpoint" in pub_off.last and "grid" in pub_off.last)
|
|
|
|
pub_on = _Pub()
|
|
Controller({"meter_source": SOURCE_HA}, None, _Store(), pub_on).publish()
|
|
check("with P1 on, p1_age is published", "p1_age" in pub_on.last)
|
|
check("...as a number, so has_state() becomes true only once we feed it",
|
|
isinstance(pub_on.last["p1_age"], float))
|
|
|
|
print()
|
|
if fails:
|
|
print(f"{len(fails)} of {total} FAILED: {', '.join(fails)}")
|
|
sys.exit(1)
|
|
print(f"{total} checks")
|
|
print("all checks passed")
|