Bootstrap material state across production runs
This commit is contained in:
@@ -23,36 +23,63 @@ class EnlyzeApiGateway:
|
||||
self._client = client
|
||||
|
||||
def get_open_production_run(self, machine_id: str) -> EnlyzeProductionRun | None:
|
||||
response = self._client.get(
|
||||
"/v2/production-runs",
|
||||
{"machine": machine_id},
|
||||
)
|
||||
runs = [run for run in self.get_production_runs(machine_id) if run.end is None]
|
||||
if len(runs) > 1:
|
||||
raise ValueError(
|
||||
"Invalid production-run response: multiple open Production Runs returned"
|
||||
)
|
||||
return runs[0] if runs else 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_production_runs(self, machine_id: str) -> list[EnlyzeProductionRun]:
|
||||
"""Retrieve all machine runs, following validated cursor pagination."""
|
||||
params = {"machine": machine_id}
|
||||
runs = []
|
||||
followed_cursors: set[str] = set()
|
||||
page = 1
|
||||
while True:
|
||||
response = self._client.get("/v2/production-runs", params)
|
||||
try:
|
||||
if not isinstance(response.body, dict):
|
||||
raise ValueError("body must be an object")
|
||||
items = response.body["data"]
|
||||
if not isinstance(items, list):
|
||||
raise ValueError("data must be a list")
|
||||
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")
|
||||
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")
|
||||
start = _parse_timestamp(item["start"])
|
||||
end = _parse_timestamp(item["end"]) if item["end"] is not None else None
|
||||
if end is not None and end < start:
|
||||
raise ValueError("run end must not precede start")
|
||||
runs.append(EnlyzeProductionRun(
|
||||
uuid=item["uuid"], machine_id=item["machine"],
|
||||
product_id=item.get("product"), production_order=item["production_order"],
|
||||
start=start, end=end,
|
||||
))
|
||||
next_cursor = None
|
||||
if "metadata" in response.body:
|
||||
metadata = response.body["metadata"]
|
||||
if not isinstance(metadata, dict) or "next_cursor" not in metadata:
|
||||
raise ValueError("metadata must be an object with a next_cursor field")
|
||||
next_cursor = metadata["next_cursor"]
|
||||
if next_cursor is not None and (
|
||||
not isinstance(next_cursor, str) or not next_cursor
|
||||
):
|
||||
raise ValueError("next_cursor must be null or a non-empty string")
|
||||
if next_cursor is None:
|
||||
return runs
|
||||
if next_cursor in followed_cursors:
|
||||
raise ValueError("next_cursor has already been followed")
|
||||
followed_cursors.add(next_cursor)
|
||||
params = {**params, "cursor": next_cursor}
|
||||
page += 1
|
||||
except (KeyError, TypeError, ValueError) as exc:
|
||||
raise ValueError(f"Invalid production-run response on page {page}: {exc}") from exc
|
||||
|
||||
def get_material_samples(
|
||||
self,
|
||||
|
||||
@@ -77,7 +77,26 @@ class MaterialPollingService:
|
||||
)
|
||||
|
||||
if saved_polling_state is None:
|
||||
initial_state = None
|
||||
initial_state = MaterialIntegrationState()
|
||||
historical_runs = sorted(
|
||||
(
|
||||
previous for previous in self._gateway.get_production_runs(self._machine_id)
|
||||
if previous.production_order == run.production_order
|
||||
and previous.uuid != run.uuid
|
||||
and previous.start <= run.start
|
||||
),
|
||||
key=lambda previous: previous.start,
|
||||
)
|
||||
for previous in historical_runs:
|
||||
if previous.end is None or previous.end > run.start:
|
||||
raise ValueError("Historical production run overlaps the current open run")
|
||||
integrated = self._integrate(
|
||||
initial_state, start=previous.start, end=previous.end,
|
||||
)
|
||||
initial_state = MaterialIntegrationState(
|
||||
cumulative_consumption_kg=integrated.cumulative_consumption_kg,
|
||||
integrated_running_seconds=integrated.integrated_running_seconds,
|
||||
)
|
||||
start = run.start
|
||||
elif saved_polling_state.run_id == run.uuid:
|
||||
initial_state = saved_polling_state.integration_state
|
||||
@@ -97,31 +116,31 @@ class MaterialPollingService:
|
||||
if now < start:
|
||||
return None
|
||||
|
||||
integrator = MaterialConsumptionIntegrator(
|
||||
self._config,
|
||||
initial_state=initial_state,
|
||||
)
|
||||
|
||||
samples = self._gateway.get_material_samples(
|
||||
machine_id=self._machine_id,
|
||||
rate_variable_id=self._rate_variable_id,
|
||||
gate_variable_id=self._gate_variable_id,
|
||||
start=start,
|
||||
end=now,
|
||||
)
|
||||
|
||||
integrator.process_many(samples)
|
||||
state = self._integrate(initial_state, start=start, end=now)
|
||||
|
||||
self._state_store.save(
|
||||
self._machine_id,
|
||||
run.production_order,
|
||||
MaterialPollingState(
|
||||
run_id=run.uuid,
|
||||
integration_state=integrator.state,
|
||||
integration_state=state,
|
||||
),
|
||||
)
|
||||
|
||||
return MaterialPollResult(
|
||||
run=run,
|
||||
state=integrator.state,
|
||||
state=state,
|
||||
)
|
||||
|
||||
def _integrate(
|
||||
self, initial_state: MaterialIntegrationState, *, start: datetime, end: datetime,
|
||||
) -> MaterialIntegrationState:
|
||||
integrator = MaterialConsumptionIntegrator(self._config, initial_state=initial_state)
|
||||
samples = self._gateway.get_material_samples(
|
||||
machine_id=self._machine_id,
|
||||
rate_variable_id=self._rate_variable_id,
|
||||
gate_variable_id=self._gate_variable_id,
|
||||
start=start,
|
||||
end=end,
|
||||
)
|
||||
return integrator.process_many(samples)
|
||||
|
||||
@@ -290,3 +290,64 @@ def test_malformed_later_page(timeseries, data, error) -> None:
|
||||
with pytest.raises(ValueError, match=f"Invalid timeseries response on page 2:.*{error}"):
|
||||
gateway.get_material_samples(**args)
|
||||
assert client.post_json.call_count == 2
|
||||
|
||||
|
||||
def test_production_runs_follow_pages_and_find_open_run() -> None:
|
||||
client = Mock()
|
||||
closed = dict(uuid="closed", machine="m", production_order="opaque-order",
|
||||
start="2026-09-01T00:00:00Z", end="2026-09-01T01:00:00Z")
|
||||
opened = {**closed, "uuid": "open", "end": None}
|
||||
pages = [
|
||||
Mock(body={"data": [closed], "metadata": {"next_cursor": "a"}}),
|
||||
Mock(body={"data": [opened], "metadata": {"next_cursor": None}}),
|
||||
]
|
||||
client.get.side_effect = pages
|
||||
gateway = EnlyzeApiGateway(client)
|
||||
runs = gateway.get_production_runs("m")
|
||||
assert [run.uuid for run in runs] == ["closed", "open"]
|
||||
assert runs[0].end == datetime(2026, 9, 1, 1, tzinfo=UTC)
|
||||
assert [call.args for call in client.get.call_args_list] == [
|
||||
("/v2/production-runs", {"machine": "m"}),
|
||||
("/v2/production-runs", {"machine": "m", "cursor": "a"}),
|
||||
]
|
||||
client.get.side_effect = pages
|
||||
assert gateway.get_open_production_run("m") == runs[1]
|
||||
|
||||
|
||||
@pytest.mark.parametrize("metadata", [
|
||||
None, [], "bad", {}, {"next_cursor": ""}, {"next_cursor": 1},
|
||||
{"next_cursor": False}, {"next_cursor": []}, {"next_cursor": {}},
|
||||
])
|
||||
def test_production_run_invalid_pagination(metadata) -> None:
|
||||
client = Mock()
|
||||
client.get.return_value.body = {"data": [], "metadata": metadata}
|
||||
with pytest.raises(ValueError, match="Invalid production-run response.*next_cursor"):
|
||||
EnlyzeApiGateway(client).get_production_runs("m")
|
||||
|
||||
|
||||
@pytest.mark.parametrize("cursors", [["a", "a"], ["a", "b", "a"]])
|
||||
def test_production_run_cyclic_pagination(cursors) -> None:
|
||||
client = Mock()
|
||||
client.get.side_effect = [
|
||||
Mock(body={"data": [], "metadata": {"next_cursor": cursor}})
|
||||
for cursor in cursors
|
||||
]
|
||||
with pytest.raises(ValueError, match="next_cursor has already been followed"):
|
||||
EnlyzeApiGateway(client).get_production_runs("m")
|
||||
assert client.get.call_count == len(cursors)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("body", [
|
||||
None, [], {}, {"data": {}}, {"data": [None]},
|
||||
{"data": [dict(uuid="r", machine="wrong", production_order="o",
|
||||
start="2026-09-03T00:00:00Z", end=None)]},
|
||||
{"data": [dict(uuid="r", machine="m", production_order="o",
|
||||
start="2026-09-03T00:00:00Z", end="2026-09-02T00:00:00Z")]},
|
||||
])
|
||||
def test_production_run_malformed_later_page(body) -> None:
|
||||
client = Mock()
|
||||
client.get.side_effect = [
|
||||
Mock(body={"data": [], "metadata": {"next_cursor": "a"}}), Mock(body=body),
|
||||
]
|
||||
with pytest.raises(ValueError, match="Invalid production-run response on page 2"):
|
||||
EnlyzeApiGateway(client).get_production_runs("m")
|
||||
|
||||
@@ -16,6 +16,7 @@ from production_analytics.service.material_polling import (
|
||||
|
||||
def test_poll_once_starts_at_run_start_without_saved_state() -> None:
|
||||
gateway = Mock()
|
||||
gateway.get_production_runs.return_value = []
|
||||
state_store = Mock()
|
||||
state_store.load.return_value = None
|
||||
|
||||
@@ -28,6 +29,7 @@ def test_poll_once_starts_at_run_start_without_saved_state() -> None:
|
||||
end=None,
|
||||
)
|
||||
gateway.get_open_production_run.return_value = run
|
||||
gateway.get_production_runs.return_value = [run]
|
||||
|
||||
samples = [
|
||||
MaterialSample(
|
||||
@@ -82,6 +84,7 @@ def test_poll_once_starts_at_run_start_without_saved_state() -> None:
|
||||
|
||||
def test_poll_once_resumes_from_saved_timestamp_without_double_counting() -> None:
|
||||
gateway = Mock()
|
||||
gateway.get_production_runs.return_value = []
|
||||
state_store = Mock()
|
||||
|
||||
run = EnlyzeProductionRun(
|
||||
@@ -135,6 +138,7 @@ def test_poll_once_resumes_from_saved_timestamp_without_double_counting() -> Non
|
||||
)
|
||||
|
||||
assert result is not None
|
||||
gateway.get_production_runs.assert_not_called()
|
||||
assert result.state.cumulative_consumption_kg == 20.0
|
||||
assert result.state.integrated_running_seconds == 20.0
|
||||
|
||||
@@ -149,6 +153,7 @@ def test_poll_once_resumes_from_saved_timestamp_without_double_counting() -> Non
|
||||
|
||||
def test_poll_once_without_open_run_does_nothing() -> None:
|
||||
gateway = Mock()
|
||||
gateway.get_production_runs.return_value = []
|
||||
state_store = Mock()
|
||||
gateway.get_open_production_run.return_value = None
|
||||
|
||||
@@ -174,6 +179,7 @@ def test_poll_once_without_open_run_does_nothing() -> None:
|
||||
|
||||
def test_new_run_of_same_order_keeps_total_but_restarts_integration_at_run_start() -> None:
|
||||
gateway = Mock()
|
||||
gateway.get_production_runs.return_value = []
|
||||
state_store = Mock()
|
||||
|
||||
run = EnlyzeProductionRun(
|
||||
@@ -228,6 +234,7 @@ def test_new_run_of_same_order_keeps_total_but_restarts_integration_at_run_start
|
||||
)
|
||||
|
||||
assert result is not None
|
||||
gateway.get_production_runs.assert_not_called()
|
||||
assert result.state.cumulative_consumption_kg == 110.0
|
||||
assert result.state.integrated_running_seconds == 370.0
|
||||
|
||||
@@ -250,6 +257,7 @@ def polling(tmp_path):
|
||||
from production_analytics.service.material_state_store import JsonMaterialStateStore
|
||||
|
||||
gateway = Mock()
|
||||
gateway.get_production_runs.return_value = []
|
||||
gateway.get_open_production_run.return_value = EnlyzeProductionRun(
|
||||
"run", "machine", None, "order", datetime(2026, 9, 3, tzinfo=UTC), None,
|
||||
)
|
||||
@@ -291,6 +299,7 @@ def test_different_order_does_not_reuse_state(polling) -> None:
|
||||
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_production_runs.assert_not_called()
|
||||
gateway.get_material_samples.assert_not_called()
|
||||
assert store.load("machine", "order") is None
|
||||
|
||||
@@ -325,3 +334,76 @@ def test_now_before_saved_timestamp_leaves_state_unchanged(polling) -> None:
|
||||
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
|
||||
|
||||
|
||||
@pytest.mark.parametrize("order", [
|
||||
"K 7-12026000815", "K 7-12025000074-K 7-12025000075",
|
||||
])
|
||||
@pytest.mark.parametrize("gap_seconds", [5, 86400])
|
||||
def test_bootstrap_replays_exact_order_segments_and_then_resumes(
|
||||
polling, order, gap_seconds,
|
||||
) -> None:
|
||||
from dataclasses import replace
|
||||
from datetime import timedelta
|
||||
|
||||
service, gateway, store = polling
|
||||
base = datetime(2026, 9, 3, tzinfo=UTC)
|
||||
runs = [
|
||||
EnlyzeProductionRun(
|
||||
f"run-{i}", "machine", None, order,
|
||||
base + timedelta(seconds=i * (10 + gap_seconds)),
|
||||
base + timedelta(seconds=i * (10 + gap_seconds) + 10) if i < 4 else None,
|
||||
)
|
||||
for i in range(5)
|
||||
]
|
||||
current = runs[-1]
|
||||
now = current.start + timedelta(seconds=10)
|
||||
gateway.get_open_production_run.return_value = current
|
||||
gateway.get_production_runs.return_value = [
|
||||
current, runs[2], runs[0], runs[3], runs[1],
|
||||
replace(runs[0], uuid="unrelated", production_order=order + "-suffix"),
|
||||
replace(runs[0], uuid="component", production_order="K 7-12025000074"),
|
||||
replace(runs[0], uuid="whitespace", production_order=order + " "),
|
||||
replace(current, uuid="future", start=now + timedelta(days=1)),
|
||||
]
|
||||
|
||||
def samples(**kwargs):
|
||||
return [
|
||||
MaterialSample(kwargs["start"], 3600, 1),
|
||||
MaterialSample(kwargs["end"], 3600, 1),
|
||||
]
|
||||
|
||||
gateway.get_material_samples.side_effect = samples
|
||||
result = service.poll_once(now=now)
|
||||
assert result.state.cumulative_consumption_kg == 50
|
||||
assert result.state.integrated_running_seconds == 50
|
||||
assert [
|
||||
(call.kwargs["start"], call.kwargs["end"])
|
||||
for call in gateway.get_material_samples.call_args_list
|
||||
] == [(run.start, run.end or now) for run in runs]
|
||||
assert store.load("machine", order) == MaterialPollingState(current.uuid, result.state)
|
||||
|
||||
gateway.get_production_runs.reset_mock()
|
||||
gateway.get_material_samples.reset_mock()
|
||||
result = service.poll_once(now=now + timedelta(seconds=10))
|
||||
gateway.get_production_runs.assert_not_called()
|
||||
assert gateway.get_material_samples.call_args.kwargs["start"] == now
|
||||
assert result.state.cumulative_consumption_kg == 60
|
||||
assert result.state.integrated_running_seconds == 60
|
||||
|
||||
|
||||
def test_bootstrap_failure_does_not_persist_partial_totals(polling) -> None:
|
||||
from dataclasses import replace
|
||||
from datetime import timedelta
|
||||
|
||||
service, gateway, store = polling
|
||||
current = gateway.get_open_production_run.return_value
|
||||
gateway.get_production_runs.return_value = [
|
||||
replace(current, uuid="old", start=current.start - timedelta(hours=1),
|
||||
end=current.start - timedelta(minutes=30)),
|
||||
current,
|
||||
]
|
||||
gateway.get_material_samples.side_effect = [[], ValueError("failed page")]
|
||||
with pytest.raises(ValueError, match="failed page"):
|
||||
service.poll_once(now=current.start + timedelta(seconds=10))
|
||||
assert store.load("machine", "order") is None
|
||||
|
||||
Reference in New Issue
Block a user