diff --git a/src/production_analytics/enlyze/gateway.py b/src/production_analytics/enlyze/gateway.py index da76a0d..08e42eb 100644 --- a/src/production_analytics/enlyze/gateway.py +++ b/src/production_analytics/enlyze/gateway.py @@ -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, diff --git a/src/production_analytics/service/material_polling.py b/src/production_analytics/service/material_polling.py index 6058977..a2e1ae5 100644 --- a/src/production_analytics/service/material_polling.py +++ b/src/production_analytics/service/material_polling.py @@ -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) diff --git a/tests/test_enlyze_gateway.py b/tests/test_enlyze_gateway.py index 2a05166..7636f4d 100644 --- a/tests/test_enlyze_gateway.py +++ b/tests/test_enlyze_gateway.py @@ -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") diff --git a/tests/test_material_polling_service.py b/tests/test_material_polling_service.py index 0df9a3b..59e4255 100644 --- a/tests/test_material_polling_service.py +++ b/tests/test_material_polling_service.py @@ -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