diff --git a/app/integrations/expose.py b/app/integrations/expose.py index 16c5795..2ab71cf 100644 --- a/app/integrations/expose.py +++ b/app/integrations/expose.py @@ -45,12 +45,15 @@ class DeviceInfo: """Grouping metadata that maps an entity to a logical HA device. ``identifiers`` corresponds to ``device.identifiers`` in HA Discovery - config — a stable tuple used to group entities under one device card. + config. HA merges device cards if *any* identifier overlaps, so every + logical device deliberately has exactly one, globally namespaced value. + ``identity`` is an independent internal seed for MQTT topics and unique + ids; it must never be inferred from the HA identifier list. ``name`` is the human-readable device name (friendly_name). """ identifiers: tuple[str, ...] - """Stable identifiers for HA device grouping (e.g. ``("modbus", "")``). + """Stable HA identifiers, normally a one-item tuple. These must not change over the device lifetime; they anchor the HA device card even when friendly_name changes. @@ -59,6 +62,14 @@ class DeviceInfo: name: str """Human-readable device name (may change; triggers HA discovery re-publish).""" + identity: Optional[str] = None + """Stable internal topic/unique-id identity, separate from HA identifiers. + + The fallback preserves compatibility for third-party providers still using + the historic ``("kind", "uuid")`` construction while first-party providers + use this field explicitly. + """ + provides_availability: bool = True """Whether this device publishes an availability ("online"/"offline") heartbeat. @@ -85,6 +96,11 @@ class DeviceInfo: ) """Return whether the source behind this device is currently usable.""" + @property + def internal_identity(self) -> str: + """Return the topic/unique-id seed without exposing HA grouping details.""" + return self.identity or self.identifiers[-1] + @dataclass class ExposableEntity: @@ -276,8 +292,9 @@ def _modbus_provider(session: Session) -> list[ExposableEntity]: for device in devices: device_info = DeviceInfo( - identifiers=("modbus", device.uuid), + identifiers=(f"home-automation:modbus:{device.uuid}",), name=device.friendly_name, + identity=device.uuid, ) # Load the profile to get metric metadata. @@ -407,9 +424,10 @@ def _energy_cost_provider(session: Session) -> list[ExposableEntity]: HA device identity (换表 → 新 sensor) -------------------------------------- - ``identifiers[1]`` is set to the active meter's **uuid** (not the fixed - string ``"energy-cost"``). ``ha_discovery.py`` uses ``identifiers[1]`` as - the MQTT node_id and as part of the ``unique_id`` for every entity. + The independent internal identity is set to the active meter's **uuid**. + ``ha_discovery.py`` uses that internal identity as the MQTT node_id and as + part of the ``unique_id`` for every entity; HA receives one separate, + namespaced device identifier. Declaring a new active electricity meter produces a new uuid → new node_id / unique_id → HA creates a brand-new sensor, cleanly isolating post-swap data. @@ -460,9 +478,9 @@ def _energy_cost_provider(session: Session) -> list[ExposableEntity]: DeviceInfo identifiers ---------------------- - **Two-element tuple** ``("energy-cost", meter.uuid)`` so that - ``ha_discovery.py``'s ``entity.device.identifiers[1]`` resolves to the - meter uuid (used as the MQTT node_id and unique_id seed throughout). + One-element namespaced tuple ``("home-automation:energy-cost:",)``. + The meter uuid used for MQTT topics and unique IDs is carried independently + in ``DeviceInfo.identity``. """ from app.models.energy import EnergyCostPeriod, Meter # local import to avoid circular from sqlalchemy import select @@ -499,14 +517,15 @@ def _energy_cost_provider(session: Session) -> list[ExposableEntity]: currency = latest_period.currency # --- Shared DeviceInfo anchored to the active meter's uuid --- - # identifiers[1] = meter.uuid drives the MQTT node_id and unique_id in - # ha_discovery.py. Swapping the meter produces a new uuid → new HA sensor. + # The internal identity is meter.uuid, while the HA identifier is a single + # namespaced value. Swapping the meter produces a new HA sensor/card. # provides_availability=False: the energy-cost device has only sensors and no # online/offline heartbeat, so its entities must be "always available" in HA. # (Otherwise HA shows them unavailable despite state being published.) device_info = DeviceInfo( - identifiers=("energy-cost", active_meter.uuid), + identifiers=(f"home-automation:energy-cost:{active_meter.uuid}",), name=active_meter.label, + identity=active_meter.uuid, provides_availability=False, ) @@ -967,7 +986,8 @@ def _m8_energy_provider(session: Session) -> list[ExposableEntity]: for source in sources: source_info = DeviceInfo( - identifiers=("meter-source", source.uuid), name=source.name, + identifiers=(f"home-automation:meter-source:{source.uuid}",), name=source.name, + identity=source.uuid, availability_id=source.uuid, availability_getter=lambda sess, source_id=source.id: _source_online_by_id(sess, source_id), ) @@ -994,7 +1014,8 @@ def _m8_energy_provider(session: Session) -> list[ExposableEntity]: continue binding, channel, source = bound info = DeviceInfo( - identifiers=("meter", meter.uuid), name=meter.label, + identifiers=(f"home-automation:meter:{meter.uuid}",), name=meter.label, + identity=meter.uuid, # Keep this opaque and Meter-anchored. In particular, do not use a # source UUID here: two channels of one source can be independently # stale/invalid and must not overwrite each other's availability. @@ -1021,15 +1042,20 @@ def _m8_energy_provider(session: Session) -> list[ExposableEntity]: heating, hot_water = active_by_commodity.get("heating"), active_by_commodity.get("hot_water") if heating is not None and hot_water is not None: + # Keep the v1.6.1 canonical identity/key seed stable. Discovery topic + # segments are sanitised independently by ha_discovery, so this dot is + # never emitted in a topic while existing toggle rows remain usable. identity = ".".join(sorted((heating.uuid, hot_water.uuid))) currency = _thermal_currency(session) cost_info = DeviceInfo( - identifiers=("thermal-cost", identity), name="Thermal Energy Cost", + identifiers=(f"home-automation:thermal-cost:{identity}",), name="Thermal Energy Cost", + identity=identity, provides_availability=False, ) labels = { "heating": "Heating", "hot_water_heating": "Hot Water Heating", "water": "Water", - "water_tax": "Water Tax", "fixed": "Fixed", "all_in": "All-in", + "water_tax": "Water Tax", "hot_water_total": "Hot Water", + "fixed": "Fixed", "all_in": "All-in", } for suffix, window in (("total", None), ("today", "today")): for metric, label in labels.items(): @@ -1203,6 +1229,8 @@ def _thermal_cost_getter(metric: str, window: str | None) -> Callable[[Session], return result["breakdown"]["hot_water"] if metric == "water_tax": return result["breakdown"]["hot_water_tax"] + if metric == "hot_water_total": + return result["breakdown"]["hot_water_heating"] + result["breakdown"]["hot_water"] return result["breakdown"][metric] return _getter diff --git a/app/integrations/homeassistant.py b/app/integrations/homeassistant.py index d371a17..72e21f4 100644 --- a/app/integrations/homeassistant.py +++ b/app/integrations/homeassistant.py @@ -2,10 +2,14 @@ from __future__ import annotations import json import logging +from time import monotonic from dataclasses import dataclass, field from typing import Any from urllib import error, parse, request +from websockets.exceptions import WebSocketException +from websockets.sync.client import connect + from app.config import Settings logger = logging.getLogger(__name__) @@ -57,6 +61,81 @@ class HomeAssistantClient: self._post_json(f"/api/webhook/{webhook_id}", body, operation="trigger_webhook") + def discovery_registry_bindings(self, unique_ids: set[str]) -> dict[str, set[str]]: + """Return HA device identifiers currently bound to MQTT unique IDs. + + This is deliberately a read-only WebSocket query. MQTT only confirms + broker receipt; the entity/device registries are the authoritative HA + observation that a discovery unload/re-add was actually processed. + """ + self._require_config() + if not unique_ids: + return {} + try: + deadline = monotonic() + self.timeout_seconds + with connect(self._websocket_url(), open_timeout=self.timeout_seconds, + close_timeout=self.timeout_seconds) as websocket: + greeting = json.loads(self._websocket_recv(websocket, deadline)) + if greeting.get("type") != "auth_required": + raise HomeAssistantRequestError("Unexpected Home Assistant WebSocket greeting") + websocket.send(json.dumps({"type": "auth", "access_token": self.settings.home_assistant_auth_token})) + auth = json.loads(self._websocket_recv(websocket, deadline)) + if auth.get("type") != "auth_ok": + raise HomeAssistantRequestError("Home Assistant WebSocket authentication failed") + entities = self._websocket_command(websocket, 1, "config/entity_registry/list", deadline) + devices = self._websocket_command(websocket, 2, "config/device_registry/list", deadline) + except (OSError, WebSocketException, TimeoutError, ValueError, KeyError, TypeError) as exc: + raise HomeAssistantRequestError("Home Assistant registry query failed") from exc + + devices_by_id = { + device["id"]: { + identifier[1] + for identifier in device.get("identifiers", []) + if ( + isinstance(identifier, (list, tuple)) + and len(identifier) == 2 + and identifier[0] == "mqtt" + and isinstance(identifier[1], str) + and identifier[1] + ) + } + for device in devices + if isinstance(device, dict) and isinstance(device.get("id"), str) + } + return { + entity["unique_id"]: devices_by_id.get(entity.get("device_id"), set()) + for entity in entities + if entity.get("platform") == "mqtt" and entity.get("unique_id") in unique_ids + } + + @staticmethod + def _websocket_recv(websocket: Any, deadline: float) -> str: + remaining = deadline - monotonic() + if remaining <= 0: + raise TimeoutError("Home Assistant WebSocket registry query timed out") + return websocket.recv(timeout=remaining) + + @classmethod + def _websocket_command( + cls, websocket: Any, message_id: int, command: str, deadline: float + ) -> list[dict[str, Any]]: + websocket.send(json.dumps({"id": message_id, "type": command})) + while True: + response = json.loads(cls._websocket_recv(websocket, deadline)) + if response.get("id") != message_id: + continue + if not response.get("success"): + raise HomeAssistantRequestError(f"Home Assistant WebSocket {command} failed") + return response.get("result", []) + + def _websocket_url(self) -> str: + parsed = parse.urlsplit(self.settings.home_assistant_base_url) + if parsed.scheme not in {"http", "https"} or not parsed.netloc: + raise HomeAssistantConfigError("HOME_ASSISTANT_BASE_URL must be an HTTP(S) URL") + scheme = "wss" if parsed.scheme == "https" else "ws" + path = f"{parsed.path.rstrip('/')}/api/websocket" + return parse.urlunsplit((scheme, parsed.netloc, path, "", "")) + def _require_config(self) -> None: if self.is_configured(): return diff --git a/app/integrations/mqtt.py b/app/integrations/mqtt.py index 5bf111f..e3920f1 100644 --- a/app/integrations/mqtt.py +++ b/app/integrations/mqtt.py @@ -174,23 +174,33 @@ class MqttManager: *, retain: bool = False, qos: int = 0, - ) -> None: + ) -> bool: """Publish a message to *topic*. - If the client is not connected the call is silently skipped. - Errors are logged but do not raise. + Return whether paho accepted the message for publication. If the + client is not connected, or paho raises/rejects it, return ``False``; + errors still do not escape this best-effort boundary. """ with self._lock: client = self._client if client is None or not self._connected: logger.debug("MQTT publish skipped — not connected (topic=%s).", topic) - return + return False try: - client.publish(topic, payload=payload, qos=qos, retain=retain) + result = client.publish(topic, payload=payload, qos=qos, retain=retain) + # paho returns MQTTMessageInfo with an integer rc. Keep fakes and + # alternate clients compatible by treating an absent/non-int rc as + # accepted after the call itself succeeded. + rc = getattr(result, "rc", None) + if isinstance(rc, int) and rc != mqtt.MQTT_ERR_SUCCESS: + logger.warning("MQTT publish rejected (topic=%s, rc=%s).", topic, rc) + return False + return True except Exception: logger.exception("MQTT publish error (topic=%s).", topic) + return False def subscribe(self, topic: str, handler: Callable[[bytes], None]) -> None: """Register *handler* to be called when a message arrives on *topic*. diff --git a/app/main.py b/app/main.py index 9bfee89..e3c50b6 100644 --- a/app/main.py +++ b/app/main.py @@ -38,7 +38,7 @@ from app.services.config_page import build_runtime_settings, seed_missing_config from app.services.dsmr_ingest import apply_dsmr_subscription from app.services.public_ip import check_public_ipv4_and_notify from app.services.modbus_poll import poll_all_enabled_devices, BASE_POLL_TICK_SECONDS -from app.services.ha_discovery import publish_discovery, publish_states +from app.services.ha_discovery import initialize_legacy_thermal_cleanup, publish_discovery, publish_states from app.services.tibber_prices import run_tibber_refresh_best_effort from app.services.energy_cost import compute_closed_periods from app.services.meter_cost import compute_closed_periods as compute_closed_meter_cost_periods @@ -203,6 +203,7 @@ def ensure_auth_db_ready() -> None: initialize_auth_schema(session, get_settings()) seed_missing_config_from_bootstrap(session, get_settings()) sync_app_hostname_from_bootstrap(session, get_settings()) + initialize_legacy_thermal_cleanup(session) except AppDatabaseAdoptionError as exc: raise RuntimeError(str(exc)) from exc except AuthBootstrapError as exc: diff --git a/app/services/ha_discovery.py b/app/services/ha_discovery.py index 3ce6e29..f5aa488 100644 --- a/app/services/ha_discovery.py +++ b/app/services/ha_discovery.py @@ -34,16 +34,30 @@ from __future__ import annotations import json import logging +import re +from datetime import UTC, datetime from typing import Any from sqlalchemy.orm import Session +from sqlalchemy import inspect from app.integrations.expose import ExposableEntity, build_catalog +from app.models.config import AppConfigEntry from app.integrations.mqtt import mqtt_manager from app.services.config_page import build_runtime_settings logger = logging.getLogger(__name__) +_HA_SEGMENT_RE = re.compile(r"[^A-Za-z0-9_-]+") +# These durable, versioned application metadata entries prevent the minute +# publisher from repeatedly emitting already accepted known-bad v1.6.1 thermal +# topics. A failed publish is not acknowledged and remains retryable across +# restarts. They are intentionally application metadata, rather than broker +# state. An empty retained publish to an illegal v1.6.1 topic cannot be used +# as an acknowledgement: HA rejects that topic before it processes its payload. +_LEGACY_THERMAL_CLEANUP_KEY = "HA_DISCOVERY_LEGACY_THERMAL_CLEANUP_V1" +_REGISTRY_REPAIR_KEY = "HA_DISCOVERY_REGISTRY_REPAIR_V2" + # --------------------------------------------------------------------------- # Internal helpers @@ -59,14 +73,25 @@ def _should_publish(settings: Any) -> bool: ) -def _node_id(device_uuid: str) -> str: - """Stable MQTT node_id derived from device uuid (replace hyphens for safety).""" - return device_uuid.replace("-", "_") +def _safe_segment(value: str) -> str: + """Return a deterministic HA discovery topic segment. + + Home Assistant accepts only ``[A-Za-z0-9_-]+`` for node/object ids. + Replacing each run of other characters keeps ordinary UUID/key output + readable while ensuring composite identities cannot produce invalid topics. + """ + result = _HA_SEGMENT_RE.sub("_", value).strip("_") + return result or "home_automation" + + +def _node_id(identity: str) -> str: + """Stable, strictly valid MQTT node_id from an internal identity.""" + return _safe_segment(identity.replace("-", "_")) def _object_id(entity: ExposableEntity) -> str: - """Stable MQTT object_id derived from entity key (dots + hyphens → underscores).""" - return entity.key.replace(".", "_").replace("-", "_") + """Stable, strictly valid MQTT object_id derived from the entity key.""" + return _safe_segment(entity.key.replace("-", "_")) def _discovery_topic(entity: ExposableEntity, prefix: str) -> str: @@ -74,7 +99,7 @@ def _discovery_topic(entity: ExposableEntity, prefix: str) -> str: Format: ``////config`` """ - node = _node_id(entity.device.identifiers[1]) # device uuid + node = _node_id(entity.device.internal_identity) obj = _object_id(entity) return f"{prefix}/{entity.component}/{node}/{obj}/config" @@ -84,7 +109,7 @@ def _state_topic(entity: ExposableEntity, prefix: str) -> str: Format: ``////state`` """ - node = _node_id(entity.device.identifiers[1]) + node = _node_id(entity.device.internal_identity) obj = _object_id(entity) return f"{prefix}/{entity.component}/{node}/{obj}/state" @@ -104,16 +129,133 @@ def _availability_id(entity: ExposableEntity) -> str: M8 meters deliberately retain their own UUID as HA node/unique identity, while their availability is supplied by a MeterSource UUID. """ - return entity.device.availability_id or entity.device.identifiers[1] + return entity.device.availability_id or entity.device.internal_identity def _unique_id(entity: ExposableEntity) -> str: """Stable unique_id — device uuid + metric key (never from mutable fields).""" - device_uuid = entity.device.identifiers[1] + device_uuid = entity.device.internal_identity # entity.key is already "modbus.." — use it as the seed return f"{device_uuid}_{entity.key.replace('.', '_')}" +def _migration_state(session: Session, key: str, default: str) -> str: + """Read a versioned discovery-repair acknowledgement from ``app_config``.""" + entry = session.query(AppConfigEntry).filter(AppConfigEntry.key == key).one_or_none() + return entry.value if entry is not None else default + + +def _has_repair_store(session: Session) -> bool: + """Whether this session uses the app schema (not a lightweight unit DB).""" + return inspect(session.get_bind()).has_table("app_config") + + +def _set_migration_state(session: Session, key: str, value: str) -> None: + """Durably acknowledge a completed discovery-repair step. + + The state is operational metadata only; it never alters meters, readings, + contracts, costs, or expose toggles. + """ + entry = session.query(AppConfigEntry).filter(AppConfigEntry.key == key).one_or_none() + if entry is None: + session.add(AppConfigEntry(key=key, value=value, updated_at=datetime.now(UTC))) + else: + entry.value = value + entry.updated_at = datetime.now(UTC) + session.commit() + + +def _migration_json(session: Session, key: str) -> dict[str, Any]: + """Read a versioned JSON migration ledger, treating corrupt values as pending.""" + raw = _migration_state(session, key, "{}") + try: + value = json.loads(raw) + except json.JSONDecodeError: + logger.warning("Ignoring malformed HA discovery migration ledger %s", key) + return {} + return value if isinstance(value, dict) else {} + + +def _set_migration_json(session: Session, key: str, value: dict[str, Any]) -> None: + _set_migration_state(session, key, json.dumps(value, sort_keys=True)) + + +def _frozen_legacy_inventory(ledger: dict[str, Any]) -> list[str] | None: + """Return a structurally valid immutable cleanup inventory, if present. + + The ledger predates a schema migration, so it must tolerate both an absent + entry and old ``topics``-only values. An inventory becomes immutable only + once it is a list of non-empty topic strings; anything else is compatibility + data to be frozen during startup, not state a publisher may reinterpret. + """ + value = ledger.get("inventory") + if not isinstance(value, list) or not all(isinstance(topic, str) and topic for topic in value): + return None + return list(dict.fromkeys(value)) + + +def _publish_ok(topic: str, payload: str | bytes, *, retain: bool = True) -> bool: + """Publish best-effort while treating only an explicit ``False`` as failure. + + Production ``MqttManager.publish`` returns ``bool``. Accepting ``None`` + keeps existing third-party/mock publishers best-effort compatible. + """ + return mqtt_manager.publish(topic, payload, retain=retain) is not False + + +def initialize_legacy_thermal_cleanup(session: Session) -> None: + """Create the legacy-cleanup ledger before expose toggles can be edited. + + A newly installed database has no v1.6.1 retained thermal topics. Marking + that case complete during startup prevents a later first-time toggle from + manufacturing an old, illegal topic. Conversely, an upgraded database + freezes its exact enabled legacy topics before the UI can change a toggle. + """ + if not _has_repair_store(session): + return + + ledger = _migration_json(session, _LEGACY_THERMAL_CLEANUP_KEY) + if ledger.get("complete") is True: + return + + inventory = _frozen_legacy_inventory(ledger) + if inventory is not None: + # A previously frozen non-empty inventory must never be expanded or + # recomputed from mutable toggle/prefix state. A zero-item inventory + # is just the fresh-install terminal state written by an older build. + if not inventory: + _set_migration_json( + session, + _LEGACY_THERMAL_CLEANUP_KEY, + {"complete": True, "inventory": [], "topics": []}, + ) + return + + # This deliberately propagates. Startup runs before the expose UI can + # write a toggle, so swallowing an enumeration or durable-write error would + # create a window in which the cleanup scope could be changed or lost. + legacy_entities = _legacy_thermal_entities(session) + from app.config import get_settings + + settings = build_runtime_settings(session, get_settings()) + inventory = list(dict.fromkeys( + _legacy_discovery_topic(entity, settings.ha_discovery_prefix) + for entity in legacy_entities + )) + _set_migration_json( + session, + _LEGACY_THERMAL_CLEANUP_KEY, + { + "complete": not inventory, + "inventory": inventory, + "topics": sorted({ + topic for topic in ledger.get("topics", []) + if isinstance(topic, str) and topic in inventory + }), + }, + ) + + # --------------------------------------------------------------------------- # Public: build discovery payload # --------------------------------------------------------------------------- @@ -220,34 +362,90 @@ def publish_discovery(session: Session) -> None: logger.exception("publish_discovery: failed to build catalog; aborting") return - # Meter UUIDs are intentionally identity-changing epochs. Discovery config - # is retained, so clear only the precisely enumerable old M8 identities; - # never wildcard a provider/topic and risk removing another source's card. - try: - stale_entities = _stale_m8_entities(session) - except Exception: - logger.exception("publish_discovery: unable to enumerate old M8 identities") - stale_entities = [] - for old_entity in stale_entities: + # Clear the *illegal* v1.6.1 thermal keys once. Startup freezes the exact + # inventory while the legacy toggle state still describes v1.6.1. Thus a + # later UI toggle/prefix/meter change cannot erase or expand this repair. + # ``topics`` is only durable success progress; ``inventory`` is immutable. + # Startup freezes compatibility ledgers before the UI can mutate the + # toggle/prefix inputs. The publisher only consumes that frozen snapshot. + repaired_keys: set[str] = set() + if _has_repair_store(session): + ledger = _migration_json(session, _LEGACY_THERMAL_CLEANUP_KEY) + if not ledger.get("complete"): + inventory = _frozen_legacy_inventory(ledger) + if inventory is None: + logger.error( + "publish_discovery: legacy cleanup ledger was not frozen at startup; refusing cleanup" + ) + + if inventory is not None: + completed = { + topic for topic in ledger.get("topics", []) + if isinstance(topic, str) and topic in inventory + } + for topic in (topic for topic in inventory if topic not in completed): + try: + if _publish_ok(topic, b""): + completed.add(topic) + # Commit each accepted illegal-topic tombstone before + # attempting another one: a crash must not replay it. + _set_migration_json( + session, + _LEGACY_THERMAL_CLEANUP_KEY, + { + "complete": False, + "inventory": inventory, + "topics": sorted(completed), + }, + ) + else: + logger.warning("publish_discovery: broker rejected v1.6.1 cleanup for %s", topic) + except Exception: + logger.exception("publish_discovery: unable to clear v1.6.1 topic %s", topic) + if set(inventory).issubset(completed): + _set_migration_json( + session, + _LEGACY_THERMAL_CLEANUP_KEY, + {"complete": True, "inventory": inventory, "topics": sorted(completed)}, + ) + + # Current-format stale topics are legal and deliberately remain retryable: + # they cover post-repair meter swaps and include all 14 thermal metrics. try: - old_topic, _ = build_discovery_payload(old_entity, discovery_prefix, state_prefix) - mqtt_manager.publish(old_topic, b"", retain=True) + stale_entities = _stale_m8_entities(session) except Exception: - logger.exception("publish_discovery: unable to clear old identity %r", old_entity.key) + logger.exception("publish_discovery: unable to enumerate stale M8 identities") + stale_entities = [] + for old_entity in stale_entities: + try: + if not _publish_ok(_discovery_topic(old_entity, discovery_prefix), b""): + logger.warning("publish_discovery: broker rejected stale cleanup for %r", old_entity.key) + except Exception: + logger.exception("publish_discovery: unable to clear stale identity %r", old_entity.key) + + # v1.6.1 used the same legal topics/unique_ids for Meter, Source, Modbus + # and electricity-cost entities, but attached their registry entries to + # wrongly merged HA devices. HA does not move those entries when only a + # discovery ``device`` block changes. One durable unload -> re-add cycle + # lets HA rebuild the existing unique_id entries against the new singleton + # identifiers, preserving entity keys, toggles and user customizations. + repaired_keys = _run_registry_repair(session, catalog, discovery_prefix, state_prefix, settings) for entry in catalog: entity = entry.entity + if entity.key in repaired_keys: + continue try: topic, config = build_discovery_payload(entity, discovery_prefix, state_prefix) if entry.enabled: payload = json.dumps(config) - mqtt_manager.publish(topic, payload, retain=True) + _publish_ok(topic, payload) logger.debug( "publish_discovery: published config for %r → %s", entity.key, topic ) else: # Clear retained config for disabled entities. - mqtt_manager.publish(topic, b"", retain=True) + _publish_ok(topic, b"") logger.debug( "publish_discovery: cleared config for disabled entity %r → %s", entity.key, @@ -259,39 +457,163 @@ def publish_discovery(session: Session) -> None: ) -def _stale_m8_entities(session: Session) -> list[ExposableEntity]: - """Return synthetic discovery entries for superseded thermal identities. +def _is_registry_repair_entity(entity: ExposableEntity) -> bool: + """Whether a legal v1.6.1 config needs an unload/re-add device split.""" + return entity.key.startswith(("modbus.", "source.", "meter.", "energy.", "thermal_cost.")) - This is deliberately a narrow, best-effort cleanup: ended Meter UUIDs and - historically possible thermal combinations only; current identities are excluded. + +def _ha_registry_bindings(settings: Any, unique_ids: set[str]) -> dict[str, set[str]] | None: + """Read HA's registry, returning ``None`` when it is not observable.""" + from app.integrations.homeassistant import ( + HomeAssistantClient, + HomeAssistantConfigError, + HomeAssistantRequestError, + ) + + # Runtime settings from old callers/tests may not expose outbound HA + # fields. Treat absent/non-string credentials as intentionally optional. + if not isinstance(getattr(settings, "home_assistant_base_url", None), str) or not isinstance( + getattr(settings, "home_assistant_auth_token", None), str + ): + return None + client = HomeAssistantClient(settings) + if not client.is_configured(): + return None + try: + return client.discovery_registry_bindings(unique_ids) + except (HomeAssistantConfigError, HomeAssistantRequestError, TypeError, ValueError): + logger.warning("HA registry repair remains pending: HA registry is unavailable") + return None + + +def _run_registry_repair( + session: Session, + catalog: list[Any], + discovery_prefix: str, + state_prefix: str, + settings: Any, +) -> set[str]: + """Repair each v1.6.1 entity only after HA confirms every phase. + + MQTT acknowledges broker receipt, not Home Assistant processing. The + ledger therefore holds each target at ``unloaded`` until HA's entity + registry no longer contains its stable unique_id, then at ``republished`` + until HA reports the expected singleton device identifier. A partial + catalog simply adds targets on a later run; it cannot complete others. + When HA's registry is unavailable, ordinary discovery remains untouched. """ - from app.integrations.expose import DeviceInfo + targets = [entry for entry in catalog if _is_registry_repair_entity(entry.entity)] + if not targets: + return set() + ledger = _migration_json(session, _REGISTRY_REPAIR_KEY) + pending_targets = [ + entry for entry in targets if ledger.get(_unique_id(entry.entity), {}).get("phase") != "complete" + ] + if not pending_targets: + return set() + unique_ids = {_unique_id(entry.entity) for entry in pending_targets} + bindings = _ha_registry_bindings(settings, unique_ids) + if bindings is None: + return set() + blocked: set[str] = set() + for entry in pending_targets: + entity = entry.entity + unique_id = _unique_id(entity) + expected = set(entity.device.identifiers) + phase = ledger.get(unique_id, {}).get("phase", "pending") + observed = bindings.get(unique_id) + if phase == "unloaded": + if observed is not None: + # HA was disconnected or otherwise missed the retained + # tombstone. Keep it retained until HA itself confirms delete. + try: + _publish_ok(_discovery_topic(entity, discovery_prefix), b"") + except Exception: + logger.exception("publish_discovery: unable to repeat registry unload for %r", entity.key) + blocked.add(entity.key) + continue + if not entry.enabled: + _set_registry_phase(session, ledger, unique_id, "complete", expected) + continue + _publish_registry_config(session, ledger, entry, discovery_prefix, state_prefix, expected, blocked) + elif phase == "republished": + if observed == expected: + _set_registry_phase(session, ledger, unique_id, "complete", expected) + elif observed is None and entry.enabled: + _publish_registry_config(session, ledger, entry, discovery_prefix, state_prefix, expected, blocked) + else: + # A stale/wrong device binding must pass through deletion again. + try: + if _publish_ok(_discovery_topic(entity, discovery_prefix), b""): + _set_registry_phase(session, ledger, unique_id, "unloaded", expected) + except Exception: + logger.exception("publish_discovery: unable to unload wrong registry entity %r", entity.key) + blocked.add(entity.key) + elif phase != "complete": + if observed == expected or (observed is None and not entry.enabled): + _set_registry_phase(session, ledger, unique_id, "complete", expected) + elif observed is None: + _publish_registry_config(session, ledger, entry, discovery_prefix, state_prefix, expected, blocked) + else: + try: + if _publish_ok(_discovery_topic(entity, discovery_prefix), b""): + _set_registry_phase(session, ledger, unique_id, "unloaded", expected) + blocked.add(entity.key) + else: + logger.warning("publish_discovery: broker rejected registry unload for %r", entity.key) + except Exception: + logger.exception("publish_discovery: unable to unload registry entity %r", entity.key) + return blocked + + +def _set_registry_phase( + session: Session, ledger: dict[str, Any], unique_id: str, phase: str, expected: set[str] +) -> None: + ledger[unique_id] = {"phase": phase, "identifier": sorted(expected)} + _set_migration_json(session, _REGISTRY_REPAIR_KEY, ledger) + + +def _publish_registry_config( + session: Session, ledger: dict[str, Any], entry: Any, discovery_prefix: str, + state_prefix: str, expected: set[str], blocked: set[str], +) -> None: + entity = entry.entity + topic, config = build_discovery_payload(entity, discovery_prefix, state_prefix) + try: + if _publish_ok(topic, json.dumps(config)): + _set_registry_phase(session, ledger, _unique_id(entity), "republished", expected) + else: + logger.warning("publish_discovery: broker rejected registry re-add for %r", entity.key) + except Exception: + logger.exception("publish_discovery: unable to re-add registry entity %r", entity.key) + blocked.add(entity.key) + + +def _legacy_discovery_topic(entity: ExposableEntity, prefix: str) -> str: + """Return the exact v1.6.1 config topic for a known stale M8 entity. + + This intentionally preserves the old hyphen-only conversion because the + point is to remove that exact retained broker key once, never to publish a + wildcard or manufacture a new invalid topic. + """ + node = entity.device.internal_identity.replace("-", "_") + obj = entity.key.replace(".", "_").replace("-", "_") + return f"{prefix}/{entity.component}/{node}/{obj}/config" + + +def _overlapping_thermal_pairs(session: Session) -> list[tuple[Any, Any]]: + """Return only heating/water Meter epochs that could have coexisted.""" from app.models.energy import Meter from sqlalchemy import select meters = session.execute(select(Meter).where( Meter.commodity.in_(("electricity", "heating", "hot_water")) )).scalars().all() - current = {meter.commodity: meter for meter in meters if meter.ended_at is None} - old = [meter for meter in meters if meter.ended_at is not None] - entities: list[ExposableEntity] = [] - for meter in old: - info = DeviceInfo(identifiers=("meter", meter.uuid), name=meter.label) - for suffix in ("total", "today"): - entities.append(ExposableEntity( - key=f"meter.{meter.uuid}.{suffix}", component="sensor", device=info, - device_class=None, unit="", name="obsolete", - )) heatings = [meter for meter in meters if meter.commodity == "heating"] waters = [meter for meter in meters if meter.commodity == "hot_water"] - current_identity = ( - ".".join(sorted((current["heating"].uuid, current["hot_water"].uuid))) - if current.get("heating") is not None and current.get("hot_water") is not None else None - ) + pairs: list[tuple[Any, Any]] = [] for heating in heatings: for water in waters: - if heating.ended_at is None and water.ended_at is None: - continue # A thermal identity can only have been published when both Meter # epochs were current at the same instant. Do not form a Cartesian # product of historical records: that would tombstone identities @@ -302,16 +624,86 @@ def _stale_m8_entities(session: Session) -> list[ExposableEntity]: water_end is not None and heating_start >= water_end ): continue - identity = ".".join(sorted((heating.uuid, water.uuid))) - if identity == current_identity: - continue - info = DeviceInfo(identifiers=("thermal-cost", identity), name="obsolete") - for metric in ("heating", "hot_water_heating", "water", "water_tax", "fixed", "all_in"): - for suffix in ("total", "today"): - entities.append(ExposableEntity( - key=f"thermal_cost.{identity}.{metric}_{suffix}", component="sensor", device=info, - device_class=None, unit="", name="obsolete", - )) + pairs.append((heating, water)) + return pairs + + +def _thermal_cleanup_entities( + pairs: list[tuple[Any, Any]], *, include_hot_water_total: bool +) -> list[ExposableEntity]: + """Build exact synthetic config entries for known thermal identities.""" + from app.integrations.expose import DeviceInfo + + metrics = ["heating", "hot_water_heating", "water", "water_tax", "fixed", "all_in"] + if include_hot_water_total: + metrics.insert(2, "hot_water_total") + entities: list[ExposableEntity] = [] + for heating, water in pairs: + identity = ".".join(sorted((heating.uuid, water.uuid))) + info = DeviceInfo( + identifiers=(f"home-automation:thermal-cost:{identity}",), + name="obsolete", + identity=identity, + ) + for metric in metrics: + for suffix in ("total", "today"): + entities.append(ExposableEntity( + key=f"thermal_cost.{identity}.{metric}_{suffix}", component="sensor", device=info, + device_class=None, unit="", name="obsolete", + )) + return entities + + +def _legacy_thermal_entities(session: Session) -> list[ExposableEntity]: + """Return only v1.6.1 illegal topics that could retain a non-empty config. + + v1.6.1 published an illegal thermal config only while its expose toggle was + enabled. A disabled toggle published its own tombstone, so querying the + durable toggle state prevents a fresh install from manufacturing warnings + for twelve never-used illegal topics. + """ + from app.models.expose import ExposedEntityToggle + + entities = _thermal_cleanup_entities(_overlapping_thermal_pairs(session), include_hot_water_total=False) + keys = [entity.key for entity in entities] + enabled_keys = { + row.key for row in session.query(ExposedEntityToggle).filter( + ExposedEntityToggle.key.in_(keys), ExposedEntityToggle.enabled.is_(True) + ) + } + return [entity for entity in entities if entity.key in enabled_keys] + + +def _stale_m8_entities(session: Session) -> list[ExposableEntity]: + """Return *safe-format* configs that became stale after a later epoch swap. + + Unlike the v1.6.1 cleanup, this intentionally excludes the active thermal + pair and includes the post-repair ``hot_water_total`` metric (14 configs + per stale pair). All generated topics use ``_discovery_topic``. + """ + from app.integrations.expose import DeviceInfo + from app.models.energy import Meter + from sqlalchemy import select + + meters = session.execute(select(Meter).where( + Meter.commodity.in_(("electricity", "heating", "hot_water")) + )).scalars().all() + entities: list[ExposableEntity] = [] + for meter in (meter for meter in meters if meter.ended_at is not None): + info = DeviceInfo( + identifiers=(f"home-automation:meter:{meter.uuid}",), name=meter.label, identity=meter.uuid + ) + for suffix in ("total", "today"): + entities.append(ExposableEntity( + key=f"meter.{meter.uuid}.{suffix}", component="sensor", device=info, + device_class=None, unit="", name="obsolete", + )) + stale_pairs = [ + (heating, water) + for heating, water in _overlapping_thermal_pairs(session) + if heating.ended_at is not None or water.ended_at is not None + ] + entities.extend(_thermal_cleanup_entities(stale_pairs, include_hot_water_total=True)) return entities @@ -468,7 +860,7 @@ def publish_device_state(session: Session, device: Any) -> None: continue entity = entry.entity # Only process entities belonging to this device. - if entity.device.identifiers[1] != device_uuid: + if entity.device.internal_identity != device_uuid: continue # Skip the online binary_sensor itself (availability handled above). if entity.component == "binary_sensor" and "online" in entity.key: @@ -543,7 +935,7 @@ def clear_device_discovery(session: Session, device_uuid: str) -> None: for entry in catalog: entity = entry.entity # Only clear entities belonging to this device. - if entity.device.identifiers[1] != device_uuid: + if entity.device.internal_identity != device_uuid: continue try: topic, _config = build_discovery_payload(entity, discovery_prefix, state_prefix) diff --git a/docs/design/m8-warmtelink-energy.md b/docs/design/m8-warmtelink-energy.md index 7d4265a..de8b083 100644 --- a/docs/design/m8-warmtelink-energy.md +++ b/docs/design/m8-warmtelink-energy.md @@ -1214,6 +1214,10 @@ M8 收尾的前置条件。agent 不得执行、记录为已执行,或以 mock ## 14. Post-M8 lifecycle repair(M8-R08~R10) +### M8-R15 — HA Discovery identity/topic repair + +修复 HA 将多值 `device.identifiers` 任一匹配合并的风险:Source、Meter、Modbus 与成本 epoch 均使用单一完整 namespaced identifier,内部 MQTT identity 独立保存。thermal 的既有点号 identity/key 保持兼容,node/object topic segment 单独规范化;仅对 v1.6.1 实际发布过的 dot topic 做可重试、成功后抑制的精确 retained cleanup。新增默认关闭的 `Thermal Hot Water Total/Today`,金额为 `hot_water_heating + hot_water`,不含 `hot_water_tax`。 + M8 交付后的 Meter lifecycle 修复链记录在本地 `review-notes/M8-meter-lifecycle-repair-plan.md`。 它不新增 ORM / **数据库** schema 或 Alembic revision,不做启动自动修复、一次性数据脚本或历史删除;R08 虽然 更新了 API/Pydantic schema 及 OpenAPI/codegen,但没有变更 ORM 或数据库 schema。已有的 stranded binding 只能 diff --git a/docs/homeassistant-outbound.md b/docs/homeassistant-outbound.md index ab5e5f1..88cd235 100644 --- a/docs/homeassistant-outbound.md +++ b/docs/homeassistant-outbound.md @@ -55,7 +55,9 @@ Expose 框架还可以把已勾选的 Energy 实体通过 MQTT Home Assistant Discovery 发布;开关位于应用 Config 页的 HA Expose 面板,默认均为关闭。M8 增加了 source online、按 Meter UUID 锚定的累计量/today,以及 heating、hot-water-heating、water、water-tax、fixed、all-in total/today 等 thermal 实体。 - source 和 Meter identity 不依赖可变 label;换表会产生新 Meter UUID identity。 -- thermal 组合成本 identity 由当前 heating/hot_water Meter UUID 的有序组合锚定,任一换表都会产生新 identity,避免不同累计域拼接。 +- thermal 组合成本 identity 由当前 heating/hot_water Meter UUID 的有序组合锚定,任一换表都会产生新 identity,避免不同累计域拼接。`Thermal Hot Water Total/Today` 精确为 `hot_water_heating + hot_water`,不包含 `hot_water_tax`。 +- 每个 HA device 只发布一个完整、`home-automation:` namespaced identifier;内部 topic/unique-id seed 与该 identifier 分离。所有 discovery node/object segment 都会转为 HA 允许的 `[A-Za-z0-9_-]+` 字符集。 +- v1.6.1 的 dot-containing thermal retained topics 会在启动、UI 尚未能修改 toggle 前,按当时实际重叠的 Meter epoch、enabled toggle 和 runtime discovery prefix 冻结为精确清单;每个成功 topic 都会持久记账,失败/未尝试项才会重试。冻结清单完成后跨重启不再枚举或发布旧 topic;fresh install 的空清单也会立即完成。没有 wildcard,也不会清理任何应用数据。旧 device 合并修复则以已配置的 Home Assistant WebSocket entity/device registry 实际观察 unload、重建和正确 device identifier;HA 不可达时普通 discovery 继续发布,repair 保持 pending。 - availability、unit、device/state class 与 today reset 由 provider 声明;operator 应在 HA 中核对,而不应假设同名实体可跨换表连续。 - 关闭 toggle 后 retained discovery 会被清理;关闭暴露不删除 source、Meter、合同、读数或成本历史。 diff --git a/tests/test_app.py b/tests/test_app.py index aaac8d1..36b4efc 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -118,6 +118,35 @@ def test_app_start_seeds_missing_config_from_env_without_overwriting_existing_va reset_db_caches() +def test_startup_initializes_fresh_legacy_discovery_cleanup_ledger( + tmp_path, monkeypatch: pytest.MonkeyPatch +) -> None: + """The cleanup ledger exists before the expose UI can write its first toggle.""" + import app.main as main + + app_database_url = _prepare_app_db(tmp_path) + monkeypatch.setenv("APP_DATABASE_URL", app_database_url) + monkeypatch.setenv("AUTH_BOOTSTRAP_USERNAME", "admin") + monkeypatch.setenv("AUTH_BOOTSTRAP_PASSWORD", "test-password") + get_settings.cache_clear() + reset_db_caches() + + main.ensure_auth_db_ready() + + conn = sqlite3.connect(tmp_path / "app_ready.db") + try: + value = conn.execute( + "SELECT value FROM app_config WHERE key = ?", + ("HA_DISCOVERY_LEGACY_THERMAL_CLEANUP_V1",), + ).fetchone()[0] + finally: + conn.close() + assert value == '{"complete": true, "inventory": [], "topics": []}' + + get_settings.cache_clear() + reset_db_caches() + + def test_app_start_syncs_app_hostname_from_env_even_when_db_has_old_value( tmp_path, monkeypatch: pytest.MonkeyPatch ) -> None: diff --git a/tests/test_energy_expose.py b/tests/test_energy_expose.py index c315c6c..9021bf3 100644 --- a/tests/test_energy_expose.py +++ b/tests/test_energy_expose.py @@ -13,7 +13,7 @@ Coverage: 7. value_getter returns None when no non-degraded period exists. 8. MQTT not enabled → publish_states is a no-op (no raises, no publish calls). 9. Integration: build_discovery_payload on an energy_cost entity does NOT raise - IndexError (validates 2-element identifiers). + IndexError (validates one-item HA identifiers and independent internal identities). 10. Keys are stable fixed strings (not derived from mutable data or DB ids). 11. Provider registered: energy_cost entities appear alongside modbus entities in the full catalog. @@ -767,9 +767,8 @@ def test_publish_states_noop_when_not_connected() -> None: def test_build_discovery_payload_no_index_error_for_energy_entities(energy_db) -> None: """build_discovery_payload must NOT raise IndexError for energy_cost entities. - FUE-T05: identifiers is now ('energy-cost', meter.uuid). - Validates that the 2-element identifiers tuple satisfies ha_discovery.py's - requirement to access identifiers[1] as node_id. + The HA grouping identifier and internal MQTT identity are deliberately + separate; a one-item HA identifier must therefore remain sufficient. """ from app.integrations.expose import build_catalog from app.services.ha_discovery import build_discovery_payload @@ -800,15 +799,11 @@ def test_build_discovery_payload_no_index_error_for_energy_entities(energy_db) - assert len(energy_entries) == 6, "Expected 6 energy_cost entities in catalog" for entry in energy_entries: - # identifiers[1] must be the meter uuid (not "energy-cost"). - assert entry.entity.device.identifiers[1] == meter_uuid, ( - f"identifiers[1] must be meter uuid {meter_uuid!r}, " - f"got {entry.entity.device.identifiers[1]!r}" - ) - # Must not raise — specifically no IndexError from identifiers[1] + assert entry.entity.device.internal_identity == meter_uuid + assert entry.entity.device.identifiers == (f"home-automation:energy-cost:{meter_uuid}",) topic, config = build_discovery_payload(entry.entity, "homeassistant") - # node_id = identifiers[1] with hyphens → underscores + # node_id = internal meter identity with hyphens → underscores node_id = meter_uuid.replace("-", "_") assert node_id in topic, ( f"Expected meter uuid node_id {node_id!r} in discovery topic, got {topic!r}" @@ -817,13 +812,12 @@ def test_build_discovery_payload_no_index_error_for_energy_entities(energy_db) - f"Discovery topic must end with /config, got {topic!r}" ) assert "unique_id" in config - # unique_id seed is identifiers[1] (meter uuid) + entity key + # unique_id seed is the internal meter identity + entity key assert meter_uuid in config["unique_id"], ( f"unique_id must contain meter uuid, got {config['unique_id']!r}" ) assert "device" in config - assert "energy-cost" in config["device"]["identifiers"] - assert meter_uuid in config["device"]["identifiers"] + assert config["device"]["identifiers"] == [f"home-automation:energy-cost:{meter_uuid}"] def test_energy_cost_entities_omit_availability_so_ha_shows_them(energy_db) -> None: @@ -859,18 +853,18 @@ def test_energy_cost_entities_omit_availability_so_ha_shows_them(energy_db) -> N def test_energy_entity_discovery_topics_contain_correct_node_id() -> None: """Discovery topic node_id for energy entities must be derived from meter uuid. - FUE-T05: identifiers[1] is now the active meter's uuid. - ha_discovery._node_id() replaces hyphens with underscores in identifiers[1] - to build the MQTT node_id. This test verifies that the topic reflects the + The internal identity is the active meter's uuid, independent from the + singleton HA identifier. This test verifies that the topic reflects the meter uuid (not the old fixed 'energy-cost' string). """ from app.integrations.expose import DeviceInfo, ExposableEntity from app.services.ha_discovery import build_discovery_payload meter_uuid = "12345678-abcd-ef00-1234-567890abcdef" - # identifiers[1] = meter uuid — this is what FUE-T05 sets. + # ``identity`` is the meter UUID; HA grouping is a separate one-item tuple. device = DeviceInfo( - identifiers=("energy-cost", meter_uuid), + identifiers=(f"home-automation:energy-cost:{meter_uuid}",), + identity=meter_uuid, name="Test Meter", provides_availability=False, ) @@ -886,7 +880,7 @@ def test_energy_entity_discovery_topics_contain_correct_node_id() -> None: topic, config = build_discovery_payload(entity, discovery_prefix="homeassistant") - # node_id: identifiers[1] = meter_uuid, hyphens → underscores + # node_id: internal identity = meter_uuid, hyphens → underscores expected_node = meter_uuid.replace("-", "_") assert f"/{expected_node}/" in topic, ( f"Expected meter uuid node_id {expected_node!r} in topic {topic!r}" @@ -2345,12 +2339,11 @@ def test_energy_cost_provider_returns_empty_when_no_active_meter(energy_db) -> N ) -def test_energy_cost_provider_identifiers_match_meter_uuid(energy_db) -> None: - """FUE-T05 ②: with an active meter, identifiers[1] == meter.uuid. +def test_energy_cost_provider_identity_matches_meter_uuid(energy_db) -> None: + """The active meter anchors the internal identity and namespaced HA identifier. - The HA device identity is anchored to the active meter's uuid. ha_discovery.py - uses identifiers[1] as the MQTT node_id and unique_id seed; changing the active - meter (meter swap) produces a new uuid → new node_id → new HA sensor. + The HA device identity and MQTT internal identity are both anchored to the + active meter's uuid. Changing the meter produces a new node_id and HA sensor. """ from app.integrations.expose import _energy_cost_provider @@ -2367,11 +2360,8 @@ def test_energy_cost_provider_identifiers_match_meter_uuid(energy_db) -> None: assert len(entities) == 6, f"Expected 6 entities, got {len(entities)}" for entity in entities: - assert entity.device.identifiers == ("energy-cost", meter_uuid), ( - f"identifiers must be ('energy-cost', meter.uuid); " - f"expected ('energy-cost', {meter_uuid!r}), " - f"got {entity.device.identifiers!r}" - ) + assert entity.device.internal_identity == meter_uuid + assert entity.device.identifiers == (f"home-automation:energy-cost:{meter_uuid}",) assert entity.device.name == meter_label, ( f"device.name must be exactly the meter label {meter_label!r}, " f"got {entity.device.name!r}" @@ -2382,7 +2372,7 @@ def test_energy_cost_entity_keys_do_not_contain_meter_uuid(energy_db) -> None: """FUE-T05 ③: entity keys remain stable 'energy.*' strings (no uuid injected). Keys are the anchor for the toggle table; they must NOT change when the meter - changes. Only identifiers[1] (node_id / unique_id) changes on a meter swap. + changes. Only the internal identity (node_id / unique_id) changes on a meter swap. """ from app.integrations.expose import _energy_cost_provider @@ -2418,7 +2408,7 @@ def test_energy_cost_entity_keys_do_not_contain_meter_uuid(energy_db) -> None: def test_energy_cost_toggle_survives_meter_swap(energy_db) -> None: """FUE-T05 ③ (toggle stability): enabled toggle on 'energy.buy_price_now' survives meter swap. - After a meter swap the provider queries a new active meter → new identifiers[1] / + After a meter swap the provider queries a new active meter → new internal identity / unique_id / topic in HA. But the entity key stays 'energy.buy_price_now', so the existing toggle row (keyed by 'energy.buy_price_now') is still found → enabled=True. @@ -2475,11 +2465,7 @@ def test_energy_cost_toggle_survives_meter_swap(energy_db) -> None: (e for e in catalog if e.entity.key == "energy.buy_price_now"), None ) assert buy_entry is not None, "energy.buy_price_now must be present in catalog" - # identifiers[1] must now be the NEW meter's uuid - assert buy_entry.entity.device.identifiers[1] == new_meter_uuid, ( - f"After swap, identifiers[1] must be new meter uuid {new_meter_uuid!r}, " - f"got {buy_entry.entity.device.identifiers[1]!r}" - ) + assert buy_entry.entity.device.internal_identity == new_meter_uuid # Toggle state must still be enabled (key unchanged → same toggle row found) assert buy_entry.enabled is True, ( "energy.buy_price_now toggle must remain enabled after meter swap " @@ -2488,7 +2474,7 @@ def test_energy_cost_toggle_survives_meter_swap(energy_db) -> None: def test_energy_cost_identifiers_change_after_meter_swap(energy_db) -> None: - """FUE-T05 ④: after meter swap, provider produces new identifiers[1] (new meter uuid). + """After a meter swap, provider produces a new internal meter identity. Old uuid's entities are no longer produced → HA sensor for old uuid is frozen. New uuid's entities appear → HA creates fresh sensors for the new meter. @@ -2544,12 +2530,9 @@ def test_energy_cost_identifiers_change_after_meter_swap(energy_db) -> None: f"Expected 6 entities after swap, got {len(entities_after_swap)}" ) for entity in entities_after_swap: - assert entity.device.identifiers[1] == new_uuid, ( - f"After swap, identifiers[1] must be new uuid {new_uuid!r}, " - f"got {entity.device.identifiers[1]!r}" - ) + assert entity.device.internal_identity == new_uuid # Old uuid must not appear in identifiers - assert entity.device.identifiers[1] != old_uuid, ( + assert entity.device.internal_identity != old_uuid, ( f"After swap, old uuid {old_uuid!r} must not appear in identifiers" ) @@ -2809,8 +2792,22 @@ def test_m8_catalog_has_source_meter_and_thermal_entities_disabled(energy_db) -> assert entries[f"meter.{electricity.uuid}.total"].entity.unit == "kWh" assert entries[f"meter.{electricity.uuid}.total"].entity.device_class == "energy" assert entries[f"meter.{electricity.uuid}.today"].entity.state_class == "total_increasing" + device_ids = { + entries[f"source.{heating_source.uuid}.online"].entity.device.identifiers[0], + entries[f"meter.{heating.uuid}.total"].entity.device.identifiers[0], + entries[f"meter.{water.uuid}.total"].entity.device.identifiers[0], + entries[f"meter.{electricity.uuid}.total"].entity.device.identifiers[0], + } + assert len(device_ids) == 4 + assert all(len(entry.entity.device.identifiers) == 1 for entry in entries.values()) thermal = [entry for key, entry in entries.items() if key.startswith("thermal_cost.")] - assert len(thermal) == 12 + assert len(thermal) == 14 + assert {entry.entity.key.rsplit(".", 1)[-1] for entry in thermal} >= { + "hot_water_total_total", "hot_water_total_today" + } + assert {entry.entity.name for entry in thermal if ".hot_water_total_" in entry.entity.key} == { + "Thermal Hot Water Total", "Thermal Hot Water Today" + } assert all(entry.enabled is False for entry in thermal) assert all(entry.entity.unit == "EUR" and entry.entity.device_class == "monetary" for entry in thermal) @@ -2972,9 +2969,10 @@ def test_m8_thermal_today_summary_ends_at_frozen_now(energy_db) -> None: captured: dict[str, datetime] = {} result = { "period_count": 1, "fixed_cost": Decimal("0"), "total_cost": Decimal("0"), - "breakdown": {key: Decimal("0") for key in ( - "heating", "hot_water_heating", "hot_water", "hot_water_tax", - )}, + "breakdown": { + "heating": Decimal("0"), "hot_water_heating": Decimal("1.25"), + "hot_water": Decimal("2.75"), "hot_water_tax": Decimal("9.99"), + }, } def summarize_spy(_session: Session, start: datetime, end: datetime, *, now: datetime) -> dict: @@ -2994,4 +2992,7 @@ def test_m8_thermal_today_summary_ends_at_frozen_now(energy_db) -> None: entity = next(item.entity for item in build_catalog(session) if item.entity.key.endswith(".heating_today")) assert entity.value_getter(session) == Decimal("0") + hot_water_total = next(item.entity for item in build_catalog(session) + if item.entity.key.endswith(".hot_water_total_today")) + assert hot_water_total.value_getter(session) == Decimal("4.00") assert captured["end"] == now diff --git a/tests/test_expose_catalog.py b/tests/test_expose_catalog.py index 4f27aa6..78dc7f7 100644 --- a/tests/test_expose_catalog.py +++ b/tests/test_expose_catalog.py @@ -362,7 +362,21 @@ def test_build_catalog_device_grouping(expose_db): assert len(identifiers_set) == 1, ( "All entities for one device must share the same DeviceInfo identifiers" ) - assert identifiers_set.pop() == ("modbus", test_uuid) + own_identifiers = identifiers_set.pop() + assert own_identifiers == (f"home-automation:modbus:{test_uuid}",) + + # A second Modbus device must never overlap its singleton HA identifier. + other_uuid = "aaaaaaaa-0000-0000-0000-000000000099" + with Session(expose_db) as session: + _make_modbus_device(session, friendly_name="SDM120 E", uuid=other_uuid) + session.commit() + with Session(expose_db) as session: + second_catalog = build_catalog(session) + other_identifiers = { + entry.entity.device.identifiers for entry in second_catalog if other_uuid in entry.entity.key + } + assert other_identifiers == {(f"home-automation:modbus:{other_uuid}",)} + assert {own_identifiers} != other_identifiers def test_build_catalog_metric_metadata_from_profile(expose_db): diff --git a/tests/test_ha_discovery.py b/tests/test_ha_discovery.py index 6ad8272..235d5e8 100644 --- a/tests/test_ha_discovery.py +++ b/tests/test_ha_discovery.py @@ -15,6 +15,7 @@ Coverage: from __future__ import annotations +import json from datetime import datetime, timedelta, timezone from pathlib import Path from typing import Any @@ -443,14 +444,14 @@ def test_publish_discovery_enabled_vs_disabled_payload(disco_db) -> None: voltage_config_topic = f"homeassistant/sensor/{node_id}/{voltage_obj_id}/config" # The voltage entity's config topic should have non-empty JSON payload - voltage_call = next( - ((t, p, r) for t, p, r in config_calls if t == voltage_config_topic), None - ) - assert voltage_call is not None, ( + voltage_calls = [(t, p, r) for t, p, r in config_calls if t == voltage_config_topic] + assert voltage_calls, ( f"Expected config publish for voltage topic {voltage_config_topic!r}. " f"Got topics: {[t for t, _, _ in config_calls]}" ) - _t, payload, retain = voltage_call + # First v1.6.1 repair run deliberately unloads the old config before it + # re-adds the same unique_id under the corrected HA device identifier. + _t, payload, retain = voltage_calls[-1] assert payload not in (b"", "", None), "Enabled entity should get non-empty config payload" assert retain is True, "Discovery config must be retained" @@ -1378,3 +1379,497 @@ def test_stale_m8_entities_includes_ended_electricity_meter(disco_db) -> None: session.commit() keys = {entity.key for entity in _stale_m8_entities(session)} assert keys == {f"meter.{old.uuid}.total", f"meter.{old.uuid}.today"} + + +def test_discovery_uses_single_namespaced_identifier_and_safe_thermal_topic() -> None: + """HA device grouping is independent from a dot-containing internal seed.""" + from app.integrations.expose import DeviceInfo, ExposableEntity + from app.services.ha_discovery import build_discovery_payload + + identity = "aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa.bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb" + device = DeviceInfo( + identifiers=(f"home-automation:thermal-cost:{identity}",), + identity=identity, + name="Thermal Energy Cost", + provides_availability=False, + ) + entity = ExposableEntity( + key=f"thermal_cost.{identity}.heating_total", component="sensor", device=device, + device_class="monetary", unit="EUR", name="Thermal Heating Total", state_class="total", + ) + topic, config = build_discovery_payload(entity, "homeassistant") + assert config["device"]["identifiers"] == [f"home-automation:thermal-cost:{identity}"] + assert all(part.replace("-", "").replace("_", "").isalnum() + for part in topic.split("/")[2:4]) + assert "." not in topic + + +def test_legacy_thermal_cleanup_persists_each_success_and_retries_only_failure(disco_db) -> None: + """Illegal v1.6.1 cleanup never re-sends a durably accepted topic.""" + from app.integrations.expose import DeviceInfo, ExposableEntity + from app.services import ha_discovery + from app.models.config import AppConfigEntry + + identity = "aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa.bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb" + legacy = ExposableEntity( + key=f"thermal_cost.{identity}.heating_total", component="sensor", + device=DeviceInfo(identifiers=("thermal-cost", identity), identity=identity, name="obsolete"), + device_class=None, unit="", name="obsolete", + ) + second_legacy = ExposableEntity( + key=f"thermal_cost.{identity}.heating_today", component="sensor", + device=legacy.device, device_class=None, unit="", name="obsolete", + ) + settings = _make_settings() + manager = _make_mock_manager() + with ( + patch("app.services.ha_discovery.build_runtime_settings", return_value=settings), + patch("app.services.ha_discovery.mqtt_manager", manager), + patch("app.services.ha_discovery.build_catalog", return_value=[]), + patch("app.services.ha_discovery._legacy_thermal_entities", return_value=[legacy, second_legacy]), + patch("app.services.ha_discovery._stale_m8_entities", return_value=[]), + ): + with Session(disco_db) as session: + ha_discovery.initialize_legacy_thermal_cleanup(session) + manager.publish.side_effect = [True, False] + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + assert manager.publish.call_count == 2 + with Session(disco_db) as session: + progress = session.query(AppConfigEntry).filter_by( + key=ha_discovery._LEGACY_THERMAL_CLEANUP_KEY + ).one() + assert ha_discovery._legacy_discovery_topic(legacy, "homeassistant") in progress.value + + manager.publish.reset_mock() + manager.publish.side_effect = None + manager.publish.return_value = True + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + assert manager.publish.call_args_list == [ + ((ha_discovery._legacy_discovery_topic(second_legacy, "homeassistant"), b""), {"retain": True}), + ] + + # A fresh Session simulates a process restart: no illegal topic is sent. + manager.publish.reset_mock() + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + manager.publish.assert_not_called() + + +def test_legacy_thermal_enumeration_failure_does_not_advance_marker(disco_db) -> None: + """An inventory error is pending work, never an empty successful cleanup.""" + from app.models.config import AppConfigEntry + from app.services import ha_discovery + + settings = _make_settings() + manager = _make_mock_manager() + with ( + patch("app.services.ha_discovery.build_runtime_settings", return_value=settings), + patch("app.services.ha_discovery.mqtt_manager", manager), + patch("app.services.ha_discovery.build_catalog", return_value=[]), + patch("app.services.ha_discovery._legacy_thermal_entities", side_effect=RuntimeError("enumeration failed")), + patch("app.services.ha_discovery._stale_m8_entities", return_value=[]), + ): + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + with Session(disco_db) as session: + assert session.query(AppConfigEntry).filter_by( + key=ha_discovery._LEGACY_THERMAL_CLEANUP_KEY + ).one_or_none() is None + + +def test_legacy_cleanup_compatibility_freeze_keeps_previous_success_progress(disco_db) -> None: + """Startup freezes a pre-inventory ledger and preserves its success progress.""" + from app.integrations.expose import DeviceInfo, ExposableEntity + from app.models.config import AppConfigEntry + from app.services import ha_discovery + + device = DeviceInfo(identifiers=("legacy",), identity="old", name="obsolete") + first = ExposableEntity(key="thermal_cost.old.heating_total", component="sensor", device=device, + device_class=None, unit="", name="obsolete") + second = ExposableEntity(key="thermal_cost.old.heating_today", component="sensor", device=device, + device_class=None, unit="", name="obsolete") + first_topic = ha_discovery._legacy_discovery_topic(first, "homeassistant") + with Session(disco_db) as session: + session.add(AppConfigEntry( + key=ha_discovery._LEGACY_THERMAL_CLEANUP_KEY, + value=json.dumps({"complete": False, "topics": [first_topic]}), + updated_at=datetime.now(timezone.utc), + )) + session.commit() + + with ( + patch("app.services.ha_discovery.build_runtime_settings", return_value=_make_settings()), + patch("app.services.ha_discovery._legacy_thermal_entities", return_value=[first, second]), + ): + with Session(disco_db) as session: + ha_discovery.initialize_legacy_thermal_cleanup(session) + + manager = _make_mock_manager() + with ( + patch("app.services.ha_discovery.build_runtime_settings", return_value=_make_settings()), + patch("app.services.ha_discovery.mqtt_manager", manager), + patch("app.services.ha_discovery.build_catalog", return_value=[]), + patch("app.services.ha_discovery._stale_m8_entities", return_value=[]), + patch("app.services.ha_discovery._ha_registry_bindings", return_value=None), + ): + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + assert manager.publish.call_args_list == [ + ((ha_discovery._legacy_discovery_topic(second, "homeassistant"), b""), {"retain": True}), + ] + with Session(disco_db) as session: + ledger = ha_discovery._migration_json(session, ha_discovery._LEGACY_THERMAL_CLEANUP_KEY) + assert ledger == { + "complete": True, + "inventory": [first_topic, ha_discovery._legacy_discovery_topic(second, "homeassistant")], + "topics": sorted((first_topic, ha_discovery._legacy_discovery_topic(second, "homeassistant"))), + } + + +def test_legacy_cleanup_startup_marks_fresh_install_complete_and_short_circuits(disco_db) -> None: + """A later first toggle cannot make a v1.6.1 topic after fresh startup.""" + from app.models.config import AppConfigEntry + from app.services import ha_discovery + + manager = _make_mock_manager() + settings = _make_settings() + with patch("app.services.ha_discovery._legacy_thermal_entities", return_value=[]): + with Session(disco_db) as session: + ha_discovery.initialize_legacy_thermal_cleanup(session) + with Session(disco_db) as session: + ledger = session.query(AppConfigEntry).filter_by( + key=ha_discovery._LEGACY_THERMAL_CLEANUP_KEY + ).one() + assert ledger.value == '{"complete": true, "inventory": [], "topics": []}' + + # A fresh Session is equivalent to a restarted process. If a user now + # enables a thermal entity, the completed ledger prevents re-enumeration. + with ( + patch("app.services.ha_discovery.build_runtime_settings", return_value=settings), + patch("app.services.ha_discovery.mqtt_manager", manager), + patch("app.services.ha_discovery.build_catalog", return_value=[]), + patch("app.services.ha_discovery._legacy_thermal_entities") as legacy, + patch("app.services.ha_discovery._stale_m8_entities", return_value=[]), + patch("app.services.ha_discovery._ha_registry_bindings", return_value=None), + ): + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + legacy.assert_not_called() + manager.publish.assert_not_called() + + +def test_legacy_cleanup_startup_completes_empty_compat_ledger_before_later_enable(disco_db) -> None: + """An Alembic-head empty compat ledger cannot manufacture a later old topic.""" + from app.models.config import AppConfigEntry + from app.models.energy import Meter + from app.services import ha_discovery + + now = datetime.now(timezone.utc) + with Session(disco_db) as session: + heating = Meter(label="heating", commodity="heating", started_at=now, ended_at=None, + reason="initial", note=None, created_at=now) + water = Meter(label="water", commodity="hot_water", started_at=now, ended_at=None, + reason="initial", note=None, created_at=now) + session.add_all((heating, water)) + session.add(AppConfigEntry( + key=ha_discovery._LEGACY_THERMAL_CLEANUP_KEY, + value=json.dumps({"complete": False, "topics": []}), + updated_at=now, + )) + session.commit() + + settings = _make_settings(ha_discovery_prefix="startup_prefix") + with patch("app.services.ha_discovery.build_runtime_settings", return_value=settings): + with Session(disco_db) as session: + ha_discovery.initialize_legacy_thermal_cleanup(session) + with Session(disco_db) as session: + ledger = ha_discovery._migration_json(session, ha_discovery._LEGACY_THERMAL_CLEANUP_KEY) + assert ledger == {"complete": True, "inventory": [], "topics": []} + active_meters = session.query(Meter).filter(Meter.ended_at.is_(None)).all() + pair = ha_discovery._thermal_cleanup_entities( + [(next(meter for meter in active_meters if meter.commodity == "heating"), + next(meter for meter in active_meters if meter.commodity == "hot_water"))], + include_hot_water_total=False, + ) + _enable_entity(session, pair[0].key) + session.commit() + + manager = _make_mock_manager() + with ( + patch("app.services.ha_discovery.build_runtime_settings", + return_value=_make_settings(ha_discovery_prefix="changed_prefix")), + patch("app.services.ha_discovery.mqtt_manager", manager), + patch("app.services.ha_discovery.build_catalog", return_value=[]), + patch("app.services.ha_discovery._stale_m8_entities", return_value=[]), + patch("app.services.ha_discovery._ha_registry_bindings", return_value=None), + ): + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + manager.publish.assert_not_called() + + +def test_legacy_cleanup_startup_propagates_inventory_failure(disco_db) -> None: + """Startup must fail closed instead of exposing a mutable cleanup window.""" + from app.services import ha_discovery + + with patch("app.services.ha_discovery._legacy_thermal_entities", side_effect=RuntimeError("boom")): + with Session(disco_db) as session: + with pytest.raises(RuntimeError, match="boom"): + ha_discovery.initialize_legacy_thermal_cleanup(session) + + +def test_legacy_cleanup_startup_propagates_durable_write_failure(disco_db) -> None: + """A failed freeze write is fatal; UI must not open with an unfrozen ledger.""" + from app.services import ha_discovery + + with ( + patch("app.services.ha_discovery._legacy_thermal_entities", return_value=[]), + patch("app.services.ha_discovery._set_migration_json", side_effect=RuntimeError("disk full")), + ): + with Session(disco_db) as session: + with pytest.raises(RuntimeError, match="disk full"): + ha_discovery.initialize_legacy_thermal_cleanup(session) + + +def test_legacy_cleanup_startup_keeps_upgrade_with_enabled_topic_pending(disco_db) -> None: + """Existing enabled v1.6.1 inventory is not mistaken for a fresh install.""" + from app.integrations.expose import DeviceInfo, ExposableEntity + from app.models.config import AppConfigEntry + from app.services import ha_discovery + + legacy = ExposableEntity( + key="thermal_cost.old.heating_total", component="sensor", + device=DeviceInfo(identifiers=("legacy",), identity="old", name="obsolete"), + device_class=None, unit="", name="obsolete", + ) + with patch("app.services.ha_discovery._legacy_thermal_entities", return_value=[legacy]): + with Session(disco_db) as session: + ha_discovery.initialize_legacy_thermal_cleanup(session) + with Session(disco_db) as session: + ledger = session.query(AppConfigEntry).filter_by( + key=ha_discovery._LEGACY_THERMAL_CLEANUP_KEY + ).one() + assert ledger.value == ( + '{"complete": false, "inventory": ["homeassistant/sensor/old/' + 'thermal_cost_old_heating_total/config"], "topics": []}' + ) + + +def test_legacy_cleanup_freezes_startup_upgrade_inventory_across_toggle_changes(disco_db) -> None: + """Alembic-head compat ledger survives UI disable and runtime prefix changes.""" + from app.models.config import AppConfigEntry + from app.models.energy import Meter + from app.models.expose import ExposedEntityToggle + from app.services import ha_discovery + + now = datetime.now(timezone.utc) + settings = _make_settings(ha_discovery_prefix="frozen_prefix") + with Session(disco_db) as session: + heating = Meter(label="heating", commodity="heating", started_at=now, ended_at=None, + reason="initial", note=None, created_at=now) + water = Meter(label="water", commodity="hot_water", started_at=now, ended_at=None, + reason="initial", note=None, created_at=now) + session.add_all((heating, water)) + session.flush() + enabled = ha_discovery._thermal_cleanup_entities( + [(heating, water)], include_hot_water_total=False + )[:2] + for entity in enabled: + _enable_entity(session, entity.key) + first_topic = ha_discovery._legacy_discovery_topic(enabled[0], "frozen_prefix") + session.add(AppConfigEntry( + key=ha_discovery._LEGACY_THERMAL_CLEANUP_KEY, + value=json.dumps({"complete": False, "topics": [first_topic]}), + updated_at=now, + )) + session.commit() + + with patch("app.services.ha_discovery.build_runtime_settings", return_value=settings): + with Session(disco_db) as session: + ha_discovery.initialize_legacy_thermal_cleanup(session) + + with Session(disco_db) as session: + ledger = ha_discovery._migration_json(session, ha_discovery._LEGACY_THERMAL_CLEANUP_KEY) + inventory = ledger["inventory"] + assert len(inventory) == 2 + assert all(topic.startswith("frozen_prefix/") for topic in inventory) + assert ledger["topics"] == [first_topic] + # This mirrors PUT /api/expose: persist the UI change before it invokes + # publish_discovery in the same request. + toggle = session.query(ExposedEntityToggle).filter_by(key=enabled[0].key).one() + toggle.enabled = False + session.commit() + + manager = _make_mock_manager() + failed_topic = inventory[-1] + published: list[str] = [] + manager.publish.side_effect = lambda topic, _payload, **_kwargs: ( + published.append(topic) or topic != failed_topic + ) + with ( + patch("app.services.ha_discovery.build_runtime_settings", + return_value=_make_settings(ha_discovery_prefix="changed_prefix")), + patch("app.services.ha_discovery.mqtt_manager", manager), + patch("app.services.ha_discovery.build_catalog", return_value=[]), + patch("app.services.ha_discovery._stale_m8_entities", return_value=[]), + patch("app.services.ha_discovery._ha_registry_bindings", return_value=None), + ): + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + assert published == [failed_topic] + with Session(disco_db) as session: + pending = ha_discovery._migration_json(session, ha_discovery._LEGACY_THERMAL_CLEANUP_KEY) + assert pending["complete"] is False + assert pending["inventory"] == inventory + assert pending["topics"] == [first_topic] + + manager.publish.reset_mock() + manager.publish.side_effect = lambda topic, _payload, **_kwargs: published.append(topic) or True + published.clear() + with ( + patch("app.services.ha_discovery.build_runtime_settings", + return_value=_make_settings(ha_discovery_prefix="changed_prefix")), + patch("app.services.ha_discovery.mqtt_manager", manager), + patch("app.services.ha_discovery.build_catalog", return_value=[]), + patch("app.services.ha_discovery._legacy_thermal_entities") as enumerate_legacy, + patch("app.services.ha_discovery._stale_m8_entities", return_value=[]), + patch("app.services.ha_discovery._ha_registry_bindings", return_value=None), + ): + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + assert published == [failed_topic] + enumerate_legacy.assert_not_called() + + manager.publish.reset_mock() + with ( + patch("app.services.ha_discovery.build_runtime_settings", return_value=settings), + patch("app.services.ha_discovery.mqtt_manager", manager), + patch("app.services.ha_discovery.build_catalog", return_value=[]), + patch("app.services.ha_discovery._legacy_thermal_entities") as enumerate_legacy, + patch("app.services.ha_discovery._stale_m8_entities", return_value=[]), + patch("app.services.ha_discovery._ha_registry_bindings", return_value=None), + ): + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + enumerate_legacy.assert_not_called() + manager.publish.assert_not_called() + + +def test_registry_repair_waits_for_registry_observation_and_partial_catalog_recovers(disco_db) -> None: + """The V2 repair is per entity and never mistakes broker ACK for HA ACK.""" + from app.integrations.expose import CatalogEntry, DeviceInfo, ExposableEntity + from app.services import ha_discovery + + def entity(kind: str, identity: str, metric: str) -> ExposableEntity: + return ExposableEntity( + key=f"{kind}.{identity}.{metric}" if kind != "energy" else f"energy.{metric}", + component="sensor", + device=DeviceInfo( + identifiers=(f"home-automation:{kind}:{identity}",), identity=identity, name=identity, + provides_availability=False, + ), + device_class=None, unit="", name=metric, + ) + + entities = [ + *(entity("meter", f"meter-{number}", "total") for number in range(3)), + *(entity("source", f"source-{number}", "online") for number in range(2)), + *(entity("modbus", f"modbus-{number}", "voltage") for number in range(2)), + entity("energy", "electricity-epoch", "import_cost_total"), + ] + catalog = [CatalogEntry(entity=item, enabled=True) for item in entities] + manager = _make_mock_manager() + manager.publish.return_value = True + settings = _make_settings() + calls: list[tuple[str, object]] = [] + manager.publish.side_effect = lambda topic, payload, **_kwargs: calls.append((topic, payload)) or True + + bindings: dict[str, set[str]] = {ha_discovery._unique_id(item): {"old"} for item in entities} + with ( + patch("app.services.ha_discovery.build_runtime_settings", return_value=settings), + patch("app.services.ha_discovery.mqtt_manager", manager), + patch("app.services.ha_discovery.build_catalog", return_value=catalog), + patch("app.services.ha_discovery._legacy_thermal_entities", return_value=[]), + patch("app.services.ha_discovery._stale_m8_entities", return_value=[]), + patch("app.services.ha_discovery._ha_registry_bindings", side_effect=lambda _settings, ids: { + key: value for key, value in bindings.items() if key in ids + }), + ): + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + + topics = [ha_discovery._discovery_topic(item, "homeassistant") for item in entities] + assert [payload for _topic, payload in calls[:len(entities)]] == [b""] * len(entities) + assert [topic for topic, _payload in calls[:len(entities)]] == topics + assert len({ha_discovery._unique_id(item) for item in entities}) == len(entities) + + # HA confirms every tombstone; only then is each target re-added. + calls.clear() + bindings.clear() + with ( + patch("app.services.ha_discovery.build_runtime_settings", return_value=settings), + patch("app.services.ha_discovery.mqtt_manager", manager), + patch("app.services.ha_discovery.build_catalog", return_value=catalog), + patch("app.services.ha_discovery._stale_m8_entities", return_value=[]), + patch("app.services.ha_discovery._ha_registry_bindings", side_effect=lambda _settings, _ids: dict(bindings)), + ): + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + assert len(calls) == len(entities) + assert all(payload not in (b"", "", None) for _topic, payload in calls) + + +def test_registry_repair_unavailable_keeps_normal_discovery_publishing(disco_db) -> None: + """A broken optional HA WS link cannot leave legal configs tombstoned.""" + from app.integrations.expose import CatalogEntry, DeviceInfo, ExposableEntity + from app.services import ha_discovery + + entity = ExposableEntity( + key="meter.meter-1.total", component="sensor", + device=DeviceInfo(identifiers=("home-automation:meter:meter-1",), identity="meter-1", name="m1"), + device_class=None, unit="", name="total", + ) + manager = _make_mock_manager() + with ( + patch("app.services.ha_discovery.build_runtime_settings", return_value=_make_settings()), + patch("app.services.ha_discovery.mqtt_manager", manager), + patch("app.services.ha_discovery.build_catalog", return_value=[CatalogEntry(entity, True)]), + patch("app.services.ha_discovery._legacy_thermal_entities", return_value=[]), + patch("app.services.ha_discovery._stale_m8_entities", return_value=[]), + patch("app.services.ha_discovery._ha_registry_bindings", return_value=None), + ): + with Session(disco_db) as session: + ha_discovery.publish_discovery(session) + assert manager.publish.call_args.args[1] != b"" + +def test_stale_thermal_cleanup_uses_safe_topics_and_all_fourteen_metrics(disco_db) -> None: + """Later thermal meter swaps clear only ended safe-format pair configs.""" + from app.models.energy import Meter + from app.services import ha_discovery + + now = datetime.now(timezone.utc) + with Session(disco_db) as session: + old_heating = Meter(label="old heat", commodity="heating", started_at=now - timedelta(days=2), + ended_at=now - timedelta(days=1), reason="meter_swap", note=None, created_at=now) + old_water = Meter(label="old water", commodity="hot_water", started_at=now - timedelta(days=2), + ended_at=now - timedelta(days=1), reason="meter_swap", note=None, created_at=now) + active_heating = Meter(label="new heat", commodity="heating", started_at=now - timedelta(days=1), + ended_at=None, reason="meter_swap", note=None, created_at=now) + active_water = Meter(label="new water", commodity="hot_water", started_at=now - timedelta(days=1), + ended_at=None, reason="meter_swap", note=None, created_at=now) + session.add_all((old_heating, old_water, active_heating, active_water)) + session.commit() + stale = ha_discovery._stale_m8_entities(session) + + old_identity = ".".join(sorted((old_heating.uuid, old_water.uuid))) + thermal = [item for item in stale if item.key.startswith("thermal_cost.")] + assert len(thermal) == 14 + assert {item.key for item in thermal} == { + f"thermal_cost.{old_identity}.{metric}_{suffix}" + for metric in ("heating", "hot_water_heating", "hot_water_total", "water", "water_tax", "fixed", "all_in") + for suffix in ("total", "today") + } + assert all("." not in ha_discovery._discovery_topic(item, "homeassistant") for item in thermal) + assert not any(active_heating.uuid in item.key and active_water.uuid in item.key for item in thermal) diff --git a/tests/test_homeassistant.py b/tests/test_homeassistant.py index 9dd0bb4..ee87383 100644 --- a/tests/test_homeassistant.py +++ b/tests/test_homeassistant.py @@ -111,3 +111,87 @@ def test_homeassistant_client_raises_on_invalid_arguments() -> None: with pytest.raises(ValueError, match="webhook_id"): client.trigger_webhook(webhook_id="", body={}) + + +def test_discovery_registry_bindings_reads_entity_and_device_registry(monkeypatch: pytest.MonkeyPatch) -> None: + """The repair confirmation reads HA's authoritative registry bindings.""" + sent: list[dict] = [] + + class _Socket: + replies = iter(( + '{"type":"auth_required"}', + '{"type":"auth_ok"}', + '{"id":1,"success":true,"result":[{"platform":"mqtt","unique_id":"u1","device_id":"d1"}]}', + '{"id":2,"success":true,"result":[{"id":"d1","identifiers":[["mqtt","home-automation:meter:m1"]]}]}', + )) + + def __enter__(self): + return self + + def __exit__(self, *_args): + return None + + def recv(self, *, timeout=None): + assert timeout is not None + assert timeout <= 1.5 + return next(self.replies) + + def send(self, payload): + sent.append(json.loads(payload)) + + monkeypatch.setattr("app.integrations.homeassistant.connect", lambda *_args, **_kwargs: _Socket()) + bindings = HomeAssistantClient(_configured_settings()).discovery_registry_bindings({"u1", "missing"}) + + assert bindings == {"u1": {"home-automation:meter:m1"}} + assert [message.get("type") for message in sent] == [ + "auth", "config/entity_registry/list", "config/device_registry/list" + ] + + +def test_discovery_registry_bindings_uses_one_deadline_and_ignores_non_mqtt_identifiers( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Only official ["mqtt", value] pairs participate in repair matching.""" + received_timeouts: list[float] = [] + + class _Socket: + replies = iter(( + '{"type":"auth_required"}', + '{"type":"auth_ok"}', + '{"id":1,"success":true,"result":[' + '{"platform":"mqtt","unique_id":"u1","device_id":"d1"},' + '{"platform":"mqtt","unique_id":"u2","device_id":"d2"},' + '{"platform":"mqtt","unique_id":"u3","device_id":"d3"}]}', + '{"id":2,"success":true,"result":[' + '{"id":"d1","identifiers":[["mqtt","expected"],["esphome","aux"]]},' + '{"id":"d2","identifiers":[["esphome","expected"],"expected",["mqtt",3],[]]},' + '{"id":"d3","identifiers":[["mqtt","expected"],["mqtt","old"]]}]}', + )) + + def __enter__(self): return self + def __exit__(self, *_args): return None + def send(self, _payload): return None + def recv(self, *, timeout=None): + received_timeouts.append(timeout) + return next(self.replies) + + monkeypatch.setattr("app.integrations.homeassistant.connect", lambda *_args, **_kwargs: _Socket()) + bindings = HomeAssistantClient(_configured_settings()).discovery_registry_bindings({"u1", "u2", "u3"}) + + assert bindings == {"u1": {"expected"}, "u2": set(), "u3": {"expected", "old"}} + assert len(received_timeouts) == 4 + assert all(timeout is not None and 0 < timeout <= 1.5 for timeout in received_timeouts) + + +def test_discovery_registry_bindings_silent_socket_times_out(monkeypatch: pytest.MonkeyPatch) -> None: + class _Socket: + def __enter__(self): return self + def __exit__(self, *_args): return None + def send(self, _payload): return None + def recv(self, *, timeout=None): + assert timeout is not None and timeout > 0 + raise TimeoutError("silent") + + monkeypatch.setattr("app.integrations.homeassistant.connect", lambda *_args, **_kwargs: _Socket()) + with pytest.raises(HomeAssistantRequestError, match="registry query failed"): + HomeAssistantClient(_configured_settings()).discovery_registry_bindings({"u1"}) diff --git a/tests/test_mqtt_client.py b/tests/test_mqtt_client.py index f1d18c0..fc01b90 100644 --- a/tests/test_mqtt_client.py +++ b/tests/test_mqtt_client.py @@ -257,7 +257,7 @@ def test_publish_passes_topic_payload_retain_to_paho() -> None: manager._connected = True manager._client = mock_client - manager.publish("test/topic", '{"key": "value"}', retain=True, qos=1) + assert manager.publish("test/topic", '{"key": "value"}', retain=True, qos=1) is True mock_client.publish.assert_called_once_with( "test/topic", payload='{"key": "value"}', qos=1, retain=True @@ -267,8 +267,7 @@ def test_publish_passes_topic_payload_retain_to_paho() -> None: def test_publish_is_noop_when_not_connected() -> None: manager = MqttManager() # No connect — should silently skip - manager.publish("topic", "payload", retain=False) - # No exception raised + assert manager.publish("topic", "payload", retain=False) is False def test_publish_does_not_raise_on_paho_error() -> None: @@ -279,7 +278,17 @@ def test_publish_does_not_raise_on_paho_error() -> None: manager._connected = True # Must not raise - manager.publish("topic", "payload") + assert manager.publish("topic", "payload") is False + + +def test_publish_returns_false_when_paho_rejects_message() -> None: + manager = MqttManager() + mock_client = MagicMock() + mock_client.publish.return_value.rc = 1 + manager._client = mock_client + manager._connected = True + + assert manager.publish("topic", "payload") is False # ---------------------------------------------------------------------------