"""Single-cycle live material polling orchestration.""" from collections.abc import Callable 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 | None, gate_threshold: float, max_sample_gap_seconds: float, application_variable_ids: tuple[str, ...] = (), rotational_speed_variable_ids: tuple[str, ...] = (), specific_discharge_kg_per_rev_m: float | None = None, nominal_width_provider: Callable[[str], float] | None = None, ) -> None: if application_variable_ids and rotational_speed_variable_ids: raise ValueError("Material source modes are mutually exclusive") self._rotational_speed_variable_ids = rotational_speed_variable_ids self._specific_discharge = specific_discharge_kg_per_rev_m self._application_variable_ids = application_variable_ids self._nominal_width_provider = nominal_width_provider if ((application_variable_ids or rotational_speed_variable_ids) and nominal_width_provider is None): raise ValueError("width-based source requires a nominal width provider") 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 if gate_variable_id is not None else 0.0, max_sample_gap_seconds=max_sample_gap_seconds, ) 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 or now < run.start: return None width = (self._nominal_width_provider(run.production_order) if self._application_variable_ids or self._rotational_speed_variable_ids else None) saved_polling_state = self._state_store.load( self._machine_id, run.production_order, ) if saved_polling_state is 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, width=width, ) 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 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 if now < start: return None state = self._integrate(initial_state, start=start, end=now, width=width) self._state_store.save( self._machine_id, run.production_order, MaterialPollingState( run_id=run.uuid, integration_state=state, ), ) return MaterialPollResult( run=run, state=state, ) def _integrate( self, initial_state: MaterialIntegrationState, *, start: datetime, end: datetime, width: float | None = None, ) -> MaterialIntegrationState: integrator = MaterialConsumptionIntegrator(self._config, initial_state=initial_state) source_options = {} if self._application_variable_ids: source_options = dict( application_variable_ids=self._application_variable_ids, nominal_width_m=width, ) if self._rotational_speed_variable_ids: source_options = dict( rotational_speed_variable_ids=self._rotational_speed_variable_ids, specific_discharge_kg_per_rev_m=self._specific_discharge, nominal_width_m=width, ) 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, **source_options, ) return integrator.process_many(samples)