Skip to content

Commit 0813d94

Browse files
committed
Simplify zeroconf discovery.
Switch to the AsyncZeroconfBrowser. Remove unused public functions that were confusing a code bot. Fixed issue with potentially lost callbacks due to create_task/gc interactions. Improved type hinting.
1 parent ee4c50a commit 0813d94

1 file changed

Lines changed: 53 additions & 98 deletions

File tree

‎src/powersensor_local/zeroconf_devices.py‎

Lines changed: 53 additions & 98 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,8 @@
11
"""Zeroconf/mDNS-based discovery for Powersensor devices.
22
33
This module provides PowersensorZeroconfDevices, which uses continuous mDNS
4-
browsing to discover Powersensor plugs rather than the legacy one-shot UDP
5-
broadcast in PowersensorLegacyDevices.
4+
browsing to discover Powersensor plugs rather than the alternative, legacy
5+
one-shot UDP broadcast in PowersensorLegacyDevices.
66
77
The zeroconf package is an optional dependency. Install it via::
88
@@ -12,12 +12,10 @@
1212
------------
1313
PowersensorZeroconfDevices owns the full lifecycle:
1414
15-
- It starts a zeroconf ServiceBrowser that calls back on plug add/update/remove.
15+
- It starts a zeroconf AsyncServiceBrowser that calls back on plug
16+
add/update/remove.
1617
- Plug removals are debounced (default 60 s) to absorb transient disappearances
1718
such as reboots or DHCP renewals.
18-
- The public add_plug() / remove_plug() methods are the seam between discovery
19-
and the plug API lifecycle, and can also be called directly (e.g. from a
20-
test or from HA's own mDNS handler) without needing a real zeroconf instance.
2119
2220
Zeroconf instance ownership
2321
----------------------------
@@ -28,38 +26,23 @@
2826
2927
Thread safety
3028
-------------
31-
On zeroconf >= 0.32 (including Home Assistant's 0.149.x), ServiceBrowser
32-
callbacks run inside the asyncio event loop rather than on a background thread.
33-
On older zeroconf (e.g. 1.0.0, which used a select() thread), they run on a
34-
background thread.
35-
36-
``_Listener`` is therefore written to be safe in both models:
37-
38-
- ``_extract`` uses ``ServiceInfo.load_from_cache()`` rather than
39-
``Zeroconf.get_service_info()``. On >= 0.32, calling get_service_info()
40-
from inside a ServiceBrowser callback deadlocks — it blocks waiting for a
41-
DNS reply that can never arrive because it holds the event loop.
42-
load_from_cache() is synchronous, non-blocking, and explicitly threadsafe;
43-
the ServiceBrowser guarantees the cache is populated before firing the
44-
callback, so the record is always present.
45-
46-
- ``_name_to_mac`` is populated in ``add_service`` / ``update_service`` and
47-
consumed in ``remove_service``. On >= 0.32 all three callbacks run on the
48-
same event loop thread, so no locking is needed. On older versions they run
49-
on the same zeroconf background thread, so no locking is needed there either.
50-
51-
- All work that touches PowersensorZeroconfDevices state crosses the thread
52-
boundary via ``loop.call_soon_threadsafe``, making it safe regardless of
53-
which threading model the installed zeroconf uses.
29+
PowersensorZeroconfDevices uses purely async zeroconf discovery where all
30+
callbacks are fired on the event loop rather than a separate thread. Older
31+
versions of zeroconf that do not provide this guarantee are not supported.
32+
Mixing of event loops is also not supported - a single
33+
PowersensorZeroconfDevices cannot be used on more than one event loop or thread.
5434
"""
5535
from __future__ import annotations
5636

5737
import asyncio
5838
import logging
5939
import sys
60-
from typing import Any
40+
from typing import Any, Callable, Coroutine
41+
42+
from .devices import _AsyncCallback, _LogLevel, _PowersensorDevicesBase
43+
44+
_InternalCallback = Coroutine[None, None, None]
6145

62-
from .devices import _PowersensorDevicesBase, _LogLevel
6346

6447
_SERVICE_TYPE_UDP = '_powersensor._udp.local.'
6548
_SERVICE_TYPE_TCP = '_powersensor._tcp.local.'
@@ -69,6 +52,8 @@
6952

