diff --git a/src/production_analytics/service/material_polling.py b/src/production_analytics/service/material_polling.py new file mode 100644 index 0000000..002ed97 --- /dev/null +++ b/src/production_analytics/service/material_polling.py @@ -0,0 +1,122 @@ +"""Single-cycle live material polling orchestration.""" + +from dataclasses import dataclass +from datetime import datetime +from typing import Protocol + +from production_analytics.calculations.material_consumption import ( + MaterialConsumptionIntegrator, + MaterialIntegrationState, + MaterialIntegratorConfig, +) +from production_analytics.enlyze.gateway import ( + EnlyzeApiGateway, + EnlyzeProductionRun, +) + + +@dataclass(frozen=True, slots=True) +class MaterialPollingState: + run_id: str + integration_state: MaterialIntegrationState + + +class MaterialStateStore(Protocol): + def load( + self, + machine_id: str, + production_order: str, + ) -> MaterialPollingState | None: ... + + def save( + self, + machine_id: str, + production_order: str, + state: MaterialPollingState, + ) -> None: ... + + +@dataclass(frozen=True, slots=True) +class MaterialPollResult: + run: EnlyzeProductionRun + state: MaterialIntegrationState + + +class MaterialPollingService: + def __init__( + self, + *, + gateway: EnlyzeApiGateway, + state_store: MaterialStateStore, + machine_id: str, + rate_variable_id: str, + gate_variable_id: str, + gate_threshold: float, + max_sample_gap_seconds: float, + ) -> None: + self._gateway = gateway + self._state_store = state_store + self._machine_id = machine_id + self._rate_variable_id = rate_variable_id + self._gate_variable_id = gate_variable_id + self._config = MaterialIntegratorConfig( + gate_threshold=gate_threshold, + max_sample_gap_seconds=max_sample_gap_seconds, + ) + + def poll_once(self, *, now: datetime) -> MaterialPollResult | None: + run = self._gateway.get_open_production_run(self._machine_id) + if run is None: + return None + + saved_polling_state = self._state_store.load( + self._machine_id, + run.production_order, + ) + + if saved_polling_state is None: + initial_state = None + start = run.start + elif saved_polling_state.run_id == run.uuid: + initial_state = saved_polling_state.integration_state + start = ( + initial_state.last_processed_timestamp + if initial_state.last_processed_timestamp is not None + else run.start + ) + else: + previous = saved_polling_state.integration_state + initial_state = MaterialIntegrationState( + cumulative_consumption_kg=previous.cumulative_consumption_kg, + integrated_running_seconds=previous.integrated_running_seconds, + ) + start = run.start + + 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) + + self._state_store.save( + self._machine_id, + run.production_order, + MaterialPollingState( + run_id=run.uuid, + integration_state=integrator.state, + ), + ) + + return MaterialPollResult( + run=run, + state=integrator.state, + ) diff --git a/tests/test_material_polling_service.py b/tests/test_material_polling_service.py new file mode 100644 index 0000000..3da1a16 --- /dev/null +++ b/tests/test_material_polling_service.py @@ -0,0 +1,243 @@ +from datetime import UTC, datetime +from unittest.mock import Mock + +from production_analytics.calculations.material_consumption import ( + MaterialIntegrationState, + MaterialSample, +) +from production_analytics.enlyze.gateway import EnlyzeProductionRun +from production_analytics.service.material_polling import ( + MaterialPollingService, + MaterialPollingState, +) + + +def test_poll_once_starts_at_run_start_without_saved_state() -> None: + gateway = Mock() + state_store = Mock() + state_store.load.return_value = None + + run = EnlyzeProductionRun( + uuid="run-1", + machine_id="machine-1", + product_id="product-1", + production_order="ORDER-1", + start=datetime(2026, 9, 3, 5, 7, 47, tzinfo=UTC), + end=None, + ) + gateway.get_open_production_run.return_value = run + + samples = [ + MaterialSample( + timestamp=datetime(2026, 9, 3, 5, 7, 50, tzinfo=UTC), + material_rate_kg_per_hour=3600.0, + gate_value=1.0, + ), + MaterialSample( + timestamp=datetime(2026, 9, 3, 5, 8, 0, tzinfo=UTC), + material_rate_kg_per_hour=3600.0, + gate_value=1.0, + ), + ] + gateway.get_material_samples.return_value = samples + + service = MaterialPollingService( + gateway=gateway, + state_store=state_store, + machine_id="machine-1", + rate_variable_id="rate-variable", + gate_variable_id="gate-variable", + gate_threshold=0.5, + max_sample_gap_seconds=20.0, + ) + + result = service.poll_once( + now=datetime(2026, 9, 3, 5, 8, 10, tzinfo=UTC), + ) + + assert result is not None + assert result.run == run + assert result.state.cumulative_consumption_kg == 10.0 + assert result.state.integrated_running_seconds == 10.0 + + gateway.get_material_samples.assert_called_once_with( + machine_id="machine-1", + rate_variable_id="rate-variable", + gate_variable_id="gate-variable", + start=run.start, + end=datetime(2026, 9, 3, 5, 8, 10, tzinfo=UTC), + ) + + state_store.save.assert_called_once_with( + "machine-1", + run.production_order, + MaterialPollingState( + run_id=run.uuid, + integration_state=result.state, + ), + ) + + +def test_poll_once_resumes_from_saved_timestamp_without_double_counting() -> None: + gateway = Mock() + state_store = Mock() + + run = EnlyzeProductionRun( + uuid="run-1", + machine_id="machine-1", + product_id="product-1", + production_order="ORDER-1", + start=datetime(2026, 9, 3, 5, 7, 47, tzinfo=UTC), + end=None, + ) + gateway.get_open_production_run.return_value = run + + saved_state = MaterialIntegrationState( + cumulative_consumption_kg=10.0, + integrated_running_seconds=10.0, + last_processed_timestamp=datetime(2026, 9, 3, 5, 8, 0, tzinfo=UTC), + last_material_rate_kg_per_hour=3600.0, + last_gate_value=1.0, + integration_active=True, + ) + state_store.load.return_value = MaterialPollingState( + run_id=run.uuid, + integration_state=saved_state, + ) + + gateway.get_material_samples.return_value = [ + MaterialSample( + timestamp=datetime(2026, 9, 3, 5, 8, 0, tzinfo=UTC), + material_rate_kg_per_hour=3600.0, + gate_value=1.0, + ), + MaterialSample( + timestamp=datetime(2026, 9, 3, 5, 8, 10, tzinfo=UTC), + material_rate_kg_per_hour=3600.0, + gate_value=1.0, + ), + ] + + service = MaterialPollingService( + gateway=gateway, + state_store=state_store, + machine_id="machine-1", + rate_variable_id="rate-variable", + gate_variable_id="gate-variable", + gate_threshold=0.5, + max_sample_gap_seconds=20.0, + ) + + result = service.poll_once( + now=datetime(2026, 9, 3, 5, 8, 20, tzinfo=UTC), + ) + + assert result is not None + assert result.state.cumulative_consumption_kg == 20.0 + assert result.state.integrated_running_seconds == 20.0 + + gateway.get_material_samples.assert_called_once_with( + machine_id="machine-1", + rate_variable_id="rate-variable", + gate_variable_id="gate-variable", + start=saved_state.last_processed_timestamp, + end=datetime(2026, 9, 3, 5, 8, 20, tzinfo=UTC), + ) + + +def test_poll_once_without_open_run_does_nothing() -> None: + gateway = Mock() + state_store = Mock() + gateway.get_open_production_run.return_value = None + + service = MaterialPollingService( + gateway=gateway, + state_store=state_store, + machine_id="machine-1", + rate_variable_id="rate-variable", + gate_variable_id="gate-variable", + gate_threshold=0.5, + max_sample_gap_seconds=20.0, + ) + + result = service.poll_once( + now=datetime(2026, 9, 3, 5, 8, 20, tzinfo=UTC), + ) + + assert result is None + gateway.get_material_samples.assert_not_called() + state_store.load.assert_not_called() + state_store.save.assert_not_called() + + +def test_new_run_of_same_order_keeps_total_but_restarts_integration_at_run_start() -> None: + gateway = Mock() + state_store = Mock() + + run = EnlyzeProductionRun( + uuid="run-2", + machine_id="machine-1", + product_id="product-1", + production_order="ORDER-1", + start=datetime(2026, 9, 3, 8, 0, 0, tzinfo=UTC), + end=None, + ) + gateway.get_open_production_run.return_value = run + + previous_integration_state = MaterialIntegrationState( + cumulative_consumption_kg=100.0, + integrated_running_seconds=360.0, + last_processed_timestamp=datetime(2026, 9, 3, 6, 0, 0, tzinfo=UTC), + last_material_rate_kg_per_hour=3600.0, + last_gate_value=1.0, + integration_active=True, + ) + + state_store.load.return_value = MaterialPollingState( + run_id="run-1", + integration_state=previous_integration_state, + ) + + gateway.get_material_samples.return_value = [ + MaterialSample( + timestamp=datetime(2026, 9, 3, 8, 0, 10, tzinfo=UTC), + material_rate_kg_per_hour=3600.0, + gate_value=1.0, + ), + MaterialSample( + timestamp=datetime(2026, 9, 3, 8, 0, 20, tzinfo=UTC), + material_rate_kg_per_hour=3600.0, + gate_value=1.0, + ), + ] + + service = MaterialPollingService( + gateway=gateway, + state_store=state_store, + machine_id="machine-1", + rate_variable_id="rate-variable", + gate_variable_id="gate-variable", + gate_threshold=0.5, + max_sample_gap_seconds=20.0, + ) + + result = service.poll_once( + now=datetime(2026, 9, 3, 8, 0, 30, tzinfo=UTC), + ) + + assert result is not None + assert result.state.cumulative_consumption_kg == 110.0 + assert result.state.integrated_running_seconds == 370.0 + + state_store.load.assert_called_once_with( + "machine-1", + "ORDER-1", + ) + + gateway.get_material_samples.assert_called_once_with( + machine_id="machine-1", + rate_variable_id="rate-variable", + gate_variable_id="gate-variable", + start=run.start, + end=datetime(2026, 9, 3, 8, 0, 30, tzinfo=UTC), + )