M8-T11: add WarmteLink discovery and history API
This commit is contained in:
@@ -12,10 +12,12 @@ from app.api.routes.api.deps import require_csrf, require_session
|
||||
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
|
||||
from app.models.meter_source import MeterSource, MeterSourceBinding, MeterSourceChannel
|
||||
from app.models.meter_source import MeterSource, MeterSourceBinding, MeterSourceChannel, WarmteLinkReading
|
||||
from app.schemas.meter_source import (
|
||||
BindingCreate, BindingListResponse, BindingPatch, BindingResponse, ChannelReadingResponse,
|
||||
BindingCreate, BindingListResponse, BindingPatch, BindingResponse, ChannelBindingSummaryResponse,
|
||||
ChannelReadingResponse,
|
||||
ChannelReadingsResponse, CommoditiesResponse, CommodityResponse, DiscoverResponse,
|
||||
DiscoverChannelResponse,
|
||||
MeterSourceChannelListResponse, MeterSourceChannelResponse, MeterSourceCreate,
|
||||
MeterSourceListResponse, MeterSourcePatch, MeterSourceResponse, SourceConfigFieldResponse,
|
||||
SourceProfileResponse, SourceProfilesResponse,
|
||||
@@ -181,8 +183,25 @@ def discover_source(source_uuid: str, db: Session = Depends(get_db),
|
||||
_auth: AuthenticatedSession = Depends(require_session), _csrf: None = Depends(require_csrf)) -> DiscoverResponse:
|
||||
source = _source_or_404(db, source_uuid)
|
||||
if source.kind == "warmtelink_serial":
|
||||
return DiscoverResponse(requested=False, supported=False, status="not_implemented",
|
||||
detail="Serial discovery is available after the WarmteLink worker is installed.")
|
||||
if not source.enabled:
|
||||
return DiscoverResponse(
|
||||
requested=False, supported=True, status="error",
|
||||
detail="The WarmteLink source is disabled.", channels=_discover_channels(db, source),
|
||||
)
|
||||
# This merely schedules lifecycle convergence. It never opens a serial
|
||||
# descriptor or waits for a frame in the request thread; the one managed
|
||||
# worker remains the sole owner of serial I/O and can keep reconnecting.
|
||||
request = warmtelink_worker_manager.request_discovery(source.id)
|
||||
if request.completed.is_set():
|
||||
# A worker may have accepted a frame during the bounded wait.
|
||||
# Refresh only durable accepted metadata, never candidates/raw data.
|
||||
db.expire_all()
|
||||
source = _source_or_404(db, source_uuid)
|
||||
return DiscoverResponse(
|
||||
requested=request.status != "error", supported=True, status=request.status,
|
||||
request_id=request.request_id or None, detail=request.detail,
|
||||
channels=_discover_channels(db, source),
|
||||
)
|
||||
return DiscoverResponse(requested=False, supported=True, status="managed_by_runtime",
|
||||
detail="This source is discovered by its runtime subscription; no connection was opened.")
|
||||
|
||||
@@ -195,32 +214,61 @@ def source_channels(source_uuid: str, db: Session = Depends(get_db),
|
||||
items = []
|
||||
for channel in channels:
|
||||
bindings = list_bindings(db, channel_id=channel.id)
|
||||
meter_ids = [binding.meter_id for binding in bindings]
|
||||
items.append(MeterSourceChannelResponse(
|
||||
uuid=channel.uuid, label=channel.label, suggested_commodity=channel.suggested_commodity,
|
||||
unit=channel.unit, device_type=channel.device_type, latest_value=channel.latest_value,
|
||||
latest_at=channel.latest_at, latest_quality=channel.latest_quality, binding_count=len(bindings),
|
||||
bound_meter_ids=[binding.meter_id for binding in bindings],
|
||||
bound_meter_ids=meter_ids,
|
||||
binding_summary=ChannelBindingSummaryResponse(count=len(bindings), meter_ids=meter_ids),
|
||||
))
|
||||
return MeterSourceChannelListResponse(items=items, total=len(items))
|
||||
return MeterSourceChannelListResponse(items=items, total=len(items), source_status=source.status)
|
||||
|
||||
|
||||
@router.get("/sources/{source_uuid}/channels/{channel_uuid}/readings", response_model=ChannelReadingsResponse)
|
||||
def channel_readings(source_uuid: str, channel_uuid: str, limit: int = Query(default=100, ge=1, le=1000),
|
||||
start: datetime | None = None, end: datetime | None = None, db: Session = Depends(get_db),
|
||||
from_: datetime | None = Query(default=None, alias="from"),
|
||||
to: datetime | None = Query(default=None), db: Session = Depends(get_db),
|
||||
_auth: AuthenticatedSession = Depends(require_session)) -> ChannelReadingsResponse:
|
||||
source = _source_or_404(db, source_uuid)
|
||||
_channel_or_404(db, source, channel_uuid)
|
||||
if source.kind != "dsmr_mqtt":
|
||||
return ChannelReadingsResponse(items=[], total=0)
|
||||
statement = select(DsmrReading).where(DsmrReading.meter_source_id == source.id)
|
||||
if start is not None:
|
||||
statement = statement.where(DsmrReading.recorded_at >= _as_utc(start))
|
||||
if end is not None:
|
||||
statement = statement.where(DsmrReading.recorded_at < _as_utc(end))
|
||||
rows = list(db.execute(statement.order_by(DsmrReading.recorded_at.desc()).limit(limit)).scalars())
|
||||
# The DSMR payload remains available only from its legacy compatibility endpoint;
|
||||
# this generic endpoint intentionally exposes no telegram/equipment identifiers.
|
||||
return ChannelReadingsResponse(items=[ChannelReadingResponse(recorded_at=row.recorded_at) for row in rows], total=len(rows))
|
||||
channel = _channel_or_404(db, source, channel_uuid)
|
||||
if from_ is not None and to is not None and _as_utc(from_) >= _as_utc(to):
|
||||
raise HTTPException(status_code=422, detail="'from' must be earlier than 'to'.")
|
||||
if source.kind == "warmtelink_serial":
|
||||
statement = select(WarmteLinkReading).where(WarmteLinkReading.channel_id == channel.id)
|
||||
model = WarmteLinkReading
|
||||
else:
|
||||
# DSMR remains a source-level protocol history. Its channel is the
|
||||
# public electricity identity, while payload/telegram diagnostics stay
|
||||
# private to ingestion and the legacy latest endpoint.
|
||||
statement = select(DsmrReading).where(DsmrReading.meter_source_id == source.id)
|
||||
model = DsmrReading
|
||||
if from_ is not None:
|
||||
statement = statement.where(model.recorded_at >= _as_utc(from_))
|
||||
if to is not None:
|
||||
statement = statement.where(model.recorded_at < _as_utc(to))
|
||||
rows = list(db.execute(statement.order_by(model.recorded_at.asc()).limit(limit)).scalars())
|
||||
return ChannelReadingsResponse(
|
||||
items=[ChannelReadingResponse(
|
||||
recorded_at=row.recorded_at,
|
||||
value=getattr(row, "value", None), quality=getattr(row, "quality", None),
|
||||
) for row in rows],
|
||||
total=len(rows),
|
||||
)
|
||||
|
||||
|
||||
def _discover_channels(db: Session, source: MeterSource) -> list[DiscoverChannelResponse]:
|
||||
"""Return only public, accepted channel metadata for discover responses."""
|
||||
return [
|
||||
DiscoverChannelResponse(
|
||||
uuid=channel.uuid, label=channel.label, unit=channel.unit,
|
||||
latest_value=channel.latest_value, latest_at=channel.latest_at,
|
||||
latest_quality=channel.latest_quality,
|
||||
)
|
||||
for channel in db.execute(
|
||||
select(MeterSourceChannel).where(MeterSourceChannel.source_id == source.id)
|
||||
).scalars()
|
||||
]
|
||||
|
||||
|
||||
@router.get("/meters/{meter_id}/bindings", response_model=BindingListResponse)
|
||||
|
||||
@@ -64,7 +64,23 @@ class DiscoverResponse(BaseModel):
|
||||
requested: bool
|
||||
supported: bool
|
||||
status: str
|
||||
request_id: int | None = None
|
||||
detail: str | None = None
|
||||
channels: list["DiscoverChannelResponse"] = Field(default_factory=list)
|
||||
|
||||
|
||||
class DiscoverChannelResponse(BaseModel):
|
||||
uuid: str
|
||||
label: str
|
||||
unit: str
|
||||
latest_value: Decimal | None
|
||||
latest_at: datetime | None
|
||||
latest_quality: str | None
|
||||
|
||||
|
||||
class ChannelBindingSummaryResponse(BaseModel):
|
||||
count: int
|
||||
meter_ids: list[int]
|
||||
|
||||
|
||||
class CommodityResponse(BaseModel):
|
||||
@@ -88,11 +104,13 @@ class MeterSourceChannelResponse(BaseModel):
|
||||
latest_quality: str | None
|
||||
binding_count: int
|
||||
bound_meter_ids: list[int]
|
||||
binding_summary: ChannelBindingSummaryResponse
|
||||
|
||||
|
||||
class MeterSourceChannelListResponse(BaseModel):
|
||||
items: list[MeterSourceChannelResponse]
|
||||
total: int
|
||||
source_status: str
|
||||
|
||||
|
||||
class ChannelReadingResponse(BaseModel):
|
||||
|
||||
@@ -8,8 +8,8 @@ from __future__ import annotations
|
||||
|
||||
from collections.abc import Callable
|
||||
from contextlib import suppress
|
||||
from dataclasses import dataclass
|
||||
from datetime import UTC, datetime
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import UTC, datetime, timedelta
|
||||
import logging
|
||||
import threading
|
||||
from typing import Protocol
|
||||
@@ -27,6 +27,9 @@ logger = logging.getLogger(__name__)
|
||||
|
||||
_BACKOFF_SECONDS = (1, 2, 4, 8, 16, 32, 60)
|
||||
_JOIN_TIMEOUT_SECONDS = 5
|
||||
_DISCOVERY_LOCK_TIMEOUT_SECONDS = 0.05
|
||||
_DISCOVERY_WAIT_SECONDS = 0.1
|
||||
_DISCOVERY_TIMEOUT_SECONDS = 5
|
||||
|
||||
|
||||
class ReadOnlySerial(Protocol):
|
||||
@@ -69,6 +72,18 @@ class _WorkerConfig:
|
||||
config: dict
|
||||
|
||||
|
||||
@dataclass
|
||||
class DiscoveryRequest:
|
||||
"""One source-scoped request, completed only by its serial owner."""
|
||||
|
||||
request_id: int
|
||||
source_id: int
|
||||
deadline: datetime
|
||||
status: str = "pending"
|
||||
detail: str | None = None
|
||||
completed: threading.Event = field(default_factory=threading.Event)
|
||||
|
||||
|
||||
class WarmteLinkWorker:
|
||||
"""One interruptible, read-only serial loop for one meter source."""
|
||||
|
||||
@@ -88,6 +103,8 @@ class WarmteLinkWorker:
|
||||
self._clock = clock or _EventClock()
|
||||
self._serial: ReadOnlySerial | None = None
|
||||
self._serial_lock = threading.Lock()
|
||||
self._discovery_lock = threading.Lock()
|
||||
self._discoveries: list[DiscoveryRequest] = []
|
||||
# Never inherit a daemon flag from a caller's background thread: a serial
|
||||
# descriptor and its orderly shutdown must remain visible to the process.
|
||||
self._thread = threading.Thread(
|
||||
@@ -105,6 +122,35 @@ class WarmteLinkWorker:
|
||||
self._stop_event.set()
|
||||
self._close_serial()
|
||||
|
||||
def request_discovery(self, request: DiscoveryRequest) -> None:
|
||||
"""Queue a read request; this worker remains the sole serial owner."""
|
||||
with self._discovery_lock:
|
||||
self._discoveries.append(request)
|
||||
|
||||
def _finish_discoveries(self, status: str, detail: str | None = None) -> None:
|
||||
now = datetime.now(UTC)
|
||||
with self._discovery_lock:
|
||||
pending, self._discoveries = self._discoveries, []
|
||||
for request in pending:
|
||||
if request.completed.is_set():
|
||||
continue
|
||||
if request.deadline <= now and status == "completed":
|
||||
request.status, request.detail = "error", "Discovery timed out."
|
||||
else:
|
||||
request.status, request.detail = status, detail
|
||||
request.completed.set()
|
||||
|
||||
def _expire_discoveries(self) -> None:
|
||||
now = datetime.now(UTC)
|
||||
with self._discovery_lock:
|
||||
expired = [request for request in self._discoveries if request.deadline <= now]
|
||||
self._discoveries = [request for request in self._discoveries if request.deadline > now]
|
||||
for request in expired:
|
||||
if request.completed.is_set():
|
||||
continue
|
||||
request.status, request.detail = "error", "Discovery timed out."
|
||||
request.completed.set()
|
||||
|
||||
def join(self, timeout: float = _JOIN_TIMEOUT_SECONDS) -> bool:
|
||||
self._thread.join(timeout)
|
||||
return not self._thread.is_alive()
|
||||
@@ -142,6 +188,7 @@ class WarmteLinkWorker:
|
||||
return
|
||||
self._serial = device
|
||||
while not self._stop_event.is_set():
|
||||
self._expire_discoveries()
|
||||
chunk = device.read(1024)
|
||||
if not chunk:
|
||||
# ``timeout`` reads are normal, but still yield so a bad
|
||||
@@ -152,13 +199,16 @@ class WarmteLinkWorker:
|
||||
for frame in frames:
|
||||
if self._stop_event.is_set():
|
||||
break
|
||||
self._ingestor.handle_frame(
|
||||
admitted = self._ingestor.handle_frame(
|
||||
self.source_id, frame, session_factory=self._session_factory
|
||||
)
|
||||
if admitted:
|
||||
self._finish_discoveries("completed")
|
||||
# A complete frame proves transport recovery even if its
|
||||
# contents are rejected by the privacy/admission layer.
|
||||
backoff_index = 0
|
||||
except Exception:
|
||||
self._finish_discoveries("error", "WarmteLink discovery failed.")
|
||||
self._record_error("WarmteLink serial connection failed")
|
||||
delay = _BACKOFF_SECONDS[min(backoff_index, len(_BACKOFF_SECONDS) - 1)]
|
||||
backoff_index += 1
|
||||
@@ -181,6 +231,7 @@ class WarmteLinkWorkerManager:
|
||||
self._workers: dict[int, tuple[_WorkerConfig, WarmteLinkWorker]] = {}
|
||||
self._lock = threading.Lock()
|
||||
self._reapers: set[int] = set()
|
||||
self._next_discovery_id = 0
|
||||
self._shutting_down = False
|
||||
|
||||
@property
|
||||
@@ -203,6 +254,48 @@ class WarmteLinkWorkerManager:
|
||||
self._shutting_down = False
|
||||
self.reconcile()
|
||||
|
||||
def request_discovery(self, source_id: int) -> DiscoveryRequest:
|
||||
"""Ask the current source worker for one bounded read/discovery attempt.
|
||||
|
||||
This intentionally does not reconcile or open a descriptor. Lifecycle
|
||||
convergence remains separate; a request can neither replace nor stop a
|
||||
worker when an HTTP client times out or disconnects.
|
||||
"""
|
||||
now = datetime.now(UTC)
|
||||
request = DiscoveryRequest(0, source_id, now)
|
||||
if not self._lock.acquire(timeout=_DISCOVERY_LOCK_TIMEOUT_SECONDS):
|
||||
request.status, request.detail = "error", "Discovery queue is busy."
|
||||
request.completed.set()
|
||||
return request
|
||||
try:
|
||||
self._next_discovery_id += 1
|
||||
request.request_id = self._next_discovery_id
|
||||
request.deadline = now + timedelta(seconds=_DISCOVERY_TIMEOUT_SECONDS)
|
||||
if self._shutting_down:
|
||||
request.status, request.detail = "error", "WarmteLink manager is stopped."
|
||||
request.completed.set()
|
||||
elif (entry := self._workers.get(source_id)) is None:
|
||||
request.status, request.detail = "error", "WarmteLink worker is not running."
|
||||
request.completed.set()
|
||||
else:
|
||||
entry[1].request_discovery(request)
|
||||
timer = threading.Timer(_DISCOVERY_TIMEOUT_SECONDS, self._timeout_discovery, args=(request,))
|
||||
timer.daemon = True
|
||||
timer.start()
|
||||
finally:
|
||||
self._lock.release()
|
||||
# A tiny bounded wait makes an immediately available frame observable,
|
||||
# without turning an HTTP call into serial I/O or an unbounded wait.
|
||||
request.completed.wait(_DISCOVERY_WAIT_SECONDS)
|
||||
return request
|
||||
|
||||
@staticmethod
|
||||
def _timeout_discovery(request: DiscoveryRequest) -> None:
|
||||
"""Resolve a stale HTTP request without touching its healthy worker."""
|
||||
if not request.completed.is_set():
|
||||
request.status, request.detail = "error", "Discovery timed out."
|
||||
request.completed.set()
|
||||
|
||||
def _read_desired(self) -> dict[int, _WorkerConfig]:
|
||||
with self._session_factory() as session:
|
||||
return {
|
||||
|
||||
Reference in New Issue
Block a user