Skip to content
Open
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
2 changes: 1 addition & 1 deletion src/powersensor_local/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@
'Events',
'Message',
]
__version__ = "2.4.0rc2"
__version__ = "2.4.0rc3"
from .devices import PowersensorDevices, PowersensorLegacyDevices
from .legacy_discovery import LegacyDiscovery
from .plug_api import PlugApi
Expand Down
16 changes: 3 additions & 13 deletions src/powersensor_local/async_event_emitter.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,23 +31,13 @@ async def emit(self, event_name: str, *args: Any) -> None:
"""Emits an event to all registered listeners for that event type.
Additional arguments may be supplied with event as appropriate. Each
event handler is awaited before delivering the event to the next.
If an event handler raises an exception, this is funneled through
to an 'exception' event being emitted. If no 'exception' listener
is registered, or an exception handler callback raises an exception,
the exception is logged (if a logger was provided), and discarded."""
If an event handler raises an exception it is logged (if a logger was
provided), and discarded."""
if self._listeners.get(event_name) is None:
return
for callback in self._listeners[event_name]:
try:
await callback(event_name, *args)
except Exception as e:
if 'exception' not in self._listeners:
if self._logger is not None:
self._logger.exception(f"Discarding unhandled exception: {e}")
else:
for handler in self._listeners['exception']:
try:
await handler('exception', e)
except Exception as e2:
if self._logger is not None:
self._logger.exception(f"Exception handling callback raised an exception itself, discarding it: {e2}")
self._logger.exception("Logic error: exception escaped from callback: %s", e)
5 changes: 2 additions & 3 deletions src/powersensor_local/devices.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,6 @@
'average_power',
'average_power_components',
'battery_level',
'exception',
'now_relaying_for',
'radio_signal_quality',
'summation_energy',
Expand Down Expand Up @@ -73,7 +72,7 @@ class _PowersensorDevicesBase:
``event`` field. Known measurement events include:

``average_flow``, ``average_power``, ``average_power_components``,
``battery_level``, ``exception``, ``now_relaying_for``,
``battery_level``, ``now_relaying_for``,
``radio_signal_quality``, ``summation_energy``, ``summation_volume``.

When ``relay_now_relaying_for=True`` the raw ``now_relaying_for`` wire
Expand Down Expand Up @@ -187,7 +186,7 @@ async def _plug_discovered(self, mac: str, ip: str, port: int) -> None:
await self._remove_device(mac)

await self._add_device(mac, 'plug')
api = PlugApi(mac, ip, port)
api = PlugApi(mac, ip, port, 'udp', self._logger)
self._plug_apis[mac] = api
for event in _KNOWN_PLUG_EVENTS:
api.subscribe(event, self._reemit)
Expand Down
6 changes: 6 additions & 0 deletions src/powersensor_local/events.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,12 @@
import typing
import sys

from pathlib import Path

PROJECT_ROOT = str(Path(__file__).parents[1])
if PROJECT_ROOT not in sys.path:
sys.path.append(PROJECT_ROOT)

from powersensor_local.devices import PowersensorDevices
from powersensor_local.abstract_event_handler import AbstractEventHandler
from powersensor_local.xlatemsg import Event
Expand Down
17 changes: 8 additions & 9 deletions src/powersensor_local/plug_api.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
"""Interface abstraction for Powersensor plugs."""
import sys

from logging import Logger

from .async_event_emitter import AsyncEventEmitter
from .plug_listener_tcp import PlugListenerTcp
from .plug_listener_udp import PlugListenerUdp
Expand All @@ -17,7 +19,7 @@ class PlugApi(AsyncEventEmitter):
documented in xlatemsg.translate_raw_message.
"""

def __init__(self, mac: str, ip: str, port: int = 49476, proto: str = 'udp'):
def __init__(self, mac: str, ip: str, port: int = 49476, proto: str = 'udp', logger: Logger | None = None):
"""Create a :class:`PlugApi` instance for a single plug.

Parameters
Expand All @@ -32,23 +34,24 @@ def __init__(self, mac: str, ip: str, port: int = 49476, proto: str = 'udp'):
Protocol used for communication. ``'udp'`` selects :class:`PlugListenerUdp`,
while ``'tcp'`` selects :class:`PlugListenerTcp`. Any other value raises a
:class:`ValueError`.
logger : Logger, optional
If provided, enables the logging of escaped exceptions from callbacks.

Raises
------
ValueError
If *proto* is not ``'udp'`` or ``'tcp'``.
"""
super().__init__()
super().__init__(logger)
self._mac: str = mac
self._listener: PlugListenerUdp | PlugListenerTcp
if proto == 'udp':
self._listener = PlugListenerUdp(ip, port)
self._listener = PlugListenerUdp(ip, port, logger)
elif proto == 'tcp':
self._listener = PlugListenerTcp(ip, port)
self._listener = PlugListenerTcp(ip, port, logger)
else:
raise ValueError(f'Unsupported proto: {proto}')
self._listener.subscribe('message', self._on_message)
self._listener.subscribe('exception', self._on_exception)
self._seen: set[str] = set()

