From 5068c951583ffc3a916c0a212d6b586b6e13ba04 Mon Sep 17 00:00:00 2001 From: Martin Tazl Date: Fri, 4 Sep 2026 08:00:59 +0200 Subject: [PATCH] Add generic peak cycle detector --- PROJECT_KNOWLEDGE.md | 24 ++ README.md | 22 ++ config/calculations.example.yaml | 6 +- docs/architecture.md | 16 ++ docs/roadmap.md | 10 +- scripts/validate_k7_roll_length.py | 61 +++++ .../calculations/__init__.py | 9 +- .../calculations/peak_cycles.py | 157 ++++++++++++ tests/test_peak_cycles.py | 224 ++++++++++++++++++ 9 files changed, 521 insertions(+), 8 deletions(-) create mode 100644 scripts/validate_k7_roll_length.py create mode 100644 src/production_analytics/calculations/peak_cycles.py create mode 100644 tests/test_peak_cycles.py diff --git a/PROJECT_KNOWLEDGE.md b/PROJECT_KNOWLEDGE.md index 1652a46..74c67c9 100644 --- a/PROJECT_KNOWLEDGE.md +++ b/PROJECT_KNOWLEDGE.md @@ -17,6 +17,30 @@ ENLYZE and is retrieved again when historical calculations need reproduction. - Calculations are modular, testable Python implementations. Configuration declares instances; it is not a generic low-code language. +## Peak-cycle detector + +The first generic calculation is a timestamp-based `PeakCycleDetector` under +`calculations`. It is transport-independent and works for roll length, roll +weight, and similar cyclic signals. It maintains the maximum in an open cycle; +after a value falls below `drop_ratio * current_peak`, it emits that maximum +only after actual subsequent below-threshold observations support a continuous +`hold_seconds` interval. `min_peak` is a generic detector setting, along with +`drop_ratio` (0 < ratio < 1), `hold_seconds` (>= 0), and an explicit positive +`max_sample_gap_seconds`. Elapsed time, never a sample count, determines +confirmation. An interval longer than the configured maximum sample gap +restarts a reset candidate, so a timestamp gap is not +continuous-below evidence. Reset values can be negative rather than zero, +equal timestamps add no elapsed time, and backwards timestamps are rejected. +The maximum gap is required rather than inferred or defaulted, making each +source's continuity assumption explicit. + +K7 roll length, variable `bf81c547-dccf-4709-aee2-0f78366d1dfc` (m), is the +first real-world validation case. The tested capture is continuous on a +10-second grid through its latest available sample; its shorter-than-requested +result was caused by an end time queried in the future, not observed time-series +gaps. The raw capture remains ignored and deterministic tests use synthetic +fixtures only. + ## Central domain context Production orders connect machine, article/material, source interval, and diff --git a/README.md b/README.md index 42552cd..45d6876 100644 --- a/README.md +++ b/README.md @@ -62,6 +62,28 @@ The current Compose file has no application service, so it deliberately does not pass ENLYZE credentials to TimescaleDB. A future application service should use `env_file: ./secrets/enlyze.env` rather than copying secrets into Compose. +## Peak-cycle detection + +`PeakCycleDetector` is a pure calculation-domain component for roll length, +roll weight, and comparable sawtooth/batch signals. It retains the current +maximum and emits it only after the value is below +`drop_ratio * current_peak` continuously for `hold_seconds`. Configuration is +`min_peak`, `drop_ratio` (strictly between 0 and 1), non-negative +`hold_seconds`, and positive `max_sample_gap_seconds`. The hold uses elapsed +timestamps, never a sample count; an interval longer than the configured +maximum sample gap restarts a reset candidate, +so a timestamp gap alone is not evidence that a signal remained below threshold. +`max_sample_gap_seconds` is required rather than defaulted, so the source's +continuity assumption is explicit for each detector configuration. +Reset values can be negative and need not be zero. Equal timestamps are +accepted in arrival order but add no elapsed hold time; backwards timestamps +are rejected. + +The first real-world validation case is K7 roll length (m), variable +`bf81c547-dccf-4709-aee2-0f78366d1dfc`. When the ignored local capture is +present, run `python scripts/validate_k7_roll_length.py` to inspect detected +peaks without adding the raw capture to tests or version control. + `docker compose up -d timescaledb` is an optional local database design for a future persistence milestone. It is not required for the bootstrap tests. diff --git a/config/calculations.example.yaml b/config/calculations.example.yaml index 04c85e5..76b22db 100644 --- a/config/calculations.example.yaml +++ b/config/calculations.example.yaml @@ -16,6 +16,8 @@ calculations: version: "1" machine_ref: machine-to-be-verified signal_ref: process-signal-to-be-verified - below_peak_fraction: 0.75 - below_peak_duration_seconds: 60 + min_peak: 10.0 + drop_ratio: 0.75 + hold_seconds: 60.0 + max_sample_gap_seconds: 20.0 output_event: cycle_peak diff --git a/docs/architecture.md b/docs/architecture.md index 116a226..e463dd9 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -39,3 +39,19 @@ Future calculation state is stored per calculation instance and relevant context partition. It enables safe continuation (for example an open peak cycle), while the source interval recorded on results keeps a historical run reproducible by fetching ENLYZE data again. + +## Peak-cycle semantics + +The generic peak detector maintains an open cycle maximum. A reset begins when +a sample falls below `drop_ratio * current_peak`; it is confirmed only after +actual subsequent below-threshold samples support `hold_seconds` of elapsed +time. It emits only maxima at least `min_peak`, starts a fresh cycle after a +completed reset, and suppresses duplicates during the continuing low phase. +This is calculation-domain logic, with no ENLYZE transport dependency. + +The detector uses timestamps rather than a sample count. An interval longer +than the configured positive `max_sample_gap_seconds` restarts a reset candidate, +so a timestamp alone does not establish that a signal was continuously below +threshold. Equal +timestamps are processed in arrival order but contribute no elapsed time, and +backwards timestamps are rejected. Reset values are not assumed to be zero. diff --git a/docs/roadmap.md b/docs/roadmap.md index df6e7a9..69df6f4 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -26,12 +26,12 @@ Implement one configured integration calculation over a verified signal and production-order context. Include backfill, rerun/provenance behavior, and tests against sanitized fixtures. -## 4. Peak detection +## 4. Peak detection (first detector completed) -Implement the configurable sawtooth peak detector: retain the current maximum; -close a cycle after a configurable below-peak fraction remains true for a -configurable duration; emit one completed peak event; preserve open-cycle -state. +The generic configurable sawtooth detector is implemented: it retains the +current maximum, closes a cycle after a configurable below-peak fraction has +actual sample support for a configurable elapsed duration, and emits one +completed peak. Its pure state can later be persisted for incremental use. ## Later diff --git a/scripts/validate_k7_roll_length.py b/scripts/validate_k7_roll_length.py new file mode 100644 index 0000000..602d2d8 --- /dev/null +++ b/scripts/validate_k7_roll_length.py @@ -0,0 +1,61 @@ +#!/usr/bin/env python3 +"""Validate the generic peak detector against an ignored local K7 capture. + +This developer utility reads no network data and does not modify the capture. +Run it from the repository root when data/raw/enlyze/k7-roll-length-2h.raw.json +is available. +""" + +from __future__ import annotations + +import json +from datetime import datetime +from pathlib import Path + +from production_analytics.calculations import ( + PeakCycleDetector, + PeakDetectorConfig, + TimeSeriesSample, +) + + +CAPTURE = Path("data/raw/enlyze/k7-roll-length-2h.raw.json") +VARIABLE_UUID = "bf81c547-dccf-4709-aee2-0f78366d1dfc" + + +def main() -> int: + if not CAPTURE.is_file(): + print(f"K7 capture is not available: {CAPTURE}") + return 1 + + payload = json.loads(CAPTURE.read_text(encoding="utf-8")) + data = payload["response"]["body"]["data"] + value_index = data["columns"].index(VARIABLE_UUID) + samples = ( + TimeSeriesSample( + datetime.fromisoformat(record[0].replace("Z", "+00:00")), + float(record[value_index]), + ) + for record in data["records"] + if record[value_index] is not None + ) + detector = PeakCycleDetector( + PeakDetectorConfig( + min_peak=10.0, + drop_ratio=0.75, + hold_seconds=60.0, + max_sample_gap_seconds=20.0, + ) + ) + peaks = list(detector.process_many(samples)) + print(f"Detected peaks: {len(peaks)}") + for peak in peaks: + print( + f"peak={peak.value:.6f} at {peak.peak_timestamp.isoformat()} " + f"confirmed={peak.confirmed_at.isoformat()}" + ) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/src/production_analytics/calculations/__init__.py b/src/production_analytics/calculations/__init__.py index 26c27b2..f7ba102 100644 --- a/src/production_analytics/calculations/__init__.py +++ b/src/production_analytics/calculations/__init__.py @@ -1,5 +1,12 @@ """Versioned, pure calculation operators and their contracts.""" from .base import Calculation +from .peak_cycles import DetectedPeak, PeakCycleDetector, PeakDetectorConfig, TimeSeriesSample -__all__ = ["Calculation"] +__all__ = [ + "Calculation", + "DetectedPeak", + "PeakCycleDetector", + "PeakDetectorConfig", + "TimeSeriesSample", +] diff --git a/src/production_analytics/calculations/peak_cycles.py b/src/production_analytics/calculations/peak_cycles.py new file mode 100644 index 0000000..aba8b09 --- /dev/null +++ b/src/production_analytics/calculations/peak_cycles.py @@ -0,0 +1,157 @@ +"""Generic, timestamp-based detection of completed sawtooth peak cycles. + +The detector is deliberately independent of input transport and persistence. +Call :meth:`PeakCycleDetector.process` for each chronological sample, or use +``process_many`` for an iterable of :class:`TimeSeriesSample` values. +""" + +from collections.abc import Iterable, Iterator +from dataclasses import dataclass +from datetime import datetime +from math import isfinite + + +@dataclass(frozen=True, slots=True) +class TimeSeriesSample: + """One numeric observation from a chronological time series.""" + + timestamp: datetime + value: float + + +@dataclass(frozen=True, slots=True) +class PeakDetectorConfig: + """Configuration for :class:`PeakCycleDetector`.""" + + min_peak: float + drop_ratio: float + hold_seconds: float + max_sample_gap_seconds: float + + def __post_init__(self) -> None: + if not isfinite(self.min_peak): + raise ValueError("min_peak must be finite") + if not isfinite(self.drop_ratio) or not 0 < self.drop_ratio < 1: + raise ValueError("drop_ratio must be greater than 0 and less than 1") + if not isfinite(self.hold_seconds) or self.hold_seconds < 0: + raise ValueError("hold_seconds must be finite and at least 0") + if not isfinite(self.max_sample_gap_seconds) or self.max_sample_gap_seconds <= 0: + raise ValueError("max_sample_gap_seconds must be finite and greater than 0") + + +@dataclass(frozen=True, slots=True) +class DetectedPeak: + """A maximum whose following reset phase has been confirmed.""" + + value: float + peak_timestamp: datetime + confirmed_at: datetime + cycle_started_at: datetime + + +class PeakCycleDetector: + """Maintain one maximum and emit it once its reset phase is confirmed. + + A confirmation requires a sample below ``drop_ratio * current_peak`` to + start the reset phase and later observed below-threshold samples spanning + ``hold_seconds``. A reset candidate is restarted when the interval between + observations exceeds ``max_sample_gap_seconds``: elapsed time is used, but + an unobserved gap on its own never provides evidence that the signal stayed + below the threshold. + Equal timestamps are accepted in input order and contribute no elapsed + time; timestamps that move backwards are rejected. + """ + + def __init__(self, config: PeakDetectorConfig) -> None: + self.config = config + self._last_timestamp: datetime | None = None + self._cycle_started_at: datetime | None = None + self._current_peak: float | None = None + self._peak_timestamp: datetime | None = None + self._below_since: datetime | None = None + # A completed reset starts the next cycle at a low value. Do not let + # that same low/reset phase become a second cycle before it rises. + self._needs_rise_after_reset = False + self._reset_baseline: float | None = None + + def process(self, timestamp: datetime, value: float) -> DetectedPeak | None: + """Consume one sample and return a newly confirmed peak, if any.""" + if not isfinite(value): + raise ValueError("sample value must be finite") + if self._last_timestamp is not None and timestamp < self._last_timestamp: + raise ValueError("sample timestamps must not move backwards") + previous_timestamp = self._last_timestamp + self._last_timestamp = timestamp + + if self._current_peak is None: + self._start_cycle(timestamp, value, needs_rise=False) + return None + + if self._needs_rise_after_reset: + assert self._reset_baseline is not None + if value > self._reset_baseline: + self._needs_rise_after_reset = False + else: + # Keep collecting the highest reset-phase value, while still + # requiring a genuine rise before a new reset can be detected. + if value > self._current_peak: + self._current_peak = value + self._peak_timestamp = timestamp + self._reset_baseline = value + return None + + assert self._peak_timestamp is not None + assert self._cycle_started_at is not None + if value > self._current_peak: + self._current_peak = value + self._peak_timestamp = timestamp + self._below_since = None + return None + + reset_threshold = self.config.drop_ratio * self._current_peak + if value >= reset_threshold: + self._below_since = None + return None + + if self._below_since is None: + self._below_since = timestamp + elif ( + previous_timestamp is not None + and (timestamp - previous_timestamp).total_seconds() + > self.config.max_sample_gap_seconds + ): + # We have no evidence that the signal remained below threshold + # during the unobserved portion of this gap, so this sample begins + # a new candidate. + self._below_since = timestamp + return None + + elapsed_seconds = (timestamp - self._below_since).total_seconds() + if elapsed_seconds < self.config.hold_seconds: + return None + + event = None + if self._current_peak >= self.config.min_peak: + event = DetectedPeak( + value=self._current_peak, + peak_timestamp=self._peak_timestamp, + confirmed_at=timestamp, + cycle_started_at=self._cycle_started_at, + ) + self._start_cycle(timestamp, value, needs_rise=True) + return event + + def process_many(self, samples: Iterable[TimeSeriesSample]) -> Iterator[DetectedPeak]: + """Yield each peak confirmed while consuming ``samples``.""" + for sample in samples: + event = self.process(sample.timestamp, sample.value) + if event is not None: + yield event + + def _start_cycle(self, timestamp: datetime, value: float, *, needs_rise: bool) -> None: + self._cycle_started_at = timestamp + self._current_peak = value + self._peak_timestamp = timestamp + self._below_since = None + self._needs_rise_after_reset = needs_rise + self._reset_baseline = value if needs_rise else None diff --git a/tests/test_peak_cycles.py b/tests/test_peak_cycles.py new file mode 100644 index 0000000..fc06827 --- /dev/null +++ b/tests/test_peak_cycles.py @@ -0,0 +1,224 @@ +import unittest +from datetime import UTC, datetime, timedelta + +from production_analytics.calculations import ( + PeakCycleDetector, + PeakDetectorConfig, + TimeSeriesSample, +) + + +BASE = datetime(2026, 9, 4, tzinfo=UTC) + + +def samples(*values_and_seconds: tuple[float, float]) -> list[TimeSeriesSample]: + return [ + TimeSeriesSample(BASE + timedelta(seconds=seconds), value) + for value, seconds in values_and_seconds + ] + + +class PeakCycleDetectorTests(unittest.TestCase): + def detector(self, **overrides: float) -> PeakCycleDetector: + config = PeakDetectorConfig( + min_peak=10.0, + drop_ratio=0.75, + hold_seconds=60.0, + max_sample_gap_seconds=20.0, + **overrides, + ) + return PeakCycleDetector(config) + + def events(self, detector: PeakCycleDetector, sample_values: list[TimeSeriesSample]): + return list(detector.process_many(sample_values)) + + def test_simple_sawtooth_emits_one_peak_after_hold(self) -> None: + events = self.events( + self.detector(), + samples( + (2, 0), (12, 10), (20, 20), (14, 30), (13, 40), (12, 50), + (11, 60), (10, 70), (9, 80), (8, 90), + ), + ) + + self.assertEqual(len(events), 1) + self.assertEqual(events[0].value, 20) + self.assertEqual(events[0].peak_timestamp, BASE + timedelta(seconds=20)) + self.assertEqual(events[0].confirmed_at, BASE + timedelta(seconds=90)) + + def test_multiple_cycles_emit_once_each(self) -> None: + events = self.events( + self.detector(), + samples( + (15, 0), (30, 10), (5, 20), (4, 30), (4, 40), (4, 50), + (4, 60), (4, 70), (4, 80), (8, 90), (25, 100), (6, 110), + (5, 120), (5, 130), (5, 140), (5, 150), (5, 160), (5, 170), + ), + ) + + self.assertEqual([event.value for event in events], [30, 25]) + self.assertEqual( + [event.confirmed_at for event in events], + [BASE + timedelta(seconds=80), BASE + timedelta(seconds=170)], + ) + + def test_peak_below_minimum_is_not_emitted(self) -> None: + events = self.events(self.detector(), samples((2, 0), (9, 10), (1, 20), (1, 80))) + + self.assertEqual(events, []) + + def test_recovery_before_hold_cancels_reset_candidate(self) -> None: + events = self.events( + self.detector(), + samples( + (10, 0), (20, 10), (14, 20), (16, 50), (16, 80), (10, 90), + (10, 100), (10, 110), (10, 120), (10, 130), (10, 140), (10, 150), + ), + ) + + self.assertEqual(len(events), 1) + self.assertEqual(events[0].confirmed_at, BASE + timedelta(seconds=150)) + + def test_sustained_drop_emits_at_one_actual_confirmation_sample(self) -> None: + events = self.events( + self.detector(), + samples( + (20, 0), (14, 10), (13, 20), (12, 30), (12, 40), (12, 50), + (12, 60), (12, 70), + ), + ) + + self.assertEqual(len(events), 1) + self.assertEqual(events[0].confirmed_at, BASE + timedelta(seconds=70)) + + def test_prolonged_low_phase_does_not_duplicate_event(self) -> None: + events = self.events( + self.detector(), + samples( + (20, 0), (5, 10), (4, 20), (4, 30), (4, 40), (4, 50), + (4, 60), (4, 70), (3, 80), (2, 90), (-1, 100), + ), + ) + + self.assertEqual([event.value for event in events], [20]) + + def test_new_cycle_after_completed_reset(self) -> None: + events = self.events( + self.detector(), + samples( + (20, 0), (5, 10), (4, 20), (4, 30), (4, 40), (4, 50), + (4, 60), (4, 70), (6, 80), (30, 90), (10, 100), (9, 110), + (9, 120), (9, 130), (9, 140), (9, 150), (9, 160), + ), + ) + + self.assertEqual([event.value for event in events], [20, 30]) + self.assertEqual(events[1].cycle_started_at, BASE + timedelta(seconds=70)) + + def test_negative_reset_values_are_valid(self) -> None: + events = self.events( + self.detector(), + samples( + (50, 0), (-0.7, 10), (-0.5, 20), (-0.5, 30), (-0.5, 40), + (-0.5, 50), (-0.5, 60), (-0.5, 70), + ), + ) + + self.assertEqual([event.value for event in events], [50]) + + def test_regular_intervals_use_elapsed_time_not_sample_count(self) -> None: + events = self.events( + self.detector(), + samples( + (20, 0), (10, 10), (9, 20), (8, 30), (7, 40), (6, 50), + (5, 60), (4, 70), + ), + ) + + self.assertEqual(len(events), 1) + self.assertEqual(events[0].confirmed_at, BASE + timedelta(seconds=70)) + + def test_mildly_irregular_intervals_below_gap_limit_confirm_reset(self) -> None: + events = self.events( + self.detector(), + samples( + (20, 0), (10, 10), (9, 21), (8, 30), (7, 40), (6, 51), + (5, 60), (4, 70), + ), + ) + + self.assertEqual(len(events), 1) + self.assertEqual(events[0].confirmed_at, BASE + timedelta(seconds=70)) + + def test_backwards_timestamps_are_rejected(self) -> None: + detector = self.detector() + detector.process(BASE + timedelta(seconds=10), 20) + + with self.assertRaisesRegex(ValueError, "must not move backwards"): + detector.process(BASE, 10) + + def test_duplicate_timestamps_are_accepted_but_do_not_advance_hold_time(self) -> None: + events = self.events( + self.detector(), + samples( + (20, 0), (10, 10), (9, 10), (8, 20), (8, 30), (8, 40), + (8, 50), (8, 60), (8, 70), + ), + ) + + self.assertEqual(len(events), 1) + self.assertEqual(events[0].confirmed_at, BASE + timedelta(seconds=70)) + + def test_large_gap_to_first_low_sample_does_not_manufacture_confirmation(self) -> None: + events = self.events(self.detector(), samples((20, 0), (5, 3600))) + + self.assertEqual(events, []) + + def test_large_gap_during_reset_does_not_prove_continuous_below_threshold(self) -> None: + events = self.events(self.detector(), samples((20, 0), (5, 10), (4, 3600), (3, 3650))) + + self.assertEqual(events, []) + + def test_missing_observations_during_reset_restart_the_candidate(self) -> None: + detector = self.detector() + events = self.events( + detector, + samples( + (20, 0), (5, 10), (4, 20), (3, 70), (3, 80), (3, 90), + (3, 100), (3, 110), (3, 120), (2, 130), + ), + ) + + self.assertEqual(len(events), 1) + self.assertEqual(events[0].confirmed_at, BASE + timedelta(seconds=130)) + + def test_config_validation(self) -> None: + for ratio in (0.0, 1.0, -0.1, 1.1): + with self.assertRaisesRegex(ValueError, "drop_ratio"): + PeakDetectorConfig( + min_peak=1, + drop_ratio=ratio, + hold_seconds=0, + max_sample_gap_seconds=1, + ) + with self.assertRaisesRegex(ValueError, "hold_seconds"): + PeakDetectorConfig( + min_peak=1, + drop_ratio=0.5, + hold_seconds=-1, + max_sample_gap_seconds=1, + ) + with self.assertRaisesRegex(ValueError, "max_sample_gap_seconds"): + PeakDetectorConfig( + min_peak=1, + drop_ratio=0.5, + hold_seconds=0, + max_sample_gap_seconds=0, + ) + with self.assertRaisesRegex(ValueError, "min_peak"): + PeakDetectorConfig( + min_peak=float("nan"), + drop_ratio=0.5, + hold_seconds=0, + max_sample_gap_seconds=1, + )