diff options
Diffstat (limited to 'daq_system/core/acquisition.py')
| -rw-r--r-- | daq_system/core/acquisition.py | 214 |
1 files changed, 0 insertions, 214 deletions
diff --git a/daq_system/core/acquisition.py b/daq_system/core/acquisition.py deleted file mode 100644 index 7e43529..0000000 --- a/daq_system/core/acquisition.py +++ /dev/null @@ -1,214 +0,0 @@ -""" -core/acquisition.py - -Background acquisition engine. Polls all connected devices at -their configured sample rates and emits data via Qt signals. -""" - -import time -import threading -import csv -import os -from collections import deque -from datetime import datetime -from typing import Callable, Dict, List, Optional, Tuple - -from PyQt6.QtCore import QObject, pyqtSignal - -from devices.base_device import BaseDevice, DeviceStatus - - -MAX_BUFFER = 10_000 # points per channel - - -class ChannelBuffer: - """Ring buffer for one channel's time-series data.""" - def __init__(self, maxlen: int = MAX_BUFFER): - self.times: deque = deque(maxlen=maxlen) - self.values: deque = deque(maxlen=maxlen) - - def append(self, t: float, v: float): - self.times.append(t) - self.values.append(v) - - def latest(self, n: int = 1) -> Tuple[List[float], List[float]]: - ts = list(self.times)[-n:] - vs = list(self.values)[-n:] - return ts, vs - - def all(self) -> Tuple[List[float], List[float]]: - return list(self.times), list(self.values) - - def clear(self): - self.times.clear() - self.values.clear() - - def __len__(self): - return len(self.times) - - -class AcquisitionEngine(QObject): - """ - Runs a background polling thread for all registered devices. - - Signals: - new_data(device_id, channel_id, timestamp, value) - alarm_triggered(device_id, channel_id, value, kind) kind = 'low'|'high' - device_status_changed(device_id, status_str) - log_written(filepath) - """ - - new_data = pyqtSignal(str, str, float, float) - alarm_triggered = pyqtSignal(str, str, float, str) - device_status_changed = pyqtSignal(str, str) - log_written = pyqtSignal(str) - - def __init__(self, poll_interval_ms: int = 100): - super().__init__() - self.poll_interval_s = poll_interval_ms / 1000.0 - self._devices: List[BaseDevice] = [] - self._buffers: Dict[str, Dict[str, ChannelBuffer]] = {} - self._running = False - self._thread: Optional[threading.Thread] = None - self._start_time = 0.0 - - # Logging - self._logging = False - self._log_file: Optional[str] = None - self._csv_writer = None - self._csv_handle = None - - # Alarm state (prevent repeated triggers) - self._alarm_active: Dict[str, bool] = {} - - # ------------------------------------------------------------------ # - # Device management # - # ------------------------------------------------------------------ # - - def add_device(self, device: BaseDevice) -> None: - self._devices.append(device) - self._buffers[device.info.device_id] = { - ch.channel_id: ChannelBuffer() - for ch in device.info.channels - } - - def remove_device(self, device_id: str) -> None: - self._devices = [d for d in self._devices if d.info.device_id != device_id] - self._buffers.pop(device_id, None) - - def get_buffer(self, device_id: str, channel_id: str) -> Optional[ChannelBuffer]: - return self._buffers.get(device_id, {}).get(channel_id) - - # ------------------------------------------------------------------ # - # Start / Stop # - # ------------------------------------------------------------------ # - - def start(self) -> None: - if self._running: - return - self._running = True - self._start_time = time.time() - self._thread = threading.Thread(target=self._loop, daemon=True) - self._thread.start() - - def stop(self) -> None: - self._running = False - if self._thread: - self._thread.join(timeout=2.0) - self.stop_logging() - - # ------------------------------------------------------------------ # - # Logging # - # ------------------------------------------------------------------ # - - def start_logging(self, filepath: Optional[str] = None) -> str: - if filepath is None: - ts = datetime.now().strftime("%Y%m%d_%H%M%S") - os.makedirs("logs", exist_ok=True) - filepath = f"logs/daq_{ts}.csv" - self._log_file = filepath - self._csv_handle = open(filepath, "w", newline="") - # Build header - headers = ["timestamp"] - for dev in self._devices: - for ch in dev.info.channels: - headers.append(f"{dev.info.device_id}.{ch.channel_id}") - self._csv_writer = csv.writer(self._csv_handle) - self._csv_writer.writerow(headers) - self._logging = True - return filepath - - def stop_logging(self) -> None: - self._logging = False - if self._csv_handle: - try: - self._csv_handle.close() - except Exception: - pass - self._csv_handle = None - self._csv_writer = None - - # ------------------------------------------------------------------ # - # Background loop # - # ------------------------------------------------------------------ # - - def _loop(self) -> None: - while self._running: - t0 = time.time() - timestamp = t0 - self._start_time - log_row = [f"{timestamp:.3f}"] - - for dev in self._devices: - if dev.status not in (DeviceStatus.CONNECTED, DeviceStatus.SIMULATED): - for ch in dev.info.channels: - log_row.append("") - continue - try: - readings = dev.read_channels() - except Exception as e: - print(f"[Acq] Error reading {dev.info.device_id}: {e}") - readings = {} - - for ch in dev.info.channels: - val = readings.get(ch.channel_id) - if val is None: - log_row.append("") - continue - # Buffer - buf = self._buffers.get(dev.info.device_id, {}).get(ch.channel_id) - if buf is not None: - buf.append(timestamp, val) - # Signal (emit on main thread via Qt queued connection) - self.new_data.emit(dev.info.device_id, ch.channel_id, timestamp, val) - log_row.append(f"{val:.4f}") - # Alarms - self._check_alarm(dev.info.device_id, ch, val) - - if self._logging and self._csv_writer: - try: - self._csv_writer.writerow(log_row) - except Exception: - pass - - elapsed = time.time() - t0 - sleep_t = self.poll_interval_s - elapsed - if sleep_t > 0: - time.sleep(sleep_t) - - def _check_alarm(self, device_id: str, ch, val: float) -> None: - key_lo = f"{device_id}.{ch.channel_id}.low" - key_hi = f"{device_id}.{ch.channel_id}.high" - - if ch.alarm_low is not None: - if val < ch.alarm_low and not self._alarm_active.get(key_lo): - self._alarm_active[key_lo] = True - self.alarm_triggered.emit(device_id, ch.channel_id, val, "low") - elif val >= ch.alarm_low: - self._alarm_active[key_lo] = False - - if ch.alarm_high is not None: - if val > ch.alarm_high and not self._alarm_active.get(key_hi): - self._alarm_active[key_hi] = True - self.alarm_triggered.emit(device_id, ch.channel_id, val, "high") - elif val <= ch.alarm_high: - self._alarm_active[key_hi] = False |
