Compare commits
3
Commits
231e18082c
...
363be0a47e
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
363be0a47e | ||
|
|
0f7862fb0a | ||
|
|
b308b4f545 |
@@ -27,3 +27,5 @@ secrets/*
|
|||||||
.idea/
|
.idea/
|
||||||
.vscode/
|
.vscode/
|
||||||
.DS_Store
|
.DS_Store
|
||||||
|
|
||||||
|
activate.sh
|
||||||
|
|||||||
@@ -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
|
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.
|
||||||
|
|
||||||
|
## 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
|
## Peak-cycle detection
|
||||||
|
|
||||||
`PeakCycleDetector` is a pure calculation-domain component for roll length,
|
`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."""
|
"""Versioned, pure calculation operators and their contracts."""
|
||||||
|
|
||||||
from .base import Calculation
|
from .base import Calculation
|
||||||
|
from .material_consumption import (
|
||||||
|
MaterialConsumptionIntegrator,
|
||||||
|
MaterialIntegrationState,
|
||||||
|
MaterialIntegratorConfig,
|
||||||
|
MaterialSample,
|
||||||
|
)
|
||||||
from .peak_cycles import DetectedPeak, PeakCycleDetector, PeakDetectorConfig, TimeSeriesSample
|
from .peak_cycles import DetectedPeak, PeakCycleDetector, PeakDetectorConfig, TimeSeriesSample
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
"Calculation",
|
"Calculation",
|
||||||
|
"MaterialConsumptionIntegrator",
|
||||||
|
"MaterialIntegrationState",
|
||||||
|
"MaterialIntegratorConfig",
|
||||||
|
"MaterialSample",
|
||||||
"DetectedPeak",
|
"DetectedPeak",
|
||||||
"PeakCycleDetector",
|
"PeakCycleDetector",
|
||||||
"PeakDetectorConfig",
|
"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)
|
||||||
@@ -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
|
||||||
@@ -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)
|
||||||
@@ -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
|
||||||
Reference in New Issue
Block a user