""" core/signal_processor.py Signal processing pipeline engine. Provides two capabilities: 1. FILTERS — applied to a raw channel buffer before plotting: LowPass, HighPass, MovingAverage, Median, Derivative, Integral, Scale+Offset 2. DERIVED CHANNELS — virtual channels computed from one or more physical channels. Each derived channel runs a user-defined function every time new data arrives. Built-ins: velocity/acceleration from displacement, power from voltage+current, RMS, etc. Custom: arbitrary Python snippet. Architecture ------------ SignalProcessor sits between AcquisitionEngine and StripChartWidget. engine.new_data → SignalProcessor.process(dev, ch, t, val) → emits processed_data(virtual_or_real_id, ch_id, t, val) The processor maintains its own ring buffers for derived channels so the strip chart can query history just like physical channels. """ from __future__ import annotations import math import threading import traceback from collections import deque from dataclasses import dataclass, field from typing import Callable, Dict, List, Optional, Tuple, Any from PyQt6.QtCore import QObject, pyqtSignal # ── Constants ────────────────────────────────────────────────────────────── MAX_BUF = 20_000 # ══════════════════════════════════════════════════════════════════════════════ # Filter definitions # ══════════════════════════════════════════════════════════════════════════════ class FilterBase: """All filters implement __call__(value: float) -> float.""" name: str = "identity" params: dict = {} def __call__(self, value: float) -> float: return value def reset(self): pass def to_dict(self) -> dict: return {"type": self.name, **self.params} class MovingAverageFilter(FilterBase): name = "moving_average" def __init__(self, window: int = 10): self.params = {"window": window} self._buf = deque(maxlen=window) def __call__(self, v: float) -> float: self._buf.append(v) return sum(self._buf) / len(self._buf) def reset(self): self._buf.clear() class MedianFilter(FilterBase): name = "median" def __init__(self, window: int = 5): self.params = {"window": window} self._buf = deque(maxlen=window) def __call__(self, v: float) -> float: self._buf.append(v) s = sorted(self._buf) n = len(s) return s[n // 2] if n % 2 else (s[n//2 - 1] + s[n//2]) / 2 def reset(self): self._buf.clear() class LowPassFilter(FilterBase): """Exponential moving average (single-pole IIR low-pass).""" name = "low_pass" def __init__(self, alpha: float = 0.1): """alpha=0.0 → no change, 1.0 → unfiltered.""" self.params = {"alpha": alpha} self._prev = None def __call__(self, v: float) -> float: if self._prev is None: self._prev = v self._prev = self._prev + self.params["alpha"] * (v - self._prev) return self._prev def reset(self): self._prev = None class HighPassFilter(FilterBase): """Simple single-pole IIR high-pass (compliment of low-pass).""" name = "high_pass" def __init__(self, alpha: float = 0.9): self.params = {"alpha": alpha} self._prev_v = None self._prev_y = 0.0 def __call__(self, v: float) -> float: if self._prev_v is None: self._prev_v = v y = self.params["alpha"] * (self._prev_y + v - self._prev_v) self._prev_y = y self._prev_v = v return y def reset(self): self._prev_v = None; self._prev_y = 0.0 class ScaleOffsetFilter(FilterBase): """y = scale * x + offset (unit conversion, calibration).""" name = "scale_offset" def __init__(self, scale: float = 1.0, offset: float = 0.0): self.params = {"scale": scale, "offset": offset} def __call__(self, v: float) -> float: return self.params["scale"] * v + self.params["offset"] class DerivativeFilter(FilterBase): """Numerical first derivative dy/dt.""" name = "derivative" def __init__(self): self.params = {}; self._prev_v = None; self._prev_t = None def process_with_t(self, v: float, t: float) -> float: if self._prev_t is None or t == self._prev_t: self._prev_v = v; self._prev_t = t; return 0.0 dy = (v - self._prev_v) / (t - self._prev_t) self._prev_v = v; self._prev_t = t return dy def __call__(self, v: float) -> float: return 0.0 # use process_with_t for real output def reset(self): self._prev_v = None; self._prev_t = None class IntegralFilter(FilterBase): """Numerical integration (trapezoidal rule).""" name = "integral" def __init__(self): self.params = {}; self._sum = 0.0; self._prev_v = None; self._prev_t = None def process_with_t(self, v: float, t: float) -> float: if self._prev_t is not None and t != self._prev_t: self._sum += 0.5 * (v + self._prev_v) * (t - self._prev_t) self._prev_v = v; self._prev_t = t return self._sum def __call__(self, v: float) -> float: return self._sum def reset(self): self._sum = 0.0; self._prev_v = None; self._prev_t = None FILTER_CLASSES = { "moving_average": MovingAverageFilter, "median": MedianFilter, "low_pass": LowPassFilter, "high_pass": HighPassFilter, "scale_offset": ScaleOffsetFilter, "derivative": DerivativeFilter, "integral": IntegralFilter, } def filter_from_dict(d: dict) -> FilterBase: cls = FILTER_CLASSES.get(d.get("type", "")) if cls is None: return FilterBase() params = {k: v for k, v in d.items() if k != "type"} return cls(**params) # ══════════════════════════════════════════════════════════════════════════════ # Derived channel definitions # ══════════════════════════════════════════════════════════════════════════════ @dataclass class DerivedChannel: """ A virtual channel computed from one or more physical channels. kind options: "velocity" — derivative of a displacement source "acceleration" — second derivative of a displacement source "power" — voltage_source * current_source "rms" — rolling RMS of a source (window samples) "expression" — arbitrary Python expression string "function" — multi-line Python function body (def compute(...)) "custom_script" — full Python script, must define compute(inputs, t) """ channel_id: str # virtual ID, e.g. "vel_0" name: str # display name unit: str = "" color: str = "#f72585" kind: str = "expression" # see above # Source channel references [("dev_id", "ch_id"), ...] sources: List[Tuple[str, str]] = field(default_factory=list) # For built-in kinds params: Dict[str, Any] = field(default_factory=dict) # For expression / function / custom_script expression: str = "" # single-line: "x[0] * 2" script: str = "" # multi-line function body enabled: bool = True # Runtime: compiled callable (not serialised) _fn: Optional[Callable] = field(default=None, repr=False, compare=False) def compile(self) -> Optional[str]: """ Compile expression/script into self._fn. Returns None on success, or error string on failure. """ try: if self.kind == "expression": # Single-line: inputs are x (list of latest values), t (time) code = compile(f"__result__ = {self.expression}", "", "exec") def _expr_fn(inputs, t, _code=code): ns = {"x": inputs, "t": t, "math": math} exec(_code, ns) return float(ns["__result__"]) self._fn = _expr_fn elif self.kind in ("function", "custom_script"): # User provides a def compute(x, t): ... body # We wrap it in a module namespace src = self.script if not src.strip().startswith("def compute"): src = "def compute(x, t):\n" + "\n".join( " " + ln for ln in src.splitlines() ) ns: dict = {"math": math} exec(compile(src, "