2026-08-23 02:56:12 +02:00
|
|
|
"""Contract tests for M8 source/channel/binding management routes."""
|
|
|
|
|
|
|
|
|
|
from __future__ import annotations
|
|
|
|
|
|
2026-08-23 07:09:53 +02:00
|
|
|
from datetime import UTC, datetime, timedelta
|
|
|
|
|
from decimal import Decimal
|
|
|
|
|
from queue import Empty, Queue
|
|
|
|
|
from types import SimpleNamespace
|
|
|
|
|
import threading
|
|
|
|
|
import time
|
2026-08-23 02:56:12 +02:00
|
|
|
from unittest.mock import patch
|
|
|
|
|
|
|
|
|
|
from fastapi.testclient import TestClient
|
|
|
|
|
from sqlalchemy import create_engine, select
|
|
|
|
|
from sqlalchemy.orm import Session
|
|
|
|
|
|
2026-08-24 01:43:41 +02:00
|
|
|
from app.models.config import AppConfigEntry
|
2026-08-23 02:56:12 +02:00
|
|
|
from app.models.energy import DsmrReading, Meter
|
2026-08-23 07:09:53 +02:00
|
|
|
from app.models.meter_source import MeterSource, MeterSourceBinding, MeterSourceChannel, WarmteLinkReading
|
2026-08-23 02:56:12 +02:00
|
|
|
|
|
|
|
|
_CSRF = "test-csrf-token"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _login(client: TestClient) -> None:
|
|
|
|
|
assert client.post("/api/auth/login", json={"username": "admin", "password": "test-password"}).status_code == 200
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _client(auth_database):
|
|
|
|
|
from app.main import create_app
|
|
|
|
|
|
|
|
|
|
engine = create_engine(auth_database["app_url"], connect_args={"check_same_thread": False})
|
|
|
|
|
return TestClient(create_app()), engine
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _create_source(client: TestClient, *, config: dict | None = None) -> dict:
|
|
|
|
|
response = client.post("/api/energy/sources", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"name": "Synthetic DSMR", "kind": "dsmr_mqtt", "config": config or {},
|
|
|
|
|
})
|
|
|
|
|
assert response.status_code == 201
|
|
|
|
|
return response.json()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _add_channel(engine, source_uuid: str, *, key: str = "electricity") -> str:
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
source = session.execute(select(MeterSource).where(MeterSource.uuid == source_uuid)).scalar_one()
|
|
|
|
|
now = datetime.now(UTC)
|
|
|
|
|
channel = MeterSourceChannel(
|
|
|
|
|
source_id=source.id, channel_key=key, label="Electricity", unit="kWh",
|
|
|
|
|
created_at=now, updated_at=now,
|
|
|
|
|
)
|
|
|
|
|
session.add(channel)
|
|
|
|
|
session.commit()
|
|
|
|
|
return channel.uuid
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_source_profiles_and_crud_mask_secrets(auth_database):
|
|
|
|
|
client, engine = _client(auth_database)
|
|
|
|
|
with client:
|
|
|
|
|
assert client.get("/api/energy/source-profiles").status_code == 401
|
|
|
|
|
_login(client)
|
|
|
|
|
profiles = client.get("/api/energy/source-profiles")
|
|
|
|
|
assert profiles.status_code == 200
|
|
|
|
|
assert {item["kind"] for item in profiles.json()["items"]} == {"dsmr_mqtt", "warmtelink_serial"}
|
|
|
|
|
|
|
|
|
|
create = client.post("/api/energy/sources", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"name": "Synthetic DSMR", "kind": "dsmr_mqtt",
|
|
|
|
|
"config": {"username": "user", "password": "not-for-api"},
|
|
|
|
|
})
|
|
|
|
|
assert create.status_code == 201
|
|
|
|
|
source = create.json()
|
|
|
|
|
assert source["config"]["username"] == ""
|
|
|
|
|
assert source["config"]["password"] == ""
|
|
|
|
|
assert "not-for-api" not in str(source)
|
|
|
|
|
|
|
|
|
|
patch = client.patch(f"/api/energy/sources/{source['uuid']}", headers={"X-CSRF-Token": _CSRF}, json={"config": {"password": ""}})
|
|
|
|
|
assert patch.status_code == 200
|
|
|
|
|
assert patch.json()["config"]["password"] == ""
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
stored = session.execute(
|
|
|
|
|
select(MeterSource).where(MeterSource.uuid == source["uuid"])
|
|
|
|
|
).scalar_one()
|
|
|
|
|
assert stored.config["password"] == "not-for-api"
|
|
|
|
|
assert client.post(f"/api/energy/sources/{source['uuid']}/discover", headers={"X-CSRF-Token": _CSRF}).status_code == 200
|
|
|
|
|
assert client.delete(f"/api/energy/sources/{source['uuid']}", headers={"X-CSRF-Token": _CSRF}).status_code == 204
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
2026-08-23 06:33:18 +02:00
|
|
|
def test_serial_source_reconciles_only_after_commit(auth_database, monkeypatch):
|
|
|
|
|
"""Creation and mutation trigger the worker manager after a successful commit."""
|
|
|
|
|
from app.api.routes.api import meter_sources
|
|
|
|
|
|
|
|
|
|
calls = []
|
|
|
|
|
monkeypatch.setattr(meter_sources.warmtelink_worker_manager, "reconcile", lambda: calls.append(True))
|
|
|
|
|
client, engine = _client(auth_database)
|
|
|
|
|
with client:
|
|
|
|
|
_login(client)
|
|
|
|
|
baseline = len(calls)
|
|
|
|
|
response = client.post("/api/energy/sources", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"name": "Serial", "kind": "warmtelink_serial", "config": {"path": "/dev/ttyUSB0"},
|
|
|
|
|
})
|
|
|
|
|
assert response.status_code == 201
|
|
|
|
|
source_uuid = response.json()["uuid"]
|
|
|
|
|
assert len(calls) == baseline + 1
|
|
|
|
|
assert client.patch(f"/api/energy/sources/{source_uuid}", headers={"X-CSRF-Token": _CSRF}, json={"enabled": False}).status_code == 200
|
|
|
|
|
assert len(calls) == baseline + 2
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_serial_source_commit_is_not_reported_as_failed_when_reconcile_raises(auth_database, monkeypatch):
|
|
|
|
|
from app.api.routes.api import meter_sources
|
|
|
|
|
|
|
|
|
|
client, engine = _client(auth_database)
|
|
|
|
|
with client:
|
|
|
|
|
_login(client)
|
|
|
|
|
monkeypatch.setattr(
|
|
|
|
|
meter_sources.warmtelink_worker_manager, "reconcile",
|
|
|
|
|
lambda: (_ for _ in ()).throw(RuntimeError()),
|
|
|
|
|
)
|
|
|
|
|
response = client.post("/api/energy/sources", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"name": "Serial", "kind": "warmtelink_serial", "config": {"path": "/dev/ttyUSB0"},
|
|
|
|
|
})
|
|
|
|
|
assert response.status_code == 201
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
2026-08-24 01:43:41 +02:00
|
|
|
def test_dsmr_crud_reconciles_runtime_and_invalidates_retained_callbacks(auth_database, monkeypatch):
|
|
|
|
|
"""CRUD converges real DSMR ownership, not merely a mocked reconcile call."""
|
|
|
|
|
from app.api.routes.api import meter_sources
|
|
|
|
|
from app.services import dsmr_ingest
|
|
|
|
|
|
|
|
|
|
class FakeMqttManager:
|
|
|
|
|
def __init__(self) -> None:
|
|
|
|
|
self.handlers: dict[int, dict[str, object]] = {}
|
|
|
|
|
self.replace_calls: list[tuple[int, dict[str, object]]] = []
|
|
|
|
|
self.remove_calls: list[int] = []
|
|
|
|
|
|
|
|
|
|
def replace_source(self, source_id: int, **kwargs: object) -> bool:
|
|
|
|
|
self.replace_calls.append((source_id, kwargs))
|
|
|
|
|
self.handlers[source_id] = kwargs["subscriptions"] # type: ignore[assignment]
|
|
|
|
|
kwargs["state_handler"]("connecting") # type: ignore[operator]
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
def remove_source(self, source_id: int) -> None:
|
|
|
|
|
self.remove_calls.append(source_id)
|
|
|
|
|
self.handlers.pop(source_id, None)
|
|
|
|
|
|
|
|
|
|
def source_is_active(self, source_id: int) -> bool:
|
|
|
|
|
return source_id in self.handlers
|
|
|
|
|
|
|
|
|
|
mqtt = FakeMqttManager()
|
|
|
|
|
warmtelink_calls: list[None] = []
|
|
|
|
|
monkeypatch.setattr("app.integrations.mqtt.mqtt_manager", mqtt)
|
|
|
|
|
monkeypatch.setattr(
|
|
|
|
|
meter_sources.warmtelink_worker_manager, "reconcile", lambda: warmtelink_calls.append(None),
|
|
|
|
|
)
|
|
|
|
|
monkeypatch.setattr(dsmr_ingest, "_subscriptions", {})
|
|
|
|
|
monkeypatch.setattr(dsmr_ingest, "_subscription_client_ids", {})
|
|
|
|
|
monkeypatch.setattr(dsmr_ingest, "_subscription_tokens", {})
|
|
|
|
|
monkeypatch.setattr(dsmr_ingest, "_tariffs", {})
|
|
|
|
|
|
|
|
|
|
client, engine = _client(auth_database)
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
session.add(AppConfigEntry(
|
|
|
|
|
key="MQTT_CLIENT_ID", value="home-automation-api-test", updated_at=datetime.now(UTC),
|
|
|
|
|
))
|
|
|
|
|
session.commit()
|
|
|
|
|
payload = b'{"timestamp":"2030-01-01T00:00:00Z"}'
|
|
|
|
|
with client:
|
|
|
|
|
_login(client)
|
|
|
|
|
baseline_warmtelink_calls = len(warmtelink_calls)
|
|
|
|
|
created = client.post("/api/energy/sources", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"name": "Managed DSMR", "kind": "dsmr_mqtt",
|
|
|
|
|
"config": {"broker_host": "broker.test", "topic": "meter/first"},
|
|
|
|
|
})
|
|
|
|
|
assert created.status_code == 201
|
|
|
|
|
assert created.json()["status"] == "connecting"
|
|
|
|
|
source_uuid = created.json()["uuid"]
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
source_id = session.scalar(select(MeterSource.id).where(MeterSource.uuid == source_uuid))
|
|
|
|
|
assert source_id is not None
|
|
|
|
|
first_handler = mqtt.handlers[source_id]["meter/first"]
|
|
|
|
|
assert mqtt.replace_calls[-1][1]["base_client_id"] == "home-automation-api-test"
|
|
|
|
|
|
|
|
|
|
updated = client.patch(
|
|
|
|
|
f"/api/energy/sources/{source_uuid}", headers={"X-CSRF-Token": _CSRF},
|
|
|
|
|
json={"config": {"topic": "meter/second"}},
|
|
|
|
|
)
|
|
|
|
|
assert updated.status_code == 200
|
|
|
|
|
assert updated.json()["status"] == "connecting"
|
|
|
|
|
assert mqtt.remove_calls == [source_id]
|
|
|
|
|
assert "meter/second" in mqtt.handlers[source_id]
|
|
|
|
|
first_handler(payload) # type: ignore[operator]
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
assert session.query(DsmrReading).filter_by(meter_source_id=source_id).count() == 0
|
|
|
|
|
|
|
|
|
|
retained_handler = mqtt.handlers[source_id]["meter/second"]
|
|
|
|
|
disabled = client.patch(
|
|
|
|
|
f"/api/energy/sources/{source_uuid}", headers={"X-CSRF-Token": _CSRF}, json={"enabled": False},
|
|
|
|
|
)
|
|
|
|
|
assert disabled.status_code == 200
|
|
|
|
|
assert disabled.json()["status"] == "unknown"
|
|
|
|
|
assert source_id not in mqtt.handlers
|
|
|
|
|
retained_handler(payload) # type: ignore[operator]
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
assert session.query(DsmrReading).filter_by(meter_source_id=source_id).count() == 0
|
|
|
|
|
|
|
|
|
|
assert client.patch(
|
|
|
|
|
f"/api/energy/sources/{source_uuid}", headers={"X-CSRF-Token": _CSRF}, json={"enabled": True},
|
|
|
|
|
).status_code == 200
|
|
|
|
|
delete_handler = mqtt.handlers[source_id]["meter/second"]
|
|
|
|
|
assert client.delete(f"/api/energy/sources/{source_uuid}", headers={"X-CSRF-Token": _CSRF}).status_code == 204
|
|
|
|
|
assert source_id not in mqtt.handlers
|
|
|
|
|
delete_handler(payload) # type: ignore[operator]
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
assert session.get(MeterSource, source_id) is None
|
|
|
|
|
assert session.query(DsmrReading).filter_by(meter_source_id=source_id).count() == 0
|
|
|
|
|
assert len(warmtelink_calls) == baseline_warmtelink_calls + 5
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_source_crud_runtime_failures_do_not_hide_commits(auth_database, monkeypatch):
|
|
|
|
|
from app.api.routes.api import meter_sources
|
|
|
|
|
|
|
|
|
|
calls: list[str] = []
|
|
|
|
|
client, engine = _client(auth_database)
|
|
|
|
|
with client:
|
|
|
|
|
_login(client)
|
|
|
|
|
monkeypatch.setattr(
|
|
|
|
|
meter_sources.warmtelink_worker_manager, "reconcile",
|
|
|
|
|
lambda: calls.append("warmtelink") or (_ for _ in ()).throw(RuntimeError()),
|
|
|
|
|
)
|
|
|
|
|
monkeypatch.setattr(
|
|
|
|
|
meter_sources, "apply_dsmr_subscription",
|
|
|
|
|
lambda _settings: calls.append("dsmr") or (_ for _ in ()).throw(RuntimeError()),
|
|
|
|
|
)
|
|
|
|
|
created = client.post("/api/energy/sources", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"name": "Durable DSMR", "kind": "dsmr_mqtt", "config": {},
|
|
|
|
|
})
|
|
|
|
|
assert created.status_code == 201
|
|
|
|
|
source_uuid = created.json()["uuid"]
|
|
|
|
|
assert calls == ["warmtelink", "dsmr"]
|
|
|
|
|
assert client.patch(
|
|
|
|
|
f"/api/energy/sources/{source_uuid}", headers={"X-CSRF-Token": _CSRF}, json={"enabled": False},
|
|
|
|
|
).status_code == 200
|
|
|
|
|
assert client.delete(f"/api/energy/sources/{source_uuid}", headers={"X-CSRF-Token": _CSRF}).status_code == 204
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
assert session.scalar(select(MeterSource.id).where(MeterSource.uuid == source_uuid)) is None
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
2026-08-23 02:56:12 +02:00
|
|
|
def test_binding_routes_and_atomic_meter_declaration(auth_database):
|
|
|
|
|
client, engine = _client(auth_database)
|
|
|
|
|
with client:
|
|
|
|
|
_login(client)
|
|
|
|
|
source_response = client.post("/api/energy/sources", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"name": "Synthetic DSMR", "kind": "dsmr_mqtt", "config": {},
|
|
|
|
|
})
|
|
|
|
|
source_uuid = source_response.json()["uuid"]
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
source = session.query(MeterSource).filter_by(uuid=source_uuid).one()
|
|
|
|
|
now = datetime.now(UTC)
|
|
|
|
|
channel = MeterSourceChannel(
|
|
|
|
|
source_id=source.id, channel_key="electricity", label="Electricity", unit="kWh",
|
|
|
|
|
created_at=now, updated_at=now,
|
|
|
|
|
)
|
|
|
|
|
session.add(channel)
|
|
|
|
|
session.commit()
|
|
|
|
|
channel_uuid = channel.uuid
|
|
|
|
|
|
|
|
|
|
declaration = {
|
|
|
|
|
"label": "Bound meter", "started_at": "2030-01-01T00:00:00Z", "reason": "initial",
|
|
|
|
|
"commodity": "electricity", "source_channel_uuid": channel_uuid,
|
|
|
|
|
}
|
|
|
|
|
created = client.post("/api/energy/meters", headers={"X-CSRF-Token": _CSRF}, json=declaration)
|
|
|
|
|
assert created.status_code == 201
|
|
|
|
|
assert created.json()["bindings"][0]["source_channel_uuid"] == channel_uuid
|
|
|
|
|
meter_id = created.json()["id"]
|
|
|
|
|
assert client.get(f"/api/energy/meters/{meter_id}/bindings").json()["total"] == 1
|
|
|
|
|
assert client.get(f"/api/energy/sources/{source_uuid}/channels").json()["items"][0]["binding_count"] == 1
|
|
|
|
|
|
|
|
|
|
invalid = dict(declaration, label="Must roll back", started_at="2030-02-01T00:00:00Z", source_channel_uuid="missing-channel")
|
|
|
|
|
assert client.post("/api/energy/meters", headers={"X-CSRF-Token": _CSRF}, json=invalid).status_code == 404
|
|
|
|
|
assert client.get("/api/energy/meters").json()["total"] == 1
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
original = session.get(Meter, meter_id)
|
|
|
|
|
assert original is not None
|
|
|
|
|
assert original.ended_at is None
|
|
|
|
|
assert session.execute(select(Meter).where(Meter.label == "Must roll back")).scalar_one_or_none() is None
|
|
|
|
|
assert session.execute(select(MeterSourceBinding).where(MeterSourceBinding.meter_id == meter_id)).scalars().all()
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_source_and_binding_error_contracts_csrf_timezone_and_dsmr_compatibility(auth_database, monkeypatch):
|
|
|
|
|
"""Exercise the public error boundary without opening serial or MQTT I/O."""
|
|
|
|
|
from zoneinfo import ZoneInfo
|
|
|
|
|
|
|
|
|
|
from app.services import timezone as timezone_service
|
|
|
|
|
|
|
|
|
|
monkeypatch.setattr(timezone_service, "local_tz", lambda: ZoneInfo("Europe/Amsterdam"))
|
|
|
|
|
client, engine = _client(auth_database)
|
|
|
|
|
with client:
|
|
|
|
|
# All management reads require a session and mutations require CSRF.
|
|
|
|
|
assert client.get("/api/energy/commodities").status_code == 401
|
|
|
|
|
_login(client)
|
|
|
|
|
assert client.post("/api/energy/sources", json={
|
|
|
|
|
"name": "No CSRF", "kind": "dsmr_mqtt", "config": {},
|
|
|
|
|
}).status_code == 403
|
|
|
|
|
assert client.post("/api/energy/sources", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"name": "Invalid", "kind": "dsmr_mqtt", "config": {"unexpected": True},
|
|
|
|
|
}).status_code == 422
|
|
|
|
|
assert client.get("/api/energy/sources/does-not-exist").status_code == 404
|
|
|
|
|
|
|
|
|
|
source = client.post("/api/energy/sources", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"name": "DSMR source", "kind": "dsmr_mqtt", "config": {"password": "stored-secret"},
|
|
|
|
|
})
|
|
|
|
|
assert source.status_code == 201
|
|
|
|
|
source_uuid = source.json()["uuid"]
|
|
|
|
|
assert "stored-secret" not in client.get(f"/api/energy/sources/{source_uuid}").text
|
|
|
|
|
assert client.get("/api/energy/commodities").json()["items"] == [
|
|
|
|
|
{"key": "electricity", "unit": "kWh", "capabilities": ["meter", "binding", "cost"]},
|
|
|
|
|
{"key": "heating", "unit": "GJ", "capabilities": ["meter", "binding"]},
|
|
|
|
|
{"key": "hot_water", "unit": "m³", "capabilities": ["meter", "binding"]},
|
|
|
|
|
]
|
|
|
|
|
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
source_model = session.query(MeterSource).filter_by(uuid=source_uuid).one()
|
|
|
|
|
now = datetime.now(UTC)
|
|
|
|
|
channel = MeterSourceChannel(
|
|
|
|
|
source_id=source_model.id, channel_key="electricity", label="Electricity", unit="kWh",
|
|
|
|
|
created_at=now, updated_at=now,
|
|
|
|
|
)
|
|
|
|
|
session.add(channel)
|
|
|
|
|
session.add(DsmrReading(
|
|
|
|
|
meter_source_id=source_model.id, telegram_id=7, recorded_at=now,
|
|
|
|
|
payload={"compatibility": "latest"},
|
|
|
|
|
))
|
|
|
|
|
session.commit()
|
|
|
|
|
channel_uuid = channel.uuid
|
|
|
|
|
|
|
|
|
|
# Retained channels prohibit deletion; there is no cascade escape hatch.
|
|
|
|
|
assert client.delete(f"/api/energy/sources/{source_uuid}", headers={"X-CSRF-Token": _CSRF}).status_code == 409
|
|
|
|
|
assert client.get(f"/api/energy/sources/{source_uuid}/channels/not-a-channel/readings").status_code == 404
|
|
|
|
|
readings = client.get(f"/api/energy/sources/{source_uuid}/channels/{channel_uuid}/readings")
|
|
|
|
|
assert readings.status_code == 200
|
|
|
|
|
assert readings.json()["total"] == 1
|
2026-08-23 07:09:53 +02:00
|
|
|
assert readings.json()["items"] == [{"recorded_at": now.isoformat().replace("+00:00", ""), "value": None, "quality": None}]
|
2026-08-23 02:56:12 +02:00
|
|
|
assert "telegram_id" not in readings.text
|
|
|
|
|
latest = client.get("/api/energy/dsmr/latest")
|
|
|
|
|
assert latest.status_code == 200
|
|
|
|
|
assert latest.json()["payload"] == {"compatibility": "latest"}
|
|
|
|
|
|
|
|
|
|
# This declaration is deliberately retroactive for timezone coverage;
|
|
|
|
|
# mock the billing sweep so the API contract test stays bounded.
|
|
|
|
|
with patch("app.api.routes.api.meters.recompute_range", return_value=0):
|
|
|
|
|
meter = client.post("/api/energy/meters", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"label": "Local-time meter", "started_at": "2025-01-02T00:00:00", "reason": "initial",
|
|
|
|
|
})
|
|
|
|
|
assert meter.status_code == 201
|
|
|
|
|
meter_id = meter.json()["id"]
|
|
|
|
|
binding = client.post(f"/api/energy/meters/{meter_id}/bindings", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"source_channel_uuid": channel_uuid, "started_at": "2025-01-02T00:00:00",
|
|
|
|
|
})
|
|
|
|
|
assert binding.status_code == 201
|
|
|
|
|
binding_body = binding.json()
|
|
|
|
|
localized_start = datetime.fromisoformat(binding_body["started_at"].replace("Z", "+00:00"))
|
|
|
|
|
# SQLite returns UTC columns without tzinfo; retain the UTC clock instant
|
|
|
|
|
# regardless of that transport detail.
|
|
|
|
|
assert localized_start.replace(tzinfo=None) == datetime(2025, 1, 1, 23, 0, 0)
|
|
|
|
|
assert client.post(f"/api/energy/meters/{meter_id}/bindings", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"source_channel_uuid": channel_uuid, "started_at": "2025-01-02T00:00:00",
|
|
|
|
|
}).status_code == 422
|
|
|
|
|
assert client.patch(f"/api/energy/bindings/{binding_body['uuid']}", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"ended_at": "2025-01-02T00:00:00",
|
|
|
|
|
}).status_code == 422
|
|
|
|
|
assert client.patch(f"/api/energy/bindings/{binding_body['uuid']}", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"ended_at": "2025-01-03T00:00:00",
|
|
|
|
|
}).status_code == 200
|
|
|
|
|
assert client.patch("/api/energy/bindings/not-a-binding", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"ended_at": "2025-01-03T00:00:00",
|
|
|
|
|
}).status_code == 404
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_management_reads_require_auth_and_mutations_require_csrf(auth_database):
|
|
|
|
|
"""Every M8-T06 management route enforces the session/CSRF contract."""
|
|
|
|
|
client, engine = _client(auth_database)
|
|
|
|
|
with client:
|
|
|
|
|
for path in (
|
|
|
|
|
"/api/energy/source-profiles",
|
|
|
|
|
"/api/energy/commodities",
|
|
|
|
|
"/api/energy/sources",
|
|
|
|
|
"/api/energy/sources/missing",
|
|
|
|
|
"/api/energy/sources/missing/channels",
|
|
|
|
|
"/api/energy/sources/missing/channels/missing/readings",
|
|
|
|
|
"/api/energy/meters/1/bindings",
|
|
|
|
|
):
|
|
|
|
|
assert client.get(path).status_code == 401
|
|
|
|
|
|
|
|
|
|
_login(client)
|
|
|
|
|
assert client.post("/api/energy/sources", json={
|
|
|
|
|
"name": "No CSRF", "kind": "dsmr_mqtt", "config": {},
|
|
|
|
|
}).status_code == 403
|
|
|
|
|
source = _create_source(client)
|
|
|
|
|
channel_uuid = _add_channel(engine, source["uuid"])
|
|
|
|
|
meter = client.post("/api/energy/meters", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"label": "CSRF meter", "started_at": "2030-01-01T00:00:00Z", "reason": "initial",
|
|
|
|
|
})
|
|
|
|
|
assert meter.status_code == 201
|
|
|
|
|
meter_id = meter.json()["id"]
|
|
|
|
|
binding = client.post(f"/api/energy/meters/{meter_id}/bindings", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"source_channel_uuid": channel_uuid, "started_at": "2030-01-01T00:00:00Z",
|
|
|
|
|
})
|
|
|
|
|
assert binding.status_code == 201
|
|
|
|
|
|
|
|
|
|
assert client.patch(f"/api/energy/sources/{source['uuid']}", json={"name": "blocked"}).status_code == 403
|
|
|
|
|
assert client.delete(f"/api/energy/sources/{source['uuid']}").status_code == 403
|
|
|
|
|
assert client.post(f"/api/energy/sources/{source['uuid']}/discover").status_code == 403
|
|
|
|
|
assert client.post(f"/api/energy/meters/{meter_id}/bindings", json={
|
|
|
|
|
"source_channel_uuid": channel_uuid, "started_at": "2031-01-01T00:00:00Z",
|
|
|
|
|
}).status_code == 403
|
|
|
|
|
assert client.patch(f"/api/energy/bindings/{binding.json()['uuid']}", json={}).status_code == 403
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_source_channel_binding_response_contract_and_discover_capabilities(auth_database):
|
|
|
|
|
client, engine = _client(auth_database)
|
|
|
|
|
with client:
|
|
|
|
|
_login(client)
|
|
|
|
|
source = _create_source(client, config={"username": "private-user", "password": "private-secret"})
|
|
|
|
|
channel_uuid = _add_channel(engine, source["uuid"])
|
2026-08-24 01:02:23 +02:00
|
|
|
with Session(engine) as session:
|
|
|
|
|
source_model = session.query(MeterSource).filter_by(uuid=source["uuid"]).one()
|
|
|
|
|
source_model.status = "online"
|
|
|
|
|
source_model.last_seen_at = datetime.now(UTC)
|
|
|
|
|
source_model.last_error = None
|
|
|
|
|
session.commit()
|
2026-08-23 02:56:12 +02:00
|
|
|
listed = client.get("/api/energy/sources")
|
|
|
|
|
assert listed.status_code == 200
|
|
|
|
|
assert listed.json()["total"] >= 1
|
|
|
|
|
source_item = next(item for item in listed.json()["items"] if item["uuid"] == source["uuid"])
|
|
|
|
|
assert source_item["uuid"] == source["uuid"]
|
|
|
|
|
detail = client.get(f"/api/energy/sources/{source['uuid']}")
|
|
|
|
|
assert detail.status_code == 200
|
|
|
|
|
assert detail.json()["uuid"] == source["uuid"]
|
|
|
|
|
for body in (listed.json(), detail.json()):
|
|
|
|
|
rendered = str(body)
|
|
|
|
|
assert "private-secret" not in rendered
|
|
|
|
|
assert "channel_key" not in rendered
|
|
|
|
|
assert "fingerprint" not in rendered
|
|
|
|
|
|
|
|
|
|
discovered = client.post(f"/api/energy/sources/{source['uuid']}/discover", headers={"X-CSRF-Token": _CSRF})
|
|
|
|
|
assert discovered.status_code == 200
|
|
|
|
|
assert discovered.json() == {
|
|
|
|
|
"requested": False, "supported": True, "status": "managed_by_runtime",
|
2026-08-23 07:09:53 +02:00
|
|
|
"request_id": None,
|
2026-08-23 02:56:12 +02:00
|
|
|
"detail": "This source is discovered by its runtime subscription; no connection was opened.",
|
2026-08-23 07:09:53 +02:00
|
|
|
"channels": [],
|
2026-08-23 02:56:12 +02:00
|
|
|
}
|
|
|
|
|
channels = client.get(f"/api/energy/sources/{source['uuid']}/channels")
|
|
|
|
|
assert channels.status_code == 200
|
|
|
|
|
channel = channels.json()["items"][0]
|
|
|
|
|
assert channel["uuid"] == channel_uuid
|
|
|
|
|
assert set(channel) == {
|
|
|
|
|
"uuid", "label", "suggested_commodity", "unit", "device_type", "latest_value",
|
2026-08-23 07:09:53 +02:00
|
|
|
"latest_at", "latest_quality", "binding_count", "bound_meter_ids", "binding_summary",
|
2026-08-23 02:56:12 +02:00
|
|
|
}
|
2026-08-24 01:02:23 +02:00
|
|
|
assert channels.json()["source_status"] == "online"
|
2026-08-23 02:56:12 +02:00
|
|
|
|
|
|
|
|
meter = client.post("/api/energy/meters", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"label": "Contract meter", "started_at": "2030-01-01T00:00:00Z", "reason": "initial",
|
|
|
|
|
})
|
|
|
|
|
assert meter.status_code == 201
|
|
|
|
|
binding = client.post(f"/api/energy/meters/{meter.json()['id']}/bindings", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"source_channel_uuid": channel_uuid, "started_at": "2030-01-01T00:00:00Z",
|
|
|
|
|
})
|
|
|
|
|
assert binding.status_code == 201
|
|
|
|
|
binding_item = client.get(f"/api/energy/meters/{meter.json()['id']}/bindings").json()["items"][0]
|
|
|
|
|
assert binding_item["uuid"] == binding.json()["uuid"]
|
|
|
|
|
assert binding_item["source_uuid"] == source["uuid"]
|
|
|
|
|
assert binding_item["source_channel_uuid"] == channel_uuid
|
|
|
|
|
assert set(binding_item) == {
|
|
|
|
|
"uuid", "meter_id", "source_channel_uuid", "source_uuid", "started_at", "ended_at",
|
|
|
|
|
"created_at", "updated_at",
|
|
|
|
|
}
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_binding_patch_omitted_null_and_adjacent_half_open_boundaries(auth_database):
|
|
|
|
|
client, engine = _client(auth_database)
|
|
|
|
|
with client:
|
|
|
|
|
_login(client)
|
|
|
|
|
source = _create_source(client)
|
|
|
|
|
channel_uuid = _add_channel(engine, source["uuid"])
|
|
|
|
|
meter = client.post("/api/energy/meters", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"label": "Timeline meter", "started_at": "2030-01-01T00:00:00Z", "reason": "initial",
|
|
|
|
|
})
|
|
|
|
|
assert meter.status_code == 201
|
|
|
|
|
meter_id = meter.json()["id"]
|
|
|
|
|
first = client.post(f"/api/energy/meters/{meter_id}/bindings", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"source_channel_uuid": channel_uuid, "started_at": "2030-01-01T00:00:00Z",
|
|
|
|
|
"ended_at": "2030-02-01T00:00:00Z",
|
|
|
|
|
})
|
|
|
|
|
assert first.status_code == 201
|
|
|
|
|
first_uuid = first.json()["uuid"]
|
|
|
|
|
corrected = client.patch(f"/api/energy/bindings/{first_uuid}", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"started_at": "2030-01-02T00:00:00Z",
|
|
|
|
|
})
|
|
|
|
|
assert corrected.status_code == 200
|
|
|
|
|
assert corrected.json()["ended_at"] == "2030-02-01T00:00:00"
|
|
|
|
|
unchanged = client.patch(f"/api/energy/bindings/{first_uuid}", headers={"X-CSRF-Token": _CSRF}, json={})
|
|
|
|
|
assert unchanged.status_code == 200
|
|
|
|
|
assert unchanged.json()["ended_at"] == "2030-02-01T00:00:00"
|
|
|
|
|
reopened = client.patch(f"/api/energy/bindings/{first_uuid}", headers={"X-CSRF-Token": _CSRF}, json={"ended_at": None})
|
|
|
|
|
assert reopened.status_code == 200
|
|
|
|
|
assert reopened.json()["ended_at"] is None
|
|
|
|
|
reclosed = client.patch(f"/api/energy/bindings/{first_uuid}", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"ended_at": "2030-02-01T00:00:00Z",
|
|
|
|
|
})
|
|
|
|
|
assert reclosed.status_code == 200
|
|
|
|
|
adjacent = client.post(f"/api/energy/meters/{meter_id}/bindings", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"source_channel_uuid": channel_uuid, "started_at": "2030-02-01T00:00:00Z",
|
|
|
|
|
})
|
|
|
|
|
assert adjacent.status_code == 201
|
|
|
|
|
engine.dispose()
|
2026-08-23 07:09:53 +02:00
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_warmtelink_discover_and_minute_history_are_bounded_and_private(auth_database, monkeypatch):
|
|
|
|
|
"""Discover delegates to the manager; readings expose accepted minute samples only."""
|
|
|
|
|
from app.api.routes.api import meter_sources
|
|
|
|
|
|
|
|
|
|
requested: list[int] = []
|
|
|
|
|
monkeypatch.setattr(
|
|
|
|
|
meter_sources.warmtelink_worker_manager, "request_discovery",
|
|
|
|
|
lambda source_id: requested.append(source_id) or SimpleNamespace(
|
|
|
|
|
status="completed", request_id=1, detail=None, completed=SimpleNamespace(is_set=lambda: False),
|
|
|
|
|
),
|
|
|
|
|
)
|
|
|
|
|
monkeypatch.setattr(meter_sources.warmtelink_worker_manager, "reconcile", lambda: None)
|
|
|
|
|
client, engine = _client(auth_database)
|
|
|
|
|
with client:
|
|
|
|
|
_login(client)
|
|
|
|
|
created = client.post("/api/energy/sources", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"name": "WarmteLink", "kind": "warmtelink_serial", "config": {"path": "/dev/fake"},
|
|
|
|
|
})
|
|
|
|
|
assert created.status_code == 201
|
|
|
|
|
source_uuid = created.json()["uuid"]
|
|
|
|
|
now = datetime(2030, 1, 1, 12, 0, 30, tzinfo=UTC)
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
source = session.execute(select(MeterSource).where(MeterSource.uuid == source_uuid)).scalar_one()
|
|
|
|
|
source.status = "online"
|
|
|
|
|
channel = MeterSourceChannel(
|
|
|
|
|
source_id=source.id, channel_key="heating", label="Heating", unit="GJ",
|
|
|
|
|
latest_value=Decimal("7.002"), latest_at=now, latest_quality="unverifiable",
|
|
|
|
|
created_at=now, updated_at=now,
|
|
|
|
|
)
|
|
|
|
|
session.add(channel)
|
|
|
|
|
session.flush()
|
|
|
|
|
session.add_all([
|
|
|
|
|
WarmteLinkReading(
|
|
|
|
|
channel_id=channel.id, recorded_at=now - timedelta(minutes=1), received_at=now,
|
|
|
|
|
value=Decimal("7.001"), unit="GJ", quality="unverifiable", equipment_fingerprint="masked",
|
|
|
|
|
),
|
|
|
|
|
WarmteLinkReading(
|
|
|
|
|
channel_id=channel.id, recorded_at=now, received_at=now,
|
|
|
|
|
value=Decimal("7.002"), unit="GJ", quality="unverifiable", equipment_fingerprint="masked",
|
|
|
|
|
),
|
|
|
|
|
])
|
|
|
|
|
session.commit()
|
|
|
|
|
channel_uuid = channel.uuid
|
|
|
|
|
|
|
|
|
|
discover = client.post(f"/api/energy/sources/{source_uuid}/discover", headers={"X-CSRF-Token": _CSRF})
|
|
|
|
|
assert discover.status_code == 200
|
|
|
|
|
assert discover.json()["status"] == "completed"
|
|
|
|
|
assert requested and "fingerprint" not in discover.text and "channel_key" not in discover.text
|
|
|
|
|
assert discover.json()["channels"][0]["latest_quality"] == "unverifiable"
|
|
|
|
|
|
|
|
|
|
history = client.get(
|
|
|
|
|
f"/api/energy/sources/{source_uuid}/channels/{channel_uuid}/readings",
|
|
|
|
|
params={"from": "2030-01-01T11:59:00Z", "to": "2030-01-01T12:01:00Z", "limit": 1},
|
|
|
|
|
)
|
|
|
|
|
assert history.status_code == 200
|
|
|
|
|
assert history.json()["total"] == 1
|
|
|
|
|
assert history.json()["items"] == [{
|
|
|
|
|
"recorded_at": "2030-01-01T11:59:30", "value": "7.001", "quality": "unverifiable",
|
|
|
|
|
}]
|
|
|
|
|
assert client.get(
|
|
|
|
|
f"/api/energy/sources/{source_uuid}/channels/{channel_uuid}/readings",
|
|
|
|
|
params={"from": "2030-01-01T12:01:00Z", "to": "2030-01-01T12:00:00Z"},
|
|
|
|
|
).status_code == 422
|
|
|
|
|
assert client.get(
|
|
|
|
|
f"/api/energy/sources/{source_uuid}/channels/{channel_uuid}/readings", params={"limit": 0}
|
|
|
|
|
).status_code == 422
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_warmtelink_discover_auth_csrf_disabled_and_source_ownership(auth_database, monkeypatch):
|
|
|
|
|
from app.api.routes.api import meter_sources
|
|
|
|
|
|
|
|
|
|
monkeypatch.setattr(
|
|
|
|
|
meter_sources.warmtelink_worker_manager, "request_discovery",
|
|
|
|
|
lambda _source_id: SimpleNamespace(
|
|
|
|
|
status="pending", request_id=1, detail=None, completed=SimpleNamespace(is_set=lambda: False),
|
|
|
|
|
),
|
|
|
|
|
)
|
|
|
|
|
monkeypatch.setattr(meter_sources.warmtelink_worker_manager, "reconcile", lambda: None)
|
|
|
|
|
client, engine = _client(auth_database)
|
|
|
|
|
with client:
|
|
|
|
|
serial = client.post("/api/energy/sources", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"name": "Serial", "kind": "warmtelink_serial", "enabled": False, "config": {"path": "/dev/fake"},
|
|
|
|
|
})
|
|
|
|
|
assert serial.status_code == 401 # no session yet
|
|
|
|
|
_login(client)
|
|
|
|
|
serial = client.post("/api/energy/sources", headers={"X-CSRF-Token": _CSRF}, json={
|
|
|
|
|
"name": "Serial", "kind": "warmtelink_serial", "enabled": False, "config": {"path": "/dev/fake"},
|
|
|
|
|
})
|
|
|
|
|
other = _create_source(client)
|
|
|
|
|
channel_uuid = _add_channel(engine, other["uuid"])
|
|
|
|
|
assert client.post(f"/api/energy/sources/{serial.json()['uuid']}/discover").status_code == 403
|
|
|
|
|
disabled = client.post(
|
|
|
|
|
f"/api/energy/sources/{serial.json()['uuid']}/discover", headers={"X-CSRF-Token": _CSRF}
|
|
|
|
|
)
|
|
|
|
|
assert disabled.status_code == 200 and disabled.json()["status"] == "error"
|
|
|
|
|
assert client.get(
|
|
|
|
|
f"/api/energy/sources/{serial.json()['uuid']}/channels/{channel_uuid}/readings"
|
|
|
|
|
).status_code == 404
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_discovery_manager_is_source_scoped_and_never_replaces_a_worker(auth_database):
|
|
|
|
|
"""The real manager queues requests on one fake read-only serial owner."""
|
|
|
|
|
from app.services.warmtelink_worker import WarmteLinkWorkerManager
|
|
|
|
|
|
|
|
|
|
engine = create_engine(auth_database["app_url"], connect_args={"check_same_thread": False})
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
now = datetime.now(UTC)
|
|
|
|
|
source = MeterSource(
|
|
|
|
|
name="Serial", kind="warmtelink_serial", enabled=True, config={"path": "/dev/fake"},
|
|
|
|
|
created_at=now, updated_at=now,
|
|
|
|
|
)
|
|
|
|
|
session.add(source)
|
|
|
|
|
session.commit()
|
|
|
|
|
source_id = source.id
|
|
|
|
|
|
|
|
|
|
class FakeReadOnlyWorker:
|
|
|
|
|
instances: list["FakeReadOnlyWorker"] = []
|
|
|
|
|
|
|
|
|
|
def __init__(self, _source_id, _config, **_kwargs):
|
|
|
|
|
self.requests = []
|
|
|
|
|
self.thread = SimpleNamespace(is_alive=lambda: True)
|
|
|
|
|
self.__class__.instances.append(self)
|
|
|
|
|
|
|
|
|
|
def start(self):
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
def stop(self):
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
def join(self, timeout=5):
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
def request_discovery(self, request):
|
|
|
|
|
self.requests.append(request)
|
|
|
|
|
|
|
|
|
|
manager = WarmteLinkWorkerManager(
|
|
|
|
|
session_factory=lambda: Session(engine), worker_factory=FakeReadOnlyWorker,
|
|
|
|
|
)
|
|
|
|
|
manager.reconcile()
|
|
|
|
|
assert manager.worker_count == 1
|
|
|
|
|
first = manager.request_discovery(source_id)
|
|
|
|
|
assert first.status == "pending"
|
|
|
|
|
worker = FakeReadOnlyWorker.instances[0]
|
|
|
|
|
assert len(worker.requests) == 1
|
|
|
|
|
# A client timing out/cancelling leaves the queued request and worker alone;
|
|
|
|
|
# completing it later cannot open another descriptor or create bindings.
|
|
|
|
|
worker.requests[0].status = "completed"
|
|
|
|
|
worker.requests[0].completed.set()
|
|
|
|
|
results = []
|
|
|
|
|
threads = [threading.Thread(target=lambda: results.append(manager.request_discovery(source_id))) for _ in range(2)]
|
|
|
|
|
for thread in threads:
|
|
|
|
|
thread.start()
|
|
|
|
|
for thread in threads:
|
|
|
|
|
thread.join()
|
|
|
|
|
assert manager.worker_count == 1
|
|
|
|
|
assert len(FakeReadOnlyWorker.instances) == 1
|
|
|
|
|
assert len(worker.requests) == 3
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
assert session.query(Meter).count() == 0
|
|
|
|
|
assert session.query(MeterSourceBinding).count() == 0
|
|
|
|
|
manager.shutdown()
|
|
|
|
|
engine.dispose()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_real_warmtelink_discovery_waits_for_admission_and_keeps_its_serial_owner(
|
|
|
|
|
auth_database, monkeypatch,
|
|
|
|
|
):
|
|
|
|
|
"""A rejected candidate is neither a discovery success nor exposed metadata."""
|
|
|
|
|
from app.integrations.p1 import dsmr_crc16
|
2026-08-24 04:13:11 +02:00
|
|
|
from app.services.warmtelink_ingest import WarmteLinkIngestor
|
2026-08-23 07:09:53 +02:00
|
|
|
from app.services import warmtelink_worker
|
2026-08-24 04:13:11 +02:00
|
|
|
from app.services.warmtelink_worker import WarmteLinkWorker, WarmteLinkWorkerManager
|
2026-08-23 07:09:53 +02:00
|
|
|
|
|
|
|
|
engine = create_engine(auth_database["app_url"], connect_args={"check_same_thread": False})
|
2026-08-24 04:13:11 +02:00
|
|
|
received_at = datetime(2026, 8, 22, 10, 0, 30, tzinfo=UTC) # 12:00:30 Europe/Amsterdam.
|
2026-08-23 07:09:53 +02:00
|
|
|
with Session(engine) as session:
|
|
|
|
|
source = MeterSource(
|
|
|
|
|
name="Serial", kind="warmtelink_serial", enabled=True, config={"path": "/dev/fake"},
|
2026-08-24 04:13:11 +02:00
|
|
|
created_at=received_at, updated_at=received_at,
|
2026-08-23 07:09:53 +02:00
|
|
|
)
|
|
|
|
|
session.add(source)
|
|
|
|
|
session.commit()
|
|
|
|
|
source_id = source.id
|
|
|
|
|
|
|
|
|
|
def frame(second: int, *, crc: bool = False) -> bytes:
|
|
|
|
|
body = (
|
|
|
|
|
b"/WARMTE\r\n"
|
|
|
|
|
+ f"0-0:1.0.0(2608221200{second:02d}S)\r\n".encode()
|
|
|
|
|
+ b"0-0:96.1.1(REDACTED)\r\n"
|
|
|
|
|
+ b"0-1:24.1.0(006)\r\n"
|
|
|
|
|
+ b"0-1:96.1.0(REDACTED)\r\n"
|
|
|
|
|
+ f"0-1:24.2.1(2608221200{second:02d}S)(5.900*m3)\r\n".encode()
|
|
|
|
|
+ b"0-2:24.1.0(012)\r\n"
|
|
|
|
|
+ b"0-2:96.1.0(REDACTED)\r\n"
|
|
|
|
|
+ f"0-2:24.2.1(2608221200{second:02d}S)(0.017*GJ)\r\n".encode()
|
|
|
|
|
)
|
|
|
|
|
payload = body + b"!"
|
|
|
|
|
return payload + (f"{dsmr_crc16(payload):04X}".encode() if crc else b"") + b"\r\n"
|
|
|
|
|
|
|
|
|
|
class FakeReadOnlySerial:
|
|
|
|
|
instances: list["FakeReadOnlySerial"] = []
|
|
|
|
|
|
|
|
|
|
def __init__(self):
|
|
|
|
|
self.frames: Queue[bytes] = Queue()
|
|
|
|
|
self.closed = False
|
|
|
|
|
self.__class__.instances.append(self)
|
|
|
|
|
|
|
|
|
|
def read(self, _size: int = 1) -> bytes:
|
|
|
|
|
try:
|
|
|
|
|
return self.frames.get(timeout=0.01)
|
|
|
|
|
except Empty:
|
|
|
|
|
return b""
|
|
|
|
|
|
|
|
|
|
def close(self) -> None:
|
|
|
|
|
self.closed = True
|
|
|
|
|
|
|
|
|
|
monkeypatch.setattr(warmtelink_worker, "_DISCOVERY_TIMEOUT_SECONDS", 0.15)
|
|
|
|
|
monkeypatch.setattr(warmtelink_worker, "_DISCOVERY_WAIT_SECONDS", 0.02)
|
|
|
|
|
manager = WarmteLinkWorkerManager(
|
|
|
|
|
session_factory=lambda: Session(engine), serial_factory=lambda _config: FakeReadOnlySerial(),
|
2026-08-24 04:13:11 +02:00
|
|
|
worker_factory=lambda source_id, config, **kwargs: WarmteLinkWorker(
|
|
|
|
|
source_id,
|
|
|
|
|
config,
|
|
|
|
|
ingestor=WarmteLinkIngestor(clock=lambda: received_at),
|
|
|
|
|
**kwargs,
|
|
|
|
|
),
|
2026-08-23 07:09:53 +02:00
|
|
|
)
|
|
|
|
|
manager.reconcile()
|
|
|
|
|
serial = FakeReadOnlySerial.instances[0]
|
|
|
|
|
|
|
|
|
|
request = manager.request_discovery(source_id)
|
|
|
|
|
assert request.status == "pending"
|
|
|
|
|
serial.frames.put(frame(0)) # First unverifiable candidate is not admitted.
|
|
|
|
|
time.sleep(0.04)
|
|
|
|
|
assert not request.completed.is_set()
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
assert session.query(MeterSourceChannel).filter_by(source_id=source_id).count() == 0
|
|
|
|
|
assert session.query(MeterSourceBinding).count() == 0
|
|
|
|
|
|
|
|
|
|
serial.frames.put(frame(10)) # Strictly continuous successor admits both channels.
|
|
|
|
|
assert request.completed.wait(1)
|
|
|
|
|
assert request.status == "completed"
|
|
|
|
|
with Session(engine) as session:
|
|
|
|
|
assert session.query(MeterSourceChannel).filter_by(source_id=source_id).count() == 2
|
|
|
|
|
assert session.query(MeterSourceBinding).count() == 0
|
|
|
|
|
|
|
|
|
|
rejected = manager.request_discovery(source_id)
|
|
|
|
|
assert rejected.status == "pending"
|
|
|
|
|
serial.frames.put(b"/malformed!\r\n")
|
|
|
|
|
assert rejected.completed.wait(1)
|
|
|
|
|
assert rejected.status == "error" and rejected.detail == "Discovery timed out."
|
|
|
|
|
assert manager.worker_count == 1 and len(FakeReadOnlySerial.instances) == 1
|
|
|
|
|
|
|
|
|
|
recovered = manager.request_discovery(source_id)
|
|
|
|
|
serial.frames.put(frame(20, crc=True))
|
|
|
|
|
assert recovered.completed.wait(1)
|
|
|
|
|
assert recovered.status == "completed"
|
|
|
|
|
assert manager.worker_count == 1 and len(FakeReadOnlySerial.instances) == 1
|
|
|
|
|
manager.shutdown()
|
|
|
|
|
assert serial.closed
|
|
|
|
|
engine.dispose()
|