Compare commits
2
Commits
be8d76507a
...
73a6e60450
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
73a6e60450 | ||
|
|
0210ced9e3 |
@@ -4,11 +4,13 @@ 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).
|
||||
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
|
||||
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
|
||||
max_sample_gap_seconds: 20.0
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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);
|
||||
|
||||
+240
-50
@@ -1,60 +1,250 @@
|
||||
# Bento 1 fresh bentonite consumption
|
||||
# Bento 1 fresh bentonite KPIs
|
||||
|
||||
`config/bento1-material-consumption.yaml` defines
|
||||
`bento1-fresh-bentonite-consumption`. This KPI estimates **fresh bentonite
|
||||
consumption** from the two fresh-spreader setpoints in g/m². The third spreader
|
||||
uses recycled/recovered bentonite and is deliberately excluded; it is neither
|
||||
missing data nor estimated. This is not total bentonite deposited on the product,
|
||||
and setpoints are not a measurement of actual mass flow.
|
||||
|
||||
The generic `area_application` source sums the configured application signals
|
||||
and converts them to kg/h using `sum(g/m²) × nominal_width_m × speed_m/min × 60 / 1000`.
|
||||
The existing previous-value integrator then accumulates kilograms only while
|
||||
speed is strictly greater than 0.3 m/min. `direct_mass_rate` remains the default
|
||||
for existing configurations, including K7. The persistent state schema and
|
||||
production-order bootstrap/resume behavior are unchanged; each disjoint run
|
||||
starts with a new sample baseline while retaining order totals.
|
||||
## Three distinct quantities
|
||||
|
||||
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.
|
||||
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)**.
|
||||
|
||||
Run with:
|
||||
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
|
||||
production-analytics run material-poll \
|
||||
--config config/bento1-material-consumption.yaml \
|
||||
--erp-secrets-file secrets/erp.env
|
||||
docker compose exec -T timescaledb psql -U production_analytics -d production_analytics \
|
||||
-v ON_ERROR_STOP=1 < db/schema.sql
|
||||
```
|
||||
|
||||
`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.
|
||||
Restart the configured Bento consumption runner to enable process persistence,
|
||||
and supervise a separate generic efficiency runner:
|
||||
|
||||
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.
|
||||
```sh
|
||||
production-analytics run material-efficiency \
|
||||
--config config/bento1-material-efficiency.yaml \
|
||||
--secrets-file secrets/enlyze.env --erp-secrets-file secrets/erp.env
|
||||
```
|
||||
|
||||
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.
|
||||
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**
|
||||
from actual spreader roll speeds. Both needle rolls apply sequentially across the
|
||||
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 **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 105.80625 kg/min
|
||||
(6348.375 kg/h), giving 17.634375 kg in ten seconds.
|
||||
|
||||
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.
|
||||
|
||||
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 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.
|
||||
|
||||
## 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;
|
||||
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
|
||||
```
|
||||
|
||||
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.
|
||||
|
||||
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.
|
||||
|
||||
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.
|
||||
|
||||
@@ -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.
|
||||
@@ -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
|
||||
@@ -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,8 +56,12 @@ 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
|
||||
calibration_ref: str | None = None
|
||||
erp_workplace: str = ""
|
||||
production_order_format: str = ""
|
||||
process_application_calculation_id: str | None = None
|
||||
|
||||
|
||||
def load_material_calculation(
|
||||
@@ -79,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):
|
||||
@@ -86,7 +91,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 +103,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,34 +119,81 @@ 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 "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"),
|
||||
"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 "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")
|
||||
|
||||
@@ -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
|
||||
@@ -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)
|
||||
@@ -119,3 +120,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
|
||||
|
||||
@@ -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,29 @@ 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,
|
||||
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:
|
||||
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 +134,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,21 +142,39 @@ 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)
|
||||
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
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -38,12 +38,14 @@ class MaterialEfficiencyService:
|
||||
self, erp: WorkplaceStatusReader, materials: MaterialSnapshotRepository, *,
|
||||
workplace: str, machine_id: str, calculation_id: str, format_template: str,
|
||||
erp_timezone: ZoneInfo | None = None,
|
||||
output_calculation_id: str | None = None,
|
||||
) -> None:
|
||||
self.erp = erp
|
||||
self.materials = materials
|
||||
self.workplace = workplace
|
||||
self.machine_id = machine_id
|
||||
self.calculation_id = calculation_id
|
||||
self.output_calculation_id = output_calculation_id or calculation_id
|
||||
self.format_template = format_template
|
||||
self.erp_timezone = erp_timezone
|
||||
|
||||
@@ -58,7 +60,7 @@ class MaterialEfficiencyService:
|
||||
Configuration, timestamp and adapter contract errors raise ValueError.
|
||||
Database failures propagate, rather than being treated as missing data.
|
||||
"""
|
||||
if status.workplace != self.workplace:
|
||||
if status.workplace.strip().casefold() != self.workplace.strip().casefold():
|
||||
raise ValueError("ERP workplace does not match configured workplace")
|
||||
order = build_enlyze_production_order(status.production_order, self.format_template)
|
||||
quantity = status.good_quantity_m2
|
||||
@@ -88,7 +90,8 @@ class MaterialEfficiencyService:
|
||||
return None
|
||||
return MaterialEfficiencySnapshot(
|
||||
workplace=status.workplace, machine_id=self.machine_id,
|
||||
calculation_id=self.calculation_id, erp_production_order=status.production_order,
|
||||
calculation_id=self.output_calculation_id,
|
||||
erp_production_order=status.production_order,
|
||||
enlyze_production_order=order, article_number=status.article_number,
|
||||
article_description=status.article_description,
|
||||
nominal_width_m=extract_nominal_width_m(status.article_description),
|
||||
|
||||
@@ -45,10 +45,10 @@ def build_material_efficiency_runner(
|
||||
if (
|
||||
not isinstance(document, dict)
|
||||
or not required <= document.keys()
|
||||
or document.keys() - required - {"poll_interval_seconds"}
|
||||
or document.keys() - required - {"poll_interval_seconds", "output_calculation_id"}
|
||||
):
|
||||
raise CalculationConfigError("Invalid material efficiency configuration fields")
|
||||
for name in required:
|
||||
for name in required | ({"output_calculation_id"} & document.keys()):
|
||||
value = document[name]
|
||||
if not isinstance(value, str) or not value.strip() or "\x00" in value:
|
||||
raise CalculationConfigError(f"{name} must be a non-empty string without NUL bytes")
|
||||
@@ -76,6 +76,7 @@ def build_material_efficiency_runner(
|
||||
calculation_id=document["calculation_id"],
|
||||
format_template=document["production_order_format"],
|
||||
erp_timezone=timezone,
|
||||
output_calculation_id=document.get("output_calculation_id"),
|
||||
),
|
||||
PostgresMaterialEfficiencyWriter(postgres),
|
||||
poll_interval_seconds=interval,
|
||||
|
||||
@@ -14,6 +14,7 @@ from production_analytics.enlyze.gateway import (
|
||||
EnlyzeApiGateway,
|
||||
EnlyzeProductionRun,
|
||||
)
|
||||
from production_analytics.service.postgres_material_application import MaterialApplicationWriter
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
@@ -51,23 +52,41 @@ 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,
|
||||
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
|
||||
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 +98,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,
|
||||
@@ -100,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,
|
||||
@@ -125,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,
|
||||
@@ -143,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 = {}
|
||||
@@ -151,6 +171,13 @@ 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,
|
||||
)
|
||||
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,
|
||||
@@ -159,4 +186,13 @@ class MaterialPollingService:
|
||||
end=end,
|
||||
**source_options,
|
||||
)
|
||||
return integrator.process_many(samples)
|
||||
state = integrator.process_many(samples)
|
||||
if self._process_application_calculation_id is not None:
|
||||
assert self._application_writer is not None and run is not None
|
||||
# Persist before advancing the checkpoint, so failed writes can be retried.
|
||||
self._application_writer.write(
|
||||
calculation_id=self._process_application_calculation_id,
|
||||
machine_id=self._machine_id, production_order=run.production_order,
|
||||
run_id=run.uuid, samples=samples,
|
||||
)
|
||||
return state
|
||||
|
||||
@@ -26,6 +26,9 @@ from production_analytics.service.postgres_material import (
|
||||
PostgresMaterialSnapshotWriter,
|
||||
PostgresSettings,
|
||||
)
|
||||
from production_analytics.service.postgres_material_application import (
|
||||
PostgresMaterialApplicationWriter,
|
||||
)
|
||||
|
||||
|
||||
class _ReportingStateStore:
|
||||
@@ -73,7 +76,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 +85,17 @@ 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,
|
||||
)
|
||||
if calculation.process_application_calculation_id is not None:
|
||||
source_options.update(
|
||||
process_application_calculation_id=calculation.process_application_calculation_id,
|
||||
application_writer=PostgresMaterialApplicationWriter(postgres_settings),
|
||||
)
|
||||
service = MaterialPollingService(
|
||||
gateway=EnlyzeApiGateway(ExplorationClient(settings)),
|
||||
state_store=_ReportingStateStore(JsonMaterialStateStore(state_directory)),
|
||||
|
||||
@@ -0,0 +1,53 @@
|
||||
"""Generic process application snapshots, timestamped at the source sample."""
|
||||
|
||||
from collections.abc import Sequence
|
||||
from typing import Protocol
|
||||
|
||||
import psycopg
|
||||
|
||||
from production_analytics.calculations.config import finite_number
|
||||
from production_analytics.calculations.material_consumption import MaterialSample
|
||||
from production_analytics.service.postgres_material import PostgresSettings
|
||||
|
||||
|
||||
class MaterialApplicationWriter(Protocol):
|
||||
def write(
|
||||
self, *, calculation_id: str, machine_id: str, production_order: str,
|
||||
run_id: str, samples: Sequence[MaterialSample],
|
||||
) -> None: ...
|
||||
|
||||
|
||||
class PostgresMaterialApplicationWriter:
|
||||
def __init__(self, settings: PostgresSettings) -> None:
|
||||
self.settings = settings
|
||||
|
||||
def write(
|
||||
self, *, calculation_id: str, machine_id: str, production_order: str,
|
||||
run_id: str, samples: Sequence[MaterialSample],
|
||||
) -> None:
|
||||
rows = []
|
||||
for sample in samples:
|
||||
if sample.application_g_m2 is None:
|
||||
continue
|
||||
if sample.timestamp.utcoffset() is None:
|
||||
raise ValueError("Application timestamp must be timezone-aware")
|
||||
finite_number(sample.application_g_m2, "application_g_m2")
|
||||
rows.append((sample.timestamp, calculation_id, machine_id, production_order,
|
||||
run_id, sample.application_g_m2))
|
||||
if not rows:
|
||||
return
|
||||
with psycopg.connect(
|
||||
host=self.settings.host, port=self.settings.port, dbname=self.settings.dbname,
|
||||
user=self.settings.user, password=self.settings.password, connect_timeout=10,
|
||||
options="-c statement_timeout=10000",
|
||||
) as connection:
|
||||
with connection.cursor() as cursor:
|
||||
cursor.executemany(
|
||||
"""INSERT INTO material_application_snapshots
|
||||
(timestamp, calculation_id, machine_id, production_order,
|
||||
run_id, application_g_m2)
|
||||
VALUES (%s, %s, %s, %s, %s, %s)
|
||||
ON CONFLICT (calculation_id, machine_id, production_order, timestamp)
|
||||
DO NOTHING""",
|
||||
rows,
|
||||
)
|
||||
@@ -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 * 2.75 * 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 * 2.75 * 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)
|
||||
|
||||
@@ -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
|
||||
@@ -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'])
|
||||
@@ -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)
|
||||
|
||||
@@ -0,0 +1,122 @@
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from pathlib import Path
|
||||
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]['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
|
||||
|
||||
|
||||
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._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 == 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'
|
||||
erp_settings.from_secret_file.assert_called_once_with(tmp_path / 'erp.env')
|
||||
Reference in New Issue
Block a user