From d5acb04b88373d33b038bb59945fb5ab8b4f543b Mon Sep 17 00:00:00 2001 From: Christian Kolset Date: Mon, 20 Apr 2026 16:55:57 -0600 Subject: V8 --- core/signal_processor.py | 464 +++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 464 insertions(+) create mode 100644 core/signal_processor.py (limited to 'core/signal_processor.py') diff --git a/core/signal_processor.py b/core/signal_processor.py new file mode 100644 index 0000000..4161d4d --- /dev/null +++ b/core/signal_processor.py @@ -0,0 +1,464 @@ +""" +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, "