Add Bento 1 fresh bentonite consumption

This commit is contained in:
2026-09-06 13:29:06 +02:00
parent ed0e20a04b
commit be8d76507a
11 changed files with 430 additions and 14 deletions
@@ -1,6 +1,6 @@
"""Strict configuration for executable material-consumption instances."""
from dataclasses import dataclass, fields
from dataclasses import MISSING, dataclass, fields
from math import isfinite
from pathlib import Path
@@ -54,6 +54,10 @@ class MaterialCalculationConfig:
output_metric: str
output_unit: str
group_by: str
source_mode: str = "direct_mass_rate"
application_signal_refs: tuple[str, ...] = ()
erp_workplace: str = ""
production_order_format: str = ""
def load_material_calculation(
@@ -71,7 +75,9 @@ def load_material_calculation(
entries = document["calculations"]
if not isinstance(entries, list) or not entries:
raise CalculationConfigError("calculations must be a non-empty list")
required = {field.name for field in fields(MaterialCalculationConfig)}
allowed = {field.name for field in fields(MaterialCalculationConfig)}
required = {field.name for field in fields(MaterialCalculationConfig)
if field.default is MISSING}
instances = {}
for index, entry in enumerate(entries):
prefix = f"calculations[{index}]"
@@ -79,10 +85,13 @@ def load_material_calculation(
raise CalculationConfigError(f"{prefix} must be a mapping")
if "type" in entry and entry["type"] != "material_consumption":
raise CalculationConfigError(f"{prefix}.type: only 'material_consumption' is supported")
area = entry.get("source_mode", "direct_mass_rate") == "area_application"
if area:
entry.setdefault("rate_signal_ref", "unused")
missing = required - entry.keys()
if missing:
raise CalculationConfigError(f"{prefix}: missing fields: {', '.join(sorted(missing))}")
if entry.keys() - required:
if entry.keys() - allowed:
raise CalculationConfigError(f"{prefix}: unknown fields (check spelling)")
for name in sorted(required - {"gate_threshold", "max_sample_gap_seconds"}):
if not isinstance(entry[name], str) or not entry[name].strip():
@@ -98,6 +107,32 @@ def load_material_calculation(
entry["max_sample_gap_seconds"] = finite_number(
entry["max_sample_gap_seconds"], f"{prefix}.max_sample_gap_seconds", positive=True,
)
mode = entry.get("source_mode", "direct_mass_rate")
if not isinstance(mode, str) or mode not in {"direct_mass_rate", "area_application"}:
raise CalculationConfigError("Unsupported material source_mode")
if area:
refs = entry.get("application_signal_refs")
if (not isinstance(refs, list) or not refs
or any(not isinstance(ref, str) or not ref.strip() for ref in refs)
or len(set(refs)) != len(refs)):
raise CalculationConfigError(
"application_signal_refs must be unique signal strings"
)
if entry["rate_signal_ref"] != "unused":
raise CalculationConfigError("area_application cannot specify rate_signal_ref")
for name in ("erp_workplace", "production_order_format"):
if not isinstance(entry.get(name), str) or not entry[name].strip():
raise CalculationConfigError(f"area_application requires {name}")
from production_analytics.context import build_enlyze_production_order
try:
build_enlyze_production_order("0", entry["production_order_format"])
except ValueError as exc:
raise CalculationConfigError("Invalid production_order_format") from exc
entry["application_signal_refs"] = tuple(refs)
elif any(name in entry for name in (
"application_signal_refs", "erp_workplace", "production_order_format",
)):
raise CalculationConfigError("Area context fields require area_application")
if entry["id"] in instances:
raise CalculationConfigError("Calculation ids must be unique")
instances[entry["id"]] = MaterialCalculationConfig(**entry)
@@ -102,3 +102,20 @@ class MaterialConsumptionIntegrator:
for sample in samples:
self.process(sample)
return self.state
def area_application_rate_kg_per_hour(
applications_g_m2: Iterable[float], nominal_width_m: float, line_speed_m_min: float,
) -> float:
"""Normalize summed area application to the integrator's canonical kg/h."""
applications = tuple(applications_g_m2)
if not applications or not all(isfinite(value) for value in applications):
raise ValueError("application values must be non-empty and finite")
if not isfinite(nominal_width_m) or nominal_width_m <= 0:
raise ValueError("nominal width must be finite and positive")
if not isfinite(line_speed_m_min):
raise ValueError("line speed must be finite")
rate = sum(applications) * nominal_width_m * line_speed_m_min * 60 / 1000
if not isfinite(rate):
raise ValueError("derived material rate must be finite")
return rate
+2
View File
@@ -112,6 +112,7 @@ def _parser() -> argparse.ArgumentParser:
material.add_argument("--poll-interval-seconds", type=float, default=10.0)
material.add_argument("--state-directory", type=Path, default=Path("data/state/material"))
material.add_argument("--secrets-file", type=Path, default=Path("secrets/enlyze.env"))
material.add_argument("--erp-secrets-file", type=Path, default=Path("secrets/erp.env"))
efficiency = runners.add_parser("material-efficiency", help="Persist ERP feedback KPI points")
efficiency.add_argument("--config", type=Path, required=True)
efficiency.add_argument("--secrets-file", type=Path, default=Path("secrets/enlyze.env"))
@@ -143,6 +144,7 @@ def _run_material(args: argparse.Namespace) -> int:
calculation,
poll_interval_seconds=args.poll_interval_seconds,
state_directory=args.state_directory,
erp_secrets_file=args.erp_secrets_file,
secrets_file=args.secrets_file,
)
runner.run()
+20 -9
View File
@@ -4,7 +4,10 @@ from dataclasses import dataclass
from datetime import UTC, datetime
from math import isfinite
from production_analytics.calculations.material_consumption import MaterialSample
from production_analytics.calculations.material_consumption import (
MaterialSample,
area_application_rate_kg_per_hour,
)
from production_analytics.enlyze.exploration import ExplorationClient
@@ -89,20 +92,21 @@ class EnlyzeApiGateway:
gate_variable_id: str,
start: datetime,
end: datetime,
application_variable_ids: tuple[str, ...] = (),
nominal_width_m: float | None = None,
) -> list[MaterialSample]:
for name, value in (("start", start), ("end", end)):
if value.tzinfo is None or value.utcoffset() is None:
raise ValueError(f"{name} must be timezone-aware")
if end < start:
raise ValueError("end must not precede start")
source_ids = application_variable_ids or (rate_variable_id,)
variable_ids = list(dict.fromkeys((*source_ids, gate_variable_id)))
request_body = {
"machine": machine_id,
"start": start.astimezone(UTC).isoformat(),
"end": end.astimezone(UTC).isoformat(),
"variables": [
{"uuid": rate_variable_id},
{"uuid": gate_variable_id},
],
"variables": [{"uuid": variable_id} for variable_id in variable_ids],
}
samples = []
followed_cursors: set[str] = set()
@@ -114,11 +118,11 @@ class EnlyzeApiGateway:
columns = data["columns"]
if not isinstance(columns, list):
raise ValueError("columns must be a list")
for required in ("time", rate_variable_id, gate_variable_id):
for required in ("time", *variable_ids):
if columns.count(required) != 1:
raise ValueError(f"required column {required!r} must occur exactly once")
time_index = columns.index("time")
rate_index = columns.index(rate_variable_id)
source_indices = [columns.index(ref) for ref in source_ids]
gate_index = columns.index(gate_variable_id)
if not isinstance(data["records"], list):
raise ValueError("records must be a list")
@@ -126,12 +130,19 @@ class EnlyzeApiGateway:
try:
if not isinstance(record, list) or len(record) != len(columns):
raise ValueError("record must match columns")
rate, gate = record[rate_index], record[gate_index]
for value in (rate, gate):
sources = [record[i] for i in source_indices]
gate = record[gate_index]
for value in (*sources, gate):
if isinstance(value, bool) or not isinstance(value, (int, float)):
raise ValueError("rate and gate must be numeric")
if not isfinite(value):
raise ValueError("rate and gate must be finite")
if application_variable_ids:
if nominal_width_m is None:
raise ValueError("area application requires nominal width")
rate = area_application_rate_kg_per_hour(sources, nominal_width_m, gate)
else:
rate = sources[0]
samples.append(MaterialSample(
timestamp=_parse_timestamp(record[time_index]),
material_rate_kg_per_hour=float(rate), gate_value=float(gate),
@@ -0,0 +1,24 @@
"""ERP nominal-width resolution for an explicitly matched production order."""
from production_analytics.context import build_enlyze_production_order, extract_nominal_width_m
from production_analytics.erp import ErpWorkplaceStatusGateway
class ErpNominalWidthProvider:
def __init__(
self, gateway: ErpWorkplaceStatusGateway, *, workplace: str, format_template: str,
) -> None:
self.gateway = gateway
self.workplace = workplace
self.format_template = format_template
def __call__(self, production_order: str) -> float:
status = self.gateway.get_current_workplace_status(self.workplace)
if (status is None or status.workplace != self.workplace
or build_enlyze_production_order(status.production_order, self.format_template)
!= production_order):
raise ValueError("ERP context does not match the production order")
width = extract_nominal_width_m(status.article_description)
if width is None:
raise ValueError("ERP article description has no unambiguous nominal width")
return width
@@ -1,5 +1,6 @@
"""Single-cycle live material polling orchestration."""
from collections.abc import Callable
from dataclasses import dataclass
from datetime import datetime
from typing import Protocol
@@ -53,7 +54,13 @@ class MaterialPollingService:
gate_variable_id: str,
gate_threshold: float,
max_sample_gap_seconds: float,
application_variable_ids: tuple[str, ...] = (),
nominal_width_provider: Callable[[str], float] | None = None,
) -> None:
self._application_variable_ids = application_variable_ids
self._nominal_width_provider = nominal_width_provider
if application_variable_ids and nominal_width_provider is None:
raise ValueError("area application requires a nominal width provider")
self._gateway = gateway
self._state_store = state_store
self._machine_id = machine_id
@@ -71,6 +78,8 @@ class MaterialPollingService:
if run is None or now < run.start:
return None
width = (self._nominal_width_provider(run.production_order)
if self._application_variable_ids else None)
saved_polling_state = self._state_store.load(
self._machine_id,
run.production_order,
@@ -91,7 +100,7 @@ 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,
initial_state, start=previous.start, end=previous.end, width=width,
)
initial_state = MaterialIntegrationState(
cumulative_consumption_kg=integrated.cumulative_consumption_kg,
@@ -116,7 +125,7 @@ class MaterialPollingService:
if now < start:
return None
state = self._integrate(initial_state, start=start, end=now)
state = self._integrate(initial_state, start=start, end=now, width=width)
self._state_store.save(
self._machine_id,
@@ -134,13 +143,20 @@ class MaterialPollingService:
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,
)
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)
@@ -14,6 +14,8 @@ from production_analytics.enlyze.exploration import (
load_secret_file,
)
from production_analytics.enlyze.gateway import EnlyzeApiGateway
from production_analytics.erp import ErpSettings, ErpWorkplaceStatusGateway
from production_analytics.service.material_context import ErpNominalWidthProvider
from production_analytics.service.material_polling import (
MaterialPollingService,
MaterialPollingState,
@@ -48,6 +50,7 @@ class _ReportingStateStore:
def build_material_runner(
calculation: MaterialCalculationConfig, *, poll_interval_seconds: float,
state_directory: Path, secrets_file: Path,
erp_secrets_file: Path = Path("secrets/erp.env"),
) -> MaterialPollingRunner:
finite_number(poll_interval_seconds, "poll interval", positive=True)
environment = dict(os.environ)
@@ -69,6 +72,16 @@ def build_material_runner(
f"State directory startup check failed: {type(exc).__name__}"
) from exc
postgres_settings = PostgresSettings.from_environment(environment)
source_options = {}
if calculation.source_mode == "area_application":
source_options = dict(
application_variable_ids=calculation.application_signal_refs,
nominal_width_provider=ErpNominalWidthProvider(
ErpWorkplaceStatusGateway(ErpSettings.from_secret_file(erp_secrets_file)),
workplace=calculation.erp_workplace,
format_template=calculation.production_order_format,
),
)
service = MaterialPollingService(
gateway=EnlyzeApiGateway(ExplorationClient(settings)),
state_store=_ReportingStateStore(JsonMaterialStateStore(state_directory)),
@@ -77,6 +90,7 @@ def build_material_runner(
gate_variable_id=calculation.gate_signal_ref,
gate_threshold=calculation.gate_threshold,
max_sample_gap_seconds=calculation.max_sample_gap_seconds,
**source_options,
)
return MaterialPollingRunner(
service, calculation_id=calculation.id,