Compare commits

...
3 Commits
Author SHA1 Message Date
admin 363be0a47e Add ENLYZE live data gateway 2026-09-05 05:19:44 +02:00
admin 0f7862fb0a Add persistent material integration state 2026-09-05 05:07:18 +02:00
admin b308b4f545 Add incremental material consumption integrator 2026-09-05 04:27:27 +02:00
10 changed files with 736 additions and 0 deletions
+2
View File
@@ -27,3 +27,5 @@ secrets/*
.idea/
.vscode/
.DS_Store
activate.sh
+34
View File
@@ -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,104 @@
"""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,
initial_state: MaterialIntegrationState | None = None,
) -> None:
self._config = config
self._state = initial_state or 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,69 @@
"""Serialization helpers for persistent material-integration state."""
from datetime import UTC, datetime
from production_analytics.calculations.material_consumption import (
MaterialIntegrationState,
)
def material_state_to_dict(state: MaterialIntegrationState) -> dict[str, object]:
"""Convert integration state to a JSON-serializable dictionary."""
return {
"cumulative_consumption_kg": state.cumulative_consumption_kg,
"integrated_running_seconds": state.integrated_running_seconds,
"last_processed_timestamp": (
state.last_processed_timestamp.astimezone(UTC).isoformat()
if state.last_processed_timestamp is not None
else None
),
"last_material_rate_kg_per_hour": state.last_material_rate_kg_per_hour,
"last_gate_value": state.last_gate_value,
"integration_active": state.integration_active,
}
def material_state_from_dict(data: dict[str, object]) -> MaterialIntegrationState:
"""Restore integration state from a JSON-compatible dictionary."""
raw_timestamp = data["last_processed_timestamp"]
timestamp = (
datetime.fromisoformat(str(raw_timestamp)).astimezone(UTC)
if raw_timestamp is not None
else None
)
return MaterialIntegrationState(
cumulative_consumption_kg=float(data["cumulative_consumption_kg"]),
integrated_running_seconds=float(data["integrated_running_seconds"]),
last_processed_timestamp=timestamp,
last_material_rate_kg_per_hour=(
float(data["last_material_rate_kg_per_hour"])
if data["last_material_rate_kg_per_hour"] is not None
else None
),
last_gate_value=(
float(data["last_gate_value"])
if data["last_gate_value"] is not None
else None
),
integration_active=bool(data["integration_active"]),
)
def save_material_state(path: str, state: MaterialIntegrationState) -> None:
"""Persist integration state as UTF-8 JSON."""
import json
with open(path, "w", encoding="utf-8") as f:
json.dump(material_state_to_dict(state), f, indent=2)
def load_material_state(path: str) -> MaterialIntegrationState:
"""Load integration state from UTF-8 JSON."""
import json
with open(path, encoding="utf-8") as f:
data = json.load(f)
return material_state_from_dict(data)
@@ -0,0 +1,89 @@
"""Verified ENLYZE adapter for production-run and timeseries access."""
from dataclasses import dataclass
from datetime import UTC, datetime
from production_analytics.calculations.material_consumption import MaterialSample
from production_analytics.enlyze.exploration import ExplorationClient
@dataclass(frozen=True, slots=True)
class EnlyzeProductionRun:
uuid: str
machine_id: str
product_id: str | None
production_order: str
start: datetime
end: datetime | None
class EnlyzeApiGateway:
def __init__(self, client: ExplorationClient) -> None:
self._client = client
def get_open_production_run(self, machine_id: str) -> EnlyzeProductionRun | None:
response = self._client.get(
"/v2/production-runs",
{"machine": machine_id},
)
for item in response.body.get("data", []):
if item.get("end") is not None:
continue
return EnlyzeProductionRun(
uuid=str(item["uuid"]),
machine_id=str(item["machine"]),
product_id=(
str(item["product"])
if item.get("product") is not None
else None
),
production_order=str(item["production_order"]),
start=_parse_timestamp(item["start"]),
end=None,
)
return None
def get_material_samples(
self,
*,
machine_id: str,
rate_variable_id: str,
gate_variable_id: str,
start: datetime,
end: datetime,
) -> list[MaterialSample]:
response = self._client.post_json(
"/v2/timeseries",
{
"machine": machine_id,
"start": start.astimezone(UTC).isoformat(),
"end": end.astimezone(UTC).isoformat(),
"variables": [
{"uuid": rate_variable_id},
{"uuid": gate_variable_id},
],
},
)
data = response.body["data"]
columns = data["columns"]
time_index = columns.index("time")
rate_index = columns.index(rate_variable_id)
gate_index = columns.index(gate_variable_id)
return [
MaterialSample(
timestamp=_parse_timestamp(record[time_index]),
material_rate_kg_per_hour=float(record[rate_index]),
gate_value=float(record[gate_index]),
)
for record in data["records"]
]
def _parse_timestamp(value: str) -> datetime:
return datetime.fromisoformat(value.replace("Z", "+00:00")).astimezone(UTC)
+133
View File
@@ -0,0 +1,133 @@
from datetime import UTC, datetime
from unittest.mock import Mock
from production_analytics.enlyze.gateway import EnlyzeApiGateway
def test_get_open_production_run_returns_current_run() -> None:
client = Mock()
client.get.return_value.body = {
"data": [
{
"uuid": "closed-run",
"machine": "machine-1",
"product": "product-a",
"production_order": "ORDER-1",
"start": "2026-09-03T04:00:00Z",
"end": "2026-09-03T05:00:00Z",
},
{
"uuid": "open-run",
"machine": "machine-1",
"product": "product-b",
"production_order": "ORDER-2",
"start": "2026-09-03T05:07:47Z",
"end": None,
},
]
}
gateway = EnlyzeApiGateway(client)
run = gateway.get_open_production_run("machine-1")
assert run is not None
assert run.uuid == "open-run"
assert run.machine_id == "machine-1"
assert run.product_id == "product-b"
assert run.production_order == "ORDER-2"
assert run.start == datetime(2026, 9, 3, 5, 7, 47, tzinfo=UTC)
assert run.end is None
client.get.assert_called_once_with(
"/v2/production-runs",
{"machine": "machine-1"},
)
def test_get_open_production_run_returns_none_without_open_run() -> None:
client = Mock()
client.get.return_value.body = {
"data": [
{
"uuid": "closed-run",
"machine": "machine-1",
"product": "product-a",
"production_order": "ORDER-1",
"start": "2026-09-03T04:00:00Z",
"end": "2026-09-03T05:00:00Z",
}
]
}
gateway = EnlyzeApiGateway(client)
assert gateway.get_open_production_run("machine-1") is None
def test_get_material_samples_requests_rate_and_gate_together() -> None:
client = Mock()
client.post_json.return_value.body = {
"data": {
"columns": ["time", "rate-variable", "gate-variable"],
"records": [
["2026-09-03T05:09:50Z", 1028.5, 2.1],
["2026-09-03T05:10:00Z", 1029.0, 2.2],
],
}
}
gateway = EnlyzeApiGateway(client)
samples = gateway.get_material_samples(
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
start=datetime(2026, 9, 3, 5, 9, 50, tzinfo=UTC),
end=datetime(2026, 9, 3, 5, 10, 10, tzinfo=UTC),
)
assert len(samples) == 2
assert samples[0].timestamp == datetime(2026, 9, 3, 5, 9, 50, tzinfo=UTC)
assert samples[0].material_rate_kg_per_hour == 1028.5
assert samples[0].gate_value == 2.1
assert samples[1].material_rate_kg_per_hour == 1029.0
assert samples[1].gate_value == 2.2
client.post_json.assert_called_once_with(
"/v2/timeseries",
{
"machine": "machine-1",
"start": "2026-09-03T05:09:50+00:00",
"end": "2026-09-03T05:10:10+00:00",
"variables": [
{"uuid": "rate-variable"},
{"uuid": "gate-variable"},
],
},
)
def test_get_material_samples_uses_column_names_not_fixed_positions() -> None:
client = Mock()
client.post_json.return_value.body = {
"data": {
"columns": ["time", "gate-variable", "rate-variable"],
"records": [
["2026-09-03T05:09:50Z", 2.1, 1028.5],
],
}
}
gateway = EnlyzeApiGateway(client)
samples = gateway.get_material_samples(
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
start=datetime(2026, 9, 3, 5, 9, 50, tzinfo=UTC),
end=datetime(2026, 9, 3, 5, 10, 0, tzinfo=UTC),
)
assert samples[0].material_rate_kg_per_hour == 1028.5
assert samples[0].gate_value == 2.1
+133
View File
@@ -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)
+97
View File
@@ -0,0 +1,97 @@
import json
from datetime import UTC, datetime
from production_analytics.calculations.material_consumption import (
MaterialConsumptionIntegrator,
MaterialIntegrationState,
MaterialIntegratorConfig,
MaterialSample,
)
from production_analytics.calculations.material_state import (
load_material_state,
material_state_from_dict,
material_state_to_dict,
save_material_state,
)
def test_material_state_dict_roundtrip() -> None:
original = MaterialIntegrationState(
cumulative_consumption_kg=51.42714233398436,
integrated_running_seconds=180.0,
last_processed_timestamp=datetime(2026, 9, 3, 5, 10, 50, tzinfo=UTC),
last_material_rate_kg_per_hour=1028.5428466796875,
last_gate_value=2.1673481464385986,
integration_active=True,
)
restored = material_state_from_dict(material_state_to_dict(original))
assert restored == original
def test_material_state_json_file_roundtrip(tmp_path) -> None:
original = MaterialIntegrationState(
cumulative_consumption_kg=12.5,
integrated_running_seconds=90.0,
last_processed_timestamp=datetime(2026, 9, 3, 5, 10, 50, tzinfo=UTC),
last_material_rate_kg_per_hour=500.0,
last_gate_value=2.0,
integration_active=True,
)
path = tmp_path / "material-state.json"
save_material_state(str(path), original)
restored = load_material_state(str(path))
assert restored == original
raw = json.loads(path.read_text(encoding="utf-8"))
assert raw["last_processed_timestamp"] == "2026-09-03T05:10:50+00:00"
def test_integrator_restore_continues_without_double_counting(tmp_path) -> None:
config = MaterialIntegratorConfig(
gate_threshold=0.5,
max_sample_gap_seconds=20.0,
)
samples = [
MaterialSample(
timestamp=datetime(2026, 9, 3, 5, 9, 40, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
MaterialSample(
timestamp=datetime(2026, 9, 3, 5, 9, 50, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
MaterialSample(
timestamp=datetime(2026, 9, 3, 5, 10, 0, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
]
continuous = MaterialConsumptionIntegrator(config)
continuous.process_many(samples)
before_restart = MaterialConsumptionIntegrator(config)
before_restart.process_many(samples[:2])
path = tmp_path / "material-state.json"
save_material_state(str(path), before_restart.state)
restored = MaterialConsumptionIntegrator(
config,
initial_state=load_material_state(str(path)),
)
# Polling deliberately starts inclusively at the last processed timestamp.
restored.process_many(samples[1:])
assert restored.state == continuous.state
assert restored.state.cumulative_consumption_kg == 20.0
assert restored.state.integrated_running_seconds == 20.0