Harden live material polling state
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
"""Serialization helpers for persistent material-integration state."""
|
||||
|
||||
from datetime import UTC, datetime
|
||||
from math import isfinite
|
||||
|
||||
from production_analytics.calculations.material_consumption import (
|
||||
MaterialIntegrationState,
|
||||
@@ -9,6 +10,11 @@ from production_analytics.calculations.material_consumption import (
|
||||
|
||||
def material_state_to_dict(state: MaterialIntegrationState) -> dict[str, object]:
|
||||
"""Convert integration state to a JSON-serializable dictionary."""
|
||||
if state.last_processed_timestamp is not None and (
|
||||
state.last_processed_timestamp.tzinfo is None
|
||||
or state.last_processed_timestamp.utcoffset() is None
|
||||
):
|
||||
raise ValueError("last_processed_timestamp must be timezone-aware")
|
||||
return {
|
||||
"cumulative_consumption_kg": state.cumulative_consumption_kg,
|
||||
"integrated_running_seconds": state.integrated_running_seconds,
|
||||
@@ -25,7 +31,33 @@ def material_state_to_dict(state: MaterialIntegrationState) -> dict[str, object]
|
||||
|
||||
def material_state_from_dict(data: dict[str, object]) -> MaterialIntegrationState:
|
||||
"""Restore integration state from a JSON-compatible dictionary."""
|
||||
if not isinstance(data, dict):
|
||||
raise ValueError("integration_state must be an object")
|
||||
for key in (
|
||||
"cumulative_consumption_kg", "integrated_running_seconds",
|
||||
"last_material_rate_kg_per_hour", "last_gate_value",
|
||||
):
|
||||
value = data[key]
|
||||
if value is None and key.startswith("last_"):
|
||||
continue
|
||||
if isinstance(value, bool) or not isinstance(value, (int, float)) or not isfinite(value):
|
||||
raise ValueError(f"{key} must be a finite number")
|
||||
if data["integrated_running_seconds"] < 0:
|
||||
raise ValueError("integrated_running_seconds must not be negative")
|
||||
if not isinstance(data["integration_active"], bool):
|
||||
raise ValueError("integration_active must be a boolean")
|
||||
raw_timestamp = data["last_processed_timestamp"]
|
||||
if raw_timestamp is not None:
|
||||
if not isinstance(raw_timestamp, str):
|
||||
raise ValueError("last_processed_timestamp must be an ISO timestamp")
|
||||
parsed = datetime.fromisoformat(raw_timestamp)
|
||||
if parsed.tzinfo is None or parsed.utcoffset() is None:
|
||||
raise ValueError("last_processed_timestamp must be timezone-aware")
|
||||
if data["integration_active"] and (
|
||||
raw_timestamp is None or data["last_material_rate_kg_per_hour"] is None
|
||||
or data["last_gate_value"] is None
|
||||
):
|
||||
raise ValueError("active integration requires a complete temporal baseline")
|
||||
|
||||
timestamp = (
|
||||
datetime.fromisoformat(str(raw_timestamp)).astimezone(UTC)
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
from dataclasses import dataclass
|
||||
from datetime import UTC, datetime
|
||||
from math import isfinite
|
||||
|
||||
from production_analytics.calculations.material_consumption import MaterialSample
|
||||
from production_analytics.enlyze.exploration import ExplorationClient
|
||||
@@ -27,24 +28,31 @@ class EnlyzeApiGateway:
|
||||
{"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
|
||||
try:
|
||||
items = response.body["data"]
|
||||
if not isinstance(items, list):
|
||||
raise ValueError("data must be a list")
|
||||
runs = []
|
||||
for item in items:
|
||||
if not isinstance(item, dict) or "end" not in item:
|
||||
raise ValueError("run must be an object with an end field")
|
||||
if item["end"] is not None:
|
||||
continue
|
||||
for key in ("uuid", "machine", "production_order"):
|
||||
if not isinstance(item[key], str) or not item[key]:
|
||||
raise ValueError(f"{key} must be a non-empty string")
|
||||
if item["machine"] != machine_id:
|
||||
raise ValueError("run machine does not match requested machine")
|
||||
runs.append(EnlyzeProductionRun(
|
||||
uuid=item["uuid"], machine_id=item["machine"],
|
||||
product_id=item.get("product"), production_order=item["production_order"],
|
||||
start=_parse_timestamp(item["start"]), end=None,
|
||||
))
|
||||
if len(runs) > 1:
|
||||
raise ValueError("multiple open Production Runs returned")
|
||||
return runs[0] if runs else None
|
||||
except (KeyError, TypeError, ValueError) as exc:
|
||||
raise ValueError(f"Invalid production-run response: {exc}") from exc
|
||||
|
||||
def get_material_samples(
|
||||
self,
|
||||
@@ -55,6 +63,11 @@ class EnlyzeApiGateway:
|
||||
start: datetime,
|
||||
end: datetime,
|
||||
) -> list[MaterialSample]:
|
||||
for name, value in (("start", start), ("end", end)):
|
||||
if value.tzinfo is None or value.utcoffset() is None:
|
||||
raise ValueError(f"{name} must be timezone-aware")
|
||||
if end < start:
|
||||
raise ValueError("end must not precede start")
|
||||
response = self._client.post_json(
|
||||
"/v2/timeseries",
|
||||
{
|
||||
@@ -68,22 +81,45 @@ class EnlyzeApiGateway:
|
||||
},
|
||||
)
|
||||
|
||||
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"]
|
||||
]
|
||||
try:
|
||||
data = response.body["data"]
|
||||
columns = data["columns"]
|
||||
if not isinstance(columns, list):
|
||||
raise ValueError("columns must be a list")
|
||||
for required in ("time", rate_variable_id, gate_variable_id):
|
||||
if columns.count(required) != 1:
|
||||
raise ValueError(f"required column {required!r} must occur exactly once")
|
||||
time_index = columns.index("time")
|
||||
rate_index = columns.index(rate_variable_id)
|
||||
gate_index = columns.index(gate_variable_id)
|
||||
if not isinstance(data["records"], list):
|
||||
raise ValueError("records must be a list")
|
||||
samples = []
|
||||
for index, record in enumerate(data["records"]):
|
||||
try:
|
||||
if not isinstance(record, list) or len(record) != len(columns):
|
||||
raise ValueError("record must match columns")
|
||||
rate, gate = record[rate_index], record[gate_index]
|
||||
for value in (rate, gate):
|
||||
if isinstance(value, bool) or not isinstance(value, (int, float)):
|
||||
raise ValueError("rate and gate must be numeric")
|
||||
if not isfinite(value):
|
||||
raise ValueError("rate and gate must be finite")
|
||||
samples.append(MaterialSample(
|
||||
timestamp=_parse_timestamp(record[time_index]),
|
||||
material_rate_kg_per_hour=float(rate), gate_value=float(gate),
|
||||
))
|
||||
except (TypeError, ValueError) as exc:
|
||||
raise ValueError(f"malformed record {index}: {exc}") from exc
|
||||
return samples
|
||||
except (KeyError, TypeError, ValueError) as exc:
|
||||
raise ValueError(f"Invalid timeseries response: {exc}") from exc
|
||||
|
||||
|
||||
def _parse_timestamp(value: str) -> datetime:
|
||||
return datetime.fromisoformat(value.replace("Z", "+00:00")).astimezone(UTC)
|
||||
if not isinstance(value, str):
|
||||
raise ValueError("timestamp must be an ISO string")
|
||||
timestamp = datetime.fromisoformat(value.replace("Z", "+00:00"))
|
||||
if timestamp.tzinfo is None or timestamp.utcoffset() is None:
|
||||
raise ValueError("timestamp must be timezone-aware")
|
||||
return timestamp.astimezone(UTC)
|
||||
|
||||
@@ -65,8 +65,10 @@ class MaterialPollingService:
|
||||
)
|
||||
|
||||
def poll_once(self, *, now: datetime) -> MaterialPollResult | None:
|
||||
if now.tzinfo is None or now.utcoffset() is None:
|
||||
raise ValueError("now must be timezone-aware")
|
||||
run = self._gateway.get_open_production_run(self._machine_id)
|
||||
if run is None:
|
||||
if run is None or now < run.start:
|
||||
return None
|
||||
|
||||
saved_polling_state = self._state_store.load(
|
||||
@@ -92,6 +94,9 @@ class MaterialPollingService:
|
||||
)
|
||||
start = run.start
|
||||
|
||||
if now < start:
|
||||
return None
|
||||
|
||||
integrator = MaterialConsumptionIntegrator(
|
||||
self._config,
|
||||
initial_state=initial_state,
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
"""JSON-backed persistence for material polling state."""
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
import tempfile
|
||||
from pathlib import Path
|
||||
|
||||
from production_analytics.calculations.material_state import (
|
||||
material_state_from_dict,
|
||||
material_state_to_dict,
|
||||
)
|
||||
from production_analytics.service.material_polling import MaterialPollingState
|
||||
|
||||
|
||||
class JsonMaterialStateStore:
|
||||
def __init__(self, directory: str | Path) -> None:
|
||||
self._directory = Path(directory)
|
||||
|
||||
def load(
|
||||
self,
|
||||
machine_id: str,
|
||||
production_order: str,
|
||||
) -> MaterialPollingState | None:
|
||||
path = self._path(machine_id, production_order)
|
||||
if not path.exists():
|
||||
return None
|
||||
|
||||
try:
|
||||
with path.open(encoding="utf-8") as f:
|
||||
data = json.load(f)
|
||||
if not isinstance(data, dict) or not isinstance(data.get("run_id"), str):
|
||||
raise ValueError("run_id must be a string")
|
||||
if not data["run_id"]:
|
||||
raise ValueError("run_id must not be empty")
|
||||
return MaterialPollingState(
|
||||
run_id=data["run_id"],
|
||||
integration_state=material_state_from_dict(data["integration_state"]),
|
||||
)
|
||||
except (ValueError, KeyError, TypeError) as exc:
|
||||
raise ValueError(f"Invalid material polling state in {path}: {exc}") from exc
|
||||
|
||||
def save(
|
||||
self,
|
||||
machine_id: str,
|
||||
production_order: str,
|
||||
state: MaterialPollingState,
|
||||
) -> None:
|
||||
self._directory.mkdir(parents=True, exist_ok=True)
|
||||
path = self._path(machine_id, production_order)
|
||||
|
||||
payload = {
|
||||
"run_id": state.run_id,
|
||||
"integration_state": material_state_to_dict(state.integration_state),
|
||||
}
|
||||
|
||||
# Validate before touching the primary file, including non-finite values.
|
||||
material_state_from_dict(payload["integration_state"])
|
||||
if not isinstance(state.run_id, str) or not state.run_id:
|
||||
raise ValueError("run_id must be a non-empty string")
|
||||
temporary_path = None
|
||||
try:
|
||||
with tempfile.NamedTemporaryFile(
|
||||
mode="w", encoding="utf-8", dir=self._directory,
|
||||
prefix=".material-state-", suffix=".tmp", delete=False,
|
||||
) as f:
|
||||
temporary_path = Path(f.name)
|
||||
json.dump(payload, f, indent=2, allow_nan=False)
|
||||
f.flush()
|
||||
os.fsync(f.fileno())
|
||||
os.replace(temporary_path, path)
|
||||
finally:
|
||||
if temporary_path is not None:
|
||||
temporary_path.unlink(missing_ok=True)
|
||||
|
||||
def _path(self, machine_id: str, production_order: str) -> Path:
|
||||
# Hash the structured pair: sanitizing alone aliases distinct identifiers.
|
||||
identity = json.dumps([machine_id, production_order], ensure_ascii=True)
|
||||
digest = hashlib.sha256(identity.encode("utf-8")).hexdigest()
|
||||
return self._directory / f"{digest}.json"
|
||||
Reference in New Issue
Block a user