move new pvdatastream in place
This commit is contained in:
@@ -1,95 +1,48 @@
|
||||
class PvDataStream:
|
||||
from time import sleep
|
||||
from slic.utils.hastyepics import get_pv as PV
|
||||
|
||||
def __init__(self, ID, name=None):
|
||||
self.ID = ID
|
||||
self._pv = PV(ID)
|
||||
from .buffer import BufferInfinite, BufferFinite
|
||||
from .timer import Timer
|
||||
|
||||
|
||||
class PVDataStream:
|
||||
|
||||
def __init__(self, name, wait_time=0.1):
|
||||
self.name = name
|
||||
self.wait_time = wait_time
|
||||
self.pv = PV(name)
|
||||
self.running = False
|
||||
|
||||
|
||||
def record(self, n=None, seconds=None):
|
||||
pv = self.pv
|
||||
|
||||
def acquire(self, hold=False, **kwargs):
|
||||
_acquire = lambda: self.collect(**kwargs)
|
||||
return Task(_acquire, hold=hold)
|
||||
if n is None:
|
||||
buf = BufferInfinite()
|
||||
else:
|
||||
current = pv.get()
|
||||
buf = BufferFinite.from_example(n, current)
|
||||
|
||||
def on_value_change(value=None, **kwargs):
|
||||
buf.append(value)
|
||||
if buf.is_full:
|
||||
self.stop()
|
||||
|
||||
pv.add_callback(callback=on_value_change)
|
||||
|
||||
def collect(self, seconds=None, samples=None):
|
||||
if seconds is None and samples is None:
|
||||
raise ValueError("Either a time interval or number of samples need to be defined.")
|
||||
self.running = True
|
||||
tim = Timer(seconds)
|
||||
while self.running and not tim.is_done:
|
||||
sleep(self.wait_time)
|
||||
|
||||
if not hasattr(self, "_collection"):
|
||||
self._accumulate = {"n_cb": None}
|
||||
self.stop()
|
||||
|
||||
self._pv.callbacks.pop(self._collection["n_cb"], None)
|
||||
return buf.data
|
||||
|
||||
self._collection = {"done": False}
|
||||
|
||||
self.data_collected = []
|
||||
|
||||
if seconds:
|
||||
stopcond = mk_collect_seconds_stopcond(seconds)
|
||||
elif samples:
|
||||
stopcond = mk_collect_samples_stopcond(samples)
|
||||
|
||||
def addData(value=None, **kwargs):
|
||||
self.data_collected.append(value)
|
||||
if stopcond():
|
||||
self._pv.callbacks.pop(self._collection["n_cb"])
|
||||
self._collection["done"] = True
|
||||
|
||||
self._collection["n_cb"] = self._pv.add_callback(addData)
|
||||
|
||||
while not self._collection["done"]:
|
||||
sleep(0.005)
|
||||
|
||||
return self.data_collected
|
||||
|
||||
|
||||
|
||||
def mk_collect_seconds_stopcond(self, seconds):
|
||||
self._collection["start_time"] = time()
|
||||
self._collection["seconds"] = seconds
|
||||
return lambda: (time() - self._collection["start_time"]) > self._collection["seconds"]
|
||||
|
||||
|
||||
|
||||
def mk_collect_samples_stopcond(self, samples):
|
||||
self._collection["samples"] = samples
|
||||
return lambda: len(self.data_collected) >= self._collection["samples"]
|
||||
|
||||
|
||||
|
||||
def accumulate(self, n_buffer):
|
||||
if not hasattr(self, "_accumulate"):
|
||||
self._accumulate = {"n_cb": None}
|
||||
|
||||
self._pv.callbacks.pop(self._accumulate["n_cb"], None)
|
||||
|
||||
self._accumulate["n_buffer"] = n_buffer
|
||||
self._accumulate["ix"] = 0
|
||||
|
||||
shape = (n_buffer * 2, self._pv.count)
|
||||
self._data = np.squeeze(np.zeros(shape)) * np.nan
|
||||
|
||||
def addData(value=None, **kwargs):
|
||||
index = self._accumulate["ix"] + 1
|
||||
n_buffer = self._accumulate["n_buffer"]
|
||||
self._accumulate["ix"] = index % n_buffer
|
||||
|
||||
start = self._accumulate["ix"]
|
||||
step = self._accumulate["n_buffer"]
|
||||
self._data[start :: step] = value
|
||||
|
||||
self._accumulate["n_cb"] = self._pv.add_callback(addData)
|
||||
|
||||
|
||||
|
||||
def get_data(self):
|
||||
start = self._accumulate["ix"] + 1
|
||||
stop = start + self._accumulate["n_buffer"]
|
||||
return self._data[start:stop]
|
||||
|
||||
data = property(get_data)
|
||||
def stop(self):
|
||||
self.pv.clear_callbacks()
|
||||
self.running = False
|
||||
|
||||
|
||||
|
||||
|
||||
@@ -1,48 +0,0 @@
|
||||
from time import sleep
|
||||
from slic.utils.hastyepics import get_pv as PV
|
||||
|
||||
from .buffer import BufferInfinite, BufferFinite
|
||||
from .timer import Timer
|
||||
|
||||
|
||||
class PVDataStream:
|
||||
|
||||
def __init__(self, name, wait_time=0.1):
|
||||
self.name = name
|
||||
self.wait_time = wait_time
|
||||
self.pv = PV(name)
|
||||
self.running = False
|
||||
|
||||
|
||||
def record(self, n=None, seconds=None):
|
||||
pv = self.pv
|
||||
|
||||
if n is None:
|
||||
buf = BufferInfinite()
|
||||
else:
|
||||
current = pv.get()
|
||||
buf = BufferFinite.from_example(n, current)
|
||||
|
||||
def on_value_change(value=None, **kwargs):
|
||||
buf.append(value)
|
||||
if buf.is_full:
|
||||
self.stop()
|
||||
|
||||
pv.add_callback(callback=on_value_change)
|
||||
|
||||
self.running = True
|
||||
tim = Timer(seconds)
|
||||
while self.running and not tim.is_done:
|
||||
sleep(self.wait_time)
|
||||
|
||||
self.stop()
|
||||
|
||||
return buf.data
|
||||
|
||||
|
||||
def stop(self):
|
||||
self.pv.clear_callbacks()
|
||||
self.running = False
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user