Add incremental material consumption integrator
This commit is contained in:
@@ -62,6 +62,40 @@ 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.
|
||||
|
||||
## Material-consumption integration
|
||||
|
||||
`MaterialConsumptionIntegrator` in `calculations` accepts `MaterialSample`
|
||||
values (timezone-aware timestamp, material rate in kg/h, numeric gate value).
|
||||
Configure a strict `gate_value > gate_threshold` condition and an explicit,
|
||||
positive `max_sample_gap_seconds`. Each interval uses the preceding sample's
|
||||
rate and gate, converting elapsed seconds to hours to accumulate kg. A longer
|
||||
gap contributes neither consumption nor running time and resets the baseline
|
||||
to the newer sample. No time before the first or after the last sample is inferred.
|
||||
|
||||
Call `process(sample)` or `process_many(samples)` on the same instance for live
|
||||
input, replay, or successive chunks. Both use the same calculation. The immutable
|
||||
`state` snapshot exposes cumulative kg, integrated running seconds, last timestamp
|
||||
(UTC), last rate, last gate, and whether the latest gate is active. State restoration
|
||||
and persistence are not implemented yet. Naive timestamps and non-finite values
|
||||
are rejected; backwards timestamps raise without changing state. Duplicate
|
||||
timestamps add no consumption but replace the baseline in arrival order.
|
||||
Finite negative material rates are currently accepted and decrease cumulative
|
||||
consumption during affected integrated intervals. This is intentional generic
|
||||
behavior for now; machine-specific validation or clamping may be added later
|
||||
at the input/adapter layer if required by process semantics.
|
||||
|
||||
Run `python scripts/validate_k7_material_consumption.py` from the repository root
|
||||
with the ignored `k7-00842-throughput.raw.json` and `k7-00842-speed.raw.json`
|
||||
captures in `data/raw/enlyze/`. The utility joins common non-null timestamps
|
||||
without filling values and rejects unordered or duplicate capture records.
|
||||
K7-specific signal UUIDs are confined to the utility: `Stundenleistung Anlage`
|
||||
is gated by `Geschwindigkeit Gesamtanlage > 0.5 m/min`. With the explicit
|
||||
20-second validation gap limit (`--max-sample-gap-seconds` to override),
|
||||
2,672 common samples yield **5.205555556 h** and **5,180.127150811 kg** for run
|
||||
00842, matching the previous manual calculation's rounded results. Synthetic
|
||||
tests verify equivalence to manual interval integration without local captures
|
||||
or live ENLYZE access. Production polling, Grafana, and Bento 1 are not implemented.
|
||||
|
||||
## Peak-cycle detection
|
||||
|
||||
`PeakCycleDetector` is a pure calculation-domain component for roll length,
|
||||
|
||||
@@ -0,0 +1,65 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Replay common timestamps from ignored K7 run 00842 captures, without network access."""
|
||||
|
||||
import argparse
|
||||
import json
|
||||
from datetime import UTC, datetime
|
||||
from pathlib import Path
|
||||
|
||||
from production_analytics.calculations import (
|
||||
MaterialConsumptionIntegrator,
|
||||
MaterialIntegratorConfig,
|
||||
MaterialSample,
|
||||
)
|
||||
|
||||
RATE_UUID = "c9d06af5-f6d6-4ede-b6c4-5a98bac77129" # Stundenleistung Anlage [kg/h]
|
||||
GATE_UUID = "823867bb-f5d2-40eb-b875-657155addfd0" # Geschwindigkeit Gesamtanlage [m/min]
|
||||
CAPTURE_DIRECTORY = Path("data/raw/enlyze")
|
||||
|
||||
|
||||
def read_capture(path: Path, variable_uuid: str) -> dict[datetime, float]:
|
||||
"""Read unique, ordered observations; refuse ambiguous duplicate source records."""
|
||||
data = json.loads(path.read_text(encoding="utf-8"))["response"]["body"]["data"]
|
||||
time_index = data["columns"].index("time")
|
||||
value_index = data["columns"].index(variable_uuid)
|
||||
values = {}
|
||||
previous = None
|
||||
for record in data["records"]:
|
||||
timestamp = datetime.fromisoformat(record[time_index].replace("Z", "+00:00"))
|
||||
if timestamp.utcoffset() is None:
|
||||
raise ValueError(f"Naive timestamp in {path}")
|
||||
timestamp = timestamp.astimezone(UTC)
|
||||
if previous is not None and timestamp <= previous:
|
||||
raise ValueError(f"Capture timestamps must be strictly increasing: {path}")
|
||||
previous = timestamp
|
||||
if record[value_index] is not None:
|
||||
values[timestamp] = float(record[value_index])
|
||||
return values
|
||||
|
||||
|
||||
def main() -> int:
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument("--max-sample-gap-seconds", type=float, default=20.0)
|
||||
args = parser.parse_args()
|
||||
integrator = MaterialConsumptionIntegrator(
|
||||
MaterialIntegratorConfig(0.5, args.max_sample_gap_seconds)
|
||||
)
|
||||
try:
|
||||
rates = read_capture(CAPTURE_DIRECTORY / "k7-00842-throughput.raw.json", RATE_UUID)
|
||||
gates = read_capture(CAPTURE_DIRECTORY / "k7-00842-speed.raw.json", GATE_UUID)
|
||||
except (OSError, ValueError, KeyError) as error:
|
||||
parser.exit(1, f"Cannot replay K7 captures: {error}\n")
|
||||
# Dict insertion order preserves validated capture order. No resampling or filling.
|
||||
samples = [MaterialSample(t, rate, gates[t]) for t, rate in rates.items() if t in gates]
|
||||
if len(samples) < 2:
|
||||
parser.exit(1, "Need at least two common non-null samples\n")
|
||||
state = integrator.process_many(samples)
|
||||
print(f"Common samples: {len(samples)}")
|
||||
print(f"Maximum sample gap: {integrator.config.max_sample_gap_seconds:g} s")
|
||||
print(f"Integrated running time: {state.integrated_running_seconds / 3600:.9f} h")
|
||||
print(f"Cumulative material consumption: {state.cumulative_consumption_kg:.9f} kg")
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -1,10 +1,20 @@
|
||||
"""Versioned, pure calculation operators and their contracts."""
|
||||
|
||||
from .base import Calculation
|
||||
from .material_consumption import (
|
||||
MaterialConsumptionIntegrator,
|
||||
MaterialIntegrationState,
|
||||
MaterialIntegratorConfig,
|
||||
MaterialSample,
|
||||
)
|
||||
from .peak_cycles import DetectedPeak, PeakCycleDetector, PeakDetectorConfig, TimeSeriesSample
|
||||
|
||||
__all__ = [
|
||||
"Calculation",
|
||||
"MaterialConsumptionIntegrator",
|
||||
"MaterialIntegrationState",
|
||||
"MaterialIntegratorConfig",
|
||||
"MaterialSample",
|
||||
"DetectedPeak",
|
||||
"PeakCycleDetector",
|
||||
"PeakDetectorConfig",
|
||||
|
||||
@@ -0,0 +1,100 @@
|
||||
"""Incremental, previous-value integration shared by live input and replay."""
|
||||
|
||||
from collections.abc import Iterable
|
||||
from dataclasses import dataclass
|
||||
from datetime import UTC, datetime
|
||||
from math import isfinite
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class MaterialSample:
|
||||
"""An aware timestamp, material rate in kg/h, and numeric gate value."""
|
||||
|
||||
timestamp: datetime
|
||||
material_rate_kg_per_hour: float
|
||||
gate_value: float
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class MaterialIntegratorConfig:
|
||||
gate_threshold: float
|
||||
max_sample_gap_seconds: float
|
||||
|
||||
def __post_init__(self) -> None:
|
||||
if not isfinite(self.gate_threshold):
|
||||
raise ValueError("gate_threshold must be finite")
|
||||
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 MaterialIntegrationState:
|
||||
"""Immutable snapshot; running time excludes inactive and rejected gap intervals."""
|
||||
|
||||
cumulative_consumption_kg: float = 0.0
|
||||
integrated_running_seconds: float = 0.0
|
||||
last_processed_timestamp: datetime | None = None
|
||||
last_material_rate_kg_per_hour: float | None = None
|
||||
last_gate_value: float | None = None
|
||||
integration_active: bool = False
|
||||
|
||||
|
||||
class MaterialConsumptionIntegrator:
|
||||
"""Integrate each observed interval using its preceding rate and gate.
|
||||
|
||||
The first sample only establishes a baseline. Equal timestamps replace that
|
||||
baseline without adding time. Gaps exceeding the configured maximum also
|
||||
establish a new baseline without integration. No tail interval is inferred.
|
||||
Finite negative material rates are currently accepted without clamping and
|
||||
decrease cumulative consumption during affected integrated intervals. This
|
||||
is intentional generic behavior for now; machine-specific validation or
|
||||
clamping may be added later at the input/adapter layer if process semantics
|
||||
require it. Invalid samples raise before modifying state. State timestamps
|
||||
are always UTC.
|
||||
"""
|
||||
|
||||
def __init__(self, config: MaterialIntegratorConfig) -> None:
|
||||
self._config = config
|
||||
self._state = MaterialIntegrationState()
|
||||
|
||||
@property
|
||||
def config(self) -> MaterialIntegratorConfig:
|
||||
return self._config
|
||||
|
||||
@property
|
||||
def state(self) -> MaterialIntegrationState:
|
||||
return self._state
|
||||
|
||||
def process(self, sample: MaterialSample) -> MaterialIntegrationState:
|
||||
"""Consume one sample in arrival order and return the resulting snapshot."""
|
||||
if sample.timestamp.tzinfo is None or sample.timestamp.utcoffset() is None:
|
||||
raise ValueError("sample timestamp must be timezone-aware")
|
||||
timestamp = sample.timestamp.astimezone(UTC)
|
||||
if not isfinite(sample.material_rate_kg_per_hour) or not isfinite(sample.gate_value):
|
||||
raise ValueError("sample rate and gate value must be finite")
|
||||
previous = self.state
|
||||
consumption = previous.cumulative_consumption_kg
|
||||
running_seconds = previous.integrated_running_seconds
|
||||
if previous.last_processed_timestamp is not None:
|
||||
elapsed = (timestamp - previous.last_processed_timestamp).total_seconds()
|
||||
if elapsed < 0:
|
||||
raise ValueError("sample timestamps must not move backwards")
|
||||
if 0 < elapsed <= self.config.max_sample_gap_seconds and previous.integration_active:
|
||||
assert previous.last_material_rate_kg_per_hour is not None
|
||||
consumption += previous.last_material_rate_kg_per_hour * (elapsed / 3600.0)
|
||||
running_seconds += elapsed
|
||||
self._state = MaterialIntegrationState(
|
||||
cumulative_consumption_kg=consumption,
|
||||
integrated_running_seconds=running_seconds,
|
||||
last_processed_timestamp=timestamp,
|
||||
last_material_rate_kg_per_hour=sample.material_rate_kg_per_hour,
|
||||
last_gate_value=sample.gate_value,
|
||||
integration_active=sample.gate_value > self.config.gate_threshold,
|
||||
)
|
||||
return self.state
|
||||
|
||||
def process_many(self, samples: Iterable[MaterialSample]) -> MaterialIntegrationState:
|
||||
"""Consume an iterable without sorting; previous successful calls remain applied."""
|
||||
for sample in samples:
|
||||
self.process(sample)
|
||||
return self.state
|
||||
@@ -0,0 +1,133 @@
|
||||
from datetime import UTC, datetime, timedelta, timezone
|
||||
|
||||
import pytest
|
||||
|
||||
from production_analytics.calculations import (
|
||||
MaterialConsumptionIntegrator,
|
||||
MaterialIntegratorConfig,
|
||||
MaterialSample,
|
||||
)
|
||||
|
||||
START = datetime(2026, 8, 31, tzinfo=UTC)
|
||||
|
||||
|
||||
def sample(seconds, rate=3600.0, gate=1.0):
|
||||
return MaterialSample(START + timedelta(seconds=seconds), rate, gate)
|
||||
|
||||
|
||||
def integrator(gap=20.0):
|
||||
return MaterialConsumptionIntegrator(MaterialIntegratorConfig(0.5, gap))
|
||||
|
||||
|
||||
def test_irregular_intervals_use_previous_rate_and_no_tail():
|
||||
core = integrator()
|
||||
first = core.process(sample(0))
|
||||
assert first.cumulative_consumption_kg == 0
|
||||
state = core.process_many([sample(7, 7200), sample(20, 9999)])
|
||||
assert state.cumulative_consumption_kg == pytest.approx(33)
|
||||
assert state.integrated_running_seconds == 20
|
||||
assert first.cumulative_consumption_kg == 0 # snapshots remain stable
|
||||
assert state.last_processed_timestamp == START + timedelta(seconds=20)
|
||||
assert state.last_material_rate_kg_per_hour == 9999
|
||||
assert state.last_gate_value == 1
|
||||
assert state.integration_active
|
||||
|
||||
|
||||
def test_inactive_gate_and_strict_threshold():
|
||||
state = integrator().process_many([sample(0, gate=0.5), sample(10, gate=-1), sample(20)])
|
||||
assert state.cumulative_consumption_kg == 0
|
||||
assert state.integrated_running_seconds == 0
|
||||
|
||||
|
||||
def test_gate_transitions_use_previous_gate():
|
||||
state = integrator().process_many(
|
||||
[sample(0, gate=0), sample(10), sample(20, gate=0), sample(30), sample(40)]
|
||||
)
|
||||
assert state.cumulative_consumption_kg == pytest.approx(20)
|
||||
assert state.integrated_running_seconds == 20
|
||||
|
||||
|
||||
def test_duplicate_updates_baseline_without_consumption():
|
||||
core = integrator()
|
||||
core.process_many([sample(0), sample(10)])
|
||||
duplicate = core.process(sample(10, 7200))
|
||||
assert duplicate.cumulative_consumption_kg == 10
|
||||
assert core.process(sample(20)).cumulative_consumption_kg == 30
|
||||
core.process(sample(20, gate=0))
|
||||
assert core.process(sample(30)).cumulative_consumption_kg == 30
|
||||
|
||||
|
||||
def test_backwards_timestamp_rejected_without_state_change():
|
||||
core = integrator()
|
||||
before = core.process(sample(10))
|
||||
with pytest.raises(ValueError, match="backwards"):
|
||||
core.process(sample(0))
|
||||
assert core.state == before
|
||||
|
||||
|
||||
def test_gap_skipped_and_resumed_with_new_rate():
|
||||
core = integrator()
|
||||
core.process_many([sample(0), sample(20)]) # exactly max gap is accepted
|
||||
state = core.process(sample(41, 7200))
|
||||
assert state.cumulative_consumption_kg == 20
|
||||
assert state.integration_active
|
||||
state = core.process(sample(51))
|
||||
assert state.cumulative_consumption_kg == 40
|
||||
assert state.integrated_running_seconds == 30
|
||||
|
||||
|
||||
def test_incremental_chunks_and_replay_match_manual_intervals():
|
||||
samples = [
|
||||
sample(0),
|
||||
sample(3, 1800),
|
||||
sample(11, gate=0),
|
||||
sample(16),
|
||||
sample(50, 900),
|
||||
sample(60),
|
||||
]
|
||||
manual_seconds = sum(
|
||||
(b.timestamp - a.timestamp).total_seconds()
|
||||
for a, b in zip(samples, samples[1:])
|
||||
if a.gate_value > 0.5 and (b.timestamp - a.timestamp).total_seconds() <= 20
|
||||
)
|
||||
manual_kg = sum(
|
||||
a.material_rate_kg_per_hour * (b.timestamp - a.timestamp).total_seconds() / 3600
|
||||
for a, b in zip(samples, samples[1:])
|
||||
if a.gate_value > 0.5 and (b.timestamp - a.timestamp).total_seconds() <= 20
|
||||
)
|
||||
core = integrator()
|
||||
core.process_many(samples[:2])
|
||||
core.process(samples[2])
|
||||
core.process_many(samples[3:])
|
||||
assert core.process_many([]) == integrator().process_many(samples)
|
||||
assert core.state.cumulative_consumption_kg == pytest.approx(manual_kg)
|
||||
assert core.state.integrated_running_seconds == manual_seconds
|
||||
|
||||
|
||||
def test_timezone_normalization_and_naive_rejection():
|
||||
core = integrator()
|
||||
core.process(sample(0))
|
||||
local = (START + timedelta(seconds=10)).astimezone(timezone(timedelta(hours=2)))
|
||||
state = core.process(MaterialSample(local, 3600, 1))
|
||||
assert state.cumulative_consumption_kg == 10
|
||||
assert state.last_processed_timestamp.tzinfo is UTC
|
||||
with pytest.raises(ValueError, match="timezone-aware"):
|
||||
core.process(MaterialSample(START.replace(tzinfo=None), 3600, 1))
|
||||
assert core.state == state
|
||||
|
||||
|
||||
@pytest.mark.parametrize("rate,gate", [(float("nan"), 1), (1, float("inf"))])
|
||||
def test_nonfinite_samples_rejected(rate, gate):
|
||||
core = integrator()
|
||||
before = core.state
|
||||
with pytest.raises(ValueError, match="finite"):
|
||||
core.process(sample(0, rate, gate))
|
||||
assert core.state == before
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"threshold,gap", [(float("nan"), 20), (0.5, 0), (0.5, -1), (0.5, float("inf"))]
|
||||
)
|
||||
def test_invalid_configuration(threshold, gap):
|
||||
with pytest.raises(ValueError):
|
||||
MaterialIntegratorConfig(threshold, gap)
|
||||
Reference in New Issue
Block a user