From aac19dabf35f76d4c37f71ed74306aff98e5d4f3 Mon Sep 17 00:00:00 2001 From: timothyanderson096-ocdealcheck Date: Thu, 27 Aug 2026 16:24:53 +1000 Subject: [PATCH 1/5] Fix async Notifier callback error handling --- can/notifier.py | 40 +++++++++++++++++++++++++++++++++------- 1 file changed, 33 insertions(+), 7 deletions(-) diff --git a/can/notifier.py b/can/notifier.py index fd21a0662..70ae3fb18 100644 --- a/can/notifier.py +++ b/can/notifier.py @@ -74,7 +74,7 @@ def unregister(self, bus: BusABC, notifier: "Notifier") -> None: with self.lock: registered_pairs_to_remove: list[_BusNotifierPair] = [] for pair in self.pairs: - if pair.bus is bus and pair.notifier is notifier: + if bus is pair.bus and pair.notifier is notifier: registered_pairs_to_remove.append(pair) for pair in registered_pairs_to_remove: self.pairs.remove(pair) @@ -223,7 +223,7 @@ def _rx_thread(self, bus: BusABC) -> None: if self._loop: handle_message: Callable[[Message], Any] = functools.partial( self._loop.call_soon_threadsafe, - self._on_message_received, # type: ignore[arg-type] + self._on_message_received_with_error_handling, # type: ignore[arg-type] ) else: handle_message = self._on_message_received @@ -248,7 +248,16 @@ def _rx_thread(self, bus: BusABC) -> None: def _on_message_available(self, bus: BusABC) -> None: if msg := bus.recv(0): + self._on_message_received_with_error_handling(msg) + + def _on_message_received_with_error_handling(self, msg: Message) -> None: + try: self._on_message_received(msg) + except Exception as exc: # pylint: disable=broad-except + self.exception = exc + if not self._on_error(exc): + raise + logger.debug("suppressed exception: %s", exc) def _on_message_received(self, msg: Message) -> None: for callback in self.listeners: @@ -257,7 +266,24 @@ def _on_message_received(self, msg: Message) -> None: # Schedule coroutine and keep a reference to the task task = self._loop.create_task(res) self._tasks.add(task) - task.add_done_callback(self._tasks.discard) + task.add_done_callback(self._on_task_done) + + def _on_task_done(self, task: asyncio.Task) -> None: + self._tasks.discard(task) + if task.cancelled(): + return + + exc = task.exception() + if exc is None: + return + + self.exception = exc + if isinstance(exc, Exception): + if not self._on_error(exc): + raise exc + logger.debug("suppressed exception: %s", exc) + else: + raise exc def _on_error(self, exc: Exception) -> bool: """Calls ``on_error()`` for all listeners if they implement it. @@ -278,7 +304,7 @@ def _on_error(self, exc: Exception) -> bool: return was_handled def add_listener(self, listener: MessageRecipient) -> None: - """Add new Listener to the notification set. + """Add new Listener for notification. :param listener: Listener to be added to the list to be notified """ @@ -289,8 +315,8 @@ def remove_listener(self, listener: MessageRecipient) -> None: throws an exception if the given listener is not part of the stored listeners. - :param listener: Listener to be removed from the set to be notified - :raises ValueError: if `listener` was never added to this notifier + :param listener: Listener to be removed from the set + :raises ValueError: if `listener` was never added to the notifier """ self.listeners.remove(listener) @@ -301,7 +327,7 @@ def stopped(self) -> bool: @staticmethod def find_instances(bus: BusABC) -> tuple["Notifier", ...]: - """Find :class:`~can.Notifier` instances associated with a given CAN bus. + """Find :class:`~can.Notifier` instances associated with a given bus. This method searches the registry for the :class:`~can.Notifier` that is linked to the specified bus. If the bus is found, the From 62d5243579d8bb749c3d1b70b9e85491cf2a6242 Mon Sep 17 00:00:00 2001 From: timothyanderson096-ocdealcheck Date: Thu, 27 Aug 2026 16:25:23 +1000 Subject: [PATCH 2/5] Add async Notifier error regression tests --- test/notifier_test.py | 70 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 70 insertions(+) diff --git a/test/notifier_test.py b/test/notifier_test.py index d8512a00b..a190ab6e9 100644 --- a/test/notifier_test.py +++ b/test/notifier_test.py @@ -7,6 +7,24 @@ import can +class RaisingListener(can.Listener): + def on_message_received(self, msg: can.Message) -> None: + raise ValueError("listener failed") + + +class ErrorCollector(can.Listener): + def __init__(self, event: asyncio.Event) -> None: + self.event = event + self.errors: list[Exception] = [] + + def on_message_received(self, msg: can.Message) -> None: + pass + + def on_error(self, exc: Exception) -> None: + self.errors.append(exc) + self.event.set() + + class NotifierTest(unittest.TestCase): def test_single_bus(self): with can.Bus("test", interface="virtual", receive_own_messages=True) as bus: @@ -88,6 +106,58 @@ async def run_it(): asyncio.run(run_it()) + def test_sync_listener_error_calls_on_error_with_loop(self): + async def run_it(): + event = asyncio.Event() + collector = ErrorCollector(event) + with can.Bus( + "sync-listener-error", interface="virtual", receive_own_messages=True + ) as bus: + notifier = can.Notifier( + bus, + [RaisingListener(), collector], + 0.1, + loop=asyncio.get_running_loop(), + ) + try: + bus.send(can.Message()) + await asyncio.wait_for(event.wait(), 0.5) + self.assertEqual(len(collector.errors), 1) + self.assertIsInstance(collector.errors[0], ValueError) + self.assertIs(notifier.exception, collector.errors[0]) + finally: + notifier.stop() + + asyncio.run(run_it()) + + def test_async_listener_error_calls_on_error(self): + async def run_it(): + event = asyncio.Event() + collector = ErrorCollector(event) + + async def raising_callback(msg: can.Message) -> None: + raise RuntimeError("async listener failed") + + with can.Bus( + "async-listener-error", interface="virtual", receive_own_messages=True + ) as bus: + notifier = can.Notifier( + bus, + [raising_callback, collector], + 0.1, + loop=asyncio.get_running_loop(), + ) + try: + bus.send(can.Message()) + await asyncio.wait_for(event.wait(), 0.5) + self.assertEqual(len(collector.errors), 1) + self.assertIsInstance(collector.errors[0], RuntimeError) + self.assertIs(notifier.exception, collector.errors[0]) + finally: + notifier.stop() + + asyncio.run(run_it()) + if __name__ == "__main__": unittest.main() From 0590b6d4353257b0d45b2b8f883287a0d78b4797 Mon Sep 17 00:00:00 2001 From: timothyanderson096-ocdealcheck Date: Thu, 27 Aug 2026 16:27:03 +1000 Subject: [PATCH 3/5] Scope Notifier fix to async error handling --- can/notifier.py | 24 +++++------------------- 1 file changed, 5 insertions(+), 19 deletions(-) diff --git a/can/notifier.py b/can/notifier.py index 70ae3fb18..da736e5df 100644 --- a/can/notifier.py +++ b/can/notifier.py @@ -74,7 +74,7 @@ def unregister(self, bus: BusABC, notifier: "Notifier") -> None: with self.lock: registered_pairs_to_remove: list[_BusNotifierPair] = [] for pair in self.pairs: - if bus is pair.bus and pair.notifier is notifier: + if pair.bus is bus and pair.notifier is notifier: registered_pairs_to_remove.append(pair) for pair in registered_pairs_to_remove: self.pairs.remove(pair) @@ -165,21 +165,16 @@ def add_bus(self, bus: BusABC) -> None: :raises ValueError: If the *bus* is already assigned to an active :class:`~can.Notifier`. """ - # add bus to notifier registry Notifier._registry.register(bus, self) - - # add bus to internal bus list self._bus_list.append(bus) file_descriptor: int = -1 try: file_descriptor = bus.fileno() except NotImplementedError: - # Bus doesn't support fileno, we fall back to thread based reader pass if self._loop is not None and file_descriptor >= 0: - # Use bus file descriptor to watch for messages self._loop.add_reader(file_descriptor, self._on_message_available, bus) self._readers.append(file_descriptor) else: @@ -208,18 +203,15 @@ def stop(self, timeout: float = 5.0) -> None: if now < end_time: reader.join(end_time - now) elif self._loop: - # reader is a file descriptor self._loop.remove_reader(reader) for listener in self.listeners: if hasattr(listener, "stop"): listener.stop() - # remove bus from registry for bus in self._bus_list: Notifier._registry.unregister(bus, self) def _rx_thread(self, bus: BusABC) -> None: - # determine message handling callable early, not inside while loop if self._loop: handle_message: Callable[[Message], Any] = functools.partial( self._loop.call_soon_threadsafe, @@ -237,13 +229,10 @@ def _rx_thread(self, bus: BusABC) -> None: self.exception = exc if self._loop is not None: self._loop.call_soon_threadsafe(self._on_error, exc) - # Raise anyway raise elif not self._on_error(exc): - # If it was not handled, raise the exception here raise else: - # It was handled, so only log it logger.debug("suppressed exception: %s", exc) def _on_message_available(self, bus: BusABC) -> None: @@ -263,7 +252,6 @@ def _on_message_received(self, msg: Message) -> None: for callback in self.listeners: res = callback(msg) if res and self._loop and asyncio.iscoroutine(res): - # Schedule coroutine and keep a reference to the task task = self._loop.create_task(res) self._tasks.add(task) task.add_done_callback(self._on_task_done) @@ -272,11 +260,9 @@ def _on_task_done(self, task: asyncio.Task) -> None: self._tasks.discard(task) if task.cancelled(): return - exc = task.exception() if exc is None: return - self.exception = exc if isinstance(exc, Exception): if not self._on_error(exc): @@ -304,7 +290,7 @@ def _on_error(self, exc: Exception) -> bool: return was_handled def add_listener(self, listener: MessageRecipient) -> None: - """Add new Listener for notification. + """Add new Listener to the notification set. :param listener: Listener to be added to the list to be notified """ @@ -315,8 +301,8 @@ def remove_listener(self, listener: MessageRecipient) -> None: throws an exception if the given listener is not part of the stored listeners. - :param listener: Listener to be removed from the set - :raises ValueError: if `listener` was never added to the notifier + :param listener: Listener to be removed from the set to be notified + :raises ValueError: if `listener` was never added to this notifier """ self.listeners.remove(listener) @@ -327,7 +313,7 @@ def stopped(self) -> bool: @staticmethod def find_instances(bus: BusABC) -> tuple["Notifier", ...]: - """Find :class:`~can.Notifier` instances associated with a given bus. + """Find :class:`~can.Notifier` instances associated with a given CAN bus. This method searches the registry for the :class:`~can.Notifier` that is linked to the specified bus. If the bus is found, the From 0c474186b69977220bb19360e9487f28c84fd46e Mon Sep 17 00:00:00 2001 From: timothyanderson096-ocdealcheck Date: Thu, 27 Aug 2026 16:27:35 +1000 Subject: [PATCH 4/5] Restore unrelated Notifier lines --- can/notifier.py | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/can/notifier.py b/can/notifier.py index da736e5df..cc25a46fb 100644 --- a/can/notifier.py +++ b/can/notifier.py @@ -165,16 +165,21 @@ def add_bus(self, bus: BusABC) -> None: :raises ValueError: If the *bus* is already assigned to an active :class:`~can.Notifier`. """ + # add bus to notifier registry Notifier._registry.register(bus, self) + + # add bus to internal bus list self._bus_list.append(bus) file_descriptor: int = -1 try: file_descriptor = bus.fileno() except NotImplementedError: + # Bus doesn't support fileno, we fall back to thread based reader pass if self._loop is not None and file_descriptor >= 0: + # Use bus file descriptor to watch for messages self._loop.add_reader(file_descriptor, self._on_message_available, bus) self._readers.append(file_descriptor) else: @@ -203,15 +208,18 @@ def stop(self, timeout: float = 5.0) -> None: if now < end_time: reader.join(end_time - now) elif self._loop: + # reader is a file descriptor self._loop.remove_reader(reader) for listener in self.listeners: if hasattr(listener, "stop"): listener.stop() + # remove bus from registry for bus in self._bus_list: Notifier._registry.unregister(bus, self) def _rx_thread(self, bus: BusABC) -> None: + # determine message handling callable early, not inside while loop if self._loop: handle_message: Callable[[Message], Any] = functools.partial( self._loop.call_soon_threadsafe, @@ -229,10 +237,13 @@ def _rx_thread(self, bus: BusABC) -> None: self.exception = exc if self._loop is not None: self._loop.call_soon_threadsafe(self._on_error, exc) + # Raise anyway raise elif not self._on_error(exc): + # If it was not handled, raise the exception here raise else: + # It was handled, so only log it logger.debug("suppressed exception: %s", exc) def _on_message_available(self, bus: BusABC) -> None: @@ -252,6 +263,7 @@ def _on_message_received(self, msg: Message) -> None: for callback in self.listeners: res = callback(msg) if res and self._loop and asyncio.iscoroutine(res): + # Schedule coroutine and keep a reference to the task task = self._loop.create_task(res) self._tasks.add(task) task.add_done_callback(self._on_task_done) @@ -260,9 +272,11 @@ def _on_task_done(self, task: asyncio.Task) -> None: self._tasks.discard(task) if task.cancelled(): return + exc = task.exception() if exc is None: return + self.exception = exc if isinstance(exc, Exception): if not self._on_error(exc): From fc695130274a41a991037df5170801d1f304f063 Mon Sep 17 00:00:00 2001 From: timothyanderson096-ocdealcheck Date: Thu, 27 Aug 2026 16:28:47 +1000 Subject: [PATCH 5/5] Narrow task exception type before storing it --- can/notifier.py | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/can/notifier.py b/can/notifier.py index cc25a46fb..97d762bf6 100644 --- a/can/notifier.py +++ b/can/notifier.py @@ -276,14 +276,13 @@ def _on_task_done(self, task: asyncio.Task) -> None: exc = task.exception() if exc is None: return + if not isinstance(exc, Exception): + raise exc self.exception = exc - if isinstance(exc, Exception): - if not self._on_error(exc): - raise exc - logger.debug("suppressed exception: %s", exc) - else: + if not self._on_error(exc): raise exc + logger.debug("suppressed exception: %s", exc) def _on_error(self, exc: Exception) -> bool: """Calls ``on_error()`` for all listeners if they implement it.