Skip to content
Merged
3 changes: 3 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,9 @@ Issues = "https://github.com/SpanPanel/span-panel-api/issues"
[project.scripts]
format-markdown = "scripts.format_markdown:main"

[project.entry-points."span_panel_api.schema_adapters"]
schema_0 = "span_panel_api._impl.schema_0:SchemaZeroAdapter"

[dependency-groups]
dev = [
"pytest>=9.0.2",
Expand Down
1 change: 1 addition & 0 deletions src/span_panel_api/_impl/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
"""Internal implementation packages. Not public API."""
9 changes: 9 additions & 0 deletions src/span_panel_api/_impl/schema_0/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
"""Flat-schema adapter package (data-model-version absent)."""

from span_panel_api._impl.schema_0.adapter import SchemaZeroAdapter

# Inclusive lower bound, exclusive upper bound. The flat schema publishes no
# data-model-version, so it is treated as the synthetic version 0 range.
SUPPORTS_DATA_MODEL_VERSIONS: tuple[str, str] = (">=0", "<1.0")

__all__ = ["SUPPORTS_DATA_MODEL_VERSIONS", "SchemaZeroAdapter"]
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,8 @@
import logging
import time

from .const import HOMIE_STATE_DISCONNECTED, HOMIE_STATE_LOST, HOMIE_STATE_READY, TOPIC_PREFIX
from span_panel_api._impl.schema_0.const import TOPIC_PREFIX
from span_panel_api.mqtt.const import HOMIE_STATE_DISCONNECTED, HOMIE_STATE_LOST, HOMIE_STATE_READY

_LOGGER = logging.getLogger(__name__)

Expand Down
67 changes: 67 additions & 0 deletions src/span_panel_api/_impl/schema_0/adapter.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
"""Flat-schema (data-model-version absent) adapter.

Composes the existing accumulator + consumer and owns the flat wire format:
a single Homie device whose node ids are circuit UUIDs and capability names.
Nothing outside this package constructs a flat-schema topic.
"""

from __future__ import annotations

from collections.abc import Callable
from typing import TYPE_CHECKING

from span_panel_api._impl.schema_0.accumulator import HomiePropertyAccumulator
from span_panel_api._impl.schema_0.const import PROPERTY_SET_TOPIC_FMT, TYPE_CORE, WILDCARD_TOPIC_FMT
from span_panel_api._impl.schema_0.consumer import HomieDeviceConsumer
from span_panel_api._impl.schema_0.field_metadata import build_field_metadata

if TYPE_CHECKING:
from span_panel_api.models import FieldMetadata, HomieSchemaTypes, SpanPanelSnapshot


class SchemaZeroAdapter:
"""Parser for the flat single-device schema (firmware r202603-r202627)."""

schema_major = "schema_0"
SUPPORTS_DATA_MODEL_VERSIONS: tuple[str, str] = (">=0", "<1.0")

def __init__(self, serial_number: str, panel_size: int) -> None:
self._serial_number = serial_number
self._accumulator = HomiePropertyAccumulator(serial_number)
self._consumer = HomieDeviceConsumer(self._accumulator, panel_size)

def topics_to_subscribe(self) -> list[str]:
return [WILDCARD_TOPIC_FMT.format(serial=self._serial_number)]

def handle_message(self, topic: str, payload: str) -> None:
self._consumer.handle_message(topic, payload)

def is_ready(self) -> bool:
return self._consumer.is_ready()

def build_snapshot(self) -> SpanPanelSnapshot:
return self._consumer.build_snapshot()

def build_field_metadata(self, schema_types: HomieSchemaTypes) -> dict[str, FieldMetadata]:
return build_field_metadata(schema_types)

def circuit_nodes_missing_names(self) -> list[str]:
return self._consumer.circuit_nodes_missing_names()

def find_node_by_type(self, type_str: str) -> str | None:
return self._consumer.find_node_by_type(type_str)

def set_circuit_relay_topic(self, circuit_id: str) -> str:
return PROPERTY_SET_TOPIC_FMT.format(serial=self._serial_number, node=circuit_id, prop="relay")

def set_circuit_priority_topic(self, circuit_id: str) -> str:
return PROPERTY_SET_TOPIC_FMT.format(serial=self._serial_number, node=circuit_id, prop="shed-priority")

def set_dominant_power_source_topic(self) -> str | None:
core_node = self._consumer.find_node_by_type(TYPE_CORE)
if core_node is None:
return None
return PROPERTY_SET_TOPIC_FMT.format(serial=self._serial_number, node=core_node, prop="dominant-power-source")

def register_property_callback(self, callback: Callable[[str, str, str, str | None], None]) -> Callable[[], None]:
return self._consumer.register_property_callback(callback)
42 changes: 42 additions & 0 deletions src/span_panel_api/_impl/schema_0/const.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
"""Constants for the flat-schema (Homie v5) parsing implementation."""

# Homie v5 topic structure
HOMIE_VERSION = 5
HOMIE_DOMAIN = "ebus"
TOPIC_PREFIX = f"{HOMIE_DOMAIN}/{HOMIE_VERSION}"

# Topic patterns (serial_number substituted at runtime)
DEVICE_TOPIC_FMT = f"{TOPIC_PREFIX}/{{serial}}"
STATE_TOPIC_FMT = f"{TOPIC_PREFIX}/{{serial}}/$state"
DESCRIPTION_TOPIC_FMT = f"{TOPIC_PREFIX}/{{serial}}/$description"
PROPERTY_TOPIC_FMT = f"{TOPIC_PREFIX}/{{serial}}/{{node}}/{{prop}}"
PROPERTY_SET_TOPIC_FMT = f"{TOPIC_PREFIX}/{{serial}}/{{node}}/{{prop}}/set"
WILDCARD_TOPIC_FMT = f"{TOPIC_PREFIX}/{{serial}}/#"

# Homie type strings from schema
TYPE_CORE = "energy.ebus.device.distribution-enclosure.core"
TYPE_LUGS = "energy.ebus.device.lugs"
TYPE_LUGS_UPSTREAM = "energy.ebus.device.lugs.upstream"
TYPE_LUGS_DOWNSTREAM = "energy.ebus.device.lugs.downstream"
TYPE_CIRCUIT = "energy.ebus.device.circuit"
TYPE_BESS = "energy.ebus.device.bess"
TYPE_PV = "energy.ebus.device.pv"
TYPE_EVSE = "energy.ebus.device.evse"
TYPE_PCS = "energy.ebus.device.pcs"
TYPE_POWER_FLOWS = "energy.ebus.device.power-flows"

# Lugs direction values
LUGS_UPSTREAM = "UPSTREAM"
LUGS_DOWNSTREAM = "DOWNSTREAM"


def normalize_circuit_id(node_id: str) -> str:
"""Strip dashes from Homie UUID for entity stability."""
return node_id.replace("-", "")


def denormalize_circuit_id(circuit_id: str) -> str:
"""Restore dashes to a 32-char dashless UUID (8-4-4-4-12 format)."""
if len(circuit_id) == 32 and "-" not in circuit_id:
return f"{circuit_id[:8]}-{circuit_id[8:12]}-{circuit_id[12:16]}-{circuit_id[16:20]}-{circuit_id[20:]}"
return circuit_id
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,8 @@
import time
from typing import ClassVar

from ..models import SpanBatterySnapshot, SpanCircuitSnapshot, SpanEvseSnapshot, SpanPanelSnapshot, SpanPVSnapshot
from .accumulator import HomiePropertyAccumulator
from .const import (
from span_panel_api._impl.schema_0.accumulator import HomiePropertyAccumulator
from span_panel_api._impl.schema_0.const import (
LUGS_DOWNSTREAM,
LUGS_UPSTREAM,
TYPE_BESS,
Expand All @@ -28,6 +27,13 @@
TYPE_PV,
normalize_circuit_id,
)
from span_panel_api.models import (
SpanBatterySnapshot,
SpanCircuitSnapshot,
SpanEvseSnapshot,
SpanPanelSnapshot,
SpanPVSnapshot,
)

_LOGGER = logging.getLogger(__name__)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,7 @@

from __future__ import annotations

import logging

from ..models import FieldMetadata, HomieSchemaTypes
from .const import (
from span_panel_api._impl.schema_0.const import (
TYPE_BESS,
TYPE_CIRCUIT,
TYPE_CORE,
Expand All @@ -29,8 +26,7 @@
TYPE_POWER_FLOWS,
TYPE_PV,
)

_LOGGER = logging.getLogger(__name__)
from span_panel_api.models import FieldMetadata, HomieSchemaTypes

# ---------------------------------------------------------------------------
# Static mapping: (node_type, property_id) → snapshot field path
Expand Down Expand Up @@ -177,53 +173,3 @@ def build_field_metadata(
result[field_path] = FieldMetadata(unit=unit, datatype=datatype)

return result


def log_schema_drift(
previous: HomieSchemaTypes,
current: HomieSchemaTypes,
) -> None:
"""Log property-level differences between two schema versions.

Called by the client when the schema hash changes between connections.
All Homie-specific detail stays in this module — the integration never
sees this output, only the transport-agnostic field metadata.
"""
prev_types = set(previous.keys())
curr_types = set(current.keys())

for node_type in sorted(curr_types - prev_types):
_LOGGER.debug("Schema drift: new node type '%s'", node_type)

for node_type in sorted(prev_types - curr_types):
_LOGGER.debug("Schema drift: removed node type '%s'", node_type)

for node_type in sorted(prev_types & curr_types):
prev_props = previous[node_type]
curr_props = current[node_type]
if not isinstance(prev_props, dict) or not isinstance(curr_props, dict):
continue

for prop_id in sorted(set(curr_props) - set(prev_props)):
_LOGGER.debug("Schema drift: new property '%s/%s'", node_type, prop_id)

for prop_id in sorted(set(prev_props) - set(curr_props)):
_LOGGER.debug("Schema drift: removed property '%s/%s'", node_type, prop_id)

for prop_id in sorted(set(prev_props) & set(curr_props)):
prev_def = prev_props[prop_id]
curr_def = curr_props[prop_id]
if not isinstance(prev_def, dict) or not isinstance(curr_def, dict):
continue
for attr in ("datatype", "unit", "format"):
old_val = prev_def.get(attr)
new_val = curr_def.get(attr)
if old_val != new_val:
_LOGGER.debug(
"Schema drift: '%s/%s' %s changed: '%s' → '%s'",
node_type,
prop_id,
attr,
old_val,
new_val,
)
41 changes: 41 additions & 0 deletions src/span_panel_api/adapters.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
"""Adapter discovery via the `span_panel_api.schema_adapters` entry-point group.

Called once per process on the first create_span_client(). A venv change needs a
process restart regardless, so a process-lifetime cache is correct.
"""

from __future__ import annotations

from importlib.metadata import entry_points
import logging
from typing import TYPE_CHECKING

if TYPE_CHECKING:
from span_panel_api.protocol import SchemaAdapter

_LOGGER = logging.getLogger(__name__)
_ENTRY_POINT_GROUP = "span_panel_api.schema_adapters"
_REGISTRY: dict[str, type[SchemaAdapter]] | None = None


def discover_adapters() -> dict[str, type[SchemaAdapter]]:
"""Load and cache every adapter class registered under the entry-point group."""
global _REGISTRY # pylint: disable=global-statement # process-lifetime cache by design
if _REGISTRY is None:
registry: dict[str, type[SchemaAdapter]] = {}
for ep in entry_points(group=_ENTRY_POINT_GROUP):
if ep.name in registry:
_LOGGER.warning("Duplicate schema adapter entry point %r; keeping the first found", ep.name)
continue
try:
registry[ep.name] = ep.load()
except Exception: # pylint: disable=broad-exception-caught
_LOGGER.exception("Failed to load schema adapter entry point %r", ep.name)
_REGISTRY = registry
return _REGISTRY


def _reset_adapter_cache() -> None:
"""Test hook. Not public API."""
global _REGISTRY # pylint: disable=global-statement # test hook for the cache above
_REGISTRY = None
14 changes: 14 additions & 0 deletions src/span_panel_api/exceptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,3 +40,17 @@ class SpanPanelStaleDataError(SpanPanelError):
but data cannot be trusted right now (broker disconnected, or the Homie
device has declared $state=disconnected/lost).
"""


class SpanPanelAdapterMissingError(SpanPanelError):
"""No installed adapter covers the schema this panel publishes."""

def __init__(self, needed: str, reason: str, available: list[str]) -> None:
self.needed = needed
self.reason = reason
self.available = available
super().__init__(
f"Panel requires adapter {needed!r} (reason: {reason}); "
f"installed adapters: {sorted(available)}. "
"Update the integration or install the missing adapter package."
)
47 changes: 45 additions & 2 deletions src/span_panel_api/factory.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,17 +7,51 @@
from __future__ import annotations

import logging
import re

from .adapters import discover_adapters
from .auth import register_v2
from .detection import detect_api_version
from .exceptions import SpanPanelAuthError
from .exceptions import SpanPanelAdapterMissingError, SpanPanelAuthError
from .mqtt.client import SpanMqttClient
from .mqtt.models import MqttClientConfig
from .protocol import SchemaAdapter

_LOGGER = logging.getLogger(__name__)

_V2_CLIENT_NAME = "span-panel-api"

_DMV_PATTERN = re.compile(r"^(\d+)\.\d+(?:\.\d+)?$")


def _select_adapter_key(data_model_version: str | None) -> tuple[str, str]:
"""Tier 1 dispatch: the panel's data-model-version selects the adapter major.

Absence is the flat-schema signal — the property was introduced by the same
firmware that introduced the parent/child model, so a panel that does not
publish it is speaking the flat schema.
"""
if data_model_version is None:
return "schema_0", "data-model-version absent (flat schema)"

match = _DMV_PATTERN.match(data_model_version)
if match is None:
return "schema_0", f"unrecognised data-model-version={data_model_version!r}, assuming flat"

return (
f"schema_{int(match.group(1))}",
f"data-model-version={data_model_version!r}",
)


def _resolve_adapter_cls(key: str, reason: str) -> type[SchemaAdapter]:
"""Look up the discovered adapter class for `key`, or raise with the installed list."""
registry = discover_adapters()
adapter_cls = registry.get(key)
if adapter_cls is None:
raise SpanPanelAdapterMissingError(needed=key, reason=reason, available=sorted(registry))
return adapter_cls


async def create_span_client(
host: str,
Expand Down Expand Up @@ -68,6 +102,15 @@ async def create_span_client(
if serial_number is None:
raise SpanPanelAuthError("serial_number is required for MQTT transport but could not be determined")

client = SpanMqttClient(host, serial_number, mqtt_config, panel_http_port=port)
# Phase 0: the factory does not fetch the Homie schema, so no panel can
# report a data-model-version yet. `None` is the correct observation for
# every panel currently in the field — Phase 1 adds the fetch.
data_model_version: str | None = None
adapter_key, dispatch_reason = _select_adapter_key(data_model_version)
adapter_cls = _resolve_adapter_cls(adapter_key, dispatch_reason)

client = SpanMqttClient(host, serial_number, mqtt_config, panel_http_port=port, adapter_factory=adapter_cls)
client._data_model_version = data_model_version # pylint: disable=protected-access
client._schema_dispatch_reason = dispatch_reason # pylint: disable=protected-access
await client.connect()
return client
5 changes: 3 additions & 2 deletions src/span_panel_api/mqtt/__init__.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
"""SPAN Panel MQTT/Homie transport."""

from .accumulator import HomieLifecycle, HomiePropertyAccumulator
from span_panel_api._impl.schema_0.accumulator import HomieLifecycle, HomiePropertyAccumulator
from span_panel_api._impl.schema_0.consumer import HomieDeviceConsumer

from .async_client import AsyncMQTTClient
from .client import SpanMqttClient
from .connection import AsyncMqttBridge
from .homie import HomieDeviceConsumer
from .models import MqttClientConfig

__all__ = [
Expand Down
Loading
Loading