Add live material polling service

This commit is contained in:
2026-09-05 05:35:37 +02:00
parent 363be0a47e
commit 98a92c3739
2 changed files with 365 additions and 0 deletions
@@ -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,
)