Add persistent material integration state
This commit is contained in:
@@ -27,3 +27,5 @@ secrets/*
|
|||||||
.idea/
|
.idea/
|
||||||
.vscode/
|
.vscode/
|
||||||
.DS_Store
|
.DS_Store
|
||||||
|
|
||||||
|
activate.sh
|
||||||
|
|||||||
@@ -53,9 +53,13 @@ class MaterialConsumptionIntegrator:
|
|||||||
are always UTC.
|
are always UTC.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
def __init__(self, config: MaterialIntegratorConfig) -> None:
|
def __init__(
|
||||||
|
self,
|
||||||
|
config: MaterialIntegratorConfig,
|
||||||
|
initial_state: MaterialIntegrationState | None = None,
|
||||||
|
) -> None:
|
||||||
self._config = config
|
self._config = config
|
||||||
self._state = MaterialIntegrationState()
|
self._state = initial_state or MaterialIntegrationState()
|
||||||
|
|
||||||
@property
|
@property
|
||||||
def config(self) -> MaterialIntegratorConfig:
|
def config(self) -> MaterialIntegratorConfig:
|
||||||
|
|||||||
@@ -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,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