diff --git a/config/bento1-material-consumption.yaml b/config/bento1-material-consumption.yaml index be1581d..8431d82 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 + process_application_calculation_id: bento1-fresh-bentonite-application 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 + calibration_ref: bento1-spreader-1-2 # Transport Auszug Geschwindigkeit Istwert, m/min (production gate only). gate_signal_ref: fef41976-1103-4090-b780-eaeecc02fdfa gate_threshold: 0.3 diff --git a/config/bento1-material-efficiency.yaml b/config/bento1-material-efficiency.yaml new file mode 100644 index 0000000..ce9627b --- /dev/null +++ b/config/bento1-material-efficiency.yaml @@ -0,0 +1,7 @@ +workplace: BENTO 1 +machine_id: 5f42a4f6-9ca0-4f6f-9786-40d50a35b230 +calculation_id: bento1-fresh-bentonite-consumption +output_calculation_id: bento1-fresh-bentonite-efficiency +production_order_format: "Bento 1-{production_order}" +erp_timezone: Europe/Berlin +poll_interval_seconds: 60 diff --git a/config/material-calibrations.yaml b/config/material-calibrations.yaml new file mode 100644 index 0000000..948ea1d --- /dev/null +++ b/config/material-calibrations.yaml @@ -0,0 +1,14 @@ +calibrations: + bento1-spreader-1-2: + type: rotational_discharge + value: 2.75 + unit: kg_per_rev_m + calibrated_at: 2026-09-08 + method: gravimetric_tray + description: Bento 1 fresh-bentonite spreaders 1 and 2 + reference: + measured_application_g_m2: 4068 + line_speed_m_min: 2.3 + signal_values: + left: 1.65 + right: 1.75 diff --git a/db/schema.sql b/db/schema.sql index 53f99b7..1a69194 100644 --- a/db/schema.sql +++ b/db/schema.sql @@ -32,3 +32,17 @@ CREATE TABLE IF NOT EXISTS material_efficiency_snapshots ( CREATE INDEX IF NOT EXISTS material_efficiency_snapshots_machine_time_idx ON material_efficiency_snapshots (calculation_id, machine_id, erp_feedback_timestamp); + +-- Instantaneous process measurements; no ERP good-area denominator. +CREATE TABLE IF NOT EXISTS material_application_snapshots ( + timestamp timestamptz NOT NULL, + calculation_id text NOT NULL, + machine_id text NOT NULL, + production_order text NOT NULL, + run_id text NOT NULL, + application_g_m2 double precision NOT NULL, + PRIMARY KEY (calculation_id, machine_id, production_order, timestamp) +); + +CREATE INDEX IF NOT EXISTS material_application_snapshots_machine_time_idx + ON material_application_snapshots (calculation_id, machine_id, timestamp); diff --git a/docs/bento1-fresh-bentonite.md b/docs/bento1-fresh-bentonite.md index 7a1a013..799e1bc 100644 --- a/docs/bento1-fresh-bentonite.md +++ b/docs/bento1-fresh-bentonite.md @@ -1,4 +1,104 @@ -# Bento 1 fresh bentonite consumption +# Bento 1 fresh bentonite KPIs + + +## Three distinct quantities + +All **fresh bentonite** values exclude the third recycled/recovered-material +scatterer. Only the two actual rpm signals listed below are used; no SET signals +are introduced. The specific discharge is central calibration +data: **2.75 kg/(revolution × metre product width)**. + +A) **Instantaneous process application [g/m²]** + +`application_g_m2 = (rpm_left + rpm_right) × 2.75 × 1000 / line_speed_m_min` + +Emitted only for observed samples with line speed **strictly greater than +0.3 m/min**. Zero, near-zero, negative and threshold-equal speeds emit no point. +Width cancels between kg/min and m²/min; neither nominal width nor ERP good area +enters this formula. For 3.739 + 3.956 rpm at 8 m/min the result is 2645.15625 g/m². +This is fresh-roll process application, not FA material efficiency. + +B) **Cumulative fresh consumption [kg]**, unchanged + +`kg/min = (rpm_left + rpm_right) × nominal_width_m × 2.75` + +The existing previous-value integration sums `kg/min × elapsed_seconds / 60` +for eligible intervals, using the preceding sample's production gate. Width +continues to come from the ERP nominal-width parser. The integration algorithm, +gap handling, order checkpoint format and disjoint-run accumulation are unchanged. + +C) **FA material efficiency [g/m²]** + +`efficiency_g_m2 = cumulative_fresh_bentonite_kg × 1000 / cumulative_good_area_m2` + +The generic material-efficiency service reads the latest exact-order cumulative +snapshot **at or before the ERP feedback timestamp**, then divides by that +feedback's cumulative good area. Missing/nonpositive good area emits no point. +This ratio includes losses captured by the existing consumption model (startup +material, rejects and other consumed material not represented in good area). +It does not add consumption during intervals excluded by the validated gate or +gap rules, including stopped-line periods. Disjoint runs of one FA share the +existing order total; their cumulative snapshots must not be summed again. + +## Configuration, PostgreSQL and Grafana + +| Calculation ID | Table | Value column | Time column | +| --- | --- | --- | --- | +| `bento1-fresh-bentonite-application` | `material_application_snapshots` | `application_g_m2` | `timestamp` | +| `bento1-fresh-bentonite-consumption` | `material_consumption_snapshots` | `consumption_kg` | `timestamp` | +| `bento1-fresh-bentonite-efficiency` | `material_efficiency_snapshots` | `material_consumption_g_per_m2` | `erp_feedback_timestamp` | + +Scope Grafana queries by calculation ID and +`machine_id = '5f42a4f6-9ca0-4f6f-9786-40d50a35b230'`; optionally filter by +`production_order` (application/consumption) or `enlyze_production_order` +(efficiency). The efficiency table also stores `good_quantity_m2`, +`material_consumption_kg`, `material_consumption_kg_per_m2` and +`material_snapshot_timestamp` for auditability. + +`config/bento1-material-consumption.yaml` opts into process snapshots through +`process_application_calculation_id`. This generic option requires rotational +sources and a positive speed gate. Both outputs reuse the same configured ACT +signals, discharge factor and samples. Process points use source timestamps, +including bootstrap samples from disjoint runs; no interpolation or hold-forward +points are stored for inactive periods. Configure Grafana to leave missing +periods as gaps rather than carrying the last active value forward. + +The shared consumption poller still requires valid ERP width/order context to +complete a cycle, although the application formula itself has no width input. +Application rows are committed before the existing checkpoint advances; a failed +application write retries the window. Duplicate source timestamps are ignored by +the primary key. Existing checkpoints are retained, so earlier application +history is not automatically backfilled. K7 does not enable this output. + +`config/bento1-material-efficiency.yaml` uses the existing generic runner. +`calculation_id` selects the cumulative input; optional `output_calculation_id` +sets the distinct persisted KPI ID. Omitting it preserves existing behavior, +including K7. ERP workplace `strip().casefold()` normalization is unchanged. + +Before deploying, apply the additive, idempotent `db/schema.sql` to the existing +PostgreSQL database to create `material_application_snapshots` and its index. +The existing consumption and efficiency tables need no column migration. For +example, on the deployment host: + +```sh +docker compose exec -T timescaledb psql -U production_analytics -d production_analytics \ + -v ON_ERROR_STOP=1 < db/schema.sql +``` + +Restart the configured Bento consumption runner to enable process persistence, +and supervise a separate generic efficiency runner: + +```sh +production-analytics run material-efficiency \ + --config config/bento1-material-efficiency.yaml \ + --secrets-file secrets/enlyze.env --erp-secrets-file secrets/erp.env +``` + +No cleanup/reset of validated rpm consumption state or snapshots is required for +these additions. No database migration, live service restart or historical data +cleanup was performed as part of this implementation. + +## Validated cumulative model details `config/bento1-material-consumption.yaml` defines `bento1-fresh-bentonite-consumption`, an estimate of **fresh bentonite consumption** @@ -7,18 +107,19 @@ full web width. The recycled/recovered third spreader is excluded. 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: +Bento config sets the factor to **2.75 kg/(rev·m)** and uses: - Right ACT: `6e5d2d94-98f9-4cc1-8a88-7987c6282525` - Left ACT: `cd7385c4-337b-4759-ab32-45d65beaf190` 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. +1000 is applied. For 3.739 + 3.956 rpm and 5 m width, the rate is 105.80625 kg/min +(6348.375 kg/h), giving 17.634375 kg in ten seconds. -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 +For cumulative consumption, transport speed +`fef41976-1103-4090-b780-eaeecc02fdfa` remains exclusively the production gate: +speed must be strictly greater than 0.3 m/min. It does not multiply the mass rate. +For instantaneous application, this same speed is also the denominator. 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. @@ -43,13 +144,24 @@ 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). +The current factor is **2.75 kg/(rev*m)**, independently derived on +**2026-09-08** by the **gravimetric tray method**: measured application +**4068 g/m²**, line speed **2.3 m/min**, actual rotational signals **left 1.65** +and **right 1.75**. The derivation is `4.068 × 2.3 / (1.65 + 1.75) ≈ 2.75188`, +rounded to the configured 2.75. This supersedes the earlier 3.12 estimate. +See [central calibration configuration](material-calibrations.md) for updates +and persistence semantics. + 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 +## Historical migration from SET to rpm: review before execution + +The following records the earlier source-model migration, **not a prerequisite +cleanup for adding these KPIs**. Its deployment observations are historical. +If the validated rpm model is already deployed, retain its state and snapshots; +do not execute this historical cleanup for the KPI extension. Recheck deployment +provenance separately if that earlier migration is still outstanding. No production state or database rows were changed during this implementation, and no service start was requested. The operator reports Bento stopped; diff --git a/docs/material-calibrations.md b/docs/material-calibrations.md new file mode 100644 index 0000000..7a52bd5 --- /dev/null +++ b/docs/material-calibrations.md @@ -0,0 +1,73 @@ +# Material calibrations + +Process calibration values live in the version-controlled +`config/material-calibrations.yaml`. Calculations refer to opaque IDs; IDs have +no special meaning to the parser. The current entry is: + +```yaml +calibrations: + bento1-spreader-1-2: + type: rotational_discharge + value: 2.75 + unit: kg_per_rev_m + calibrated_at: 2026-09-08 + method: gravimetric_tray + description: Bento 1 fresh-bentonite spreaders 1 and 2 + reference: + measured_application_g_m2: 4068 + line_speed_m_min: 2.3 + signal_values: + left: 1.65 + right: 1.75 +``` + +All entry fields are required. Type, unit, method and description are non-empty +strings; value must be a finite number (booleans are rejected); calibrated_at is +a calendar date in YYYY-MM-DD format, quoted or unquoted. Reference is a non-empty +mapping whose contents preserve source-specific measurement provenance. Unknown +entry fields and duplicate YAML keys are rejected. The generic registry permits +other types and units; each consuming calculation checks its own compatibility. + +`config/bento1-material-consumption.yaml` contains +`calibration_ref: bento1-spreader-1-2`. During configuration loading, references +resolve from `material-calibrations.yaml` in the calculation file's directory, +independently of the working directory. The entire registry is validated when a +reference is used. Rotational discharge requires type `rotational_discharge`, +unit `kg_per_rev_m` and a positive value. Resolution supplies the existing +`specific_discharge_kg_per_rev_m` runtime field for both application and +consumption; the calculation algorithms are unchanged. There is no unit conversion. + +Missing/unreadable files, missing IDs, invalid metadata and incompatible type or +unit raise `CalculationConfigError` before a runner starts. Existing direct numeric +configurations remain supported; specifying both a reference and a direct factor +is rejected. References are currently supported by rotational discharge consumers. +K7 uses direct mass rate and does not load or require a calibration file. + +## Updating a calibration + +After a physical measurement, edit only the central entry's value, date, method +and reference details, and keep its ID stable. Review and version-control the +change. Additional machines can add new IDs using the same structure. A new +calculation type needs an explicit consumer compatibility contract. + +Ship the central YAML alongside the calculation YAML and restart the consuming +process to load changes. No new CLI flag, environment variable, database migration +or UI is required. There is no hot reload or historical date-based selection. +Changes affect future calculations and explicitly rebuilt calculations. Historical +persisted snapshots are **not automatically recalculated**. Existing checkpoints +retain accumulated totals, so subsequent increments use the newly loaded factor; +a complete historical rebuild requires a separately planned replay. Persisted +snapshots do not gain calibration-version provenance through this change. + +## Current Bento provenance + +On 2026-09-08, an independent gravimetric tray measurement found 4068 g/m² at +2.3 m/min with actual rotational signals left 1.65 and right 1.75. +`(4068 / 1000) × 2.3 / (1.65 + 1.75) ≈ 2.75188 kg/(rev*m)` gives the configured +rounded factor **2.75 kg/(rev*m)**. The original measurement remains recorded, +without fitting or adjusting the runtime value to reproduce it exactly. + +Only fresh-bentonite spreaders 1 and 2 are included. Before adding spreader 3, +confirm its material scope, actual signal units, independent calibration and +whether its output belongs in the fresh-material KPIs. No spreader 3 support is +included here. diff --git a/src/production_analytics/calculations/calibrations.py b/src/production_analytics/calculations/calibrations.py new file mode 100644 index 0000000..54329f1 --- /dev/null +++ b/src/production_analytics/calculations/calibrations.py @@ -0,0 +1,89 @@ +"""Generic, version-controlled calibration values and measurement provenance.""" + +from dataclasses import dataclass +from datetime import date +from pathlib import Path +from typing import Any + +import yaml + +from production_analytics.calculations.config import ( + CalculationConfigError, + _UniqueLoader, + finite_number, +) + + +@dataclass(frozen=True, slots=True) +class Calibration: + type: str + value: float + unit: str + calibrated_at: date + method: str + description: str + reference: dict[str, Any] + + +def load_calibrations(path: str | Path) -> dict[str, Calibration]: + """Validate every entry, retaining opaque IDs and free-form reference metadata.""" + try: + with Path(path).open(encoding="utf-8") as stream: + document = yaml.load(stream, Loader=_UniqueLoader) + except (OSError, UnicodeError) as exc: + raise CalculationConfigError(f"Cannot read calibration configuration: {path}") from exc + except yaml.YAMLError as exc: + raise CalculationConfigError(f"Invalid YAML in calibration configuration: {path}") from exc + if not isinstance(document, dict) or set(document) != {"calibrations"}: + raise CalculationConfigError("Calibration configuration must contain only 'calibrations'") + entries = document["calibrations"] + if not isinstance(entries, dict) or not entries: + raise CalculationConfigError("calibrations must be a non-empty mapping") + result = {} + required = {"type", "value", "unit", "calibrated_at", "method", "description", "reference"} + for identifier, entry in entries.items(): + prefix = f"calibrations[{identifier!r}]" + if not identifier.strip(): + raise CalculationConfigError("Calibration IDs must be non-empty strings") + if not isinstance(entry, dict) or set(entry) != required: + raise CalculationConfigError( + f"{prefix}: required fields are {', '.join(sorted(required))}" + ) + for name in ("type", "unit", "method", "description"): + if not isinstance(entry[name], str) or not entry[name].strip(): + raise CalculationConfigError(f"{prefix}.{name} must be a non-empty string") + calibrated_at = entry["calibrated_at"] + if isinstance(calibrated_at, str): + try: + parsed = date.fromisoformat(calibrated_at) + if parsed.isoformat() != calibrated_at: + raise ValueError + calibrated_at = parsed + except ValueError as exc: + raise CalculationConfigError(f"{prefix}.calibrated_at must be YYYY-MM-DD") from exc + if type(calibrated_at) is not date: + raise CalculationConfigError(f"{prefix}.calibrated_at must be YYYY-MM-DD") + if not isinstance(entry["reference"], dict) or not entry["reference"]: + raise CalculationConfigError(f"{prefix}.reference must be a non-empty mapping") + result[identifier] = Calibration( + **{**entry, "calibrated_at": calibrated_at, + "value": finite_number(entry["value"], f"{prefix}.value")}, + ) + return result + + +def resolve_calibration( + calibrations: dict[str, Calibration], reference: str, *, expected_type: str, expected_unit: str, +) -> Calibration: + """Require explicit compatibility; no implicit conversion or ID interpretation.""" + if not isinstance(reference, str) or not reference.strip(): + raise CalculationConfigError("calibration_ref must be a non-empty string") + if reference not in calibrations: + raise CalculationConfigError(f"Missing calibration ID: {reference!r}") + calibration = calibrations[reference] + for field, expected in (("type", expected_type), ("unit", expected_unit)): + if getattr(calibration, field) != expected: + raise CalculationConfigError( + f"Calibration {reference!r}: incompatible {field}; expected {expected!r}" + ) + return calibration diff --git a/src/production_analytics/calculations/config.py b/src/production_analytics/calculations/config.py index 87d32b7..855a551 100644 --- a/src/production_analytics/calculations/config.py +++ b/src/production_analytics/calculations/config.py @@ -58,8 +58,10 @@ class MaterialCalculationConfig: application_signal_refs: tuple[str, ...] = () rotational_speed_signal_refs: tuple[str, ...] = () specific_discharge_kg_per_rev_m: float | None = None + calibration_ref: str | None = None erp_workplace: str = "" production_order_format: str = "" + process_application_calculation_id: str | None = None def load_material_calculation( @@ -81,6 +83,7 @@ def load_material_calculation( required = {field.name for field in fields(MaterialCalculationConfig) if field.default is MISSING} instances = {} + calibrations = None for index, entry in enumerate(entries): prefix = f"calculations[{index}]" if not isinstance(entry, dict): @@ -148,6 +151,24 @@ def load_material_calculation( "erp_workplace", "production_order_format", )): raise CalculationConfigError("Width context fields require a width-based source mode") + if "calibration_ref" in entry: + if not rotational: + raise CalculationConfigError("calibration_ref requires rotational_discharge") + if "specific_discharge_kg_per_rev_m" in entry: + raise CalculationConfigError( + "Specify calibration_ref or direct specific discharge, not both" + ) + from production_analytics.calculations.calibrations import ( + load_calibrations, + resolve_calibration, + ) + if calibrations is None: + calibrations = load_calibrations(Path(path).parent / "material-calibrations.yaml") + calibration = resolve_calibration( + calibrations, entry["calibration_ref"], + expected_type="rotational_discharge", expected_unit="kg_per_rev_m", + ) + entry["specific_discharge_kg_per_rev_m"] = calibration.value if rotational: entry["specific_discharge_kg_per_rev_m"] = finite_number( entry.get("specific_discharge_kg_per_rev_m"), @@ -155,9 +176,24 @@ def load_material_calculation( ) elif "specific_discharge_kg_per_rev_m" in entry: raise CalculationConfigError("Specific discharge requires rotational_discharge") + if "process_application_calculation_id" in entry: + process_id = entry["process_application_calculation_id"] + if (not isinstance(process_id, str) or not process_id.strip() + or "\x00" in process_id or process_id == entry["id"]): + raise CalculationConfigError( + "process application ID must be a distinct non-empty string" + ) + if not rotational or not entry["gate_signal_ref"] or entry["gate_threshold"] <= 0: + raise CalculationConfigError( + "Process application requires rotational_discharge and a positive speed gate" + ) if entry["id"] in instances: raise CalculationConfigError("Calculation ids must be unique") instances[entry["id"]] = MaterialCalculationConfig(**entry) + process_ids = [c.process_application_calculation_id for c in instances.values() + if c.process_application_calculation_id is not None] + if len(set(process_ids)) != len(process_ids) or set(process_ids) & instances.keys(): + raise CalculationConfigError("Calculation ids must be unique across all outputs") if calculation_id is not None: if calculation_id not in instances: raise CalculationConfigError("Requested calculation id was not found") diff --git a/src/production_analytics/calculations/material_application.py b/src/production_analytics/calculations/material_application.py new file mode 100644 index 0000000..83c2627 --- /dev/null +++ b/src/production_analytics/calculations/material_application.py @@ -0,0 +1,25 @@ +"""Width-independent process application from actual rotational speeds.""" + +from collections.abc import Iterable +from math import isfinite + + +def rotational_application_g_m2( + rotational_speeds_rpm: Iterable[float], specific_discharge_kg_per_rev_m: float, + line_speed_m_min: float, gate_threshold: float, +) -> float | None: + 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(specific_discharge_kg_per_rev_m) or specific_discharge_kg_per_rev_m <= 0: + raise ValueError("specific discharge must be finite and positive") + if not isfinite(gate_threshold) or gate_threshold <= 0: + raise ValueError("speed gate threshold must be finite and positive") + if not isfinite(line_speed_m_min): + raise ValueError("line speed must be finite") + if line_speed_m_min <= gate_threshold: + return None + value = sum(speeds) * specific_discharge_kg_per_rev_m * 1000 / line_speed_m_min + if not isfinite(value): + raise ValueError("derived application must be finite") + return value diff --git a/src/production_analytics/calculations/material_consumption.py b/src/production_analytics/calculations/material_consumption.py index 3542680..475a5cb 100644 --- a/src/production_analytics/calculations/material_consumption.py +++ b/src/production_analytics/calculations/material_consumption.py @@ -13,6 +13,7 @@ class MaterialSample: timestamp: datetime material_rate_kg_per_hour: float gate_value: float + application_g_m2: float | None = None @dataclass(frozen=True, slots=True) diff --git a/src/production_analytics/enlyze/gateway.py b/src/production_analytics/enlyze/gateway.py index 18fc69c..6d8424a 100644 --- a/src/production_analytics/enlyze/gateway.py +++ b/src/production_analytics/enlyze/gateway.py @@ -97,6 +97,7 @@ class EnlyzeApiGateway: nominal_width_m: float | None = None, rotational_speed_variable_ids: tuple[str, ...] = (), specific_discharge_kg_per_rev_m: float | None = None, + process_application_gate_threshold: float | None = None, ) -> list[MaterialSample]: for name, value in (("start", start), ("end", end)): if value.tzinfo is None or value.utcoffset() is None: @@ -159,9 +160,21 @@ class EnlyzeApiGateway: rate = area_application_rate_kg_per_hour(sources, nominal_width_m, gate) else: rate = sources[0] + application = None + if process_application_gate_threshold is not None: + if not rotational_speed_variable_ids or gate_variable_id is None: + raise ValueError("Process application requires rpm and speed gate") + from production_analytics.calculations.material_application import ( + rotational_application_g_m2, + ) + application = rotational_application_g_m2( + sources, specific_discharge_kg_per_rev_m, gate, + process_application_gate_threshold, + ) samples.append(MaterialSample( timestamp=_parse_timestamp(record[time_index]), material_rate_kg_per_hour=float(rate), gate_value=float(gate), + application_g_m2=application, )) except (TypeError, ValueError) as exc: raise ValueError(f"malformed record {index}: {exc}") from exc diff --git a/src/production_analytics/service/material_efficiency.py b/src/production_analytics/service/material_efficiency.py index e5b4aa1..38beb30 100644 --- a/src/production_analytics/service/material_efficiency.py +++ b/src/production_analytics/service/material_efficiency.py @@ -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), diff --git a/src/production_analytics/service/material_efficiency_runtime.py b/src/production_analytics/service/material_efficiency_runtime.py index 0c4bc0c..2afcd9b 100644 --- a/src/production_analytics/service/material_efficiency_runtime.py +++ b/src/production_analytics/service/material_efficiency_runtime.py @@ -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, diff --git a/src/production_analytics/service/material_polling.py b/src/production_analytics/service/material_polling.py index a9f4909..0314a74 100644 --- a/src/production_analytics/service/material_polling.py +++ b/src/production_analytics/service/material_polling.py @@ -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 diff --git a/src/production_analytics/service/material_runtime.py b/src/production_analytics/service/material_runtime.py index 6c369de..618a902 100644 --- a/src/production_analytics/service/material_runtime.py +++ b/src/production_analytics/service/material_runtime.py @@ -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)), diff --git a/src/production_analytics/service/postgres_material_application.py b/src/production_analytics/service/postgres_material_application.py new file mode 100644 index 0000000..57388fd --- /dev/null +++ b/src/production_analytics/service/postgres_material_application.py @@ -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, + ) diff --git a/tests/test_area_material_consumption.py b/tests/test_area_material_consumption.py index 7c773d1..e050e55 100644 --- a/tests/test_area_material_consumption.py +++ b/tests/test_area_material_consumption.py @@ -177,11 +177,11 @@ def test_bento_rpm_bootstrap_and_persistent_resume(tmp_path, inter_run_gap_secon 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) + pytest.approx((3.739 + 3.956) * 5 * 2.75 * 326 / 60) ) assert result.state.integrated_running_seconds / 60 == pytest.approx(1256) assert result.state.cumulative_consumption_kg == pytest.approx( - (3.739 + 3.956) * 5 * 3.12 * 1256, + (3.739 + 3.956) * 5 * 2.75 * 1256, ) assert service().poll_once(now=runs[1].end).state == result.state assert gateway.get_production_runs.call_count == 1 diff --git a/tests/test_calibrations.py b/tests/test_calibrations.py new file mode 100644 index 0000000..229f167 --- /dev/null +++ b/tests/test_calibrations.py @@ -0,0 +1,143 @@ +from dataclasses import replace +from datetime import UTC, date, datetime, timedelta +from pathlib import Path +from unittest.mock import Mock + +import pytest +import yaml + +from production_analytics.calculations.calibrations import load_calibrations, resolve_calibration +from production_analytics.calculations.config import ( + CalculationConfigError, + load_material_calculation, +) +from production_analytics.calculations.material_application import rotational_application_g_m2 +from production_analytics.calculations.material_consumption import ( + MaterialConsumptionIntegrator, + MaterialIntegratorConfig, +) +from production_analytics.enlyze.gateway import EnlyzeApiGateway + +CALIBRATIONS = Path('config/material-calibrations.yaml') +BENTO = Path('config/bento1-material-consumption.yaml') +ID = 'bento1-spreader-1-2' + + +def resolve(entries, identifier=ID): + return resolve_calibration(entries, identifier, expected_type='rotational_discharge', + expected_unit='kg_per_rev_m') + + +def test_load_and_resolve(): + calibration = resolve(load_calibrations(CALIBRATIONS)) + assert calibration.value == 2.75 + assert calibration.calibrated_at == date(2026, 9, 8) + assert calibration.method == 'gravimetric_tray' + assert calibration.reference == { + 'measured_application_g_m2': 4068, 'line_speed_m_min': 2.3, + 'signal_values': {'left': 1.65, 'right': 1.75}, + } + assert load_material_calculation(BENTO).specific_discharge_kg_per_rev_m == 2.75 + + +@pytest.mark.parametrize('field', ['type', 'value', 'unit', 'calibrated_at', 'method', + 'description', 'reference']) +def test_missing_metadata(tmp_path, field): + document = yaml.safe_load(CALIBRATIONS.read_text()) + del document['calibrations'][ID][field] + path = tmp_path / 'material-calibrations.yaml' + path.write_text(yaml.safe_dump(document)) + with pytest.raises(CalculationConfigError, match='required fields'): + load_calibrations(path) + + +@pytest.mark.parametrize('field,value', [ + ('type', ''), ('unit', None), ('method', 12), ('description', ' '), + ('calibrated_at', '2026-02-30'), ('calibrated_at', 123), + ('reference', []), ('reference', {}), ('value', True), ('value', float('nan')), +]) +def test_invalid_metadata(tmp_path, field, value): + document = yaml.safe_load(CALIBRATIONS.read_text()) + document['calibrations'][ID][field] = value + path = tmp_path / 'material-calibrations.yaml' + path.write_text(yaml.safe_dump(document)) + with pytest.raises(CalculationConfigError, match=field): + load_calibrations(path) + + +@pytest.mark.parametrize('field,value,error', [ + ('type', 'other', 'incompatible type'), ('unit', 'kg_per_rev', 'incompatible unit'), + ('value', 0, 'greater than zero'), +]) +def test_incompatible_at_config_load(tmp_path, field, value, error): + document = yaml.safe_load(CALIBRATIONS.read_text()) + document['calibrations'][ID][field] = value + (tmp_path / CALIBRATIONS.name).write_text(yaml.safe_dump(document)) + path = tmp_path / BENTO.name + path.write_text(BENTO.read_text()) + with pytest.raises(CalculationConfigError, match=error): + load_material_calculation(path) + + +def test_missing_id_and_file(tmp_path): + with pytest.raises(CalculationConfigError, match='Missing calibration ID'): + resolve(load_calibrations(CALIBRATIONS), 'unknown') + path = tmp_path / BENTO.name + path.write_text(BENTO.read_text()) + with pytest.raises(CalculationConfigError, match='Cannot read calibration'): + load_material_calculation(path) + (tmp_path / CALIBRATIONS.name).write_text(CALIBRATIONS.read_text()) + path.write_text(BENTO.read_text().replace(ID, 'unknown')) + with pytest.raises(CalculationConfigError, match='Missing calibration ID'): + load_material_calculation(path) + + +def test_exact_direct_value_equivalence(tmp_path): + referenced = load_material_calculation(BENTO) + document = yaml.safe_load(BENTO.read_text()) + entry = document['calculations'][0] + del entry['calibration_ref'] + entry['specific_discharge_kg_per_rev_m'] = 2.75 + path = tmp_path / 'direct.yaml' + path.write_text(yaml.safe_dump(document)) + direct = load_material_calculation(path) + assert replace(referenced, calibration_ref=None) == direct + start = datetime(2026, 9, 8, tzinfo=UTC) + outputs = [] + for config in (referenced, direct): + factor = config.specific_discharge_kg_per_rev_m + application = rotational_application_g_m2([1.65, 1.75], factor, 2.3, 0.3) + client = Mock() + client.post_json.return_value.body = {'data': { + 'columns': ['time', 'left', 'right', 'speed'], + 'records': [[(start + timedelta(seconds=t)).isoformat(), 1.65, 1.75, speed] + for t, speed in [(0, 2.3), (10, 0), (20, 2.3), (50, 2.3), (60, 2.3)]], + }} + samples = EnlyzeApiGateway(client).get_material_samples( + machine_id='m', rate_variable_id='unused', gate_variable_id='speed', + rotational_speed_variable_ids=('left', 'right'), nominal_width_m=5, + specific_discharge_kg_per_rev_m=factor, start=start, + end=start + timedelta(seconds=60), + ) + integrator = MaterialConsumptionIntegrator(MaterialIntegratorConfig(0.3, 20)) + state = integrator.process_many(samples) + assert samples[0].material_rate_kg_per_hour == 2805.0 + assert state.integrated_running_seconds == 20 + outputs.append((application, samples, state)) + assert outputs[0] == outputs[1] + assert outputs[0][0] == pytest.approx(4065.217391304348) + + +def test_update_only_central_file(tmp_path): + path = tmp_path / BENTO.name + path.write_text(BENTO.read_text()) + (tmp_path / CALIBRATIONS.name).write_text(CALIBRATIONS.read_text().replace('2.75', '2.8')) + assert load_material_calculation(path).specific_discharge_kg_per_rev_m == 2.8 + + +def test_k7_without_calibration_file(tmp_path): + original = Path('config/k7-material-consumption.yaml') + path = tmp_path / original.name + path.write_text(original.read_text()) + assert load_material_calculation(path) == load_material_calculation(original) + assert load_material_calculation(path).calibration_ref is None diff --git a/tests/test_material_application.py b/tests/test_material_application.py new file mode 100644 index 0000000..5cf9a7b --- /dev/null +++ b/tests/test_material_application.py @@ -0,0 +1,158 @@ +from datetime import UTC, datetime, timedelta +from pathlib import Path +from unittest.mock import MagicMock, Mock + +import psycopg +import pytest +import yaml + +from production_analytics.calculations.config import ( + CalculationConfigError, + load_material_calculation, +) +from production_analytics.calculations.material_application import rotational_application_g_m2 +from production_analytics.calculations.material_consumption import MaterialSample +from production_analytics.enlyze.gateway import EnlyzeApiGateway, EnlyzeProductionRun +from production_analytics.service.material_polling import MaterialPollingService +from production_analytics.service.postgres_material import PostgresSettings +from production_analytics.service.postgres_material_application import ( + PostgresMaterialApplicationWriter, +) + +NOW = datetime(2026, 9, 7, tzinfo=UTC) + + +def test_realistic_application_and_generic_factor(): + assert rotational_application_g_m2([3.739, 3.956], 3.12, 8, 0.3) == pytest.approx(3001.05) + assert rotational_application_g_m2([2, 3], 1.5, 10, 0.5) == 750 + + +@pytest.mark.parametrize('speed', [-10, 0, 1e-300, 0.299999999, 0.3]) +def test_inactive_and_near_zero(speed): + assert rotational_application_g_m2([3.739, 3.956], 3.12, speed, 0.3) is None + + +@pytest.mark.parametrize('speeds,factor,speed,threshold', [ + ([], 3.12, 8, 0.3), ([float('nan')], 3.12, 8, 0.3), + ([1], 0, 8, 0.3), ([1], float('inf'), 8, 0.3), + ([1], 3.12, float('nan'), 0.3), ([1], 3.12, 8, 0), + ([1], 3.12, 8, float('inf')), ([1e308], 3.12, 8, 0.3), +]) +def test_invalid_inputs(speeds, factor, speed, threshold): + with pytest.raises(ValueError): + rotational_application_g_m2(speeds, factor, speed, threshold) + + +@pytest.mark.parametrize('width', [1, 4.85, 5, 100]) +def test_gateway_width_cancels_and_only_act_signals_requested(width): + config = load_material_calculation('config/bento1-material-consumption.yaml') + refs = [*config.rotational_speed_signal_refs, config.gate_signal_ref] + client = Mock() + client.post_json.return_value.body = {'data': { + 'columns': ['time', *refs], + 'records': [[NOW.isoformat(), 3.739, 3.956, 8], + [(NOW + timedelta(seconds=10)).isoformat(), 3.739, 3.956, 0]], + }} + samples = EnlyzeApiGateway(client).get_material_samples( + machine_id=config.machine_ref, rate_variable_id='unused', + gate_variable_id=config.gate_signal_ref, start=NOW, end=NOW + timedelta(seconds=10), + rotational_speed_variable_ids=config.rotational_speed_signal_refs, + nominal_width_m=width, specific_discharge_kg_per_rev_m=3.12, + process_application_gate_threshold=0.3, + ) + assert samples[0].application_g_m2 == pytest.approx(3001.05) + assert samples[1].application_g_m2 is None + assert samples[0].material_rate_kg_per_hour == pytest.approx(7.695 * width * 3.12 * 60) + assert client.post_json.call_args.args[1]['variables'] == [{'uuid': ref} for ref in refs] + + +@pytest.mark.parametrize('changes', [ + {'process_application_calculation_id': ''}, {'process_application_calculation_id': None}, + {'process_application_calculation_id': 3}, {'process_application_calculation_id': 'x\x00'}, + {'process_application_calculation_id': 'bento1-fresh-bentonite-consumption'}, + {'gate_signal_ref': None}, {'gate_threshold': 0}, {'gate_threshold': -0.3}, +]) +def test_process_config_validation(tmp_path, changes): + document = yaml.safe_load(Path('config/bento1-material-consumption.yaml').read_text()) + document['calculations'][0].update(changes) + path = tmp_path / 'config.yaml' + path.write_text(yaml.safe_dump(document)) + with pytest.raises(CalculationConfigError): + load_material_calculation(path) + + +def test_persistence_before_checkpoint_retries_and_disjoint_runs(): + old = EnlyzeProductionRun('old', 'machine', None, 'order', NOW, NOW + timedelta(seconds=10)) + current = EnlyzeProductionRun( + 'new', 'machine', None, 'order', NOW + timedelta(seconds=15), None, + ) + gateway = Mock() + gateway.get_open_production_run.return_value = current + gateway.get_production_runs.return_value = [old, current] + + def samples(**kwargs): + start = kwargs['start'] + assert kwargs['process_application_gate_threshold'] == 0.3 + return [MaterialSample(start, 7202.52, 8, 3001.05), + MaterialSample(start + timedelta(seconds=10), 7202.52, 8, 3001.05)] + + gateway.get_material_samples.side_effect = samples + store = Mock(load=Mock(return_value=None)) + writer = Mock() + writer.write.side_effect = [None, psycopg.OperationalError(), None, None] + service = MaterialPollingService( + gateway=gateway, state_store=store, machine_id='machine', rate_variable_id='unused', + gate_variable_id='speed', gate_threshold=0.3, max_sample_gap_seconds=20, + rotational_speed_variable_ids=('right', 'left'), specific_discharge_kg_per_rev_m=3.12, + nominal_width_provider=lambda order: 5, process_application_calculation_id='application', + application_writer=writer, + ) + with pytest.raises(psycopg.OperationalError): + service.poll_once(now=NOW + timedelta(seconds=25)) + store.save.assert_not_called() + result = service.poll_once(now=NOW + timedelta(seconds=25)) + assert result.state.cumulative_consumption_kg == pytest.approx(40.014) + assert result.state.integrated_running_seconds == 20 + assert [c.kwargs['run_id'] for c in writer.write.call_args_list] == ['old', 'new', 'old', 'new'] + store.save.assert_called_once() + + +def test_sql_persistence_idempotent_and_inactive_omitted(monkeypatch): + import sqlite3 + + database = sqlite3.connect(':memory:') + database.executescript(Path('db/schema.sql').read_text()) + connection = MagicMock() + cursor = connection.__enter__.return_value.cursor.return_value.__enter__.return_value + cursor.executemany.side_effect = lambda sql, rows: database.executemany( + sql.replace('%s', '?'), + [(row[0].isoformat(), *row[1:]) for row in rows], + ) + connect = Mock(return_value=connection) + monkeypatch.setattr(psycopg, 'connect', connect) + writer = PostgresMaterialApplicationWriter(PostgresSettings('host', 5432, 'db', 'u', 'p')) + kwargs = dict(calculation_id='application', machine_id='machine', production_order='order', + run_id='run', samples=[MaterialSample(NOW, 7202.52, 8, 3001.05), + MaterialSample(NOW + timedelta(seconds=10), 7202.52, 0)]) + writer.write(**kwargs) + writer.write(**kwargs) + assert database.execute('SELECT * FROM material_application_snapshots').fetchall() == [ + (NOW.isoformat(), 'application', 'machine', 'order', 'run', 3001.05), + ] + writer.write(**(kwargs | {'samples': [MaterialSample(NOW, 1, 0)]})) + assert connect.call_count == 2 + database.close() + + +@pytest.mark.parametrize('collision', ['output', 'consumption']) +def test_output_ids_unique_across_config_entries(tmp_path, collision): + document = yaml.safe_load(Path('config/bento1-material-consumption.yaml').read_text()) + first = document['calculations'][0] + second = dict(first, id='second-consumption') + if collision == 'consumption': + second['process_application_calculation_id'] = first['id'] + document['calculations'].append(second) + path = tmp_path / 'config.yaml' + path.write_text(yaml.safe_dump(document)) + with pytest.raises(CalculationConfigError, match='unique'): + load_material_calculation(path, first['id']) diff --git a/tests/test_material_efficiency_persistence.py b/tests/test_material_efficiency_persistence.py index fd8a36b..99ef2a8 100644 --- a/tests/test_material_efficiency_persistence.py +++ b/tests/test_material_efficiency_persistence.py @@ -445,3 +445,44 @@ def test_commit_failure_does_not_log_success(storage): ) assert output.getvalue() == "" assert errors == "Material efficiency cycle failed: OperationalError\n" + + +def test_bento_uses_generic_efficiency_and_distinct_output_id(runtime_files, storage): + _, postgres, erp = runtime_files + runner = build_material_efficiency_runner( + Path('config/bento1-material-efficiency.yaml'), + secrets_file=postgres, erp_secrets_file=erp, + ) + from production_analytics.service.material_efficiency import MaterialEfficiencyService + + assert type(runner.service) is MaterialEfficiencyService + service = runner.service + status = CurrentWorkplaceStatus( + 'Bento 1', '00123', None, None, NOW, None, 800, None, None, None, + ) + service.erp = Mock(get_current_workplace_status=Mock(return_value=status)) + service.materials = Mock(latest_at_or_before=Mock(return_value=MaterialConsumptionSnapshot( + NOW, 'bento1-fresh-bentonite-consumption', service.machine_id, 'Bento 1-00123', 'run', 2400, + ))) + snapshot = service.evaluate_current() + assert snapshot.calculation_id == 'bento1-fresh-bentonite-efficiency' + assert snapshot.material_consumption_g_per_m2 == 3000 + assert snapshot.nominal_width_m is None + service.materials.latest_at_or_before.assert_called_once_with( + calculation_id='bento1-fresh-bentonite-consumption', machine_id=service.machine_id, + production_order='Bento 1-00123', timestamp=NOW, + ) + runner.writer.write(snapshot) + row = storage[0].execute( + 'SELECT calculation_id, material_consumption_g_per_m2 FROM material_efficiency_snapshots' + ).fetchone() + assert row == ('bento1-fresh-bentonite-efficiency', 3000) + + +@pytest.mark.parametrize('value', ['', None, 12, 'bad\x00id']) +def test_invalid_output_calculation_id(runtime_files, value): + config, postgres, erp = runtime_files + document = yaml.safe_load(config.read_text()) | {'output_calculation_id': value} + config.write_text(yaml.safe_dump(document)) + with pytest.raises(CalculationConfigError, match='output_calculation_id'): + build_material_efficiency_runner(config, secrets_file=postgres, erp_secrets_file=erp) diff --git a/tests/test_rotational_material_consumption.py b/tests/test_rotational_material_consumption.py index 17ff8d6..2bfbe50 100644 --- a/tests/test_rotational_material_consumption.py +++ b/tests/test_rotational_material_consumption.py @@ -1,4 +1,5 @@ from datetime import UTC, datetime, timedelta +from pathlib import Path from unittest.mock import Mock, patch import pytest @@ -78,10 +79,13 @@ def test_invalid_config(tmp_path, field, value): def test_optional_gate_config(tmp_path): document = yaml.safe_load(open('config/bento1-material-consumption.yaml')) + del document['calculations'][0]['process_application_calculation_id'] 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)) + (tmp_path / 'material-calibrations.yaml').write_text( + Path('config/material-calibrations.yaml').read_text()) config = load_material_calculation(path) assert config.gate_signal_ref is None assert config.gate_threshold == 0 @@ -102,11 +106,16 @@ def test_bento_runtime_wiring(tmp_path, monkeypatch): ) service = runner.service assert config.source_mode == 'rotational_discharge' + assert service._process_application_calculation_id == 'bento1-fresh-bentonite-application' + from production_analytics.service.postgres_material_application import ( + PostgresMaterialApplicationWriter, + ) + assert isinstance(service._application_writer, PostgresMaterialApplicationWriter) 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._specific_discharge == 2.75 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'