From efc563512f6be072194d2f1e117f6378d65fdb9d Mon Sep 17 00:00:00 2001 From: Martin Tazl Date: Sat, 5 Sep 2026 05:50:01 +0200 Subject: [PATCH] Harden live material polling state --- README.md | 27 +++- docs/roadmap.md | 7 +- .../calculations/material_state.py | 32 ++++ src/production_analytics/enlyze/gateway.py | 104 ++++++++----- .../service/material_polling.py | 7 +- .../service/material_state_store.py | 80 ++++++++++ tests/test_enlyze_gateway.py | 71 +++++++++ tests/test_material_polling_service.py | 84 +++++++++++ tests/test_material_state_store.py | 141 ++++++++++++++++++ 9 files changed, 513 insertions(+), 40 deletions(-) create mode 100644 src/production_analytics/service/material_state_store.py create mode 100644 tests/test_material_state_store.py diff --git a/README.md b/README.md index a70be71..9005610 100644 --- a/README.md +++ b/README.md @@ -7,8 +7,9 @@ production-order context, and calculation state. ## Status -This repository includes a read-only ENLYZE API exploration CLI. It contains -no verified ENLYZE operation wrappers, database migrations, or HTTP endpoints. +This repository includes an ENLYZE exploration CLI, a production-run/timeseries +gateway, and single-cycle live material polling with JSON state persistence. +Database migrations, a continuous polling runner, and HTTP endpoints remain pending. ## Intended flow @@ -76,7 +77,7 @@ 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 +and JSON persistence are supported. 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 @@ -94,7 +95,25 @@ is gated by `Geschwindigkeit Gesamtanlage > 0.5 m/min`. With the explicit 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. +or live ENLYZE access. Grafana and Bento 1 remain pending. + +`MaterialPollingService.poll_once(now=...)` starts at the open run's start or +resumes inclusively at its last processed timestamp. Duplicate boundary samples +add no time. A new run UUID for the same order retains cumulative kg and running +seconds but resets the temporal baseline; different orders and machines have +separate state. Empty responses add no consumption. Naive `now` is rejected; +no open run or a window beginning after `now` returns `None` without saving. + +`JsonMaterialStateStore` uses deterministic SHA-256 filenames derived from the +machine/order pair, preventing path traversal and sanitized-name collisions. +It flushes and fsyncs a temporary file in the state directory before atomic +replacement; a failed write leaves the previous primary file intact. Missing +state returns `None`; malformed JSON or state raises `ValueError` with file +context. Use an ignored runtime directory and one polling writer per state key; +atomic replacement does not coordinate concurrent read-modify-write cycles. +Files from the earlier uncommitted sanitized-filename prototype are not loaded +under the new names. The gateway rejects ambiguous open runs, naive windows, +and missing columns or malformed records instead of silently skipping them. ## Peak-cycle detection diff --git a/docs/roadmap.md b/docs/roadmap.md index 1f544a6..cfdd0c3 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -13,7 +13,12 @@ Exploration established a K7 candidate mass-flow signal and speed gate; the evidence and remaining uncertainties are recorded in `docs/enlyze-api.md`. Continue to record sanitized fixtures for newly verified semantics. -## 2. Live material-consumption MVP (next implementation milestone) +## 2. Live material-consumption MVP (in progress) + +The integration engine, ENLYZE gateway, single-cycle polling service, and atomic +JSON checkpoint store are implemented and covered by synthetic tests. Next: +wire a configured continuous polling runner and expose/store derived totals +for Grafana, followed by the TimescaleDB persistence foundation below. Detect the currently active Production Run and production order, continuously ingest new ENLYZE samples, maintain persistent incremental integration state, diff --git a/src/production_analytics/calculations/material_state.py b/src/production_analytics/calculations/material_state.py index 6ada5fd..e42ac78 100644 --- a/src/production_analytics/calculations/material_state.py +++ b/src/production_analytics/calculations/material_state.py @@ -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) diff --git a/src/production_analytics/enlyze/gateway.py b/src/production_analytics/enlyze/gateway.py index d20bcfe..47077f1 100644 --- a/src/production_analytics/enlyze/gateway.py +++ b/src/production_analytics/enlyze/gateway.py @@ -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) diff --git a/src/production_analytics/service/material_polling.py b/src/production_analytics/service/material_polling.py index 002ed97..6058977 100644 --- a/src/production_analytics/service/material_polling.py +++ b/src/production_analytics/service/material_polling.py @@ -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, diff --git a/src/production_analytics/service/material_state_store.py b/src/production_analytics/service/material_state_store.py new file mode 100644 index 0000000..9fe131d --- /dev/null +++ b/src/production_analytics/service/material_state_store.py @@ -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" diff --git a/tests/test_enlyze_gateway.py b/tests/test_enlyze_gateway.py index f96fe25..16b121e 100644 --- a/tests/test_enlyze_gateway.py +++ b/tests/test_enlyze_gateway.py @@ -1,6 +1,8 @@ from datetime import UTC, datetime from unittest.mock import Mock +import pytest + from production_analytics.enlyze.gateway import EnlyzeApiGateway @@ -131,3 +133,72 @@ def test_get_material_samples_uses_column_names_not_fixed_positions() -> None: assert samples[0].material_rate_kg_per_hour == 1028.5 assert samples[0].gate_value == 2.1 + + +@pytest.fixture +def timeseries(): + client = Mock() + client.post_json.return_value.body = {"data": { + "columns": ["time", "rate", "gate"], "records": [], + }} + args = dict(machine_id="m", rate_variable_id="rate", gate_variable_id="gate", + start=datetime(2026, 9, 3, tzinfo=UTC), end=datetime(2026, 9, 4, tzinfo=UTC)) + return EnlyzeApiGateway(client), client, args + + +@pytest.mark.parametrize("name", ["start", "end"]) +def test_naive_window_rejected(timeseries, name) -> None: + gateway, client, args = timeseries + args[name] = args[name].replace(tzinfo=None) + with pytest.raises(ValueError, match="timezone-aware"): + gateway.get_material_samples(**args) + client.post_json.assert_not_called() + + +def test_reversed_window_rejected(timeseries) -> None: + gateway, client, args = timeseries + args["start"], args["end"] = args["end"], args["start"] + with pytest.raises(ValueError, match="precede"): + gateway.get_material_samples(**args) + client.post_json.assert_not_called() + + +def test_empty_timeseries(timeseries) -> None: + gateway, client, args = timeseries + assert gateway.get_material_samples(**args) == [] + + +@pytest.mark.parametrize("column", ["time", "rate", "gate"]) +def test_missing_required_column(timeseries, column) -> None: + gateway, client, args = timeseries + client.post_json.return_value.body["data"]["columns"].remove(column) + with pytest.raises(ValueError, match="required column"): + gateway.get_material_samples(**args) + + +@pytest.mark.parametrize("record", [[], {}, ["bad", 1, 1], [None, 1, 1], + ["2026-09-03T00:00:00", 1, 1], ["2026-09-03T00:00:00Z", None, 1], + ["2026-09-03T00:00:00Z", 1, float("inf")], ["2026-09-03T00:00:00Z", True, 1]]) +def test_malformed_records_fail_clearly(timeseries, record) -> None: + gateway, client, args = timeseries + client.post_json.return_value.body["data"]["records"] = [record] + with pytest.raises(ValueError, match="malformed record 0"): + gateway.get_material_samples(**args) + + +def test_multiple_open_runs_rejected() -> None: + client = Mock() + run = dict(uuid="r", machine="m", production_order="o", start="2026-09-03T00:00:00Z", + end=None) + client.get.return_value.body = {"data": [run, {**run, "uuid": "r2"}]} + with pytest.raises(ValueError, match="multiple open Production Runs"): + EnlyzeApiGateway(client).get_open_production_run("m") + + +@pytest.mark.parametrize("item", [{}, None, {"end": None}, + dict(uuid="r", machine="m", production_order="o", start="bad", end=None)]) +def test_malformed_run_fails_clearly(item) -> None: + client = Mock() + client.get.return_value.body = {"data": [item]} + with pytest.raises(ValueError, match="Invalid production-run response"): + EnlyzeApiGateway(client).get_open_production_run("m") diff --git a/tests/test_material_polling_service.py b/tests/test_material_polling_service.py index 3da1a16..0df9a3b 100644 --- a/tests/test_material_polling_service.py +++ b/tests/test_material_polling_service.py @@ -1,6 +1,8 @@ from datetime import UTC, datetime from unittest.mock import Mock +import pytest + from production_analytics.calculations.material_consumption import ( MaterialIntegrationState, MaterialSample, @@ -241,3 +243,85 @@ def test_new_run_of_same_order_keeps_total_but_restarts_integration_at_run_start start=run.start, end=datetime(2026, 9, 3, 8, 0, 30, tzinfo=UTC), ) + + +@pytest.fixture +def polling(tmp_path): + from production_analytics.service.material_state_store import JsonMaterialStateStore + + gateway = Mock() + gateway.get_open_production_run.return_value = EnlyzeProductionRun( + "run", "machine", None, "order", datetime(2026, 9, 3, tzinfo=UTC), None, + ) + gateway.get_material_samples.return_value = [] + store = JsonMaterialStateStore(tmp_path) + service = MaterialPollingService( + gateway=gateway, state_store=store, machine_id="machine", + rate_variable_id="rate", gate_variable_id="gate", gate_threshold=0.5, + max_sample_gap_seconds=20, + ) + return service, gateway, store + + +def test_empty_response_persists_state_without_inferred_consumption(polling) -> None: + service, gateway, store = polling + now = datetime(2026, 9, 3, 1, tzinfo=UTC) + assert service.poll_once(now=now).state == MaterialIntegrationState() + saved = MaterialPollingState("run", MaterialIntegrationState( + cumulative_consumption_kg=10, integrated_running_seconds=10, + last_processed_timestamp=now, last_material_rate_kg_per_hour=3600, + last_gate_value=1, integration_active=True, + )) + store.save("machine", "order", saved) + assert service.poll_once(now=now).state == saved.integration_state + assert store.load("machine", "order") == saved + + +def test_different_order_does_not_reuse_state(polling) -> None: + service, gateway, store = polling + store.save("machine", "other-order", MaterialPollingState( + "old", MaterialIntegrationState(cumulative_consumption_kg=100), + )) + result = service.poll_once(now=datetime(2026, 9, 3, 1, tzinfo=UTC)) + assert result.state == MaterialIntegrationState() + assert gateway.get_material_samples.call_args.kwargs["start"] == result.run.start + assert store.load("machine", "other-order").integration_state.cumulative_consumption_kg == 100 + + +def test_now_before_run_start_does_not_query_or_save(polling) -> None: + service, gateway, store = polling + assert service.poll_once(now=datetime(2026, 9, 2, tzinfo=UTC)) is None + gateway.get_material_samples.assert_not_called() + assert store.load("machine", "order") is None + + +def test_naive_now_is_rejected_before_io(polling) -> None: + service, gateway, store = polling + with pytest.raises(ValueError, match="now must be timezone-aware"): + service.poll_once(now=datetime(2026, 9, 3)) + gateway.get_open_production_run.assert_not_called() + + +def test_new_run_empty_response_resets_baseline_and_retains_totals(polling) -> None: + service, gateway, store = polling + store.save("machine", "order", MaterialPollingState("previous", MaterialIntegrationState( + cumulative_consumption_kg=100, integrated_running_seconds=50, + last_processed_timestamp=datetime(2026, 9, 2, 23, 59, 59, tzinfo=UTC), + last_material_rate_kg_per_hour=3600, last_gate_value=1, integration_active=True, + ))) + result = service.poll_once(now=datetime(2026, 9, 3, 1, tzinfo=UTC)) + assert result.state == MaterialIntegrationState( + cumulative_consumption_kg=100, integrated_running_seconds=50, + ) + assert store.load("machine", "order") == MaterialPollingState("run", result.state) + + +def test_now_before_saved_timestamp_leaves_state_unchanged(polling) -> None: + service, gateway, store = polling + saved = MaterialPollingState("run", MaterialIntegrationState( + last_processed_timestamp=datetime(2026, 9, 3, 2, tzinfo=UTC), + )) + store.save("machine", "order", saved) + assert service.poll_once(now=datetime(2026, 9, 3, 1, tzinfo=UTC)) is None + gateway.get_material_samples.assert_not_called() + assert store.load("machine", "order") == saved diff --git a/tests/test_material_state_store.py b/tests/test_material_state_store.py new file mode 100644 index 0000000..ac16cdf --- /dev/null +++ b/tests/test_material_state_store.py @@ -0,0 +1,141 @@ +from datetime import UTC, datetime + +import pytest + +from production_analytics.calculations.material_consumption import ( + MaterialIntegrationState, +) +from production_analytics.service.material_polling import MaterialPollingState +from production_analytics.service.material_state_store import JsonMaterialStateStore + + +def test_json_state_store_returns_none_for_missing_state(tmp_path) -> None: + store = JsonMaterialStateStore(tmp_path) + + assert store.load("machine-1", "ORDER-1") is None + + +def test_json_state_store_roundtrip(tmp_path) -> None: + store = JsonMaterialStateStore(tmp_path) + + original = MaterialPollingState( + run_id="run-1", + integration_state=MaterialIntegrationState( + cumulative_consumption_kg=123.4, + integrated_running_seconds=456.0, + last_processed_timestamp=datetime( + 2026, 9, 3, 5, 10, 50, tzinfo=UTC + ), + last_material_rate_kg_per_hour=1028.5, + last_gate_value=2.1, + integration_active=True, + ), + ) + + store.save("machine-1", "ORDER-1", original) + restored = store.load("machine-1", "ORDER-1") + + assert restored == original + + +def test_json_state_store_separates_orders(tmp_path) -> None: + store = JsonMaterialStateStore(tmp_path) + + state_a = MaterialPollingState( + run_id="run-a", + integration_state=MaterialIntegrationState( + cumulative_consumption_kg=10.0, + ), + ) + state_b = MaterialPollingState( + run_id="run-b", + integration_state=MaterialIntegrationState( + cumulative_consumption_kg=20.0, + ), + ) + + store.save("machine-1", "ORDER-1", state_a) + store.save("machine-1", "ORDER-2", state_b) + + assert store.load("machine-1", "ORDER-1") == state_a + assert store.load("machine-1", "ORDER-2") == state_b + + +def test_identifiers_are_isolated_safe_and_deterministic(tmp_path) -> None: + store = JsonMaterialStateStore(tmp_path) + pairs = [("a/b", "x"), ("a_b", "x"), ("a", "b__x"), ("a__b", "x"), + ("../..", "/tmp/escape"), ("", ""), ("a/b", "y")] + for index, pair in enumerate(pairs): + store.save(*pair, MaterialPollingState(str(index), MaterialIntegrationState())) + assert len(list(tmp_path.iterdir())) == len(pairs) + for index, pair in enumerate(pairs): + assert JsonMaterialStateStore(tmp_path).load(*pair).run_id == str(index) + assert all(path.parent == tmp_path and len(path.name) == 69 for path in tmp_path.iterdir()) + + +def test_failed_write_preserves_primary_and_cleans_temporary_file(tmp_path, monkeypatch) -> None: + import production_analytics.service.material_state_store as module + + store = JsonMaterialStateStore(tmp_path) + state = MaterialPollingState("old", MaterialIntegrationState()) + store.save("machine", "order", state) + path = next(tmp_path.iterdir()) + original = path.read_bytes() + + def fail_dump(payload, file, **kwargs): + file.write('{"partial":') + raise OSError("disk full") + + monkeypatch.setattr(module.json, "dump", fail_dump) + with pytest.raises(OSError, match="disk full"): + store.save("machine", "order", MaterialPollingState("new", MaterialIntegrationState())) + assert path.read_bytes() == original + assert list(tmp_path.iterdir()) == [path] + + +def test_replace_failure_preserves_primary(tmp_path, monkeypatch) -> None: + import production_analytics.service.material_state_store as module + + store = JsonMaterialStateStore(tmp_path) + state = MaterialPollingState("old", MaterialIntegrationState()) + store.save("m", "o", state) + path = next(tmp_path.iterdir()) + + def fail_replace(source, target): + assert source.parent == target.parent == tmp_path + assert store.load("m", "o") == state + raise OSError("replace failed") + + monkeypatch.setattr(module.os, "replace", fail_replace) + with pytest.raises(OSError, match="replace failed"): + store.save("m", "o", MaterialPollingState("new", MaterialIntegrationState())) + assert store.load("m", "o") == state + assert list(tmp_path.iterdir()) == [path] + + +@pytest.mark.parametrize("contents", ['{', '[]', '{}', '{"run_id": null}', + '{"run_id": "r", "integration_state": {}}']) +def test_corrupt_state_fails_clearly(tmp_path, contents) -> None: + store = JsonMaterialStateStore(tmp_path) + store.save("m", "o", MaterialPollingState("r", MaterialIntegrationState())) + next(tmp_path.iterdir()).write_text(contents) + with pytest.raises(ValueError, match="Invalid material polling state"): + store.load("m", "o") + + +@pytest.mark.parametrize("field,value", [ + ("integration_active", "false"), ("cumulative_consumption_kg", float("nan")), + ("integrated_running_seconds", -1), ("last_processed_timestamp", "2026-09-03T05:00:00"), + ("integration_active", True), ("last_gate_value", True), +]) +def test_malformed_integration_state_fails(tmp_path, field, value) -> None: + import json + + store = JsonMaterialStateStore(tmp_path) + store.save("m", "o", MaterialPollingState("r", MaterialIntegrationState())) + path = next(tmp_path.iterdir()) + payload = json.loads(path.read_text()) + payload["integration_state"][field] = value + path.write_text(json.dumps(payload)) + with pytest.raises(ValueError, match="Invalid material polling state"): + store.load("m", "o")