mirror of
https://github.com/bec-project/bec_widgets.git
synced 2026-07-28 22:22:59 +02:00
perf(rpc): skip registry broadcasts when nothing changed
This commit is contained in:
@@ -324,6 +324,8 @@ class BECConnector:
|
||||
super().setObjectName(name)
|
||||
self.object_name = name
|
||||
if self.rpc_register.object_is_registered(self):
|
||||
# A rename changes the serialized registry state.
|
||||
self.rpc_register.mark_broadcast_pending()
|
||||
self.rpc_register.broadcast()
|
||||
|
||||
def submit_task(self, fn, *args, on_complete: SafeSlot = None, **kwargs) -> Worker:
|
||||
|
||||
@@ -26,6 +26,7 @@ def broadcast_update(func):
|
||||
@wraps(func)
|
||||
def wrapper(self, *args, **kwargs):
|
||||
result = func(self, *args, **kwargs)
|
||||
self.mark_broadcast_pending()
|
||||
self.broadcast()
|
||||
return result
|
||||
|
||||
@@ -53,6 +54,7 @@ class RPCRegister:
|
||||
self._broadcast_on_hold = RPCRegisterBroadcast(self)
|
||||
self._lock = RLock()
|
||||
self._skip_broadcast = False
|
||||
self._broadcast_pending = True
|
||||
self._initialized = True
|
||||
self.callbacks = []
|
||||
|
||||
@@ -146,16 +148,38 @@ class RPCRegister:
|
||||
widgets = [rpc for rpc in self._rpc_register.values() if isinstance(rpc, cls)]
|
||||
return [widget.object_name for widget in widgets]
|
||||
|
||||
def mark_broadcast_pending(self):
|
||||
"""
|
||||
Mark a broadcast as pending so the next broadcast() call serializes and
|
||||
delivers the state. Mutations (add_rpc/remove_rpc, renames, callback
|
||||
registration) call this; read-only paths leave nothing pending and
|
||||
their broadcasts become no-ops.
|
||||
"""
|
||||
self._broadcast_pending = True
|
||||
|
||||
def broadcast(self):
|
||||
"""
|
||||
Broadcast the update to all the callbacks. Callbacks whose owners have
|
||||
been garbage collected — or whose owning QObject's C++ side has been
|
||||
destroyed while the Python wrapper is still referenced — are pruned
|
||||
instead of being called.
|
||||
|
||||
Broadcasting is skipped while a delayed-broadcast context holds the
|
||||
registry, and when nothing changed since the last broadcast: the
|
||||
RPC execution path broadcasts after every call, and serializing an
|
||||
unchanged registry costs 1.5–8.5 ms at 25–200 widgets (measured),
|
||||
which would otherwise be paid by every read-only RPC call.
|
||||
"""
|
||||
|
||||
if self._skip_broadcast:
|
||||
return
|
||||
if not self._broadcast_pending:
|
||||
return
|
||||
if not self.callbacks:
|
||||
# No listeners: skip the registry walk but keep the broadcast pending so the
|
||||
# first callback added later still receives the current state.
|
||||
return
|
||||
self._broadcast_pending = False
|
||||
connections = self.list_all_connections()
|
||||
dead_refs = []
|
||||
for callback_ref in list(self.callbacks):
|
||||
@@ -204,6 +228,8 @@ class RPCRegister:
|
||||
callback_ref = safe_ref(callback)
|
||||
if callback_ref not in self.callbacks:
|
||||
self.callbacks.append(callback_ref)
|
||||
# The new callback has not seen any state yet.
|
||||
self.mark_broadcast_pending()
|
||||
|
||||
def remove_callback(self, callback: Callable[[dict], None]):
|
||||
"""
|
||||
|
||||
@@ -130,3 +130,62 @@ def test_callback_does_not_keep_owner_alive(rpc_register):
|
||||
|
||||
rpc_register.broadcast() # must not raise; prunes the dead reference
|
||||
assert len(rpc_register.callbacks) == callbacks_with_owner - 1
|
||||
|
||||
|
||||
def test_broadcast_skips_when_registry_unchanged(rpc_register):
|
||||
"""Perf regression test (audit item 24): the RPC execution path
|
||||
broadcasts after every call; an unchanged registry must not be
|
||||
re-serialized and callbacks must not be re-invoked."""
|
||||
owner = _CallbackOwner()
|
||||
rpc_register.add_callback(owner.on_update)
|
||||
|
||||
rpc_register.broadcast() # pending after add_callback -> delivers
|
||||
assert len(owner.received) == 1
|
||||
|
||||
rpc_register.broadcast() # nothing changed -> skipped
|
||||
rpc_register.broadcast()
|
||||
assert len(owner.received) == 1
|
||||
|
||||
class _Probe:
|
||||
gui_id = "pending_probe"
|
||||
object_name = "pending_probe"
|
||||
|
||||
import gc
|
||||
|
||||
probe = _Probe()
|
||||
try:
|
||||
rpc_register.add_rpc(probe) # mutation -> delivered
|
||||
assert len(owner.received) == 2
|
||||
rpc_register.broadcast() # clean again -> skipped
|
||||
assert len(owner.received) == 2
|
||||
finally:
|
||||
rpc_register.remove_rpc(probe)
|
||||
gc.collect()
|
||||
|
||||
|
||||
def test_broadcast_without_callbacks_skips_registry_walk(rpc_register, monkeypatch):
|
||||
"""With no listeners, a pending broadcast must not walk/serialize the registry,
|
||||
and the broadcast stays pending so the first callback added later receives it."""
|
||||
from unittest import mock
|
||||
|
||||
walk = mock.Mock(wraps=rpc_register.list_all_connections)
|
||||
monkeypatch.setattr(rpc_register, "list_all_connections", walk)
|
||||
# The test environment may carry ambient callbacks (e.g. the RPC server's);
|
||||
# this test owns the no-listeners precondition.
|
||||
monkeypatch.setattr(rpc_register, "callbacks", [])
|
||||
|
||||
rpc_register.mark_broadcast_pending()
|
||||
rpc_register.broadcast()
|
||||
walk.assert_not_called()
|
||||
|
||||
received = []
|
||||
|
||||
class _Owner:
|
||||
def on_update(self, connections):
|
||||
received.append(connections)
|
||||
|
||||
owner = _Owner()
|
||||
rpc_register.add_callback(owner.on_update)
|
||||
rpc_register.broadcast()
|
||||
walk.assert_called_once()
|
||||
assert len(received) == 1
|
||||
|
||||
Reference in New Issue
Block a user