feat: implement broadcast rate limiting for ADC signals
CI for xtreme_bec / test (push) Successful in 31s
CI for xtreme_bec / test (push) Successful in 31s
This commit is contained in:
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user