7053
try:
7154
import zeroconf as _zc
55+
from zeroconf.asyncio import AsyncServiceBrowser
56+
7257
class PowersensorZeroconfDevices(_PowersensorDevicesBase):
7358
"""Discovers and manages Powersensor plugs via continuous mDNS browsing.
7459
@@ -136,12 +121,13 @@ def __init__(
136121
self._browser: Any = None
137122
self._listener: _Listener | None = None
138123
self._pending_removals: dict[str, asyncio.TimerHandle] = {}
124+
self._internal_callbacks: set[asyncio.Task[None]] = set()
139125

140126
# ------------------------------------------------------------------
141127
# Lifecycle
142128
# ------------------------------------------------------------------
143129

144-
async def start(self, async_event_cb) -> None:
130+
async def start(self, async_event_cb: _AsyncCallback) -> None:
145131
"""Register the event callback and start the mDNS service browser.
146132
147133
The browser is event-driven; no polling loop is started here.
@@ -157,21 +143,27 @@ async def start(self, async_event_cb) -> None:
157143

158144
if self._zc_instance is None:
159145
self._zc_instance = _zc.Zeroconf()
146+
self._zc_owned = False
160147

161-
loop = asyncio.get_running_loop()
162-
self._listener = _Listener(self, loop)
163-
self._browser = _zc.ServiceBrowser(
148+
self._listener = _Listener(self)
149+
self._browser = AsyncServiceBrowser(
164150
self._zc_instance, self._service_type, self._listener
165151
)
166152

167153
async def stop(self) -> None:
168154
"""Stop browsing, cancel pending removals, and disconnect all plugs."""
155+
# This should be a noop under sane conditions
156+
for cb in self._internal_callbacks:
157+
cb.cancel()
158+
159+
self._event_cb = None
160+
169161
for handle in list(self._pending_removals.values()):
170162
handle.cancel()
171163
self._pending_removals.clear()
172164

173165
if self._browser is not None:
174-
self._browser.cancel()
166+
await self._browser.async_cancel()
175167
self._browser = None
176168

177169
self._listener = None
@@ -183,35 +175,7 @@ async def stop(self) -> None:
183175
await super().stop()
184176

185177
# ------------------------------------------------------------------
186-
# Public discovery seam — may also be called directly
187-
# ------------------------------------------------------------------
188-
189-
def add_plug(self, mac: str, ip: str, port: int) -> None:
190-
"""Notify that a plug is present at the given address.
191-
192-
Creates or reconnects the PlugApi for this plug. Safe to call
193-
directly without a zeroconf browser (e.g. from tests, or from HA's
194-
own mDNS handler). Cancels any pending debounced removal for this MAC.
195-
196-
Must be called from the event loop thread.
197-
"""
198-
self._cancel_pending_removal(mac, source='add_plug')
199-
asyncio.get_running_loop().create_task(
200-
self._plug_discovered(mac, ip, port)
201-
)
202-
203-
def remove_plug(self, mac: str) -> None:
204-
"""Schedule a debounced removal for the given plug.
205-
206-
After ``debounce_timeout`` seconds with no re-announcement, the plug
207-
API is disconnected and a ``device_lost`` event is emitted.
208-
209-
Must be called from the event loop thread.
210-
"""
211-
self._schedule_removal(mac)
212-
213-
# ------------------------------------------------------------------
214-
# Debounce helpers (event-loop side only)
178+
# Debounce helpers
215179
# ------------------------------------------------------------------
216180

217181
def _schedule_removal(self, mac: str) -> None:
@@ -230,7 +194,7 @@ def _on_debounce_expired(self, mac: str) -> None:
230194
"""Called by the event loop when the debounce timer fires."""
231195
self._pending_removals.pop(mac, None)
232196
self._maybe_log(_LogLevel.INFO, "Plug %s still absent after debounce — removing", mac)
233-
asyncio.get_running_loop().create_task(self._plug_lost(mac))
197+
self._internal_callback(self._plug_lost(mac))
234198

235199
def _cancel_pending_removal(self, mac: str, source: str) -> None:
236200
handle = self._pending_removals.pop(mac, None)
@@ -239,45 +203,36 @@ def _cancel_pending_removal(self, mac: str, source: str) -> None:
239203
self._maybe_log(_LogLevel.DEBUG, "Cancelled pending removal for %s (%s)", mac, source)
240204

241205
# ------------------------------------------------------------------
242-
# Called from _Listener (zeroconf thread → event loop via stored loop ref)
206+
# Called from _Listener
243207
# ------------------------------------------------------------------
244208

245-
def _on_zc_add(self, mac: str, ip: str, port: int, loop: asyncio.AbstractEventLoop) -> None:
246-
loop.call_soon_threadsafe(self._cancel_pending_removal, mac, 'zeroconf add')
247-
loop.call_soon_threadsafe(
248-
lambda: loop.create_task(self._plug_discovered(mac, ip, port))
249-
)
250-
251-
def _on_zc_update(self, mac: str, ip: str, port: int, loop: asyncio.AbstractEventLoop) -> None:
252-
loop.call_soon_threadsafe(self._cancel_pending_removal, mac, 'zeroconf update')
253-
loop.call_soon_threadsafe(
254-
lambda: loop.create_task(self._plug_discovered(mac, ip, port))
255-
)
209+
async def _on_zc_add(self, mac: str, ip: str, port: int) -> None:
210+
self._cancel_pending_removal( mac, 'zeroconf add')
211+
await self._plug_discovered(mac, ip, port)
256212

257-
def _on_zc_remove(self, mac: str, loop: asyncio.AbstractEventLoop) -> None:
258-
loop.call_soon_threadsafe(self._schedule_removal, mac)
213+
async def _on_zc_update(self, mac: str, ip: str, port: int) -> None:
214+
self._cancel_pending_removal(mac, 'zeroconf update')
215+
await self._plug_discovered(mac, ip, port)
259216

217+
async def _on_zc_remove(self, mac: str) -> None:
218+
self._schedule_removal(mac)
260219

261-
class _Listener(_zc.ServiceListener):
262-
"""Zeroconf ServiceListener that forwards events to PowersensorZeroconfDevices.
220+
# ------------------------------------------------------------------
221+
# Sync to async bridge for callbacks, per asyncio docs
222+
# ------------------------------------------------------------------
263223

264-
Internal implementation detail. All ServiceListener callbacks arrive on
265-
the zeroconf event loop (>= 0.32) or background thread (< 0.32 / 1.0.0).
224+
def _internal_callback(self, coro: _InternalCallback) -> None:
225+
"""Helper to prevent gc collection of short-lived callback tasks."""
226+
task = asyncio.create_task(coro)
227+
self._internal_callbacks.add(task)
228+
task.add_done_callback(self._internal_callbacks.discard)
266229

267-
Thread safety
268-
-------------
269-
``_name_to_mac`` is populated in ``add_service`` / ``update_service`` and
270-
consumed in ``remove_service``. In both threading models all three
271-
callbacks arrive on the same thread/loop, so no locking is required.
272230

273-
The stored ``_loop`` reference is captured once at construction from the
274-
running asyncio event loop, and is used (read-only) from the callback
275-
context to schedule work back onto that loop via ``call_soon_threadsafe``.
276-
"""
231+
class _Listener(_zc.ServiceListener):
232+
"""Zeroconf ServiceListener that forwards events to PowersensorZeroconfDevices."""
277233

278-
def __init__(self, owner: PowersensorZeroconfDevices, loop: asyncio.AbstractEventLoop) -> None:
234+
def __init__(self, owner: PowersensorZeroconfDevices) -> None:
279235
self._owner = owner
280-
self._loop = loop
281236
self._name_to_mac: dict[str, str] = {}
282237

283238
def _extract(self, zc: Any, type_: str, name: str) -> tuple[str, str, int] | None:
@@ -339,7 +294,7 @@ def add_service(self, zc: Any, type_: str, name: str) -> None:
339294
return
340295
mac, ip, port = result
341296
self._name_to_mac[name] = mac
342-
self._owner._on_zc_add(mac, ip, port, self._loop)
297+
self._owner._internal_callback(self._owner._on_zc_add(mac, ip, port))
343298

344299
def update_service(self, zc: Any, type_: str, name: str) -> None:
345300
result = self._extract(zc, type_, name)
@@ -352,7 +307,7 @@ def update_service(self, zc: Any, type_: str, name: str) -> None:
352307
return
353308
mac, ip, port = result
354309
self._name_to_mac[name] = mac
355-
self._owner._on_zc_update(mac, ip, port, self._loop)
310+
self._owner._internal_callback(self._owner._on_zc_update(mac, ip, port))
356311

357312
def remove_service(self, zc: Any, type_: str, name: str) -> None:
358313
mac = self._name_to_mac.pop(name, None)
@@ -362,7 +317,7 @@ def remove_service(self, zc: Any, type_: str, name: str) -> None:
362317
"remove_service for %s: MAC not in cache — removal ignored", name,
363318
)
364319
return
365-
self._owner._on_zc_remove(mac, self._loop)
320+
self._owner._internal_callback(self._owner._on_zc_remove(mac))
366321

367322
except ImportError as exc:
368323
_zeroconf_import_error = exc

0 commit comments

Comments
 (0)