Add generic peak cycle detector

This commit is contained in:
2026-09-04 08:00:59 +02:00
parent f017d187eb
commit 5068c95158
9 changed files with 521 additions and 8 deletions
+24
View File
@@ -17,6 +17,30 @@ ENLYZE and is retrieved again when historical calculations need reproduction.
- Calculations are modular, testable Python implementations. Configuration - Calculations are modular, testable Python implementations. Configuration
declares instances; it is not a generic low-code language. 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 ## Central domain context
Production orders connect machine, article/material, source interval, and Production orders connect machine, article/material, source interval, and
+22
View File
@@ -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 not pass ENLYZE credentials to TimescaleDB. A future application service should
use `env_file: ./secrets/enlyze.env` rather than copying secrets into Compose. 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 `docker compose up -d timescaledb` is an optional local database design for a
future persistence milestone. It is not required for the bootstrap tests. future persistence milestone. It is not required for the bootstrap tests.
+4 -2
View File
@@ -16,6 +16,8 @@ calculations:
version: "1" version: "1"
machine_ref: machine-to-be-verified machine_ref: machine-to-be-verified
signal_ref: process-signal-to-be-verified signal_ref: process-signal-to-be-verified
below_peak_fraction: 0.75 min_peak: 10.0
below_peak_duration_seconds: 60 drop_ratio: 0.75
hold_seconds: 60.0
max_sample_gap_seconds: 20.0
output_event: cycle_peak output_event: cycle_peak
+16
View File
@@ -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 context partition. It enables safe continuation (for example an open peak
cycle), while the source interval recorded on results keeps a historical run cycle), while the source interval recorded on results keeps a historical run
reproducible by fetching ENLYZE data again. 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.
+5 -5
View File
@@ -26,12 +26,12 @@ Implement one configured integration calculation over a verified signal and
production-order context. Include backfill, rerun/provenance behavior, and production-order context. Include backfill, rerun/provenance behavior, and
tests against sanitized fixtures. tests against sanitized fixtures.
## 4. Peak detection ## 4. Peak detection (first detector completed)
Implement the configurable sawtooth peak detector: retain the current maximum; The generic configurable sawtooth detector is implemented: it retains the
close a cycle after a configurable below-peak fraction remains true for a current maximum, closes a cycle after a configurable below-peak fraction has
configurable duration; emit one completed peak event; preserve open-cycle actual sample support for a configurable elapsed duration, and emits one
state. completed peak. Its pure state can later be persisted for incremental use.
## Later ## Later
+61
View File
@@ -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())
@@ -1,5 +1,12 @@
"""Versioned, pure calculation operators and their contracts.""" """Versioned, pure calculation operators and their contracts."""
from .base import Calculation from .base import Calculation
from .peak_cycles import DetectedPeak, PeakCycleDetector, PeakDetectorConfig, TimeSeriesSample
__all__ = ["Calculation"] __all__ = [
"Calculation",
"DetectedPeak",
"PeakCycleDetector",
"PeakDetectorConfig",
"TimeSeriesSample",
]
@@ -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
+224
View File
@@ -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,
)