M8-R02A: reconcile DSMR runtime after source CRUD
This commit is contained in:
@@ -9,6 +9,7 @@ from sqlalchemy import select
|
|||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
from app.api.routes.api.deps import require_csrf, require_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.dependencies import get_db
|
||||||
from app.integrations.meter_sources import SourceProfileError, list_source_profiles, sanitize_source_config
|
from app.integrations.meter_sources import SourceProfileError, list_source_profiles, sanitize_source_config
|
||||||
from app.models.energy import DsmrReading, Meter
|
from app.models.energy import DsmrReading, Meter
|
||||||
@@ -23,6 +24,8 @@ from app.schemas.meter_source import (
|
|||||||
SourceProfileResponse, SourceProfilesResponse,
|
SourceProfileResponse, SourceProfilesResponse,
|
||||||
)
|
)
|
||||||
from app.services.auth import AuthenticatedSession
|
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 (
|
from app.services.meter_sources import (
|
||||||
BindingNotFoundError, ChannelNotFoundError, MeterNotFoundError,
|
BindingNotFoundError, ChannelNotFoundError, MeterNotFoundError,
|
||||||
MeterSourceError, SourceDeleteRestrictedError, SourceNotFoundError, create_binding,
|
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"])
|
router = APIRouter(prefix="/api/energy", tags=["api-energy-meter-sources"])
|
||||||
|
|
||||||
|
|
||||||
def _reconcile_warmtelink_after_commit() -> None:
|
def _reconcile_runtimes_after_commit(db: Session) -> None:
|
||||||
"""Runtime convergence is best-effort; the already committed API result wins."""
|
"""Best-effort runtime convergence after a durable source CRUD commit."""
|
||||||
try:
|
try:
|
||||||
warmtelink_worker_manager.reconcile()
|
warmtelink_worker_manager.reconcile()
|
||||||
except Exception:
|
except Exception:
|
||||||
# The manager records individual source failures itself. Do not turn a
|
# The manager records individual source failures itself. Do not turn a
|
||||||
# successful durable create/update/delete into a misleading HTTP 500.
|
# 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:
|
def _as_utc(value: datetime) -> datetime:
|
||||||
@@ -130,9 +143,8 @@ def post_source(body: MeterSourceCreate, db: Session = Depends(get_db),
|
|||||||
try:
|
try:
|
||||||
source = create_source(db, name=body.name, kind=body.kind, config=body.config, enabled=body.enabled)
|
source = create_source(db, name=body.name, kind=body.kind, config=body.config, enabled=body.enabled)
|
||||||
db.commit()
|
db.commit()
|
||||||
db.refresh(source)
|
_reconcile_runtimes_after_commit(db)
|
||||||
_reconcile_warmtelink_after_commit()
|
return _source_response(_source_or_404(db, source.uuid))
|
||||||
return _source_response(source)
|
|
||||||
except (SourceProfileError, MeterSourceError) as exc:
|
except (SourceProfileError, MeterSourceError) as exc:
|
||||||
db.rollback()
|
db.rollback()
|
||||||
raise HTTPException(status_code=422, detail=str(exc)) from exc
|
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:
|
try:
|
||||||
updated = update_source(db, source.id, name=body.name, enabled=body.enabled, config_patch=body.config)
|
updated = update_source(db, source.id, name=body.name, enabled=body.enabled, config_patch=body.config)
|
||||||
db.commit()
|
db.commit()
|
||||||
db.refresh(updated)
|
_reconcile_runtimes_after_commit(db)
|
||||||
_reconcile_warmtelink_after_commit()
|
return _source_response(_source_or_404(db, updated.uuid))
|
||||||
return _source_response(updated)
|
|
||||||
except (SourceProfileError, MeterSourceError) as exc:
|
except (SourceProfileError, MeterSourceError) as exc:
|
||||||
db.rollback()
|
db.rollback()
|
||||||
raise _binding_error(exc) from exc
|
raise _binding_error(exc) from exc
|
||||||
@@ -171,7 +182,7 @@ def remove_source(source_uuid: str, db: Session = Depends(get_db),
|
|||||||
try:
|
try:
|
||||||
delete_source(db, source.id)
|
delete_source(db, source.id)
|
||||||
db.commit()
|
db.commit()
|
||||||
_reconcile_warmtelink_after_commit()
|
_reconcile_runtimes_after_commit(db)
|
||||||
return Response(status_code=status.HTTP_204_NO_CONTENT)
|
return Response(status_code=status.HTTP_204_NO_CONTENT)
|
||||||
except SourceDeleteRestrictedError as exc:
|
except SourceDeleteRestrictedError as exc:
|
||||||
db.rollback()
|
db.rollback()
|
||||||
|
|||||||
@@ -14,6 +14,7 @@ from fastapi.testclient import TestClient
|
|||||||
from sqlalchemy import create_engine, select
|
from sqlalchemy import create_engine, select
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
|
from app.models.config import AppConfigEntry
|
||||||
from app.models.energy import DsmrReading, Meter
|
from app.models.energy import DsmrReading, Meter
|
||||||
from app.models.meter_source import MeterSource, MeterSourceBinding, MeterSourceChannel, WarmteLinkReading
|
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()
|
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):
|
def test_binding_routes_and_atomic_meter_declaration(auth_database):
|
||||||
client, engine = _client(auth_database)
|
client, engine = _client(auth_database)
|
||||||
with client:
|
with client:
|
||||||
|
|||||||
Reference in New Issue
Block a user