From 70ce00c7ada9b6589192336b475acb0aa256a162 Mon Sep 17 00:00:00 2001 From: wakonig_k Date: Tue, 23 Jun 2026 17:11:48 +0200 Subject: [PATCH] feat: implement broadcast rate limiting for ADC signals --- xtreme_bec/devices/x07ma_devices.py | 54 ++++++++++++++++++++++++----- 1 file changed, 46 insertions(+), 8 deletions(-) diff --git a/xtreme_bec/devices/x07ma_devices.py b/xtreme_bec/devices/x07ma_devices.py index efafac1..4ef5761 100644 --- a/xtreme_bec/devices/x07ma_devices.py +++ b/xtreme_bec/devices/x07ma_devices.py @@ -2,6 +2,7 @@ ophyd device classes for X07MA beamline """ +import threading import time import traceback from collections import OrderedDict @@ -341,15 +342,52 @@ class X07MAAnalogSignals(Device): ADC inputs """ - s1 = Cpt(EpicsSignalRO, "SIGNAL0", kind=Kind.hinted, auto_monitor=True) - s2 = Cpt(EpicsSignalRO, "SIGNAL1", kind=Kind.hinted, auto_monitor=True) - s3 = Cpt(EpicsSignalRO, "SIGNAL2", kind=Kind.hinted, auto_monitor=True) - s4 = Cpt(EpicsSignalRO, "SIGNAL3", kind=Kind.hinted, auto_monitor=True) - s5 = Cpt(EpicsSignalRO, "SIGNAL4", kind=Kind.hinted, auto_monitor=True) - s6 = Cpt(EpicsSignalRO, "SIGNAL5", kind=Kind.hinted, auto_monitor=True) - s7 = Cpt(EpicsSignalRO, "SIGNAL6", kind=Kind.hinted, auto_monitor=True) + SUB_VALUE = "value" + _default_sub = SUB_VALUE + + s1 = Cpt(EpicsSignalRO, "SIGNAL0", kind=Kind.hinted) + s2 = Cpt(EpicsSignalRO, "SIGNAL1", kind=Kind.hinted) + s3 = Cpt(EpicsSignalRO, "SIGNAL2", kind=Kind.hinted) + s4 = Cpt(EpicsSignalRO, "SIGNAL3", kind=Kind.hinted) + s5 = Cpt(EpicsSignalRO, "SIGNAL4", kind=Kind.hinted) + s6 = Cpt(EpicsSignalRO, "SIGNAL5", kind=Kind.hinted) + s7 = Cpt(EpicsSignalRO, "SIGNAL6", kind=Kind.hinted) norm_tey = Cpt(NormTEYSignals, name="norm_tey", kind=Kind.hinted) - norm_diode = Cpt(NormDIODESignals, name="norm_tey", kind=Kind.hinted) + norm_diode = Cpt(NormDIODESignals, name="norm_diode", kind=Kind.hinted) + + def __init__(self, *args, broadcast_rate_limit_s=0.1, **kwargs): + super().__init__(*args, **kwargs) + self._broadcast_rate_limit_s = broadcast_rate_limit_s + self._broadcast_lock = threading.Lock() + self._last_broadcast_ts = 0.0 + self._pending_broadcast = None + + for signal_name in ("s1", "s2", "s3", "s4", "s5", "s6", "s7"): + getattr(self, signal_name).subscribe(self._schedule_broadcast, run=False) + + def _schedule_broadcast(self, *args, **kwargs): + with self._broadcast_lock: + now = time.monotonic() + remaining = self._broadcast_rate_limit_s - (now - self._last_broadcast_ts) + if remaining <= 0 and self._pending_broadcast is None: + self._last_broadcast_ts = now + emit_now = True + else: + emit_now = False + if self._pending_broadcast is None: + self._pending_broadcast = threading.Timer(remaining, self._emit_broadcast) + self._pending_broadcast.daemon = True + self._pending_broadcast.start() + + if emit_now: + self._run_subs(sub_type=self.SUB_VALUE, value=self.read()) + + def _emit_broadcast(self): + with self._broadcast_lock: + self._pending_broadcast = None + self._last_broadcast_ts = time.monotonic() + + self._run_subs(sub_type=self.SUB_VALUE, value=self.read()) # Aliases # tey = s1