From a8d440655148ba25d1dc6003d675140991dc9203 Mon Sep 17 00:00:00 2001 From: wyzula-jan Date: Fri, 3 Jul 2026 19:38:33 +0200 Subject: [PATCH] feat(dispatcher): sweep subscriptions of dead owners; self-heal on connect --- bec_widgets/utils/bec_dispatcher.py | 39 +++++++++++++++++++++++++++++ 1 file changed, 39 insertions(+) diff --git a/bec_widgets/utils/bec_dispatcher.py b/bec_widgets/utils/bec_dispatcher.py index 85254874..cca3ecb5 100644 --- a/bec_widgets/utils/bec_dispatcher.py +++ b/bec_widgets/utils/bec_dispatcher.py @@ -10,6 +10,7 @@ from typing import TYPE_CHECKING, Any, DefaultDict, Hashable import louie import redis +import shiboken6 from bec_lib.client import BECClient from bec_lib.logger import bec_logger from bec_lib.redis_connector import MessageObject, RedisConnector @@ -241,6 +242,9 @@ class BECDispatcher: "released automatically and stays subscribed until disconnect_slot is called. " "Pass owner= to bind its lifetime to a widget." ) + # Self-healing: reap subscriptions whose owners died without any close event + # before adding new ones, so stale registrations never accumulate. + self.cleanup_dead_slots() qt_slot = QtThreadSafeCallback(cb=slot, cb_info=cb_info, owner=owner) if not self.client.connector.any_stream_is_registered(topics, qt_slot): if qt_slot not in self._registered_slots: @@ -343,6 +347,41 @@ class BECDispatcher: qt_slot.topics.clear() self._registered_slots.pop(qt_slot, None) + def cleanup_dead_slots(self) -> int: + """ + Release subscriptions whose callback or owner no longer exists. + + Covers every death path that never delivers a close event: widgets destroyed + together with a parent, C++ objects deleted while a Python reference lingers, + and owners that were simply garbage collected. Runs automatically whenever a + BECWidget is destroyed and on every connect_slot, so stale subscriptions never + accumulate. + + Returns: + int: The number of released slot wrappers. + """ + dead_slots = [] + for qt_slot in list(self._registered_slots.values()): + cb = qt_slot.cb + if cb is None: + # The callback itself was garbage collected. + dead_slots.append(qt_slot) + continue + owner = qt_slot.cb_owner() if qt_slot.cb_owner is not None else None + if qt_slot.cb_owner is not None and owner is None: + # The owner was garbage collected (also anchors owner-bound lambdas). + dead_slots.append(qt_slot) + continue + if isinstance(owner, QObject) and not shiboken6.isValid(owner): + # The owner's C++ object was destroyed (e.g. with its parent) while a + # Python reference keeps the wrapper alive. + dead_slots.append(qt_slot) + for qt_slot in dead_slots: + self._release_slot(qt_slot) + if dead_slots: + logger.info(f"Released {len(dead_slots)} dispatcher slot(s) with dead owners") + return len(dead_slots) + def start_cli_server(self, gui_id: str | None = None): """ Start the CLI server.