Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions examples/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,12 @@ Publishes a single eBus utility-meter device (`energy.ebus.device.utility-meter`
./utility-meter --config ./utility-meter-cfg.example.json --broker-config /path/to/broker-cfg.json
```

Add `--discover` to find the broker over mDNS (`_secure-mqtt._tcp`) instead of using the `host`/`port` in the broker config; the broker config still supplies the TLS material. Needs the `mdns` extra (`pip install 'ebus-sdk[mdns]'`):

```bash
./utility-meter --config ./utility-meter-cfg.example.json --broker-config /path/to/broker-cfg.json --discover
```

Set DOE values at runtime:

```bash
Expand Down
86 changes: 86 additions & 0 deletions examples/utility-meter
Original file line number Diff line number Diff line change
Expand Up @@ -550,6 +550,56 @@ def start_tick_thread(meter: UtilityMeter, interval_s: float) -> threading.Event
# ─── Entry point ──────────────────────────────────────────────────────────────


def _discover_broker(
timeout: float, log: logging.Logger
) -> Optional[tuple]:
"""Discover an eBus broker over mDNS (`_secure-mqtt._tcp`).

Returns (host, port) from the advertisement, preferring the spec `broker`
TXT record (the `<host>.local` name the broker's cert SAN covers) and
falling back to the SRV target. Returns None on timeout. Requires the
optional `mdns` extra (zeroconf).
"""
try:
from zeroconf import ServiceBrowser, ServiceListener, Zeroconf
except ImportError:
log.error(
"discovery needs the 'mdns' extra: pip install 'ebus-sdk[mdns]'"
)
return None

service_type = "_secure-mqtt._tcp.local."
found = threading.Event()
result: Dict[str, Any] = {}

class _Listener(ServiceListener):
def add_service(self, zc, type_, name):
info = zc.get_service_info(type_, name, timeout=3000)
if info is None:
return
broker = info.properties.get(b"broker")
result["host"] = (
broker.decode() if broker else info.server.rstrip(".")
)
result["port"] = info.port or 8883
found.set()

def update_service(self, zc, type_, name):
self.add_service(zc, type_, name)

def remove_service(self, zc, type_, name):
pass

zc = Zeroconf()
ServiceBrowser(zc, service_type, _Listener())
try:
if found.wait(timeout):
return result["host"], result["port"]
return None
finally:
zc.close()


def main():
parser = argparse.ArgumentParser(
description=(
Expand All @@ -569,6 +619,19 @@ def main():
help="Path to the MQTT broker config JSON file. "
f"Defaults to ${DEFAULT_BROKER_CFG_ENV} env var.",
)
parser.add_argument(
"--discover",
action="store_true",
help="Discover the broker over mDNS (_secure-mqtt._tcp) and use the "
"advertised host/port, overriding 'host'/'port' in the broker config. "
"The broker config still supplies the TLS material.",
)
parser.add_argument(
"--discover-timeout",
type=float,
default=8.0,
help="Seconds to wait for mDNS broker discovery (default: 8.0).",
)
parser.add_argument(
"--doe-port",
type=int,
Expand Down Expand Up @@ -604,6 +667,29 @@ def main():
sys.exit(2)
mqtt_cfg = UtilityMeterAdapter.load_broker_config(broker_cfg_path)

# Prefer an mDNS-discovered broker over the config's explicit host, per the
# eBus broker-discovery flow. The broker config still supplies the TLS
# material; discovery only fills in where to connect.
if args.discover:
endpoint = _discover_broker(args.discover_timeout, log)
if endpoint is not None:
mqtt_cfg["host"], mqtt_cfg["port"] = endpoint
log.info(
"reason=brokerDiscovered,host=%s,port=%d",
mqtt_cfg["host"],
mqtt_cfg["port"],
)
elif mqtt_cfg.get("host"):
log.warning(
"reason=brokerDiscoveryTimeout,fallbackHost=%s",
mqtt_cfg["host"],
)
else:
log.error(
"broker discovery failed and no fallback 'host' in broker config"
)
sys.exit(2)

# Resolve device identity from config
device_id = meter_cfg.get("device-id") or meter_cfg.get("info", {}).get(
"serial-number"
Expand Down