From c24b6684cc68ebc3304fa0d57cdf6ec4d64bd161 Mon Sep 17 00:00:00 2001 From: Tianyu Liu Date: Mon, 24 Aug 2026 01:43:41 +0200 Subject: [PATCH] M8-R02A: reconcile DSMR runtime after source CRUD --- app/api/routes/api/meter_sources.py | 31 ++++--- tests/test_meter_source_api.py | 126 ++++++++++++++++++++++++++++ 2 files changed, 147 insertions(+), 10 deletions(-) diff --git a/app/api/routes/api/meter_sources.py b/app/api/routes/api/meter_sources.py index 574ed27..d001dfb 100644 --- a/app/api/routes/api/meter_sources.py +++ b/app/api/routes/api/meter_sources.py @@ -9,6 +9,7 @@ from sqlalchemy import select from sqlalchemy.orm import Session from app.api.routes.api.deps import require_csrf, require_session +from app.config import get_settings from app.dependencies import get_db from app.integrations.meter_sources import SourceProfileError, list_source_profiles, sanitize_source_config from app.models.energy import DsmrReading, Meter @@ -23,6 +24,8 @@ from app.schemas.meter_source import ( SourceProfileResponse, SourceProfilesResponse, ) from app.services.auth import AuthenticatedSession +from app.services.config_page import build_runtime_settings +from app.services.dsmr_ingest import apply_dsmr_subscription from app.services.meter_sources import ( BindingNotFoundError, ChannelNotFoundError, MeterNotFoundError, MeterSourceError, SourceDeleteRestrictedError, SourceNotFoundError, create_binding, @@ -34,14 +37,24 @@ from app.services.warmtelink_worker import warmtelink_worker_manager router = APIRouter(prefix="/api/energy", tags=["api-energy-meter-sources"]) -def _reconcile_warmtelink_after_commit() -> None: - """Runtime convergence is best-effort; the already committed API result wins.""" +def _reconcile_runtimes_after_commit(db: Session) -> None: + """Best-effort runtime convergence after a durable source CRUD commit.""" try: warmtelink_worker_manager.reconcile() except Exception: # The manager records individual source failures itself. Do not turn a # successful durable create/update/delete into a misleading HTTP 500. - return + pass + try: + apply_dsmr_subscription(build_runtime_settings(db, get_settings())) + except Exception: + # DSMR owns independent source clients. Its failure must neither undo + # durable CRUD nor prevent the WarmteLink manager from converging. + pass + finally: + # DSMR health callbacks use short independent sessions. Make a CRUD + # response observe any durable status change they just committed. + db.expire_all() def _as_utc(value: datetime) -> datetime: @@ -130,9 +143,8 @@ def post_source(body: MeterSourceCreate, db: Session = Depends(get_db), try: source = create_source(db, name=body.name, kind=body.kind, config=body.config, enabled=body.enabled) db.commit() - db.refresh(source) - _reconcile_warmtelink_after_commit() - return _source_response(source) + _reconcile_runtimes_after_commit(db) + return _source_response(_source_or_404(db, source.uuid)) except (SourceProfileError, MeterSourceError) as exc: db.rollback() raise HTTPException(status_code=422, detail=str(exc)) from exc @@ -151,9 +163,8 @@ def patch_source(source_uuid: str, body: MeterSourcePatch, db: Session = Depends try: updated = update_source(db, source.id, name=body.name, enabled=body.enabled, config_patch=body.config) db.commit() - db.refresh(updated) - _reconcile_warmtelink_after_commit() - return _source_response(updated) + _reconcile_runtimes_after_commit(db) + return _source_response(_source_or_404(db, updated.uuid)) except (SourceProfileError, MeterSourceError) as exc: db.rollback() raise _binding_error(exc) from exc @@ -171,7 +182,7 @@ def remove_source(source_uuid: str, db: Session = Depends(get_db), try: delete_source(db, source.id) db.commit() - _reconcile_warmtelink_after_commit() + _reconcile_runtimes_after_commit(db) return Response(status_code=status.HTTP_204_NO_CONTENT) except SourceDeleteRestrictedError as exc: db.rollback() diff --git a/tests/test_meter_source_api.py b/tests/test_meter_source_api.py index 5d5a9f0..e686e54 100644 --- a/tests/test_meter_source_api.py +++ b/tests/test_meter_source_api.py @@ -14,6 +14,7 @@ from fastapi.testclient import TestClient from sqlalchemy import create_engine, select from sqlalchemy.orm import Session +from app.models.config import AppConfigEntry from app.models.energy import DsmrReading, Meter from app.models.meter_source import MeterSource, MeterSourceBinding, MeterSourceChannel, WarmteLinkReading @@ -122,6 +123,131 @@ def test_serial_source_commit_is_not_reported_as_failed_when_reconcile_raises(au engine.dispose() +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() + + def test_binding_routes_and_atomic_meter_declaration(auth_database): client, engine = _client(auth_database) with client: