Add central material calibration support

This commit is contained in:
2026-09-09 03:22:47 +02:00
parent 0210ced9e3
commit 73a6e60450
21 changed files with 848 additions and 23 deletions
@@ -38,12 +38,14 @@ class MaterialEfficiencyService:
self, erp: WorkplaceStatusReader, materials: MaterialSnapshotRepository, *,
workplace: str, machine_id: str, calculation_id: str, format_template: str,
erp_timezone: ZoneInfo | None = None,
output_calculation_id: str | None = None,
) -> None:
self.erp = erp
self.materials = materials
self.workplace = workplace
self.machine_id = machine_id
self.calculation_id = calculation_id
self.output_calculation_id = output_calculation_id or calculation_id
self.format_template = format_template
self.erp_timezone = erp_timezone
@@ -58,7 +60,7 @@ class MaterialEfficiencyService:
Configuration, timestamp and adapter contract errors raise ValueError.
Database failures propagate, rather than being treated as missing data.
"""
if status.workplace != self.workplace:
if status.workplace.strip().casefold() != self.workplace.strip().casefold():
raise ValueError("ERP workplace does not match configured workplace")
order = build_enlyze_production_order(status.production_order, self.format_template)
quantity = status.good_quantity_m2
@@ -88,7 +90,8 @@ class MaterialEfficiencyService:
return None
return MaterialEfficiencySnapshot(
workplace=status.workplace, machine_id=self.machine_id,
calculation_id=self.calculation_id, erp_production_order=status.production_order,
calculation_id=self.output_calculation_id,
erp_production_order=status.production_order,
enlyze_production_order=order, article_number=status.article_number,
article_description=status.article_description,
nominal_width_m=extract_nominal_width_m(status.article_description),
@@ -45,10 +45,10 @@ def build_material_efficiency_runner(
if (
not isinstance(document, dict)
or not required <= document.keys()
or document.keys() - required - {"poll_interval_seconds"}
or document.keys() - required - {"poll_interval_seconds", "output_calculation_id"}
):
raise CalculationConfigError("Invalid material efficiency configuration fields")
for name in required:
for name in required | ({"output_calculation_id"} & document.keys()):
value = document[name]
if not isinstance(value, str) or not value.strip() or "\x00" in value:
raise CalculationConfigError(f"{name} must be a non-empty string without NUL bytes")
@@ -76,6 +76,7 @@ def build_material_efficiency_runner(
calculation_id=document["calculation_id"],
format_template=document["production_order_format"],
erp_timezone=timezone,
output_calculation_id=document.get("output_calculation_id"),
),
PostgresMaterialEfficiencyWriter(postgres),
poll_interval_seconds=interval,
@@ -14,6 +14,7 @@ from production_analytics.enlyze.gateway import (
EnlyzeApiGateway,
EnlyzeProductionRun,
)
from production_analytics.service.postgres_material_application import MaterialApplicationWriter
@dataclass(frozen=True, slots=True)
@@ -58,9 +59,20 @@ class MaterialPollingService:
rotational_speed_variable_ids: tuple[str, ...] = (),
specific_discharge_kg_per_rev_m: float | None = None,
nominal_width_provider: Callable[[str], float] | None = None,
process_application_calculation_id: str | None = None,
application_writer: MaterialApplicationWriter | None = None,
) -> None:
if application_variable_ids and rotational_speed_variable_ids:
raise ValueError("Material source modes are mutually exclusive")
if process_application_calculation_id is not None and (
not rotational_speed_variable_ids or not gate_variable_id
or gate_threshold <= 0 or application_writer is None
):
raise ValueError(
"Process application requires rotational source, positive gate and writer"
)
self._process_application_calculation_id = process_application_calculation_id
self._application_writer = application_writer
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
@@ -107,7 +119,8 @@ class MaterialPollingService:
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, start=previous.start, end=previous.end,
width=width, run=previous,
)
initial_state = MaterialIntegrationState(
cumulative_consumption_kg=integrated.cumulative_consumption_kg,
@@ -132,7 +145,7 @@ class MaterialPollingService:
if now < start:
return None
state = self._integrate(initial_state, start=start, end=now, width=width)
state = self._integrate(initial_state, start=start, end=now, width=width, run=run)
self._state_store.save(
self._machine_id,
@@ -150,7 +163,7 @@ class MaterialPollingService:
def _integrate(
self, initial_state: MaterialIntegrationState, *, start: datetime, end: datetime,
width: float | None = None,
width: float | None = None, run: EnlyzeProductionRun | None = None,
) -> MaterialIntegrationState:
integrator = MaterialConsumptionIntegrator(self._config, initial_state=initial_state)
source_options = {}
@@ -163,6 +176,8 @@ class MaterialPollingService:
rotational_speed_variable_ids=self._rotational_speed_variable_ids,
specific_discharge_kg_per_rev_m=self._specific_discharge, nominal_width_m=width,
)
if self._process_application_calculation_id is not None:
source_options["process_application_gate_threshold"] = self._config.gate_threshold
samples = self._gateway.get_material_samples(
machine_id=self._machine_id,
rate_variable_id=self._rate_variable_id,
@@ -171,4 +186,13 @@ class MaterialPollingService:
end=end,
**source_options,
)
return integrator.process_many(samples)
state = integrator.process_many(samples)
if self._process_application_calculation_id is not None:
assert self._application_writer is not None and run is not None
# Persist before advancing the checkpoint, so failed writes can be retried.
self._application_writer.write(
calculation_id=self._process_application_calculation_id,
machine_id=self._machine_id, production_order=run.production_order,
run_id=run.uuid, samples=samples,
)
return state
@@ -26,6 +26,9 @@ from production_analytics.service.postgres_material import (
PostgresMaterialSnapshotWriter,
PostgresSettings,
)
from production_analytics.service.postgres_material_application import (
PostgresMaterialApplicationWriter,
)
class _ReportingStateStore:
@@ -88,6 +91,11 @@ def build_material_runner(
rotational_speed_variable_ids=calculation.rotational_speed_signal_refs,
specific_discharge_kg_per_rev_m=calculation.specific_discharge_kg_per_rev_m,
)
if calculation.process_application_calculation_id is not None:
source_options.update(
process_application_calculation_id=calculation.process_application_calculation_id,
application_writer=PostgresMaterialApplicationWriter(postgres_settings),
)
service = MaterialPollingService(
gateway=EnlyzeApiGateway(ExplorationClient(settings)),
state_store=_ReportingStateStore(JsonMaterialStateStore(state_directory)),
@@ -0,0 +1,53 @@
"""Generic process application snapshots, timestamped at the source sample."""
from collections.abc import Sequence
from typing import Protocol
import psycopg
from production_analytics.calculations.config import finite_number
from production_analytics.calculations.material_consumption import MaterialSample
from production_analytics.service.postgres_material import PostgresSettings
class MaterialApplicationWriter(Protocol):
def write(
self, *, calculation_id: str, machine_id: str, production_order: str,
run_id: str, samples: Sequence[MaterialSample],
) -> None: ...
class PostgresMaterialApplicationWriter:
def __init__(self, settings: PostgresSettings) -> None:
self.settings = settings
def write(
self, *, calculation_id: str, machine_id: str, production_order: str,
run_id: str, samples: Sequence[MaterialSample],
) -> None:
rows = []
for sample in samples:
if sample.application_g_m2 is None:
continue
if sample.timestamp.utcoffset() is None:
raise ValueError("Application timestamp must be timezone-aware")
finite_number(sample.application_g_m2, "application_g_m2")
rows.append((sample.timestamp, calculation_id, machine_id, production_order,
run_id, sample.application_g_m2))
if not rows:
return
with psycopg.connect(
host=self.settings.host, port=self.settings.port, dbname=self.settings.dbname,
user=self.settings.user, password=self.settings.password, connect_timeout=10,
options="-c statement_timeout=10000",
) as connection:
with connection.cursor() as cursor:
cursor.executemany(
"""INSERT INTO material_application_snapshots
(timestamp, calculation_id, machine_id, production_order,
run_id, application_g_m2)
VALUES (%s, %s, %s, %s, %s, %s)
ON CONFLICT (calculation_id, machine_id, production_order, timestamp)
DO NOTHING""",
rows,
)