From be8d76507af5947f3546e7b4d2ea91f9f99278a8 Mon Sep 17 00:00:00 2001 From: Martin Tazl Date: Sun, 6 Sep 2026 13:29:06 +0200 Subject: [PATCH] Add Bento 1 fresh bentonite consumption --- README.md | 2 + config/bento1-material-consumption.yaml | 19 ++ docs/bento1-fresh-bentonite.md | 60 +++++ .../calculations/config.py | 41 +++- .../calculations/material_consumption.py | 17 ++ src/production_analytics/cli/__main__.py | 2 + src/production_analytics/enlyze/gateway.py | 29 ++- .../service/material_context.py | 24 ++ .../service/material_polling.py | 20 +- .../service/material_runtime.py | 14 ++ tests/test_area_material_consumption.py | 216 ++++++++++++++++++ 11 files changed, 430 insertions(+), 14 deletions(-) create mode 100644 config/bento1-material-consumption.yaml create mode 100644 docs/bento1-fresh-bentonite.md create mode 100644 src/production_analytics/service/material_context.py create mode 100644 tests/test_area_material_consumption.py diff --git a/README.md b/README.md index a9aa2c5..17f6eb2 100644 --- a/README.md +++ b/README.md @@ -464,3 +464,5 @@ The current-status ERP source cannot backfill feedback events missed between pol or during outages, nor guarantee observation of final FA feedback before the order changes. Daily 24h reports and final order summary tables remain future work. No dashboard, aggregation or lag correction is added here. + +Bento 1: [fresh bentonite consumption configuration and scope](docs/bento1-fresh-bentonite.md). diff --git a/config/bento1-material-consumption.yaml b/config/bento1-material-consumption.yaml new file mode 100644 index 0000000..e3e0db3 --- /dev/null +++ b/config/bento1-material-consumption.yaml @@ -0,0 +1,19 @@ +# Fresh bentonite only: the recycled/recovered third spreader is excluded. +calculations: + - id: bento1-fresh-bentonite-consumption + type: material_consumption + version: "1" + machine_ref: 5f42a4f6-9ca0-4f6f-9786-40d50a35b230 + source_mode: area_application + application_signal_refs: + - 19ea65d2-bd35-4d02-89de-4495583d9026 + - d9f47615-69cf-4f97-88df-199585764491 + # Transport Auszug Geschwindigkeit Istwert, m/min (also the production gate). + gate_signal_ref: fef41976-1103-4090-b780-eaeecc02fdfa + gate_threshold: 0.3 + max_sample_gap_seconds: 20.0 + erp_workplace: Bento 1 + production_order_format: "Bento 1-{production_order}" + output_metric: material_consumption + output_unit: kg + group_by: production_order diff --git a/docs/bento1-fresh-bentonite.md b/docs/bento1-fresh-bentonite.md new file mode 100644 index 0000000..4452f67 --- /dev/null +++ b/docs/bento1-fresh-bentonite.md @@ -0,0 +1,60 @@ +# Bento 1 fresh bentonite consumption + +`config/bento1-material-consumption.yaml` defines +`bento1-fresh-bentonite-consumption`. This KPI estimates **fresh bentonite +consumption** from the two fresh-spreader setpoints in g/m². The third spreader +uses recycled/recovered bentonite and is deliberately excluded; it is neither +missing data nor estimated. This is not total bentonite deposited on the product, +and setpoints are not a measurement of actual mass flow. + +The generic `area_application` source sums the configured application signals +and converts them to kg/h using `sum(g/m²) × nominal_width_m × speed_m/min × 60 / 1000`. +The existing previous-value integrator then accumulates kilograms only while +speed is strictly greater than 0.3 m/min. `direct_mass_rate` remains the default +for existing configurations, including K7. The persistent state schema and +production-order bootstrap/resume behavior are unchanged; each disjoint run +starts with a new sample baseline while retaining order totals. + +Width comes from the existing ERP workplace status adapter and generic nominal +width parser. The ERP order must format to the ENLYZE order using +`Bento 1-{production_order}`. Missing/mismatched ERP context or an unparseable +width rejects the polling cycle without saving state. ENLYZE product metadata is +not used for width. The current ERP context supplies width for all runs of that +same order; this assumes the order's article width stays constant. Historical +orders no longer present in current ERP workplace status require a separate +historical context source for offline replay. + +Run with: + +```sh +production-analytics run material-poll \ + --config config/bento1-material-consumption.yaml \ + --erp-secrets-file secrets/erp.env +``` + +`erp_workplace: Bento 1` assumes that exact ERP workplace identifier; verify it +against the deployment's workplace view. The 20-second maximum sample gap is a +generic policy shared with K7, intended for a nominal 10-second sampling cadence. +One missing sample is tolerated: the resulting 20-second interval is integrated +using the preceding rate and gate. Gaps greater than 20 seconds are not +integrated, so prolonged or repeated gaps can undercount consumption. +Missing/nonfinite values from either configured fresh spreader reject the cycle. +Intervals exceeding the maximum gap and unobserved run tails are not inferred. + +Before implementation, the calculation was independently validated manually +against the same three real ENLYZE signals over the two runs below, yielding +approximately **147021.9 kg** of fresh bentonite for `Bento 1-12026000814`. +The post-implementation live replay could not be repeated because ENLYZE +connectivity was temporarily unavailable (DNS resolution failure), not because +of a known implementation issue. Exact production-path replay remains a +follow-up verification rather than a blocker for this milestone. + +The regression uses article 180305, description `Bfix NSP 5300, 5,00 x 40 m` +(width 5.00 m), and the supplied boundaries of runs +`b675818c-4636-4976-be6c-3e0e24da9e8a` and +`25d67424-8346-42d4-b844-75b4097382bd` for order `Bento 1-12026000814`. +It constructs synthetic samples from the supplied aggregates: 1256.0 active +minutes, 31294.6 m² and 4698.0 g/m², yielding approximately 147021.9 kg. +It verifies bootstrap across both runs and persisted resume without counting the +inter-run gap. It is an aggregate plausibility regression, not an independent +replay of recorded signals; raw reference samples were not supplied. diff --git a/src/production_analytics/calculations/config.py b/src/production_analytics/calculations/config.py index e209f89..f637361 100644 --- a/src/production_analytics/calculations/config.py +++ b/src/production_analytics/calculations/config.py @@ -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) diff --git a/src/production_analytics/calculations/material_consumption.py b/src/production_analytics/calculations/material_consumption.py index a9ddf97..5412618 100644 --- a/src/production_analytics/calculations/material_consumption.py +++ b/src/production_analytics/calculations/material_consumption.py @@ -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 diff --git a/src/production_analytics/cli/__main__.py b/src/production_analytics/cli/__main__.py index 4ce2aa9..8dd5bd5 100644 --- a/src/production_analytics/cli/__main__.py +++ b/src/production_analytics/cli/__main__.py @@ -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() diff --git a/src/production_analytics/enlyze/gateway.py b/src/production_analytics/enlyze/gateway.py index 08e42eb..2d5a764 100644 --- a/src/production_analytics/enlyze/gateway.py +++ b/src/production_analytics/enlyze/gateway.py @@ -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), diff --git a/src/production_analytics/service/material_context.py b/src/production_analytics/service/material_context.py new file mode 100644 index 0000000..388a765 --- /dev/null +++ b/src/production_analytics/service/material_context.py @@ -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 diff --git a/src/production_analytics/service/material_polling.py b/src/production_analytics/service/material_polling.py index a2e1ae5..fc9e6e7 100644 --- a/src/production_analytics/service/material_polling.py +++ b/src/production_analytics/service/material_polling.py @@ -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) diff --git a/src/production_analytics/service/material_runtime.py b/src/production_analytics/service/material_runtime.py index 1fbd61a..e11d742 100644 --- a/src/production_analytics/service/material_runtime.py +++ b/src/production_analytics/service/material_runtime.py @@ -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, diff --git a/tests/test_area_material_consumption.py b/tests/test_area_material_consumption.py new file mode 100644 index 0000000..e3dc515 --- /dev/null +++ b/tests/test_area_material_consumption.py @@ -0,0 +1,216 @@ +from dataclasses import replace +from datetime import datetime, timedelta +from types import SimpleNamespace +from unittest.mock import Mock + +import pytest +import yaml + +from production_analytics.calculations.config import ( + CalculationConfigError, + load_material_calculation, +) +from production_analytics.calculations.material_consumption import ( + MaterialConsumptionIntegrator, + MaterialIntegratorConfig, + area_application_rate_kg_per_hour, +) +from production_analytics.enlyze.gateway import EnlyzeApiGateway, EnlyzeProductionRun +from production_analytics.service.material_context import ErpNominalWidthProvider +from production_analytics.service.material_polling import MaterialPollingService +from production_analytics.service.material_state_store import JsonMaterialStateStore + + +def timestamp(value): + return datetime.fromisoformat(value.replace('Z', '+00:00')) + + +def test_area_units_and_strict_gate(): + client = Mock() + start = timestamp('2026-08-25T05:27:18Z') + client.post_json.return_value.body = {'data': { + 'columns': ['speed', 's2', 'time', 's1'], + 'records': [[speed, 2000, (start + timedelta(minutes=i)).isoformat(), 3000] + for i, speed in enumerate([2, 0.3, 0, 2, 2])], + }} + samples = EnlyzeApiGateway(client).get_material_samples( + machine_id='machine', rate_variable_id='unused', gate_variable_id='speed', + application_variable_ids=('s1', 's2'), nominal_width_m=5, + start=start, end=start + timedelta(minutes=4), + ) + state = MaterialConsumptionIntegrator(MaterialIntegratorConfig(0.3, 60)).process_many(samples) + assert state.cumulative_consumption_kg == 100 + assert state.integrated_running_seconds == 120 + assert client.post_json.call_args.args[1]['variables'] == [ + {'uuid': 's1'}, {'uuid': 's2'}, {'uuid': 'speed'}, + ] + + +@pytest.mark.parametrize('sources,width,speed', [ + ([], 5, 2), ([float('nan')], 5, 2), ([1], 0, 2), ([1], 5, float('inf')), +]) +def test_invalid_area_inputs(sources, width, speed): + with pytest.raises(ValueError): + area_application_rate_kg_per_hour(sources, width, speed) + + +@pytest.mark.parametrize('bad', [None, True, float('nan')]) +def test_missing_or_invalid_spreader_fails(bad): + client = Mock() + now = timestamp('2026-08-25T05:27:18Z') + client.post_json.return_value.body = {'data': { + 'columns': ['time', 's1', 's2', 'speed'], + 'records': [[now.isoformat(), 3000, bad, 2]], + }} + with pytest.raises(ValueError): + EnlyzeApiGateway(client).get_material_samples( + machine_id='machine', rate_variable_id='unused', gate_variable_id='speed', + application_variable_ids=('s1', 's2'), nominal_width_m=5, start=now, end=now, + ) + + +def width_provider(): + erp = Mock() + erp.get_current_workplace_status.return_value = SimpleNamespace( + workplace='Bento 1', production_order='12026000814', article_number='180305', + article_description='Bfix NSP 5300, 5,00 x 40 m', + ) + return ErpNominalWidthProvider( + erp, workplace='Bento 1', format_template='Bento 1-{production_order}', + ) + + +@pytest.mark.parametrize('field,value', [ + ('production_order', '12026000815'), ('workplace', 'K7'), + ('article_description', None), ('article_description', '180305'), +]) +def test_width_context_fails_closed(field, value): + provider = width_provider() + setattr(provider.gateway.get_current_workplace_status.return_value, field, value) + with pytest.raises(ValueError): + provider('Bento 1-12026000814') + + +def test_bento_reference_aggregate_bootstrap_and_persistent_resume(tmp_path): + """Synthetic aggregate fixture, not recorded ENLYZE samples. + + Uses reported active time/area/application within the actual run boundaries. + """ + config = load_material_calculation('config/bento1-material-consumption.yaml') + order = 'Bento 1-12026000814' + runs = [EnlyzeProductionRun( + uuid, config.machine_ref, 'external-product', order, timestamp(start), timestamp(end), + ) for uuid, start, end in [ + ('b675818c-4636-4976-be6c-3e0e24da9e8a', + '2026-08-25T05:27:18Z', '2026-08-26T04:40:58Z'), + ('25d67424-8346-42d4-b844-75b4097382bd', + '2026-08-26T05:43:39Z', '2026-08-26T05:54:05Z'), + ]] + speed = 31294.6 / (5 * 1256) + client = Mock() + + def response(path, body): + start, end = timestamp(body['start']), timestamp(body['end']) + run = next(run for run in runs if run.start <= start <= run.end) + active_seconds = 1256 * 60 - 626 if run == runs[0] else 626 + stop = run.start + timedelta(seconds=active_seconds) + times = {start, end} + t = start + while t < end: + times.add(t) + t += timedelta(seconds=10) + if start <= stop <= end: + times.add(stop) + return SimpleNamespace(body={'data': { + 'columns': ['time', *config.application_signal_refs, config.gate_signal_ref], + 'records': [[t.isoformat(), 2300, 2398, speed if t < stop else 0] + for t in sorted(times)], + }}) + + client.post_json.side_effect = response + gateway = EnlyzeApiGateway(client) + gateway.get_open_production_run = Mock(return_value=replace(runs[1], end=None)) + gateway.get_production_runs = Mock(return_value=runs) + store = JsonMaterialStateStore(tmp_path) + + def service(): + return MaterialPollingService( + gateway=gateway, state_store=store, machine_id=config.machine_ref, + rate_variable_id=config.rate_signal_ref, gate_variable_id=config.gate_signal_ref, + gate_threshold=config.gate_threshold, + max_sample_gap_seconds=config.max_sample_gap_seconds, + application_variable_ids=config.application_signal_refs, + nominal_width_provider=width_provider(), + ) + + result = service().poll_once(now=runs[1].end) + assert result.state.integrated_running_seconds / 60 == pytest.approx(1256) + assert result.state.cumulative_consumption_kg == pytest.approx(147021.9, abs=0.2) + assert result.state.integrated_running_seconds / 60 * speed * 5 == pytest.approx(31294.6) + assert service().poll_once(now=runs[1].end).state == result.state + assert gateway.get_production_runs.call_count == 1 + assert [(timestamp(call.args[1]['start']), timestamp(call.args[1]['end'])) + for call in client.post_json.call_args_list[:2]] == [ + (run.start, run.end) for run in runs] + + +def test_k7_config_defaults_remain_direct(): + config = load_material_calculation('config/k7-material-consumption.yaml') + assert config.source_mode == 'direct_mass_rate' + assert config.application_signal_refs == () + assert config.gate_threshold == 0.5 + assert config.rate_signal_ref == 'c9d06af5-f6d6-4ede-b6c4-5a98bac77129' + + +@pytest.mark.parametrize('field,value', [ + ('source_mode', 'unknown'), ('application_signal_refs', []), + ('application_signal_refs', ['same', 'same']), ('erp_workplace', ''), + ('production_order_format', '{wrong}'), ('rate_signal_ref', 'direct'), +]) +def test_invalid_area_configuration(tmp_path, field, value): + with open('config/bento1-material-consumption.yaml') as stream: + document = yaml.safe_load(stream) + document['calculations'][0][field] = value + path = tmp_path / 'config.yaml' + path.write_text(yaml.safe_dump(document)) + with pytest.raises(CalculationConfigError): + load_material_calculation(path) + + +def test_area_context_failure_does_not_touch_persistent_state(tmp_path): + config = load_material_calculation('config/bento1-material-consumption.yaml') + gateway = Mock() + now = timestamp('2026-08-25T05:27:18Z') + gateway.get_open_production_run.return_value = EnlyzeProductionRun( + 'run', config.machine_ref, None, 'Bento 1-12026000814', now, None, + ) + provider = width_provider() + provider.gateway.get_current_workplace_status.return_value = None + store = Mock() + service = MaterialPollingService( + gateway=gateway, state_store=store, machine_id=config.machine_ref, + rate_variable_id=config.rate_signal_ref, gate_variable_id=config.gate_signal_ref, + gate_threshold=config.gate_threshold, max_sample_gap_seconds=20, + application_variable_ids=config.application_signal_refs, nominal_width_provider=provider, + ) + with pytest.raises(ValueError, match='context'): + service.poll_once(now=now) + gateway.get_material_samples.assert_not_called() + store.save.assert_not_called() + + +def test_area_gateway_pagination(): + client = Mock() + now = timestamp('2026-08-25T05:27:18Z') + client.post_json.side_effect = [SimpleNamespace(body={ + 'metadata': {'next_cursor': cursor}, + 'data': {'columns': ['time', 's1', 's2', 'speed'], + 'records': [[(now + timedelta(seconds=i * 10)).isoformat(), 3000, 2000, 2]]}, + }) for i, cursor in enumerate(['next', None])] + samples = EnlyzeApiGateway(client).get_material_samples( + machine_id='machine', rate_variable_id='unused', gate_variable_id='speed', + application_variable_ids=('s1', 's2'), nominal_width_m=5, + start=now, end=now + timedelta(seconds=10), + ) + assert [sample.material_rate_kg_per_hour for sample in samples] == [3000, 3000] + assert client.post_json.call_args.args[1]['cursor'] == 'next'