From 0210ced9e3f81125d08581861b85c5ff70ed88f1 Mon Sep 17 00:00:00 2001 From: Martin Tazl Date: Mon, 7 Sep 2026 07:59:48 +0200 Subject: [PATCH] Use rotational discharge for Bento consumption --- config/bento1-material-consumption.yaml | 11 +- docs/bento1-fresh-bentonite.md | 176 +++++++++++++----- .../calculations/config.py | 44 +++-- .../calculations/material_consumption.py | 18 ++ src/production_analytics/enlyze/gateway.py | 28 ++- .../service/material_context.py | 3 +- .../service/material_polling.py | 22 ++- .../service/material_runtime.py | 8 +- tests/test_area_material_consumption.py | 69 +++++-- tests/test_rotational_material_consumption.py | 113 +++++++++++ 10 files changed, 400 insertions(+), 92 deletions(-) create mode 100644 tests/test_rotational_material_consumption.py diff --git a/config/bento1-material-consumption.yaml b/config/bento1-material-consumption.yaml index e3e0db3..be1581d 100644 --- a/config/bento1-material-consumption.yaml +++ b/config/bento1-material-consumption.yaml @@ -4,11 +4,12 @@ calculations: 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). + source_mode: rotational_discharge + rotational_speed_signal_refs: + - 6e5d2d94-98f9-4cc1-8a88-7987c6282525 + - cd7385c4-337b-4759-ab32-45d65beaf190 + specific_discharge_kg_per_rev_m: 3.12 + # Transport Auszug Geschwindigkeit Istwert, m/min (production gate only). gate_signal_ref: fef41976-1103-4090-b780-eaeecc02fdfa gate_threshold: 0.3 max_sample_gap_seconds: 20.0 diff --git a/docs/bento1-fresh-bentonite.md b/docs/bento1-fresh-bentonite.md index 4452f67..7a1a013 100644 --- a/docs/bento1-fresh-bentonite.md +++ b/docs/bento1-fresh-bentonite.md @@ -1,60 +1,138 @@ # 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. +`bento1-fresh-bentonite-consumption`, an estimate of **fresh bentonite consumption** +from actual spreader roll speeds. Both needle rolls apply sequentially across the +full web width. The recycled/recovered third spreader is excluded. -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. +The generic `rotational_discharge` source converts one or more rpm signals to the +existing integrator's kg/h: `sum(rpm) × nominal_width_m × specific_discharge_kg_per_rev_m × 60`. +Bento config sets the factor to **3.12 kg/(rev·m)** and uses: -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. +- Right ACT: `6e5d2d94-98f9-4cc1-8a88-7987c6282525` +- Left ACT: `cd7385c4-337b-4759-ab32-45d65beaf190` -Run with: +ENLYZE now returns physical rpm directly (scaling factor 1.0); no division by +1000 is applied. For 3.739 + 3.956 rpm and 5 m width, the rate is 120.042 kg/min +(7202.52 kg/h), giving 20.007 kg in ten seconds. -```sh -production-analytics run material-poll \ - --config config/bento1-material-consumption.yaml \ - --erp-secrets-file secrets/erp.env +Transport speed `fef41976-1103-4090-b780-eaeecc02fdfa` is exclusively the +production gate: speed must be strictly greater than 0.3 m/min. It does not +multiply the mass rate. The generic rotational mode also supports omitting +`gate_signal_ref` and `gate_threshold`, in which case all valid intervals are +active. Bento retains its gate. + +The two old SET sources `19ea65d2-bd35-4d02-89de-4495583d9026` and +`d9f47615-69cf-4f97-88df-199585764491` are no longer requested by Bento. +`direct_mass_rate` remains the default for K7; `area_application` remains +available for other configurations. There is no Bento branch in the integrator. + +The previous-value integration, JSON state format and order bootstrap/resume +remain unchanged. Totals accumulate across disjoint runs, with a new baseline +at each run boundary, even if the inter-run gap is under 20 seconds. Intervals +up to 20 seconds use the preceding rate and gate; longer sample gaps and +unobserved run tails are not integrated. Missing/nonfinite source values reject +the cycle. Finite negative speeds retain the existing generic signed-rate +semantics; no clamping is introduced. + +Width comes from the existing ERP nominal-width parser, with workplace +`Bento 1` and order format `Bento 1-{production_order}`. Workplace comparisons +use `strip().casefold()`; order matching remains strict. Missing/mismatched ERP +context or an unparseable width rejects the cycle without saving state. The +current ERP context supplies width for all runs of the same order, assuming +constant article width. Historical orders no longer in the ERP workplace view +need a historical width source for replay. + +The configured factor is based on the supplied independent historical checks: +3.1175 at 5.00 m (article 180305), 3.1193 and 3.1187 at 4.85 m (article 173000), +and median 3.1208 at 5.00 m (article 180005; p05–p95 3.0703–3.1558). +Tests use realistic rpm values and synthetic samples, including disjoint-run +bootstrap and persisted resume. They are not a new replay of recorded ENLYZE data. + +## Old calculation cleanup: review before execution + +No production state or database rows were changed during this implementation, +and no service start was requested. The operator reports Bento stopped; +sandbox systemd access could not independently confirm that state. +The deployed `/etc/systemd/system/production-analytics-bento1.service` specifies +`/opt/git-projects/production-analytics/data/state/material` as its state directory. + +The existing state file for machine `5f42a4f6-9ca0-4f6f-9786-40d50a35b230` +and order `Bento 1-12026000857` is exactly: + +```text +/opt/git-projects/production-analytics/data/state/material/4d6cc209e35bc58887e5b4a3816835006425429e810e0d957f651971723e8c2f.json ``` -`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. +This identity was verified using the store's SHA-256 of the JSON machine/order +pair. The file exists but its contents are not readable by the sandbox user. +The other state file, `6d11088159b8e4fd21543560a54a31a0f86faf22bbf8f16873b228d1903c23f9.json`, +has not been attributed here and must be left untouched. -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. +Old PostgreSQL snapshots are in `material_consumption_snapshots`, scoped by both +`calculation_id = 'bento1-fresh-bentonite-consumption'` and +`machine_id = '5f42a4f6-9ca0-4f6f-9786-40d50a35b230'`. +The table has no calculation-source/version provenance column. Given the stopped +service and no rpm deployment yet, existing rows in this scope belong to the old +calculation. PostgreSQL was unreachable from the sandbox, so exact row counts, +order/run membership and timestamp bounds remain unverified. Do not mistake +this predicate for a completed row inventory. Review the following output first. -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. +Run these **read-only inventory commands** on the deployment host: + +```sh +cd /opt/git-projects/production-analytics +sudo systemctl is-active production-analytics-bento1.service +sudo cat data/state/material/4d6cc209e35bc58887e5b4a3816835006425429e810e0d957f651971723e8c2f.json +sudo docker compose exec -T timescaledb psql -U production_analytics -d production_analytics -v ON_ERROR_STOP=1 <<'SQL' +BEGIN READ ONLY; +SELECT production_order, run_id, count(*) AS rows, + min(timestamp) AS first_snapshot, max(timestamp) AS last_snapshot, + min(consumption_kg), max(consumption_kg) +FROM material_consumption_snapshots +WHERE calculation_id = 'bento1-fresh-bentonite-consumption' + AND machine_id = '5f42a4f6-9ca0-4f6f-9786-40d50a35b230' +GROUP BY production_order, run_id ORDER BY first_snapshot; +SELECT * FROM material_consumption_snapshots +WHERE calculation_id = 'bento1-fresh-bentonite-consumption' + AND machine_id = '5f42a4f6-9ca0-4f6f-9786-40d50a35b230' +ORDER BY production_order, timestamp; +COMMIT; +SQL +``` + +After reviewing the inventory, and **only with cleanup authorization**, keep the +service stopped and execute this backup-and-delete transaction. The backup table +intentionally has no `IF NOT EXISTS`: a repeated invocation fails rather than +reusing an older backup. Only rows with the backed-up primary keys are deleted. +K7 rows are outside the predicate. + +```sh +cd /opt/git-projects/production-analytics +sudo docker compose exec -T timescaledb psql -U production_analytics -d production_analytics -v ON_ERROR_STOP=1 <<'SQL' +BEGIN; +CREATE TABLE bento1_material_snapshots_before_rpm_20260907 AS +SELECT * FROM material_consumption_snapshots +WHERE calculation_id = 'bento1-fresh-bentonite-consumption' + AND machine_id = '5f42a4f6-9ca0-4f6f-9786-40d50a35b230'; +DELETE FROM material_consumption_snapshots AS s +USING bento1_material_snapshots_before_rpm_20260907 AS old +WHERE s.calculation_id = old.calculation_id AND s.machine_id = old.machine_id + AND s.production_order = old.production_order AND s.timestamp = old.timestamp; +COMMIT; +SQL +sudo mv -n -- \ + data/state/material/4d6cc209e35bc58887e5b4a3816835006425429e810e0d957f651971723e8c2f.json \ + data/state/material/4d6cc209e35bc58887e5b4a3816835006425429e810e0d957f651971723e8c2f.json.before-rpm-20260907 +``` + +Verify the original JSON path is absent before restarting; `mv -n` preserves an +existing backup and will not overwrite it. If inventory reveals other Bento +orders, derive their paths with `JsonMaterialStateStore._path(machine, order)` +and review any existing files before moving them too. Do not clear the shared +directory. Complete both state and snapshot cleanup before starting the new +code: the state schema deliberately does not detect a changed source model. + +A later start bootstraps only the currently open order using its current ERP +width. It does not automatically rebuild snapshots for closed historical +orders. Reconstructing those requires a separately reviewed historical replay. diff --git a/src/production_analytics/calculations/config.py b/src/production_analytics/calculations/config.py index f637361..87d32b7 100644 --- a/src/production_analytics/calculations/config.py +++ b/src/production_analytics/calculations/config.py @@ -48,7 +48,7 @@ class MaterialCalculationConfig: version: str machine_ref: str rate_signal_ref: str - gate_signal_ref: str + gate_signal_ref: str | None gate_threshold: float max_sample_gap_seconds: float output_metric: str @@ -56,6 +56,8 @@ class MaterialCalculationConfig: group_by: str source_mode: str = "direct_mass_rate" application_signal_refs: tuple[str, ...] = () + rotational_speed_signal_refs: tuple[str, ...] = () + specific_discharge_kg_per_rev_m: float | None = None erp_workplace: str = "" production_order_format: str = "" @@ -86,7 +88,11 @@ def load_material_calculation( 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: + rotational = entry.get("source_mode") == "rotational_discharge" + if rotational: + entry.setdefault("gate_signal_ref", None) + entry.setdefault("gate_threshold", 0.0) + if area or rotational: entry.setdefault("rate_signal_ref", "unused") missing = required - entry.keys() if missing: @@ -94,6 +100,8 @@ def load_material_calculation( if entry.keys() - allowed: raise CalculationConfigError(f"{prefix}: unknown fields (check spelling)") for name in sorted(required - {"gate_threshold", "max_sample_gap_seconds"}): + if name == "gate_signal_ref" and rotational and entry[name] is None: + continue if not isinstance(entry[name], str) or not entry[name].strip(): raise CalculationConfigError(f"{prefix}.{name} must be a non-empty string") for name, expected in ( @@ -108,31 +116,45 @@ def load_material_calculation( 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"}: + if not isinstance(mode, str) or mode not in { + "direct_mass_rate", "area_application", "rotational_discharge", + }: raise CalculationConfigError("Unsupported material source_mode") - if area: - refs = entry.get("application_signal_refs") + if area or rotational: + refs_name = "rotational_speed_signal_refs" if rotational else "application_signal_refs" + other = "application_signal_refs" if rotational else "rotational_speed_signal_refs" + if other in entry: + raise CalculationConfigError(f"{mode} cannot specify {other}") + refs = entry.get(refs_name) 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" + f"{refs_name} must be unique signal strings" ) if entry["rate_signal_ref"] != "unused": - raise CalculationConfigError("area_application cannot specify rate_signal_ref") + raise CalculationConfigError(f"{mode} 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}") + raise CalculationConfigError(f"{mode} 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) + entry[refs_name] = tuple(refs) elif any(name in entry for name in ( - "application_signal_refs", "erp_workplace", "production_order_format", + "application_signal_refs", "rotational_speed_signal_refs", + "erp_workplace", "production_order_format", )): - raise CalculationConfigError("Area context fields require area_application") + raise CalculationConfigError("Width context fields require a width-based source mode") + if rotational: + entry["specific_discharge_kg_per_rev_m"] = finite_number( + entry.get("specific_discharge_kg_per_rev_m"), + "specific_discharge_kg_per_rev_m", positive=True, + ) + elif "specific_discharge_kg_per_rev_m" in entry: + raise CalculationConfigError("Specific discharge requires rotational_discharge") 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 5412618..3542680 100644 --- a/src/production_analytics/calculations/material_consumption.py +++ b/src/production_analytics/calculations/material_consumption.py @@ -119,3 +119,21 @@ def area_application_rate_kg_per_hour( if not isfinite(rate): raise ValueError("derived material rate must be finite") return rate + + +def rotational_discharge_rate_kg_per_hour( + rotational_speeds_rpm: Iterable[float], nominal_width_m: float, + specific_discharge_kg_per_rev_m: float, +) -> float: + """Convert full-width roll speeds and specific discharge to canonical kg/h.""" + speeds = tuple(rotational_speeds_rpm) + if not speeds or not all(isfinite(value) for value in speeds): + raise ValueError("rotational speeds 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(specific_discharge_kg_per_rev_m) or specific_discharge_kg_per_rev_m <= 0: + raise ValueError("specific discharge must be finite and positive") + rate = sum(speeds) * nominal_width_m * specific_discharge_kg_per_rev_m * 60 + if not isfinite(rate): + raise ValueError("derived material rate must be finite") + return rate diff --git a/src/production_analytics/enlyze/gateway.py b/src/production_analytics/enlyze/gateway.py index 2d5a764..18fc69c 100644 --- a/src/production_analytics/enlyze/gateway.py +++ b/src/production_analytics/enlyze/gateway.py @@ -7,6 +7,7 @@ from math import isfinite from production_analytics.calculations.material_consumption import ( MaterialSample, area_application_rate_kg_per_hour, + rotational_discharge_rate_kg_per_hour, ) from production_analytics.enlyze.exploration import ExplorationClient @@ -89,19 +90,28 @@ class EnlyzeApiGateway: *, machine_id: str, rate_variable_id: str, - gate_variable_id: str, + gate_variable_id: str | None, start: datetime, end: datetime, application_variable_ids: tuple[str, ...] = (), nominal_width_m: float | None = None, + rotational_speed_variable_ids: tuple[str, ...] = (), + specific_discharge_kg_per_rev_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))) + if application_variable_ids and rotational_speed_variable_ids: + raise ValueError("Material source modes are mutually exclusive") + if gate_variable_id is None and not rotational_speed_variable_ids: + raise ValueError("This source requires a gate signal") + source_ids = ( + rotational_speed_variable_ids or application_variable_ids or (rate_variable_id,) + ) + gate_ids = (gate_variable_id,) if gate_variable_id else () + variable_ids = list(dict.fromkeys((*source_ids, *gate_ids))) request_body = { "machine": machine_id, "start": start.astimezone(UTC).isoformat(), @@ -123,7 +133,7 @@ class EnlyzeApiGateway: raise ValueError(f"required column {required!r} must occur exactly once") time_index = columns.index("time") source_indices = [columns.index(ref) for ref in source_ids] - gate_index = columns.index(gate_variable_id) + gate_index = columns.index(gate_variable_id) if gate_variable_id else None if not isinstance(data["records"], list): raise ValueError("records must be a list") for index, record in enumerate(data["records"]): @@ -131,13 +141,19 @@ class EnlyzeApiGateway: if not isinstance(record, list) or len(record) != len(columns): raise ValueError("record must match columns") sources = [record[i] for i in source_indices] - gate = record[gate_index] + gate = record[gate_index] if gate_index is not None else 1.0 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 rotational_speed_variable_ids: + if nominal_width_m is None or specific_discharge_kg_per_rev_m is None: + raise ValueError("rotational discharge requires width and factor") + rate = rotational_discharge_rate_kg_per_hour( + sources, nominal_width_m, specific_discharge_kg_per_rev_m, + ) + elif 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) diff --git a/src/production_analytics/service/material_context.py b/src/production_analytics/service/material_context.py index 388a765..6cc26aa 100644 --- a/src/production_analytics/service/material_context.py +++ b/src/production_analytics/service/material_context.py @@ -14,7 +14,8 @@ class ErpNominalWidthProvider: 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 + if (status is None + or status.workplace.strip().casefold() != self.workplace.strip().casefold() or build_enlyze_production_order(status.production_order, self.format_template) != production_order): raise ValueError("ERP context does not match the production order") diff --git a/src/production_analytics/service/material_polling.py b/src/production_analytics/service/material_polling.py index fc9e6e7..a9f4909 100644 --- a/src/production_analytics/service/material_polling.py +++ b/src/production_analytics/service/material_polling.py @@ -51,23 +51,30 @@ class MaterialPollingService: state_store: MaterialStateStore, machine_id: str, rate_variable_id: str, - gate_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 and nominal_width_provider is None: - raise ValueError("area application requires a 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, + gate_threshold=gate_threshold if gate_variable_id is not None else 0.0, max_sample_gap_seconds=max_sample_gap_seconds, ) @@ -79,7 +86,7 @@ class MaterialPollingService: return None width = (self._nominal_width_provider(run.production_order) - if self._application_variable_ids else None) + 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, @@ -151,6 +158,11 @@ class MaterialPollingService: 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, diff --git a/src/production_analytics/service/material_runtime.py b/src/production_analytics/service/material_runtime.py index e11d742..6c369de 100644 --- a/src/production_analytics/service/material_runtime.py +++ b/src/production_analytics/service/material_runtime.py @@ -73,7 +73,7 @@ def build_material_runner( ) from exc postgres_settings = PostgresSettings.from_environment(environment) source_options = {} - if calculation.source_mode == "area_application": + if calculation.source_mode in {"area_application", "rotational_discharge"}: source_options = dict( application_variable_ids=calculation.application_signal_refs, nominal_width_provider=ErpNominalWidthProvider( @@ -82,6 +82,12 @@ def build_material_runner( format_template=calculation.production_order_format, ), ) + if calculation.source_mode == "rotational_discharge": + source_options.pop("application_variable_ids") + source_options.update( + rotational_speed_variable_ids=calculation.rotational_speed_signal_refs, + specific_discharge_kg_per_rev_m=calculation.specific_discharge_kg_per_rev_m, + ) service = MaterialPollingService( gateway=EnlyzeApiGateway(ExplorationClient(settings)), state_store=_ReportingStateStore(JsonMaterialStateStore(state_directory)), diff --git a/tests/test_area_material_consumption.py b/tests/test_area_material_consumption.py index e3dc515..7c773d1 100644 --- a/tests/test_area_material_consumption.py +++ b/tests/test_area_material_consumption.py @@ -80,6 +80,31 @@ def width_provider(): ) +@pytest.mark.parametrize('configured,returned', [ + ('Bento 1', 'BENTO 1'), + ('Bento 1', ' \tBENTO 1\n'), + (' \tBento 1\n', 'BENTO 1'), + ('K7', ' k7 '), +]) +def test_width_context_normalizes_workplace(configured, returned): + provider = width_provider() + provider.workplace = configured + provider.gateway.get_current_workplace_status.return_value.workplace = returned + assert provider('Bento 1-12026000814') == 5.0 + + +@pytest.mark.parametrize('production_order', [ + 'Bento 1-12026000815', + 'BENTO 1-12026000814', + ' Bento 1-12026000814 ', +]) +def test_width_context_keeps_production_order_strict(production_order): + provider = width_provider() + provider.gateway.get_current_workplace_status.return_value.workplace = 'BENTO 1' + with pytest.raises(ValueError, match='ERP context does not match the production order'): + provider(production_order) + + @pytest.mark.parametrize('field,value', [ ('production_order', '12026000815'), ('workplace', 'K7'), ('article_description', None), ('article_description', '180305'), @@ -91,10 +116,11 @@ def test_width_context_fails_closed(field, value): provider('Bento 1-12026000814') -def test_bento_reference_aggregate_bootstrap_and_persistent_resume(tmp_path): +@pytest.mark.parametrize('inter_run_gap_seconds', [5, 3761]) +def test_bento_rpm_bootstrap_and_persistent_resume(tmp_path, inter_run_gap_seconds): """Synthetic aggregate fixture, not recorded ENLYZE samples. - Uses reported active time/area/application within the actual run boundaries. + Uses realistic rpm values and synthetic active time within actual run boundaries. """ config = load_material_calculation('config/bento1-material-consumption.yaml') order = 'Bento 1-12026000814' @@ -106,7 +132,11 @@ def test_bento_reference_aggregate_bootstrap_and_persistent_resume(tmp_path): ('25d67424-8346-42d4-b844-75b4097382bd', '2026-08-26T05:43:39Z', '2026-08-26T05:54:05Z'), ]] - speed = 31294.6 / (5 * 1256) + runs[1] = replace( + runs[1], start=runs[0].end + timedelta(seconds=inter_run_gap_seconds), + end=runs[0].end + timedelta(seconds=inter_run_gap_seconds + 626), + ) + speed = 2.0 client = Mock() def response(path, body): @@ -122,8 +152,8 @@ def test_bento_reference_aggregate_bootstrap_and_persistent_resume(tmp_path): 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] + 'columns': ['time', *config.rotational_speed_signal_refs, config.gate_signal_ref], + 'records': [[t.isoformat(), 3.739, 3.956, speed if t < stop else 0] for t in sorted(times)], }}) @@ -139,35 +169,44 @@ def test_bento_reference_aggregate_bootstrap_and_persistent_resume(tmp_path): 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, + rotational_speed_variable_ids=config.rotational_speed_signal_refs, + specific_discharge_kg_per_rev_m=config.specific_discharge_kg_per_rev_m, nominal_width_provider=width_provider(), ) + partial = service().poll_once(now=runs[1].start + timedelta(seconds=300)) result = service().poll_once(now=runs[1].end) + assert result.state.cumulative_consumption_kg - partial.state.cumulative_consumption_kg == ( + pytest.approx((3.739 + 3.956) * 5 * 3.12 * 326 / 60) + ) 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 result.state.cumulative_consumption_kg == pytest.approx( + (3.739 + 3.956) * 5 * 3.12 * 1256, + ) 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] + (runs[0].start, runs[0].end), + (runs[1].start, runs[1].start + timedelta(seconds=300)), + ] 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.rotational_speed_signal_refs == () 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', ''), + ('source_mode', 'unknown'), ('rotational_speed_signal_refs', []), + ('rotational_speed_signal_refs', ['same', 'same']), ('erp_workplace', ''), ('production_order_format', '{wrong}'), ('rate_signal_ref', 'direct'), ]) -def test_invalid_area_configuration(tmp_path, field, value): +def test_invalid_rotational_configuration(tmp_path, field, value): with open('config/bento1-material-consumption.yaml') as stream: document = yaml.safe_load(stream) document['calculations'][0][field] = value @@ -177,7 +216,7 @@ def test_invalid_area_configuration(tmp_path, field, value): load_material_calculation(path) -def test_area_context_failure_does_not_touch_persistent_state(tmp_path): +def test_rotational_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') @@ -191,7 +230,9 @@ def test_area_context_failure_does_not_touch_persistent_state(tmp_path): 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, + rotational_speed_variable_ids=config.rotational_speed_signal_refs, + specific_discharge_kg_per_rev_m=config.specific_discharge_kg_per_rev_m, + nominal_width_provider=provider, ) with pytest.raises(ValueError, match='context'): service.poll_once(now=now) diff --git a/tests/test_rotational_material_consumption.py b/tests/test_rotational_material_consumption.py new file mode 100644 index 0000000..17ff8d6 --- /dev/null +++ b/tests/test_rotational_material_consumption.py @@ -0,0 +1,113 @@ +from datetime import UTC, datetime, timedelta +from unittest.mock import Mock, patch + +import pytest +import yaml + +from production_analytics.calculations.config import ( + CalculationConfigError, + load_material_calculation, +) +from production_analytics.calculations.material_consumption import ( + MaterialConsumptionIntegrator, + MaterialIntegratorConfig, + rotational_discharge_rate_kg_per_hour, +) +from production_analytics.enlyze.gateway import EnlyzeApiGateway +from production_analytics.service.material_runtime import build_material_runner + + +@pytest.mark.parametrize('speeds,width,factor,expected', [ + ([3.739, 3.956], 5, 3.12, 7202.52), + ([3.739, 3.956], 4.85, 3.12, 6986.4444), + ([2], 4, 1.5, 720), + ([2, 3, 4], 4, 1.5, 3240), +]) +def test_conversion(speeds, width, factor, expected): + assert rotational_discharge_rate_kg_per_hour(speeds, width, factor) == pytest.approx(expected) + + +@pytest.mark.parametrize('speeds,width,factor', [ + ([], 5, 3.12), ([float('nan')], 5, 3.12), ([1], 0, 3.12), + ([1], float('inf'), 3.12), ([1], 5, 0), ([1], 5, float('nan')), + ([1e308], 5, 3.12), +]) +def test_invalid_conversion(speeds, width, factor): + with pytest.raises(ValueError): + rotational_discharge_rate_kg_per_hour(speeds, width, factor) + + +@pytest.mark.parametrize('gate_id', ['speed', None]) +def test_gate_and_gap(gate_id): + start = datetime(2026, 8, 25, tzinfo=UTC) + times = [0, 10, 20, 40, 61, 71, 81] + gates = [2, 0.3, 0, 8, 1, 1, 1] + client = Mock() + client.post_json.return_value.body = {'data': { + 'columns': ['time', 'right', 'left'] + ([gate_id] if gate_id else []), + 'records': [[(start + timedelta(seconds=t)).isoformat(), 3.739, 3.956] + + ([gate] if gate_id else []) for t, gate in zip(times, gates, strict=True)], + }} + samples = EnlyzeApiGateway(client).get_material_samples( + machine_id='m', rate_variable_id='unused', gate_variable_id=gate_id, + rotational_speed_variable_ids=('right', 'left'), nominal_width_m=5, + specific_discharge_kg_per_rev_m=3.12, start=start, end=start + timedelta(seconds=81), + ) + assert all(s.material_rate_kg_per_hour == pytest.approx(7202.52) for s in samples) + state = MaterialConsumptionIntegrator(MaterialIntegratorConfig(0.3, 20)).process_many(samples) + seconds = 30 if gate_id else 60 + assert state.integrated_running_seconds == seconds + assert state.cumulative_consumption_kg == pytest.approx(120.042 * seconds / 60) + assert client.post_json.call_args.args[1]['variables'] == [ + {'uuid': ref} for ref in ['right', 'left'] + ([gate_id] if gate_id else [])] + + +@pytest.mark.parametrize('field,value', [ + ('specific_discharge_kg_per_rev_m', None), ('specific_discharge_kg_per_rev_m', 0), + ('specific_discharge_kg_per_rev_m', True), ('specific_discharge_kg_per_rev_m', float('inf')), + ('application_signal_refs', ['set']), ('rotational_speed_signal_refs', ['']), +]) +def test_invalid_config(tmp_path, field, value): + document = yaml.safe_load(open('config/bento1-material-consumption.yaml')) + 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_optional_gate_config(tmp_path): + document = yaml.safe_load(open('config/bento1-material-consumption.yaml')) + del document['calculations'][0]['gate_signal_ref'] + del document['calculations'][0]['gate_threshold'] + path = tmp_path / 'config.yaml' + path.write_text(yaml.safe_dump(document)) + config = load_material_calculation(path) + assert config.gate_signal_ref is None + assert config.gate_threshold == 0 + + +def test_bento_runtime_wiring(tmp_path, monkeypatch): + monkeypatch.setenv('ENLYZE_BASE_URL', 'https://example.invalid/api/') + for key, value in dict(HOST='localhost', PORT='5432', DB='analytics', + USER='user', PASSWORD='test').items(): + monkeypatch.setenv(f'POSTGRES_{key}', value) + secrets = tmp_path / 'secret.env' + secrets.write_text('ENLYZE_API_KEY=test\n') + config = load_material_calculation('config/bento1-material-consumption.yaml') + with patch('production_analytics.service.material_runtime.ErpSettings') as erp_settings: + runner = build_material_runner( + config, poll_interval_seconds=10, state_directory=tmp_path / 'state', + secrets_file=secrets, erp_secrets_file=tmp_path / 'erp.env', + ) + service = runner.service + assert config.source_mode == 'rotational_discharge' + assert service._rotational_speed_variable_ids == ( + '6e5d2d94-98f9-4cc1-8a88-7987c6282525', 'cd7385c4-337b-4759-ab32-45d65beaf190', + ) + assert service._application_variable_ids == () + assert service._specific_discharge == 3.12 + assert service._gate_variable_id == 'fef41976-1103-4090-b780-eaeecc02fdfa' + assert service._config == MaterialIntegratorConfig(0.3, 20) + assert service._nominal_width_provider.workplace == 'Bento 1' + erp_settings.from_secret_file.assert_called_once_with(tmp_path / 'erp.env')