""" 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: derivative/second_derivative 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 re 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) def register_filter_class(name: str, cls) -> None: """Register a plugin-provided filter class by type name.""" FILTER_CLASSES[name] = cls def unregister_filter_class(name: str) -> None: """Remove a previously registered plugin filter class.""" FILTER_CLASSES.pop(name, None) # ══════════════════════════════════════════════════════════════════════════════ # Helpers # ══════════════════════════════════════════════════════════════════════════════ def _make_src_names(sources: List[Tuple[str, str]]) -> List[str]: """ Build a list of Python-safe variable names from (device_id, channel_id) pairs. Uses channel_id alone when unambiguous; prefixes device_id when two sources share the same channel_id from different devices. """ seen: Dict[str, str] = {} # ch_id -> dev_id of first occurrence names = [] for dev_id, ch_id in sources: if ch_id in seen and seen[ch_id] != dev_id: raw = f"{dev_id}_{ch_id}" else: seen[ch_id] = dev_id raw = ch_id safe = re.sub(r"\W", "_", raw) if safe and safe[0].isdigit(): safe = "_" + safe names.append(safe or "_ch") return names # ══════════════════════════════════════════════════════════════════════════════ # Derived channel definitions # ══════════════════════════════════════════════════════════════════════════════ @dataclass class DerivedChannel: """ A virtual channel computed from one or more physical channels. kind options: "derivative" — derivative of a source (dy/dt) "second_derivative" — second derivative of a source (d²y/dt²) "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 and per-channel script state (not serialised) _fn: Optional[Callable] = field(default=None, repr=False, compare=False) _exec_state: dict = field(default_factory=dict, repr=False, compare=False) def compile(self) -> Optional[str]: """ Compile expression/script into self._fn. Returns None on success, or error string on failure. Namespace available in all expression/script kinds: x — list of source values in declaration order — each source channel_id as a named variable (x[0] == first source name) t — elapsed time in seconds math — Python math module vars — shared SignalProcessor._script_vars dict (read/write) state — per-channel persistent dict (survives between evaluations) """ try: src_names = _make_src_names(self.sources) if self.kind == "expression": code = compile(f"__result__ = {self.expression}", "", "exec") _state = self._exec_state def _expr_fn(inputs, t, sv, _code=code, _names=src_names, _st=_state): ns = {"x": inputs, "t": t, "math": math, "vars": sv, "state": _st} for i, name in enumerate(_names): if i < len(inputs): ns[name] = inputs[i] exec(_code, ns) return float(ns["__result__"]) self._fn = _expr_fn elif self.kind in ("function", "custom_script"): src = self.script if not src.strip().startswith("def compute"): src = "def compute(x, t):\n" + "\n".join( " " + ln for ln in src.splitlines() ) _exec_ns: dict = {"math": math} exec(compile(src, "