def connect(self) -> None:
Expand Down Expand Up @@ -93,10 +96,6 @@ async def _on_message(self, _: str, message: Message) -> None:
for name, ev in evs.items():
await self.emit(name, ev)

async def _on_exception(self, _: str, e: Exception) -> None:
"""Propagates exceptions from the plug listener."""
await self.emit('exception', e)

@property
def ip_address(self) -> str:
"""
Expand Down
5 changes: 3 additions & 2 deletions src/powersensor_local/plug_listener_tcp.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
import sys

from asyncio import StreamReader, StreamWriter
from logging import Logger

from powersensor_local.async_event_emitter import AsyncEventEmitter

Expand All @@ -23,7 +24,7 @@ class PlugListenerTcp(AsyncEventEmitter):
The event handlers must be async.
"""

def __init__(self, ip: str, port: int = 49476):
def __init__(self, ip: str, port: int = 49476, logger: Logger|None = None):
"""
Create a :class:`PlugListenerTcp` bound to the given IP address.

Expand All @@ -34,7 +35,7 @@ def __init__(self, ip: str, port: int = 49476):
port : int, optional
TCP port used by the plug (default ``49476``).
"""
super().__init__()
super().__init__(logger = logger)
self._ip: str = ip
self._port: int = port
self._task: asyncio.Task[None] | None = None
Expand Down
5 changes: 3 additions & 2 deletions src/powersensor_local/plug_listener_udp.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
import sys

from asyncio import TimerHandle
from logging import Logger
from typing import Any, Coroutine

from powersensor_local.async_event_emitter import AsyncEventEmitter
Expand All @@ -27,7 +28,7 @@ class PlugListenerUdp(AsyncEventEmitter, asyncio.DatagramProtocol):
The event handlers must be async.
"""

def __init__(self, ip: str, port: int = 49476):
def __init__(self, ip: str, port: int = 49476, logger: Logger|None = None):
"""
Create a :class:`PlugListenerUdp` bound to the given IP address.

Expand All @@ -38,7 +39,7 @@ def __init__(self, ip: str, port: int = 49476):
port : int, optional
UDP port used by the plug (default ``49476``).
"""
super().__init__()
super().__init__(logger = logger)
self._ip: str = ip
self._port: int = port
self._backoff: int = 0 # exponential backoff
Expand Down
9 changes: 7 additions & 2 deletions src/powersensor_local/plugevents.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
"""Utility script for accessing the plug api from a single network-local
Powersensor device. Intended for advanced debugging use only."""

import logging
import sys
from pathlib import Path

Expand All @@ -15,6 +16,8 @@
from powersensor_local.abstract_event_handler import AbstractEventHandler
from powersensor_local.xlatemsg import Message

LOGGER = logging.getLogger(__name__)

async def print_event_and_message(event: str, message: Message) -> None:
"""Callback for printing event data."""
print(event, message)
Expand All @@ -37,9 +40,11 @@ async def main(self) -> None:
# Signal handler for Ctrl+C
self.register_sigint_handler()

plug = PlugApi(sys.argv[1], sys.argv[2], int(*sys.argv[3:3]))
port = int(*sys.argv[3:3])
if port == 0:
port = 49476
plug = PlugApi(sys.argv[1], sys.argv[2], port, 'udp', LOGGER)
known_evs = [
'exception',
'average_flow',
'average_power',
'average_power_components',
Expand Down
1 change: 0 additions & 1 deletion src/powersensor_local/rawplug.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,6 @@ async def main(self) -> None:
else:
print('Unsupported protocol:', self._protocol)
sys.exit(1)
self.plug.subscribe('exception', print_message_ignore_event)
self.plug.subscribe('message', print_message_ignore_event)
self.plug.subscribe('connecting', print_event)
self.plug.subscribe('connecting', print_event)
Expand Down
9 changes: 6 additions & 3 deletions src/powersensor_local/virtual_household.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

import sys
from dataclasses import dataclass
from logging import Logger
from typing import Optional

from .async_event_emitter import AsyncEventEmitter
Expand Down Expand Up @@ -138,9 +139,9 @@ class VirtualHousehold(AsyncEventEmitter):
field to take note of summation resets.
"""

def __init__(self, with_solar: bool):
def __init__(self, with_solar: bool, logger: Logger|None = None):
"""Constructor.
with_solar True if it's already known that solar exists. Will be
with_solar: True if it's already known that solar exists. Will be
automatically enabled upon encountering a solar event during
processing, but until such a time may generate incorrect values
for home usage. Similarly, if this is set to True but no solar
Expand All @@ -156,8 +157,10 @@ def __init__(self, with_solar: bool):
would be generating incorrect data until such a time the solar
sensor is recharged. It is vastly preferable to have the system
show no data than show incorrect data.
logger: An optional logger to capture leaked exceptions from event
callbacks.
"""
super().__init__()
super().__init__(logger)
self._expect_solar = with_solar
self._summation = self.SummationInfo(0, 0, 0, 0)
self._counters = self.Counters(0, 0, 0, 0, 0)
Expand Down
15 changes: 13 additions & 2 deletions src/powersensor_local/zeroconf_devices.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,9 @@
import asyncio
import logging
import sys

from asyncio import Task
from logging import Logger
from typing import Any, Callable, Coroutine

from .devices import _AsyncCallback, _LogLevel, _PowersensorDevicesBase
Expand Down Expand Up @@ -92,7 +95,7 @@ def __init__(
service_type: str = _SERVICE_TYPE_UDP,
debounce_timeout: float = _DEBOUNCE_DEFAULT_S,
relay_now_relaying_for: bool = False,
logger: 'logging.Logger | None' = None,
logger: Logger|None = None,
) -> None:
"""Initialise.

Expand All @@ -118,6 +121,7 @@ def __init__(
self._zc_owned = zeroconf_instance is None # True → we close it in stop()
self._service_type = service_type
self._debounce_seconds = debounce_timeout
self._cb_logger = logger
self._browser: Any = None
self._listener: _Listener | None = None
self._pending_removals: dict[str, asyncio.TimerHandle] = {}
Expand Down Expand Up @@ -223,8 +227,15 @@ async def _on_zc_remove(self, mac: str) -> None:
def _internal_callback(self, coro: _InternalCallback) -> None:
"""Helper to prevent gc collection of short-lived callback tasks."""
task = asyncio.create_task(coro)

def cleanup(task: Task[None]) -> None:
self._internal_callbacks.discard(task)
e = task.exception()
if e is not None and not task.cancelled():
self._maybe_log(_LogLevel.ERROR, 'Exception escaped callback: %s', e)

self._internal_callbacks.add(task)
task.add_done_callback(self._internal_callbacks.discard)
task.add_done_callback(cleanup)


class _Listener(_zc.ServiceListener):
Expand Down
51 changes: 1 addition & 50 deletions tests/test_async_event_emitter.py
Original file line number Diff line number Diff line change
Expand Up @@ -66,59 +66,10 @@ async def test_argument_passing(emitter: AsyncEventEmitter) -> None:
async def test_exception_unhandled() -> None:
logger = MagicMock()
emitter = AsyncEventEmitter(logger)
mock = AsyncMock()
emitter.subscribe('e', mock)
mock.side_effect = KeyError('oops')
await emitter.emit('e')
mock.assert_called_once()
logger.exception.assert_called_once_with("Discarding unhandled exception: 'oops'")


@pytest.mark.asyncio
async def test_exception_handler(emitter: AsyncEventEmitter) -> None:
mock = AsyncMock()
emitter.subscribe('e', mock)
e = KeyError('oops')
mock.side_effect = e
mock_exc = AsyncMock()
emitter.subscribe('exception', mock_exc)
await emitter.emit('e')
mock.assert_called_once()
mock_exc.assert_called_once_with('exception', e)


@pytest.mark.asyncio
async def test_exception_handler_exception() -> None:
logger = MagicMock()
emitter = AsyncEventEmitter(logger)
mock = AsyncMock()
emitter.subscribe('e', mock)
mock.side_effect = KeyError('oops')
mock_exc = AsyncMock()
emitter.subscribe('exception', mock_exc)
mock_exc.side_effect = ValueError('doh')
await emitter.emit('e')
mock.assert_called_once()
mock_exc.assert_called_once()
logger.exception.assert_called_once_with("Exception handling callback raised an exception itself, discarding it: doh")


@pytest.mark.asyncio
async def test_multiple_exception_handlers() -> None:
logger = MagicMock()
emitter = AsyncEventEmitter(logger)
trigger = AsyncMock()
e = ValueError('overflow')
trigger.side_effect = e
emitter.subscribe('e', trigger)
handlers = [ AsyncMock() for _ in range(5) ]
for handler in handlers:
emitter.subscribe('exception', handler)
bad_handlers = handlers[1::2] # pick every other
for handler in bad_handlers:
handler.side_effect = e
await emitter.emit('e')
trigger.assert_called_once()
for handler in handlers:
handler.assert_called_once()
assert(logger.exception.call_count == len(bad_handlers))
logger.exception.assert_called_once_with("Logic error: exception escaped from callback: %s", e)
Loading