Compare commits

..
6 Commits
32 changed files with 3114 additions and 41 deletions
+176 -3
View File
@@ -249,9 +249,127 @@ semantics. `feedback_timestamp` is preserved exactly, including a naive timezone
It represents the latest ERP feedback, and the ERP export is delayed relative to It represents the latest ERP feedback, and the ERP export is delayed relative to
live process data; the adapter makes no freshness inference. live process data; the adapter makes no freshness inference.
This milestone adds no kg/m² or material-efficiency calculation, ERP-to-ENLYZE Automated ERP tests use mocks and require no live connectivity.
order mapping, persistence, polling, CLI, or Grafana integration. Automated ERP
tests use mocks and require no live connectivity. ## ERP context normalization
Two pure helpers in `production_analytics.context` prepare context for later KPIs:
- `build_enlyze_production_order(erp_production_order, format_template) -> str`
uses an explicit template configured by the caller per machine/workplace.
The template requires exactly one literal `{production_order}` placeholder and
no other braces; invalid templates raise `ValueError`.
Outer ERP-order whitespace is stripped and the
remaining order must contain only ASCII digits. Leading zeros are preserved.
Empty/invalid orders raise `ValueError`. There is no default machine rule.
- `extract_nominal_width_m(article_description: str | None) -> float | None`
is generic across machines and conservatively reads a single
`<width> x <length> m` pair, accepting decimal
comma/point and variable whitespace. For example, `Stex R 1501 C (PR) 5,80 x 50 m`
yields `5.8`. Earlier article numbers and trailing descriptive text are ignored.
Missing, malformed, multiple or chained dimension patterns return `None`, as do
non-positive/non-finite dimensions. Signed/scientific notation and other units
are deliberately unsupported; both dimensions must be positive plain numbers.
K7 currently uses the explicit template `K 7-{production_order}`:
```python
k7_order_template = "K 7-{production_order}"
enlyze_order = build_enlyze_production_order("12026000815", k7_order_template)
# "K 7-12026000815"
```
Future machines may supply different verified templates without changing the
generic helper. Template selection belongs to machine/workplace configuration
or adapters; the helper contains no workplace lookup or machine-specific branch.
No other machine rule is introduced here.
ENLYZE production-order identifiers remain opaque everywhere else, with exact
comparison. These helpers never reverse-parse or split ENLYZE orders, including
combined identifiers such as `K 7-12025000074-K 7-12025000075`; passing such an
identifier to the ERP mapping function is rejected.
Nominal finished-product width currently comes from ERP article-description
parsing because neither verified DWH view (`dbo.GRAFANA_WORKPLACE_STATUS` and
`dbo.GRAFANA_PRODUCTION_CONFIRMATION`) has a dedicated width column. ENLYZE
`wLg1MeasuringWidth` (`Produktbreite`) is explicitly not used as K7 nominal
finished-product width:
it belongs to the MAHLO measurement system and observed values differ from nominal
article widths. Structured ERP product master data would be preferred when available.
These helpers are independent of database access. The service below composes them
for material efficiency; no background processing is attached.
## Feedback-aligned material efficiency
`MaterialEfficiencyService` in `service.material_efficiency` combines ERP good area
with cumulative material consumption for one exact production order. Its
`evaluate_current()` reads the existing ERP workplace adapter once; `evaluate(status)`
evaluates an already supplied `CurrentWorkplaceStatus`. Both return an immutable
`MaterialEfficiencySnapshot` or `None` when inputs are unavailable.
`PostgresMaterialSnapshotRepository(settings).latest_at_or_before(...)` in
`service.postgres_material` takes keyword arguments `calculation_id`, `machine_id`,
`production_order`, and an aware `timestamp`. A parameterized query selects the
latest snapshot **at or before ERP feedback time**, matching calculation, machine,
and mapped order exactly. It returns `MaterialConsumptionSnapshot` or `None`.
Run IDs do not restrict this lookup: stored consumption is already cumulative
across the order's runs. The service neither reintegrates nor sums snapshots and
does not bridge run boundaries.
The result includes workplace/machine/calculation identifiers, both order identifiers,
article context, nominal width, the UTC ERP feedback instant and selected material timestamp, good area, aligned
consumption, `material_consumption_kg_per_m2 = consumption_kg / good_quantity_m2`,
and `material_consumption_g_per_m2 = material_consumption_kg_per_m2 * 1000`.
Nominal width comes from generic ERP description parsing and is context only;
ERP good square metres remain authoritative even when width cannot be parsed.
Missing ERP status, missing aligned material, missing/non-positive/non-finite good
area, invalid consumption (negative or non-finite), or non-finite computed ratios
produce `None`. Zero consumption is valid. Configuration errors, invalid timestamp
alignment, and database failures raise rather than masquerading as missing data.
All machine/workplace settings are supplied explicitly. Example using the existing
K7 material calculation (with previously constructed ERP gateway and PG settings):
```python
from zoneinfo import ZoneInfo
from production_analytics.service.material_efficiency import MaterialEfficiencyService
from production_analytics.service.postgres_material import PostgresMaterialSnapshotRepository
service = MaterialEfficiencyService(
erp_gateway, PostgresMaterialSnapshotRepository(postgres_settings),
workplace="K7",
machine_id="c220f95c-a65e-4cb7-99b7-0626d6c7508c",
calculation_id="k7-fiber-consumption",
format_template="K 7-{production_order}",
# Supply only after verifying the ERP timestamp's source timezone:
erp_timezone=ZoneInfo("Europe/Berlin"),
)
result = service.evaluate_current()
```
Aware ERP timestamps are compared as UTC instants. Naive ERP timestamps require
explicit `erp_timezone`; there is no inferred system/database timezone. Ambiguous
or nonexistent local times during DST transitions are rejected. The result carries the
ERP feedback instant normalized to UTC and the selected aware material timestamp.
The original ERP status object remains unchanged.
This alignment avoids knowingly including future material consumption, but does
not eliminate ERP roll-feedback timing uncertainty (feedback may be roughly one
roll ahead or otherwise offset). Live values are plausibility indicators; final
production-order values are more meaningful as relative timing error diminishes.
There is no interpolation, lag correction, smoothing, freshness threshold, or
estimated timestamp. Historical bootstrap snapshots are usable only when their
stored timestamps satisfy the cutoff; a later cumulative total cannot reconstruct
an earlier value. Reliable final evaluation requires retaining final ERP feedback;
the current-workplace adapter alone does not provide historical completed orders.
The service itself returns results in memory without modifying material snapshots.
The standalone persistence runner is described below. Daily per-machine 24h reporting
and aggregation remain future work. Repository SQL tests use an in-memory SQLite
fixture with driver transport adaptation; they require no live PostgreSQL.
## Peak-cycle detection ## Peak-cycle detection
@@ -293,3 +411,58 @@ The normal unit test suite mocks PostgreSQL and requires no live database.
See [PROJECT_KNOWLEDGE.md](PROJECT_KNOWLEDGE.md) for durable project context See [PROJECT_KNOWLEDGE.md](PROJECT_KNOWLEDGE.md) for durable project context
and [docs/roadmap.md](docs/roadmap.md) for the implementation sequence. and [docs/roadmap.md](docs/roadmap.md) for the implementation sequence.
### Feedback-driven material-efficiency persistence
Apply the updated `db/schema.sql` using the PostgreSQL setup command above, then run
this independent foreground process:
```bash
production-analytics run material-efficiency --config config/k7-material-efficiency.yaml
```
`poll_interval_seconds` in YAML defaults to 60 seconds (fixed delay after each
cycle). Workplace, machine, calculation, order format and `erp_timezone` are also
configured in YAML; K7 uses `Europe/Berlin`. PostgreSQL settings use the existing
exported `POSTGRES_*` variables and `--secrets-file` (default `secrets/enlyze.env`);
ERP uses `--erp-secrets-file` (default `secrets/erp.env`). No additional secrets are
needed, and `.env` is not loaded automatically.
`PostgresMaterialEfficiencyWriter.write(snapshot) -> bool` validates aware
timestamps and finite numbers, then inserts all KPI fields in a short transaction.
`material_efficiency_snapshots` retains one immutable point per
`(calculation_id, machine_id, enlyze_production_order, erp_feedback_timestamp)`.
The writer uses `RETURNING 1` to return `True` for a new insert and `False` for a
duplicate. Only new inserts produce a flushed stdout line with feedback time,
workplace, order, good m², material kg and g/m². Unavailable evaluations and
duplicates remain silent.
Repeated evaluations attempt `ON CONFLICT DO NOTHING`; they never update history,
including after a restart. Missing/invalid KPI inputs produce no row and can be
retried on a later cycle. Nullable article context remains SQL NULL.
Persistence frequency follows ERP feedback changes, not ENLYZE sample frequency.
This runner is independent of the 10-second material polling process. Live values
are plausibility indicators; final production-order (FA) values are more meaningful.
Grafana can read PostgreSQL alone:
```sql
SELECT erp_feedback_timestamp AS "time", material_consumption_g_per_m2
FROM material_efficiency_snapshots
WHERE $__timeFilter(erp_feedback_timestamp)
AND calculation_id = 'k7-fiber-consumption'
AND machine_id = 'c220f95c-a65e-4cb7-99b7-0626d6c7508c'
ORDER BY erp_feedback_timestamp;
```
ERP read errors and transient PostgreSQL connection/operational errors are reported
using error classes only and retried after the configured delay. Each database
operation opens a fresh connection. Configuration, schema/programming and invalid
timestamp contract errors terminate with a nonzero CLI exit; ambiguous/nonexistent
DST feedback remains rejected. Ctrl-C stops the foreground runner cleanly.
The current-status ERP source cannot backfill feedback events missed between polls
or during outages, nor guarantee observation of final FA feedback before the order
changes. Daily 24h reports and final order summary tables remain future work.
No dashboard, aggregation or lag correction is added here.
Bento 1: [fresh bentonite consumption configuration and scope](docs/bento1-fresh-bentonite.md).
+21
View File
@@ -0,0 +1,21 @@
# Fresh bentonite only: the recycled/recovered third spreader is excluded.
calculations:
- id: bento1-fresh-bentonite-consumption
type: material_consumption
version: "1"
machine_ref: 5f42a4f6-9ca0-4f6f-9786-40d50a35b230
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
erp_workplace: Bento 1
production_order_format: "Bento 1-{production_order}"
output_metric: material_consumption
output_unit: kg
group_by: production_order
+7
View File
@@ -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
+6
View File
@@ -0,0 +1,6 @@
workplace: K7
machine_id: c220f95c-a65e-4cb7-99b7-0626d6c7508c
calculation_id: k7-fiber-consumption
production_order_format: "K 7-{production_order}"
erp_timezone: Europe/Berlin
poll_interval_seconds: 60
+14
View File
@@ -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
+35
View File
@@ -11,3 +11,38 @@ CREATE TABLE IF NOT EXISTS material_consumption_snapshots (
-- Machine time series across orders; the primary key covers a specific order. -- Machine time series across orders; the primary key covers a specific order.
CREATE INDEX IF NOT EXISTS material_consumption_snapshots_machine_time_idx CREATE INDEX IF NOT EXISTS material_consumption_snapshots_machine_time_idx
ON material_consumption_snapshots (calculation_id, machine_id, timestamp); ON material_consumption_snapshots (calculation_id, machine_id, timestamp);
CREATE TABLE IF NOT EXISTS material_efficiency_snapshots (
erp_feedback_timestamp timestamptz NOT NULL,
material_snapshot_timestamp timestamptz NOT NULL,
calculation_id text NOT NULL,
machine_id text NOT NULL,
workplace text NOT NULL,
erp_production_order text NOT NULL,
enlyze_production_order text NOT NULL,
article_number text NULL,
article_description text NULL,
nominal_width_m double precision NULL,
good_quantity_m2 double precision NOT NULL,
material_consumption_kg double precision NOT NULL,
material_consumption_kg_per_m2 double precision NOT NULL,
material_consumption_g_per_m2 double precision NOT NULL,
PRIMARY KEY (calculation_id, machine_id, enlyze_production_order, erp_feedback_timestamp)
);
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);
+250
View File
@@ -0,0 +1,250 @@
# 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**
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.
+73
View File
@@ -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
@@ -1,6 +1,6 @@
"""Strict configuration for executable material-consumption instances.""" """Strict configuration for executable material-consumption instances."""
from dataclasses import dataclass, fields from dataclasses import MISSING, dataclass, fields
from math import isfinite from math import isfinite
from pathlib import Path from pathlib import Path
@@ -48,12 +48,20 @@ class MaterialCalculationConfig:
version: str version: str
machine_ref: str machine_ref: str
rate_signal_ref: str rate_signal_ref: str
gate_signal_ref: str gate_signal_ref: str | None
gate_threshold: float gate_threshold: float
max_sample_gap_seconds: float max_sample_gap_seconds: float
output_metric: str output_metric: str
output_unit: str output_unit: str
group_by: str 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( def load_material_calculation(
@@ -71,20 +79,32 @@ def load_material_calculation(
entries = document["calculations"] entries = document["calculations"]
if not isinstance(entries, list) or not entries: if not isinstance(entries, list) or not entries:
raise CalculationConfigError("calculations must be a non-empty list") raise CalculationConfigError("calculations must be a non-empty list")
required = {field.name for field in fields(MaterialCalculationConfig)} allowed = {field.name for field in fields(MaterialCalculationConfig)}
required = {field.name for field in fields(MaterialCalculationConfig)
if field.default is MISSING}
instances = {} instances = {}
calibrations = None
for index, entry in enumerate(entries): for index, entry in enumerate(entries):
prefix = f"calculations[{index}]" prefix = f"calculations[{index}]"
if not isinstance(entry, dict): if not isinstance(entry, dict):
raise CalculationConfigError(f"{prefix} must be a mapping") raise CalculationConfigError(f"{prefix} must be a mapping")
if "type" in entry and entry["type"] != "material_consumption": if "type" in entry and entry["type"] != "material_consumption":
raise CalculationConfigError(f"{prefix}.type: only 'material_consumption' is supported") raise CalculationConfigError(f"{prefix}.type: only 'material_consumption' is supported")
area = entry.get("source_mode", "direct_mass_rate") == "area_application"
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() missing = required - entry.keys()
if missing: if missing:
raise CalculationConfigError(f"{prefix}: missing fields: {', '.join(sorted(missing))}") raise CalculationConfigError(f"{prefix}: missing fields: {', '.join(sorted(missing))}")
if entry.keys() - required: if entry.keys() - allowed:
raise CalculationConfigError(f"{prefix}: unknown fields (check spelling)") raise CalculationConfigError(f"{prefix}: unknown fields (check spelling)")
for name in sorted(required - {"gate_threshold", "max_sample_gap_seconds"}): 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(): if not isinstance(entry[name], str) or not entry[name].strip():
raise CalculationConfigError(f"{prefix}.{name} must be a non-empty string") raise CalculationConfigError(f"{prefix}.{name} must be a non-empty string")
for name, expected in ( for name, expected in (
@@ -98,9 +118,82 @@ def load_material_calculation(
entry["max_sample_gap_seconds"] = finite_number( entry["max_sample_gap_seconds"] = finite_number(
entry["max_sample_gap_seconds"], f"{prefix}.max_sample_gap_seconds", positive=True, 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", "rotational_discharge",
}:
raise CalculationConfigError("Unsupported material source_mode")
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(
f"{refs_name} must be unique signal strings"
)
if entry["rate_signal_ref"] != "unused":
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"{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[refs_name] = tuple(refs)
elif any(name in entry for name in (
"application_signal_refs", "rotational_speed_signal_refs",
"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"),
"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: if entry["id"] in instances:
raise CalculationConfigError("Calculation ids must be unique") raise CalculationConfigError("Calculation ids must be unique")
instances[entry["id"]] = MaterialCalculationConfig(**entry) 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 is not None:
if calculation_id not in instances: if calculation_id not in instances:
raise CalculationConfigError("Requested calculation id was not found") 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 timestamp: datetime
material_rate_kg_per_hour: float material_rate_kg_per_hour: float
gate_value: float gate_value: float
application_g_m2: float | None = None
@dataclass(frozen=True, slots=True) @dataclass(frozen=True, slots=True)
@@ -102,3 +103,38 @@ class MaterialConsumptionIntegrator:
for sample in samples: for sample in samples:
self.process(sample) self.process(sample)
return self.state return self.state
def area_application_rate_kg_per_hour(
applications_g_m2: Iterable[float], nominal_width_m: float, line_speed_m_min: float,
) -> float:
"""Normalize summed area application to the integrator's canonical kg/h."""
applications = tuple(applications_g_m2)
if not applications or not all(isfinite(value) for value in applications):
raise ValueError("application values must be non-empty and finite")
if not isfinite(nominal_width_m) or nominal_width_m <= 0:
raise ValueError("nominal width must be finite and positive")
if not isfinite(line_speed_m_min):
raise ValueError("line speed must be finite")
rate = sum(applications) * nominal_width_m * line_speed_m_min * 60 / 1000
if not isfinite(rate):
raise ValueError("derived material rate must be finite")
return rate
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
+98 -18
View File
@@ -37,10 +37,18 @@ def _parser() -> argparse.ArgumentParser:
raw.add_argument("--query", action="append", type=_query_item, default=[], metavar="KEY=VALUE") raw.add_argument("--query", action="append", type=_query_item, default=[], metavar="KEY=VALUE")
raw.add_argument("--pretty", action="store_true", help="Pretty-print sanitized JSON") raw.add_argument("--pretty", action="store_true", help="Pretty-print sanitized JSON")
raw.add_argument("--verbose", action="store_true", help="Show safe request progress on stderr") raw.add_argument("--verbose", action="store_true", help="Show safe request progress on stderr")
raw.add_argument("--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/") raw.add_argument(
raw.add_argument("--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/") "--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/"
raw.add_argument("--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS) )
raw.add_argument("--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS) raw.add_argument(
"--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/"
)
raw.add_argument(
"--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS
)
raw.add_argument(
"--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS
)
raw.add_argument( raw.add_argument(
"--secrets-file", "--secrets-file",
type=Path, type=Path,
@@ -52,18 +60,44 @@ def _parser() -> argparse.ArgumentParser:
timeseries.add_argument("--start", required=True, help="ISO 8601 datetime with timezone") timeseries.add_argument("--start", required=True, help="ISO 8601 datetime with timezone")
timeseries.add_argument("--end", required=True, help="ISO 8601 datetime with timezone") timeseries.add_argument("--end", required=True, help="ISO 8601 datetime with timezone")
timeseries.add_argument("--variable", required=True, help="Variable UUID") timeseries.add_argument("--variable", required=True, help="Variable UUID")
timeseries.add_argument("--resampling-interval", type=int, help="Seconds; schema range is 10..604800") timeseries.add_argument(
"--resampling-interval", type=int, help="Seconds; schema range is 10..604800"
)
timeseries.add_argument( timeseries.add_argument(
"--resampling-method", "--resampling-method",
choices=["first", "last", "max", "min", "count", "sum", "avg", "median", "std", "q5", "q25", "q75", "q95"], choices=[
"first",
"last",
"max",
"min",
"count",
"sum",
"avg",
"median",
"std",
"q5",
"q25",
"q75",
"q95",
],
help="Optional schema-defined method for this variable", help="Optional schema-defined method for this variable",
) )
timeseries.add_argument("--pretty", action="store_true", help="Pretty-print sanitized JSON") timeseries.add_argument("--pretty", action="store_true", help="Pretty-print sanitized JSON")
timeseries.add_argument("--verbose", action="store_true", help="Show safe request progress on stderr") timeseries.add_argument(
timeseries.add_argument("--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/") "--verbose", action="store_true", help="Show safe request progress on stderr"
timeseries.add_argument("--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/") )
timeseries.add_argument("--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS) timeseries.add_argument(
timeseries.add_argument("--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS) "--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/"
)
timeseries.add_argument(
"--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/"
)
timeseries.add_argument(
"--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS
)
timeseries.add_argument(
"--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS
)
timeseries.add_argument( timeseries.add_argument(
"--secrets-file", "--secrets-file",
type=Path, type=Path,
@@ -78,6 +112,11 @@ def _parser() -> argparse.ArgumentParser:
material.add_argument("--poll-interval-seconds", type=float, default=10.0) material.add_argument("--poll-interval-seconds", type=float, default=10.0)
material.add_argument("--state-directory", type=Path, default=Path("data/state/material")) material.add_argument("--state-directory", type=Path, default=Path("data/state/material"))
material.add_argument("--secrets-file", type=Path, default=Path("secrets/enlyze.env")) material.add_argument("--secrets-file", type=Path, default=Path("secrets/enlyze.env"))
material.add_argument("--erp-secrets-file", type=Path, default=Path("secrets/erp.env"))
efficiency = runners.add_parser("material-efficiency", help="Persist ERP feedback KPI points")
efficiency.add_argument("--config", type=Path, required=True)
efficiency.add_argument("--secrets-file", type=Path, default=Path("secrets/enlyze.env"))
efficiency.add_argument("--erp-secrets-file", type=Path, default=Path("secrets/erp.env"))
return parser return parser
@@ -102,8 +141,11 @@ def _run_material(args: argparse.Namespace) -> int:
try: try:
calculation = load_material_calculation(args.config, args.calculation_id) calculation = load_material_calculation(args.config, args.calculation_id)
runner = build_material_runner( runner = build_material_runner(
calculation, poll_interval_seconds=args.poll_interval_seconds, calculation,
state_directory=args.state_directory, secrets_file=args.secrets_file, poll_interval_seconds=args.poll_interval_seconds,
state_directory=args.state_directory,
erp_secrets_file=args.erp_secrets_file,
secrets_file=args.secrets_file,
) )
runner.run() runner.run()
except (CalculationConfigError, ConfigurationError) as exc: except (CalculationConfigError, ConfigurationError) as exc:
@@ -117,9 +159,35 @@ def _run_material(args: argparse.Namespace) -> int:
return 0 return 0
def _run_material_efficiency(args: argparse.Namespace) -> int:
from production_analytics.calculations.config import CalculationConfigError
from production_analytics.enlyze.exploration import ConfigurationError
from production_analytics.service.material_efficiency_runtime import (
build_material_efficiency_runner,
)
try:
build_material_efficiency_runner(
args.config,
secrets_file=args.secrets_file,
erp_secrets_file=args.erp_secrets_file,
).run()
except (CalculationConfigError, ConfigurationError) as exc:
print(f"Config/startup error: {exc}", file=sys.stderr)
return 2
except Exception as exc:
print(f"Startup/runtime error: {type(exc).__name__}", file=sys.stderr)
return 1
except KeyboardInterrupt:
return 0
return 0
def main(argv: Sequence[str] | None = None) -> int: def main(argv: Sequence[str] | None = None) -> int:
args = _parser().parse_args(argv) args = _parser().parse_args(argv)
if args.namespace == "run": if args.namespace == "run":
if args.command == "material-efficiency":
return _run_material_efficiency(args)
return _run_material(args) return _run_material(args)
if args.namespace != "enlyze" or args.command not in {"raw", "timeseries"}: if args.namespace != "enlyze" or args.command not in {"raw", "timeseries"}:
return 2 return 2
@@ -133,7 +201,8 @@ def main(argv: Sequence[str] | None = None) -> int:
operation = "GET" if args.command == "raw" else "POST" operation = "GET" if args.command == "raw" else "POST"
target = args.path if args.command == "raw" else "/v2/timeseries" target = args.path if args.command == "raw" else "/v2/timeseries"
print( print(
f"Requesting {operation} {target} (timeout={settings.timeout_seconds:g}s; authentication {authentication}).", f"Requesting {operation} {target} (timeout={settings.timeout_seconds:g}s; "
f"authentication {authentication}).",
file=sys.stderr, file=sys.stderr,
) )
client = ExplorationClient(settings) client = ExplorationClient(settings)
@@ -141,8 +210,13 @@ def main(argv: Sequence[str] | None = None) -> int:
response = client.get(args.path, dict(args.query)) response = client.get(args.path, dict(args.query))
request = {"method": "GET", "path": response.path} request = {"method": "GET", "path": response.path}
else: else:
if args.resampling_interval is not None and not 10 <= args.resampling_interval <= 604800: if (
raise ExplorationError("--resampling-interval must be between 10 and 604800 seconds.") args.resampling_interval is not None
and not 10 <= args.resampling_interval <= 604800
):
raise ExplorationError(
"--resampling-interval must be between 10 and 604800 seconds."
)
variable: dict[str, str] = {"uuid": args.variable} variable: dict[str, str] = {"uuid": args.variable}
if args.resampling_method: if args.resampling_method:
variable["resampling_method"] = args.resampling_method variable["resampling_method"] = args.resampling_method
@@ -170,9 +244,15 @@ def main(argv: Sequence[str] | None = None) -> int:
indent = 2 if args.pretty or args.save_fixture else None indent = 2 if args.pretty or args.save_fixture else None
print(json.dumps(payload, indent=indent, sort_keys=True)) print(json.dumps(payload, indent=indent, sort_keys=True))
if args.save_fixture: if args.save_fixture:
print(f"Sanitized fixture written to {_write_fixture(args.fixture_dir, args.save_fixture, payload)}") print(
"Sanitized fixture written to "
f"{_write_fixture(args.fixture_dir, args.save_fixture, payload)}"
)
if args.save_raw: if args.save_raw:
print(f"Raw response written to {_write_fixture(args.raw_dir, args.save_raw, raw_payload)}") print(
"Raw response written to "
f"{_write_fixture(args.raw_dir, args.save_raw, raw_payload)}"
)
except ExplorationError as error: except ExplorationError as error:
print(f"ENLYZE exploration error: {error}") print(f"ENLYZE exploration error: {error}")
return 1 return 1
@@ -0,0 +1,62 @@
"""Deterministic ERP context helpers, independent of database and timeseries access."""
import math
import re
__all__ = ["build_enlyze_production_order", "extract_nominal_width_m"]
# Capture malformed numeric tokens too, so a valid-looking suffix cannot become width.
_NUMBER_TOKEN = r"[+-]?(?:[0-9]+(?:[.,][0-9]+)*|inf(?:inity)?|nan)(?:[eE][+-]?[0-9]+)?"
_DIMENSIONS = re.compile(
rf"(?<![\w.,+\-/])(?P<width>{_NUMBER_TOKEN})\s*x\s*"
rf"(?P<length>{_NUMBER_TOKEN})\s*m(?![\w/²³^])",
re.IGNORECASE,
)
_PLAIN_NUMBER = re.compile(r"[0-9]+(?:[.,][0-9]+)?")
def build_enlyze_production_order(erp_production_order: str, format_template: str) -> str:
"""Format a single ERP order using explicit caller-supplied configuration.
The template must contain exactly one literal {production_order} placeholder
and no other braces. Outer order whitespace is stripped; the order
must otherwise contain ASCII digits only (leading zeros are preserved).
This is not a parser for ENLYZE identifiers, combined or otherwise.
"""
order = erp_production_order.strip()
if not order or not order.isascii() or not order.isdigit():
raise ValueError("ERP production order must be a non-empty ASCII digit string")
placeholder = "{production_order}"
literal = format_template.replace(placeholder, "")
if format_template.count(placeholder) != 1 or "{" in literal or "}" in literal:
raise ValueError("Format template must contain exactly one {production_order} placeholder")
return format_template.replace(placeholder, order)
def extract_nominal_width_m(article_description: str | None) -> float | None:
"""Read one unambiguous positive '<width> x <length> m' pair from ERP text.
Decimal comma/point and variable whitespace are supported. Multiple pairs,
dimension chains, signed/scientific notation and invalid dimensions fail
closed.
"""
if article_description is None:
return None
matches = list(_DIMENSIONS.finditer(article_description))
if len(matches) != 1:
return None
match = matches[0]
# Do not mistake the tail of a three-dimensional expression for a width pair.
if re.search(r"[x×]\s*$", article_description[:match.start()], re.IGNORECASE):
return None
if re.match(r"\s*[x×]", article_description[match.end():], re.IGNORECASE):
return None
values = []
for token in (match["width"], match["length"]):
if _PLAIN_NUMBER.fullmatch(token) is None:
return None
value = float(token.replace(",", "."))
if not math.isfinite(value) or value <= 0:
return None
values.append(value)
return values[0]
+51 -11
View File
@@ -4,7 +4,11 @@ from dataclasses import dataclass
from datetime import UTC, datetime from datetime import UTC, datetime
from math import isfinite from math import isfinite
from production_analytics.calculations.material_consumption import MaterialSample from production_analytics.calculations.material_consumption import (
MaterialSample,
area_application_rate_kg_per_hour,
rotational_discharge_rate_kg_per_hour,
)
from production_analytics.enlyze.exploration import ExplorationClient from production_analytics.enlyze.exploration import ExplorationClient
@@ -86,23 +90,34 @@ class EnlyzeApiGateway:
*, *,
machine_id: str, machine_id: str,
rate_variable_id: str, rate_variable_id: str,
gate_variable_id: str, gate_variable_id: str | None,
start: datetime, start: datetime,
end: 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]: ) -> list[MaterialSample]:
for name, value in (("start", start), ("end", end)): for name, value in (("start", start), ("end", end)):
if value.tzinfo is None or value.utcoffset() is None: if value.tzinfo is None or value.utcoffset() is None:
raise ValueError(f"{name} must be timezone-aware") raise ValueError(f"{name} must be timezone-aware")
if end < start: if end < start:
raise ValueError("end must not precede start") raise ValueError("end must not precede start")
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 = { request_body = {
"machine": machine_id, "machine": machine_id,
"start": start.astimezone(UTC).isoformat(), "start": start.astimezone(UTC).isoformat(),
"end": end.astimezone(UTC).isoformat(), "end": end.astimezone(UTC).isoformat(),
"variables": [ "variables": [{"uuid": variable_id} for variable_id in variable_ids],
{"uuid": rate_variable_id},
{"uuid": gate_variable_id},
],
} }
samples = [] samples = []
followed_cursors: set[str] = set() followed_cursors: set[str] = set()
@@ -114,27 +129,52 @@ class EnlyzeApiGateway:
columns = data["columns"] columns = data["columns"]
if not isinstance(columns, list): if not isinstance(columns, list):
raise ValueError("columns must be a list") raise ValueError("columns must be a list")
for required in ("time", rate_variable_id, gate_variable_id): for required in ("time", *variable_ids):
if columns.count(required) != 1: if columns.count(required) != 1:
raise ValueError(f"required column {required!r} must occur exactly once") raise ValueError(f"required column {required!r} must occur exactly once")
time_index = columns.index("time") time_index = columns.index("time")
rate_index = columns.index(rate_variable_id) source_indices = [columns.index(ref) for ref in source_ids]
gate_index = columns.index(gate_variable_id) gate_index = columns.index(gate_variable_id) if gate_variable_id else None
if not isinstance(data["records"], list): if not isinstance(data["records"], list):
raise ValueError("records must be a list") raise ValueError("records must be a list")
for index, record in enumerate(data["records"]): for index, record in enumerate(data["records"]):
try: try:
if not isinstance(record, list) or len(record) != len(columns): if not isinstance(record, list) or len(record) != len(columns):
raise ValueError("record must match columns") raise ValueError("record must match columns")
rate, gate = record[rate_index], record[gate_index] sources = [record[i] for i in source_indices]
for value in (rate, gate): 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)): if isinstance(value, bool) or not isinstance(value, (int, float)):
raise ValueError("rate and gate must be numeric") raise ValueError("rate and gate must be numeric")
if not isfinite(value): if not isfinite(value):
raise ValueError("rate and gate must be finite") raise ValueError("rate and gate must be finite")
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( samples.append(MaterialSample(
timestamp=_parse_timestamp(record[time_index]), timestamp=_parse_timestamp(record[time_index]),
material_rate_kg_per_hour=float(rate), gate_value=float(gate), material_rate_kg_per_hour=float(rate), gate_value=float(gate),
application_g_m2=application,
)) ))
except (TypeError, ValueError) as exc: except (TypeError, ValueError) as exc:
raise ValueError(f"malformed record {index}: {exc}") from exc raise ValueError(f"malformed record {index}: {exc}") from exc
@@ -0,0 +1,25 @@
"""ERP nominal-width resolution for an explicitly matched production order."""
from production_analytics.context import build_enlyze_production_order, extract_nominal_width_m
from production_analytics.erp import ErpWorkplaceStatusGateway
class ErpNominalWidthProvider:
def __init__(
self, gateway: ErpWorkplaceStatusGateway, *, workplace: str, format_template: str,
) -> None:
self.gateway = gateway
self.workplace = workplace
self.format_template = format_template
def __call__(self, production_order: str) -> float:
status = self.gateway.get_current_workplace_status(self.workplace)
if (status is None
or status.workplace.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")
width = extract_nominal_width_m(status.article_description)
if width is None:
raise ValueError("ERP article description has no unambiguous nominal width")
return width
@@ -0,0 +1,119 @@
"""Feedback-aligned, machine-independent production-order material efficiency."""
from dataclasses import dataclass
from datetime import UTC, datetime
from math import isfinite
from typing import Protocol
from zoneinfo import ZoneInfo
from production_analytics.context import build_enlyze_production_order, extract_nominal_width_m
from production_analytics.erp import CurrentWorkplaceStatus
from production_analytics.service.postgres_material import MaterialSnapshotRepository
class WorkplaceStatusReader(Protocol):
def get_current_workplace_status(self, workplace: str) -> CurrentWorkplaceStatus | None: ...
@dataclass(frozen=True, slots=True)
class MaterialEfficiencySnapshot:
workplace: str
machine_id: str
calculation_id: str
erp_production_order: str
enlyze_production_order: str
article_number: str | None
article_description: str | None
nominal_width_m: float | None
erp_feedback_timestamp: datetime
material_snapshot_timestamp: datetime
good_quantity_m2: float
material_consumption_kg: float
material_consumption_kg_per_m2: float
material_consumption_g_per_m2: float
class MaterialEfficiencyService:
def __init__(
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
def evaluate_current(self) -> MaterialEfficiencySnapshot | None:
"""Read current ERP feedback once; no scheduling or mutable KPI history."""
status = self.erp.get_current_workplace_status(self.workplace)
return None if status is None else self.evaluate(status)
def evaluate(self, status: CurrentWorkplaceStatus) -> MaterialEfficiencySnapshot | None:
"""Evaluate supplied feedback; missing/invalid numeric inputs yield None.
Configuration, timestamp and adapter contract errors raise ValueError.
Database failures propagate, rather than being treated as missing data.
"""
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
if quantity is None or not isfinite(quantity) or quantity <= 0:
return None
cutoff = _feedback_instant(status.feedback_timestamp, self.erp_timezone)
material = self.materials.latest_at_or_before(
calculation_id=self.calculation_id, machine_id=self.machine_id,
production_order=order, timestamp=cutoff,
)
if material is None:
return None
if (
material.calculation_id != self.calculation_id
or material.machine_id != self.machine_id
or material.production_order != order
or material.timestamp.utcoffset() is None
or material.timestamp > cutoff
):
raise ValueError("Material repository returned an unaligned snapshot")
consumption = material.consumption_kg
if not isfinite(consumption) or consumption < 0:
return None
kg_per_m2 = consumption / quantity
g_per_m2 = kg_per_m2 * 1000
if not isfinite(kg_per_m2) or not isfinite(g_per_m2):
return None
return MaterialEfficiencySnapshot(
workplace=status.workplace, machine_id=self.machine_id,
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),
erp_feedback_timestamp=cutoff,
material_snapshot_timestamp=material.timestamp,
good_quantity_m2=quantity, material_consumption_kg=consumption,
material_consumption_kg_per_m2=kg_per_m2,
material_consumption_g_per_m2=g_per_m2,
)
def _feedback_instant(timestamp: datetime, timezone: ZoneInfo | None) -> datetime:
if timestamp.utcoffset() is not None:
return timestamp.astimezone(UTC)
if timezone is None:
raise ValueError("Naive ERP feedback timestamp requires explicit erp_timezone")
# Round trips reject DST gaps; two distinct instants indicate a DST overlap.
instants = set()
for fold in (0, 1):
instant = timestamp.replace(tzinfo=timezone, fold=fold).astimezone(UTC)
if instant.astimezone(timezone).replace(tzinfo=None) == timestamp:
instants.add(instant)
if len(instants) != 1:
raise ValueError("ERP feedback timestamp is ambiguous or nonexistent in erp_timezone")
return instants.pop()
@@ -0,0 +1,60 @@
"""Independent foreground ERP feedback polling; PostgreSQL owns deduplication."""
import sys
import time
from collections.abc import Callable
from typing import TextIO
import psycopg
from production_analytics.calculations.config import finite_number
from production_analytics.erp import ErpReadError
from production_analytics.service.material_efficiency import MaterialEfficiencyService
from production_analytics.service.postgres_material_efficiency import MaterialEfficiencyWriter
class MaterialEfficiencyRunner:
def __init__(
self,
service: MaterialEfficiencyService,
writer: MaterialEfficiencyWriter,
*,
poll_interval_seconds: float = 60,
sleep: Callable[[float], None] = time.sleep,
stderr: TextIO | None = None,
stdout: TextIO | None = None,
) -> None:
self.service = service
self.writer = writer
self.interval = finite_number(poll_interval_seconds, "poll interval", positive=True)
self.sleep = sleep
self.stderr = stderr if stderr is not None else sys.stderr
self.stdout = stdout if stdout is not None else sys.stdout
def run(self) -> None:
try:
while True:
try:
snapshot = self.service.evaluate_current()
if snapshot is not None and self.writer.write(snapshot):
print(
f"{snapshot.erp_feedback_timestamp.isoformat()} "
f"workplace={snapshot.workplace!r} "
f"production_order={snapshot.enlyze_production_order!r} "
f"good_m2={snapshot.good_quantity_m2:.3f} "
f"material_kg={snapshot.material_consumption_kg:.3f} "
f"g_per_m2={snapshot.material_consumption_g_per_m2:.3f}",
file=self.stdout,
flush=True,
)
except (ErpReadError, psycopg.OperationalError, psycopg.InterfaceError) as exc:
# Never include driver messages, SQL, credentials or response bodies.
print(
f"Material efficiency cycle failed: {type(exc).__name__}",
file=self.stderr,
flush=True,
)
# Fixed delay even after unavailable evaluations or transient failures.
self.sleep(self.interval)
except KeyboardInterrupt:
return
@@ -0,0 +1,83 @@
"""Small strict YAML configuration and composition for KPI persistence."""
import os
from pathlib import Path
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
import yaml
from production_analytics.calculations.config import (
CalculationConfigError,
_UniqueLoader,
finite_number,
)
from production_analytics.context import build_enlyze_production_order
from production_analytics.enlyze.exploration import ConfigurationError, load_secret_file
from production_analytics.erp import ErpSettings, ErpWorkplaceStatusGateway
from production_analytics.service.material_efficiency import MaterialEfficiencyService
from production_analytics.service.material_efficiency_runner import MaterialEfficiencyRunner
from production_analytics.service.postgres_material import (
PostgresMaterialSnapshotRepository,
PostgresSettings,
)
from production_analytics.service.postgres_material_efficiency import (
PostgresMaterialEfficiencyWriter,
)
def build_material_efficiency_runner(
config: Path,
*,
secrets_file: Path = Path("secrets/enlyze.env"),
erp_secrets_file: Path = Path("secrets/erp.env"),
) -> MaterialEfficiencyRunner:
try:
document = yaml.load(config.read_text(encoding="utf-8"), Loader=_UniqueLoader)
except (OSError, UnicodeError, yaml.YAMLError):
raise CalculationConfigError("Cannot read material efficiency YAML configuration") from None
required = {
"workplace",
"machine_id",
"calculation_id",
"production_order_format",
"erp_timezone",
}
if (
not isinstance(document, dict)
or not required <= document.keys()
or document.keys() - required - {"poll_interval_seconds", "output_calculation_id"}
):
raise CalculationConfigError("Invalid material efficiency configuration fields")
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")
try:
timezone = ZoneInfo(document["erp_timezone"])
build_enlyze_production_order("0", document["production_order_format"])
except (ValueError, ZoneInfoNotFoundError):
raise CalculationConfigError("Invalid ERP timezone or production order format") from None
interval = finite_number(
document.get("poll_interval_seconds", 60), "poll interval", positive=True
)
environment = dict(os.environ)
try:
environment.update(load_secret_file(secrets_file))
except Exception:
raise ConfigurationError("PostgreSQL settings file could not be loaded") from None
postgres = PostgresSettings.from_environment(environment)
erp = ErpSettings.from_secret_file(erp_secrets_file)
return MaterialEfficiencyRunner(
MaterialEfficiencyService(
ErpWorkplaceStatusGateway(erp),
PostgresMaterialSnapshotRepository(postgres),
workplace=document["workplace"],
machine_id=document["machine_id"],
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,
)
@@ -1,5 +1,6 @@
"""Single-cycle live material polling orchestration.""" """Single-cycle live material polling orchestration."""
from collections.abc import Callable
from dataclasses import dataclass from dataclasses import dataclass
from datetime import datetime from datetime import datetime
from typing import Protocol from typing import Protocol
@@ -13,6 +14,7 @@ from production_analytics.enlyze.gateway import (
EnlyzeApiGateway, EnlyzeApiGateway,
EnlyzeProductionRun, EnlyzeProductionRun,
) )
from production_analytics.service.postgres_material_application import MaterialApplicationWriter
@dataclass(frozen=True, slots=True) @dataclass(frozen=True, slots=True)
@@ -50,17 +52,41 @@ class MaterialPollingService:
state_store: MaterialStateStore, state_store: MaterialStateStore,
machine_id: str, machine_id: str,
rate_variable_id: str, rate_variable_id: str,
gate_variable_id: str, gate_variable_id: str | None,
gate_threshold: float, gate_threshold: float,
max_sample_gap_seconds: 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: ) -> 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 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._gateway = gateway
self._state_store = state_store self._state_store = state_store
self._machine_id = machine_id self._machine_id = machine_id
self._rate_variable_id = rate_variable_id self._rate_variable_id = rate_variable_id
self._gate_variable_id = gate_variable_id self._gate_variable_id = gate_variable_id
self._config = MaterialIntegratorConfig( 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, max_sample_gap_seconds=max_sample_gap_seconds,
) )
@@ -71,6 +97,8 @@ class MaterialPollingService:
if run is None or now < run.start: if run is None or now < run.start:
return None return None
width = (self._nominal_width_provider(run.production_order)
if self._application_variable_ids or self._rotational_speed_variable_ids else None)
saved_polling_state = self._state_store.load( saved_polling_state = self._state_store.load(
self._machine_id, self._machine_id,
run.production_order, run.production_order,
@@ -92,6 +120,7 @@ class MaterialPollingService:
raise ValueError("Historical production run overlaps the current open run") raise ValueError("Historical production run overlaps the current open run")
integrated = self._integrate( integrated = self._integrate(
initial_state, start=previous.start, end=previous.end, initial_state, start=previous.start, end=previous.end,
width=width, run=previous,
) )
initial_state = MaterialIntegrationState( initial_state = MaterialIntegrationState(
cumulative_consumption_kg=integrated.cumulative_consumption_kg, cumulative_consumption_kg=integrated.cumulative_consumption_kg,
@@ -116,7 +145,7 @@ class MaterialPollingService:
if now < start: if now < start:
return None return None
state = self._integrate(initial_state, start=start, end=now) state = self._integrate(initial_state, start=start, end=now, width=width, run=run)
self._state_store.save( self._state_store.save(
self._machine_id, self._machine_id,
@@ -134,13 +163,36 @@ class MaterialPollingService:
def _integrate( def _integrate(
self, initial_state: MaterialIntegrationState, *, start: datetime, end: datetime, self, initial_state: MaterialIntegrationState, *, start: datetime, end: datetime,
width: float | None = None, run: EnlyzeProductionRun | None = None,
) -> MaterialIntegrationState: ) -> MaterialIntegrationState:
integrator = MaterialConsumptionIntegrator(self._config, initial_state=initial_state) integrator = MaterialConsumptionIntegrator(self._config, initial_state=initial_state)
source_options = {}
if self._application_variable_ids:
source_options = dict(
application_variable_ids=self._application_variable_ids, nominal_width_m=width,
)
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( samples = self._gateway.get_material_samples(
machine_id=self._machine_id, machine_id=self._machine_id,
rate_variable_id=self._rate_variable_id, rate_variable_id=self._rate_variable_id,
gate_variable_id=self._gate_variable_id, gate_variable_id=self._gate_variable_id,
start=start, start=start,
end=end, 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
@@ -14,6 +14,8 @@ from production_analytics.enlyze.exploration import (
load_secret_file, load_secret_file,
) )
from production_analytics.enlyze.gateway import EnlyzeApiGateway from production_analytics.enlyze.gateway import EnlyzeApiGateway
from production_analytics.erp import ErpSettings, ErpWorkplaceStatusGateway
from production_analytics.service.material_context import ErpNominalWidthProvider
from production_analytics.service.material_polling import ( from production_analytics.service.material_polling import (
MaterialPollingService, MaterialPollingService,
MaterialPollingState, MaterialPollingState,
@@ -24,6 +26,9 @@ from production_analytics.service.postgres_material import (
PostgresMaterialSnapshotWriter, PostgresMaterialSnapshotWriter,
PostgresSettings, PostgresSettings,
) )
from production_analytics.service.postgres_material_application import (
PostgresMaterialApplicationWriter,
)
class _ReportingStateStore: class _ReportingStateStore:
@@ -48,6 +53,7 @@ class _ReportingStateStore:
def build_material_runner( def build_material_runner(
calculation: MaterialCalculationConfig, *, poll_interval_seconds: float, calculation: MaterialCalculationConfig, *, poll_interval_seconds: float,
state_directory: Path, secrets_file: Path, state_directory: Path, secrets_file: Path,
erp_secrets_file: Path = Path("secrets/erp.env"),
) -> MaterialPollingRunner: ) -> MaterialPollingRunner:
finite_number(poll_interval_seconds, "poll interval", positive=True) finite_number(poll_interval_seconds, "poll interval", positive=True)
environment = dict(os.environ) environment = dict(os.environ)
@@ -69,6 +75,27 @@ def build_material_runner(
f"State directory startup check failed: {type(exc).__name__}" f"State directory startup check failed: {type(exc).__name__}"
) from exc ) from exc
postgres_settings = PostgresSettings.from_environment(environment) postgres_settings = PostgresSettings.from_environment(environment)
source_options = {}
if calculation.source_mode in {"area_application", "rotational_discharge"}:
source_options = dict(
application_variable_ids=calculation.application_signal_refs,
nominal_width_provider=ErpNominalWidthProvider(
ErpWorkplaceStatusGateway(ErpSettings.from_secret_file(erp_secrets_file)),
workplace=calculation.erp_workplace,
format_template=calculation.production_order_format,
),
)
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( service = MaterialPollingService(
gateway=EnlyzeApiGateway(ExplorationClient(settings)), gateway=EnlyzeApiGateway(ExplorationClient(settings)),
state_store=_ReportingStateStore(JsonMaterialStateStore(state_directory)), state_store=_ReportingStateStore(JsonMaterialStateStore(state_directory)),
@@ -77,6 +104,7 @@ def build_material_runner(
gate_variable_id=calculation.gate_signal_ref, gate_variable_id=calculation.gate_signal_ref,
gate_threshold=calculation.gate_threshold, gate_threshold=calculation.gate_threshold,
max_sample_gap_seconds=calculation.max_sample_gap_seconds, max_sample_gap_seconds=calculation.max_sample_gap_seconds,
**source_options,
) )
return MaterialPollingRunner( return MaterialPollingRunner(
service, calculation_id=calculation.id, service, calculation_id=calculation.id,
@@ -1,4 +1,4 @@
"""PostgreSQL settings and derived material snapshot inserts.""" """PostgreSQL settings and derived material snapshot reads/writes."""
from collections.abc import Mapping from collections.abc import Mapping
from dataclasses import dataclass, field from dataclasses import dataclass, field
@@ -41,6 +41,54 @@ class MaterialSnapshotWriter(Protocol):
) -> None: ... ) -> None: ...
@dataclass(frozen=True, slots=True)
class MaterialConsumptionSnapshot:
timestamp: datetime
calculation_id: str
machine_id: str
production_order: str
run_id: str
consumption_kg: float
class MaterialSnapshotRepository(Protocol):
def latest_at_or_before(
self, *, calculation_id: str, machine_id: str, production_order: str,
timestamp: datetime,
) -> MaterialConsumptionSnapshot | None: ...
class PostgresMaterialSnapshotRepository:
def __init__(self, settings: PostgresSettings) -> None:
self.settings = settings
def latest_at_or_before(
self, *, calculation_id: str, machine_id: str, production_order: str,
timestamp: datetime,
) -> MaterialConsumptionSnapshot | None:
"""Read one exact-order cumulative snapshot, never later than the aware cutoff."""
if timestamp.utcoffset() is None:
raise ValueError("Material snapshot cutoff must be timezone-aware")
import psycopg
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:
row = connection.execute(
"""SELECT timestamp, calculation_id, machine_id, production_order,
run_id, consumption_kg
FROM material_consumption_snapshots
WHERE calculation_id = %s AND machine_id = %s
AND production_order = %s AND timestamp <= %s
ORDER BY timestamp DESC
LIMIT 1""",
(calculation_id, machine_id, production_order, timestamp),
).fetchone()
return None if row is None else MaterialConsumptionSnapshot(*row)
class PostgresMaterialSnapshotWriter: class PostgresMaterialSnapshotWriter:
def __init__(self, settings: PostgresSettings) -> None: def __init__(self, settings: PostgresSettings) -> None:
self.settings = settings self.settings = settings
@@ -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,
)
@@ -0,0 +1,71 @@
"""Immutable PostgreSQL history of feedback-aligned KPI points."""
from typing import Protocol
import psycopg
from production_analytics.calculations.config import finite_number
from production_analytics.service.material_efficiency import MaterialEfficiencySnapshot
from production_analytics.service.postgres_material import PostgresSettings
class MaterialEfficiencyWriter(Protocol):
def write(self, snapshot: MaterialEfficiencySnapshot) -> bool: ...
class PostgresMaterialEfficiencyWriter:
def __init__(self, settings: PostgresSettings) -> None:
self.settings = settings
def write(self, snapshot: MaterialEfficiencySnapshot) -> bool:
for name in ("erp_feedback_timestamp", "material_snapshot_timestamp"):
if getattr(snapshot, name).utcoffset() is None:
raise ValueError(f"{name} must be timezone-aware")
for name in (
"good_quantity_m2",
"material_consumption_kg",
"material_consumption_kg_per_m2",
"material_consumption_g_per_m2",
):
finite_number(getattr(snapshot, name), name)
if snapshot.nominal_width_m is not None:
finite_number(snapshot.nominal_width_m, "nominal_width_m")
# Context exit commits/rolls back and closes; later cycles reconnect after outages.
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:
row = connection.execute(
"""INSERT INTO material_efficiency_snapshots (
erp_feedback_timestamp, material_snapshot_timestamp, calculation_id,
machine_id, workplace, erp_production_order, enlyze_production_order,
article_number, article_description, nominal_width_m, good_quantity_m2,
material_consumption_kg, material_consumption_kg_per_m2,
material_consumption_g_per_m2
) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (calculation_id, machine_id, enlyze_production_order,
erp_feedback_timestamp) DO NOTHING
RETURNING 1""",
(
snapshot.erp_feedback_timestamp,
snapshot.material_snapshot_timestamp,
snapshot.calculation_id,
snapshot.machine_id,
snapshot.workplace,
snapshot.erp_production_order,
snapshot.enlyze_production_order,
snapshot.article_number,
snapshot.article_description,
snapshot.nominal_width_m,
snapshot.good_quantity_m2,
snapshot.material_consumption_kg,
snapshot.material_consumption_kg_per_m2,
snapshot.material_consumption_g_per_m2,
),
).fetchone()
return row is not None
+257
View File
@@ -0,0 +1,257 @@
from dataclasses import replace
from datetime import datetime, timedelta
from types import SimpleNamespace
from unittest.mock import Mock
import pytest
import yaml
from production_analytics.calculations.config import (
CalculationConfigError,
load_material_calculation,
)
from production_analytics.calculations.material_consumption import (
MaterialConsumptionIntegrator,
MaterialIntegratorConfig,
area_application_rate_kg_per_hour,
)
from production_analytics.enlyze.gateway import EnlyzeApiGateway, EnlyzeProductionRun
from production_analytics.service.material_context import ErpNominalWidthProvider
from production_analytics.service.material_polling import MaterialPollingService
from production_analytics.service.material_state_store import JsonMaterialStateStore
def timestamp(value):
return datetime.fromisoformat(value.replace('Z', '+00:00'))
def test_area_units_and_strict_gate():
client = Mock()
start = timestamp('2026-08-25T05:27:18Z')
client.post_json.return_value.body = {'data': {
'columns': ['speed', 's2', 'time', 's1'],
'records': [[speed, 2000, (start + timedelta(minutes=i)).isoformat(), 3000]
for i, speed in enumerate([2, 0.3, 0, 2, 2])],
}}
samples = EnlyzeApiGateway(client).get_material_samples(
machine_id='machine', rate_variable_id='unused', gate_variable_id='speed',
application_variable_ids=('s1', 's2'), nominal_width_m=5,
start=start, end=start + timedelta(minutes=4),
)
state = MaterialConsumptionIntegrator(MaterialIntegratorConfig(0.3, 60)).process_many(samples)
assert state.cumulative_consumption_kg == 100
assert state.integrated_running_seconds == 120
assert client.post_json.call_args.args[1]['variables'] == [
{'uuid': 's1'}, {'uuid': 's2'}, {'uuid': 'speed'},
]
@pytest.mark.parametrize('sources,width,speed', [
([], 5, 2), ([float('nan')], 5, 2), ([1], 0, 2), ([1], 5, float('inf')),
])
def test_invalid_area_inputs(sources, width, speed):
with pytest.raises(ValueError):
area_application_rate_kg_per_hour(sources, width, speed)
@pytest.mark.parametrize('bad', [None, True, float('nan')])
def test_missing_or_invalid_spreader_fails(bad):
client = Mock()
now = timestamp('2026-08-25T05:27:18Z')
client.post_json.return_value.body = {'data': {
'columns': ['time', 's1', 's2', 'speed'],
'records': [[now.isoformat(), 3000, bad, 2]],
}}
with pytest.raises(ValueError):
EnlyzeApiGateway(client).get_material_samples(
machine_id='machine', rate_variable_id='unused', gate_variable_id='speed',
application_variable_ids=('s1', 's2'), nominal_width_m=5, start=now, end=now,
)
def width_provider():
erp = Mock()
erp.get_current_workplace_status.return_value = SimpleNamespace(
workplace='Bento 1', production_order='12026000814', article_number='180305',
article_description='Bfix NSP 5300, 5,00 x 40 m',
)
return ErpNominalWidthProvider(
erp, workplace='Bento 1', format_template='Bento 1-{production_order}',
)
@pytest.mark.parametrize('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'),
])
def test_width_context_fails_closed(field, value):
provider = width_provider()
setattr(provider.gateway.get_current_workplace_status.return_value, field, value)
with pytest.raises(ValueError):
provider('Bento 1-12026000814')
@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 realistic rpm values and synthetic active time within actual run boundaries.
"""
config = load_material_calculation('config/bento1-material-consumption.yaml')
order = 'Bento 1-12026000814'
runs = [EnlyzeProductionRun(
uuid, config.machine_ref, 'external-product', order, timestamp(start), timestamp(end),
) for uuid, start, end in [
('b675818c-4636-4976-be6c-3e0e24da9e8a',
'2026-08-25T05:27:18Z', '2026-08-26T04:40:58Z'),
('25d67424-8346-42d4-b844-75b4097382bd',
'2026-08-26T05:43:39Z', '2026-08-26T05:54:05Z'),
]]
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):
start, end = timestamp(body['start']), timestamp(body['end'])
run = next(run for run in runs if run.start <= start <= run.end)
active_seconds = 1256 * 60 - 626 if run == runs[0] else 626
stop = run.start + timedelta(seconds=active_seconds)
times = {start, end}
t = start
while t < end:
times.add(t)
t += timedelta(seconds=10)
if start <= stop <= end:
times.add(stop)
return SimpleNamespace(body={'data': {
'columns': ['time', *config.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)],
}})
client.post_json.side_effect = response
gateway = EnlyzeApiGateway(client)
gateway.get_open_production_run = Mock(return_value=replace(runs[1], end=None))
gateway.get_production_runs = Mock(return_value=runs)
store = JsonMaterialStateStore(tmp_path)
def service():
return MaterialPollingService(
gateway=gateway, state_store=store, machine_id=config.machine_ref,
rate_variable_id=config.rate_signal_ref, gate_variable_id=config.gate_signal_ref,
gate_threshold=config.gate_threshold,
max_sample_gap_seconds=config.max_sample_gap_seconds,
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(
(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]] == [
(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'), ('rotational_speed_signal_refs', []),
('rotational_speed_signal_refs', ['same', 'same']), ('erp_workplace', ''),
('production_order_format', '{wrong}'), ('rate_signal_ref', 'direct'),
])
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
path = tmp_path / 'config.yaml'
path.write_text(yaml.safe_dump(document))
with pytest.raises(CalculationConfigError):
load_material_calculation(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')
gateway.get_open_production_run.return_value = EnlyzeProductionRun(
'run', config.machine_ref, None, 'Bento 1-12026000814', now, None,
)
provider = width_provider()
provider.gateway.get_current_workplace_status.return_value = None
store = Mock()
service = MaterialPollingService(
gateway=gateway, state_store=store, machine_id=config.machine_ref,
rate_variable_id=config.rate_signal_ref, gate_variable_id=config.gate_signal_ref,
gate_threshold=config.gate_threshold, max_sample_gap_seconds=20,
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)
gateway.get_material_samples.assert_not_called()
store.save.assert_not_called()
def test_area_gateway_pagination():
client = Mock()
now = timestamp('2026-08-25T05:27:18Z')
client.post_json.side_effect = [SimpleNamespace(body={
'metadata': {'next_cursor': cursor},
'data': {'columns': ['time', 's1', 's2', 'speed'],
'records': [[(now + timedelta(seconds=i * 10)).isoformat(), 3000, 2000, 2]]},
}) for i, cursor in enumerate(['next', None])]
samples = EnlyzeApiGateway(client).get_material_samples(
machine_id='machine', rate_variable_id='unused', gate_variable_id='speed',
application_variable_ids=('s1', 's2'), nominal_width_m=5,
start=now, end=now + timedelta(seconds=10),
)
assert [sample.material_rate_kg_per_hour for sample in samples] == [3000, 3000]
assert client.post_json.call_args.args[1]['cursor'] == 'next'
+143
View File
@@ -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
+75
View File
@@ -0,0 +1,75 @@
import pytest
from production_analytics.context import (
build_enlyze_production_order,
extract_nominal_width_m,
)
K7_ORDER_TEMPLATE = "K 7-{production_order}"
@pytest.mark.parametrize("order, expected", [
("12026000815", "K 7-12026000815"),
(" \t12026000815\n", "K 7-12026000815"),
("00123", "K 7-00123"),
])
def test_k7_mapping(order, expected):
assert build_enlyze_production_order(order, K7_ORDER_TEMPLATE) == expected
@pytest.mark.parametrize("order", [
"", " \t", "K 7-12026000815", "K 7-12025000074-K 7-12025000075",
"12025000074-12025000075", "prefix12026000815", "12026000815suffix",
"12026 000815", "123", "-123",
])
def test_mapping_rejects_non_erp_identifiers(order):
with pytest.raises(ValueError, match="ERP production order"):
build_enlyze_production_order(order, K7_ORDER_TEMPLATE)
@pytest.mark.parametrize("template, expected", [
("EXAMPLE/{production_order}/finished", "EXAMPLE/00123/finished"),
("{production_order}", "00123"),
("{production_order}-example", "00123-example"),
])
def test_mapping_uses_explicit_template(template, expected):
assert build_enlyze_production_order(" 00123 ", template) == expected
@pytest.mark.parametrize("template", [
"", "fixed-order", "{order}", "{production_order}-{production_order}",
"{production_order:>20}", "{production_order!r}", "{production_order}{unknown}",
"{{production_order}}",
])
def test_mapping_rejects_invalid_template(template):
with pytest.raises(ValueError, match="Format template"):
build_enlyze_production_order("12026000815", template)
@pytest.mark.parametrize("description, expected", [
("Stex R 1501 C (PR) 5,80 x 50 m", 5.8),
("Stex R 401, 6,00 x 90 m", 6.0),
("Stex R 1501 5.80 x 50 m", 5.8),
("5,80x50 m", 5.8),
("5,80x50m", 5.8),
("Article 212520 grade 1501 5,80\t x\t50 m", 5.8),
("5,80 x 50 m coated (batch 42)", 5.8),
("(5,80 x 50 m), coated", 5.8),
("6 x 90 m", 6.0),
])
def test_width(description, expected):
assert extract_nominal_width_m(description) == expected
@pytest.mark.parametrize("description", [
None, "", "Stex R 1501", "5,80 m", "5,80 x 50", "5,80 x 50 cm",
"5,80 x 50 mm", "5,80 x 50 m²", "5,80 x 50 m2", "5,80 x 50 m/min",
"5,80 x 50 metres", "5,80,2 x 50 m", "5.80.2 x 50 m",
"5,80 x ? m", "5,80 x 50 m or 6,00 x 90 m", "2 x 5,80 x 50 m",
"5,80 x 50 m x 2", "0 x 50 m", "-5,80 x 50 m", "+5,80 x 50 m",
"nan x 50 m", "inf x 50 m", "Infinity x 50 m", "1e309 x 50 m",
"9" * 400 + " x 50 m", "5,80 x 0 m", "5,80 x -50 m", "5,80 x nan m",
"1/5,80 x 50 m", "abc5,80 x 50 m",
])
def test_width_fails_safely(description):
assert extract_nominal_width_m(description) is None
+158
View File
@@ -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'])
+186
View File
@@ -0,0 +1,186 @@
from dataclasses import FrozenInstanceError, replace
from datetime import UTC, datetime
from unittest.mock import Mock
from zoneinfo import ZoneInfo
import pytest
from production_analytics.erp import CurrentWorkplaceStatus
from production_analytics.service.material_efficiency import MaterialEfficiencyService
from production_analytics.service.postgres_material import MaterialConsumptionSnapshot
CALC = 'k7-fiber-consumption'
MACHINE = 'c220f95c-a65e-4cb7-99b7-0626d6c7508c'
START = datetime(2026, 9, 4, 10, tzinfo=UTC)
FEEDBACK = START.replace(minute=5)
STATUS = CurrentWorkplaceStatus(
'K7', '12026000815', '212520', 'Stex R 1501 C (PR) 5,80 x 50 m',
FEEDBACK, None, 800, None, None, None,
)
MATERIAL = MaterialConsumptionSnapshot(START, CALC, MACHINE, 'K 7-12026000815', 'run', 1000)
def service(repository=None, **config):
return MaterialEfficiencyService(
Mock(get_current_workplace_status=Mock(return_value=STATUS)),
repository if repository is not None else Mock(latest_at_or_before=Mock(
return_value=MATERIAL,
)),
**dict(workplace='K7', machine_id=MACHINE, calculation_id=CALC,
format_template='K 7-{production_order}', **config),
)
def test_valid_current_feedback():
subject = service()
result = subject.evaluate_current()
assert result.material_consumption_kg_per_m2 == 1.25
assert result.material_consumption_g_per_m2 == 1250
assert result.nominal_width_m == 5.8
assert result.good_quantity_m2 == 800
assert result.material_consumption_kg == 1000
assert result.erp_feedback_timestamp == FEEDBACK
assert result.material_snapshot_timestamp == START
assert (result.workplace, result.machine_id, result.calculation_id) == ('K7', MACHINE, CALC)
assert (result.erp_production_order, result.enlyze_production_order) == (
'12026000815', 'K 7-12026000815',
)
assert (result.article_number, result.article_description) == (
STATUS.article_number, STATUS.article_description,
)
subject.erp.get_current_workplace_status.assert_called_once_with('K7')
subject.materials.latest_at_or_before.assert_called_once_with(
calculation_id=CALC, machine_id=MACHINE, production_order='K 7-12026000815',
timestamp=FEEDBACK,
)
with pytest.raises(FrozenInstanceError):
result.good_quantity_m2 = 5
@pytest.mark.parametrize('quantity', [None, 0, -1, float('nan'), float('inf'), -float('inf')])
def test_invalid_quantity(quantity):
subject = service()
assert subject.evaluate(replace(STATUS, good_quantity_m2=quantity)) is None
subject.materials.latest_at_or_before.assert_not_called()
@pytest.mark.parametrize('description', [None, 'unstructured article'])
def test_width_is_optional_context(description):
result = service().evaluate(replace(STATUS, article_description=description))
assert result.nominal_width_m is None
assert result.material_consumption_kg_per_m2 == 1.25
def test_unavailable_sources():
subject = service()
subject.erp.get_current_workplace_status.return_value = None
assert subject.evaluate_current() is None
subject.materials.latest_at_or_before.return_value = None
assert subject.evaluate(STATUS) is None
@pytest.mark.parametrize('consumption, quantity', [
(float('nan'), 800), (float('inf'), 800), (-1, 800),
(1e308, 1e-308), (1e308, 1),
])
def test_invalid_consumption_and_overflow(consumption, quantity):
subject = service()
subject.materials.latest_at_or_before.return_value = replace(
MATERIAL, consumption_kg=consumption,
)
assert subject.evaluate(replace(STATUS, good_quantity_m2=quantity)) is None
def test_zero_consumption_is_valid():
subject = service()
subject.materials.latest_at_or_before.return_value = replace(MATERIAL, consumption_kg=0)
assert subject.evaluate(STATUS).material_consumption_kg_per_m2 == 0
@pytest.mark.parametrize('changes', [
dict(timestamp=START.replace(minute=10)), dict(timestamp=START.replace(tzinfo=None)),
dict(production_order='K 7-12026000815-K 7-12026000816'),
dict(machine_id='other'), dict(calculation_id='other'),
])
def test_repository_contract_is_checked(changes):
subject = service()
subject.materials.latest_at_or_before.return_value = replace(MATERIAL, **changes)
with pytest.raises(ValueError, match='unaligned'):
subject.evaluate(STATUS)
def test_other_explicit_configuration():
repository = Mock()
repository.latest_at_or_before.return_value = replace(
MATERIAL, machine_id='example-machine', production_order='ORDER/00123',
)
subject = MaterialEfficiencyService(
Mock(), repository, workplace='example', machine_id='example-machine',
calculation_id=CALC, format_template='ORDER/{production_order}',
)
result = subject.evaluate(replace(STATUS, workplace='example', production_order=' 00123 '))
assert result.enlyze_production_order == 'ORDER/00123'
def test_wrong_workplace_and_combined_order_rejected():
with pytest.raises(ValueError, match='workplace'):
service().evaluate(replace(STATUS, workplace='other'))
with pytest.raises(ValueError, match='ERP production order'):
service().evaluate(replace(STATUS, production_order='K 7-123-K 7-456'))
def test_temporal_regression_and_successive_feedback():
later = replace(
MATERIAL, timestamp=START.replace(minute=10), consumption_kg=1100, run_id='run2',
)
repository = Mock()
repository.latest_at_or_before.side_effect = lambda **kw: max(
(row for row in [MATERIAL, later] if row.timestamp <= kw['timestamp']),
key=lambda row: row.timestamp, default=None,
)
subject = service(repository)
first = subject.evaluate(STATUS)
second = subject.evaluate(replace(
STATUS, feedback_timestamp=later.timestamp, good_quantity_m2=1000,
))
assert first.material_snapshot_timestamp == START
assert first.material_consumption_kg == 1000
assert first.material_consumption_kg_per_m2 == 1.25
assert second.material_snapshot_timestamp == later.timestamp
assert second.material_consumption_kg == 1100
assert second.material_consumption_kg_per_m2 == 1.1
assert first.good_quantity_m2 == 800
def test_naive_feedback_requires_explicit_timezone_and_returns_instant():
naive = datetime(2026, 9, 4, 12, 5)
status = replace(STATUS, feedback_timestamp=naive)
with pytest.raises(ValueError, match='explicit erp_timezone'):
service().evaluate(status)
subject = service(erp_timezone=ZoneInfo('Europe/Berlin'))
result = subject.evaluate(status)
assert result.erp_feedback_timestamp == FEEDBACK
assert status.feedback_timestamp == naive
assert subject.materials.latest_at_or_before.call_args.kwargs['timestamp'] == FEEDBACK
@pytest.mark.parametrize('timestamp', [datetime(2026, 3, 29, 2, 30), datetime(2026, 10, 25, 2, 30)])
def test_dst_gap_and_overlap_rejected(timestamp):
with pytest.raises(ValueError, match='ambiguous or nonexistent'):
service(erp_timezone=ZoneInfo('Europe/Berlin')).evaluate(
replace(STATUS, feedback_timestamp=timestamp),
)
def test_aware_feedback_returns_utc():
timestamp = FEEDBACK.astimezone(ZoneInfo('Europe/Berlin'))
result = service().evaluate(replace(STATUS, feedback_timestamp=timestamp))
assert result.erp_feedback_timestamp == timestamp
assert result.erp_feedback_timestamp.tzinfo is UTC
def test_database_failure_propagates():
subject = service()
subject.materials.latest_at_or_before.side_effect = RuntimeError('unavailable')
with pytest.raises(RuntimeError, match='unavailable'):
subject.evaluate(STATUS)
@@ -0,0 +1,488 @@
import io
import sqlite3
from dataclasses import replace
from datetime import UTC, datetime
from pathlib import Path
from unittest.mock import MagicMock, Mock
import psycopg
import pytest
import yaml
from production_analytics.calculations.config import CalculationConfigError
from production_analytics.cli.__main__ import main
from production_analytics.erp import CurrentWorkplaceStatus, ErpReadError
from production_analytics.service.material_efficiency import (
MaterialEfficiencySnapshot,
)
from production_analytics.service.material_efficiency_runner import MaterialEfficiencyRunner
from production_analytics.service.material_efficiency_runtime import (
build_material_efficiency_runner,
)
from production_analytics.service.postgres_material import (
MaterialConsumptionSnapshot,
PostgresSettings,
)
from production_analytics.service.postgres_material_efficiency import (
PostgresMaterialEfficiencyWriter,
)
NOW = datetime(2026, 9, 4, 10, tzinfo=UTC)
SNAPSHOT = MaterialEfficiencySnapshot(
"example",
"machine",
"calc",
"00123",
" ORDER/00123'; -- ",
"article",
"fabric 5 x 50 m",
5,
NOW,
NOW,
800,
1000,
1.25,
1250,
)
SETTINGS = PostgresSettings("localhost", 5432, "analytics", "writer", "secret")
@pytest.fixture
def storage(monkeypatch):
# Real schema and INSERT semantics, with only driver transport adapted for SQLite.
database = sqlite3.connect(":memory:")
database.executescript(Path("db/schema.sql").read_text())
connection = MagicMock()
def execute(sql, params):
assert sql.lstrip().startswith("INSERT INTO")
assert "SELECT" not in sql.upper()
assert "RETURNING 1" in sql
assert sql.count("%s") == 14
assert SNAPSHOT.enlyze_production_order not in sql
return database.execute(
sql.replace("%s", "?"),
tuple(value.isoformat() if isinstance(value, datetime) else value for value in params),
)
connection.__enter__.return_value.execute.side_effect = execute
connect = Mock(return_value=connection)
monkeypatch.setattr(psycopg, "connect", connect)
yield database, connect, connection
database.close()
@pytest.mark.parametrize("nullable", [False, True])
def test_all_fields_and_immutable_event(storage, nullable):
database, connect, connection = storage
snapshot = (
replace(SNAPSHOT, article_number=None, article_description=None, nominal_width_m=None)
if nullable
else SNAPSHOT
)
writer = PostgresMaterialEfficiencyWriter(SETTINGS)
assert writer.write(snapshot) is True
inserted = writer.write(
replace(
snapshot,
material_consumption_kg=9999,
material_snapshot_timestamp=NOW.replace(minute=1),
)
)
assert inserted is False
assert connection.__enter__.return_value.execute.call_count == 2
assert database.execute("SELECT * FROM material_efficiency_snapshots").fetchall() == [
(
NOW.isoformat(),
NOW.isoformat(),
"calc",
"machine",
"example",
"00123",
SNAPSHOT.enlyze_production_order,
snapshot.article_number,
snapshot.article_description,
snapshot.nominal_width_m,
800,
1000,
1.25,
1250,
)
]
assert connect.call_count == 2
assert connect.call_args.kwargs["connect_timeout"] == 10
connection.__exit__.assert_called_with(None, None, None)
@pytest.mark.parametrize(
"field,value",
[
("calculation_id", "other"),
("machine_id", "other"),
("enlyze_production_order", SNAPSHOT.enlyze_production_order.strip()),
("enlyze_production_order", SNAPSHOT.enlyze_production_order + "-combined"),
("erp_feedback_timestamp", NOW.replace(minute=1)),
],
)
def test_exact_identity(storage, field, value):
writer = PostgresMaterialEfficiencyWriter(SETTINGS)
writer.write(SNAPSHOT)
writer.write(replace(SNAPSHOT, **{field: value}))
assert storage[0].execute("SELECT count(*) FROM material_efficiency_snapshots").fetchone() == (
2,
)
@pytest.mark.parametrize("field", ["erp_feedback_timestamp", "material_snapshot_timestamp"])
def test_naive_rejected(storage, field):
with pytest.raises(ValueError, match="timezone-aware"):
PostgresMaterialEfficiencyWriter(SETTINGS).write(
replace(SNAPSHOT, **{field: NOW.replace(tzinfo=None)}),
)
storage[1].assert_not_called()
@pytest.mark.parametrize(
"field",
[
"nominal_width_m",
"good_quantity_m2",
"material_consumption_kg",
"material_consumption_kg_per_m2",
"material_consumption_g_per_m2",
],
)
@pytest.mark.parametrize("value", [float("nan"), float("inf"), -float("inf")])
def test_nonfinite_rejected(storage, field, value):
with pytest.raises(ValueError, match="finite"):
PostgresMaterialEfficiencyWriter(SETTINGS).write(replace(SNAPSHOT, **{field: value}))
storage[1].assert_not_called()
def run_cycles(service, writer, cycles, interval=60, stdout=None):
sleep = Mock(side_effect=[None] * (cycles - 1) + [KeyboardInterrupt])
errors = io.StringIO()
MaterialEfficiencyRunner(
service, writer, poll_interval_seconds=interval, sleep=sleep, stderr=errors, stdout=stdout
).run()
assert sleep.call_args_list == [((interval,),)] * cycles
return errors.getvalue()
def test_runner_history_and_none(storage):
service = Mock(
evaluate_current=Mock(
side_effect=[
None,
SNAPSHOT,
SNAPSHOT,
replace(SNAPSHOT, erp_feedback_timestamp=NOW.replace(minute=15)),
]
)
)
output = io.StringIO()
run_cycles(service, PostgresMaterialEfficiencyWriter(SETTINGS), 4, 73, stdout=output)
lines = output.getvalue().splitlines()
assert len(lines) == 2
assert lines[0].startswith(NOW.isoformat())
assert lines[1].startswith(NOW.replace(minute=15).isoformat())
assert storage[1].call_count == 3
assert storage[0].execute("SELECT count(*) FROM material_efficiency_snapshots").fetchone() == (
2,
)
@pytest.mark.parametrize("error", [ErpReadError, psycopg.OperationalError, psycopg.InterfaceError])
def test_read_recovers(error):
service = Mock(evaluate_current=Mock(side_effect=[error("secret"), SNAPSHOT]))
writer = Mock()
errors = run_cycles(service, writer, 2)
writer.write.assert_called_once_with(SNAPSHOT)
assert error.__name__ in errors and "secret" not in errors
def test_write_recovers_with_fresh_transaction(storage):
database, connect, connection = storage
connect.side_effect = [psycopg.OperationalError("secret"), connection]
errors = run_cycles(
Mock(evaluate_current=Mock(return_value=SNAPSHOT)),
PostgresMaterialEfficiencyWriter(SETTINGS),
2,
)
assert "OperationalError" in errors and "secret" not in errors
assert connect.call_count == 2
assert database.execute("SELECT count(*) FROM material_efficiency_snapshots").fetchone() == (1,)
@pytest.mark.parametrize("error", [ValueError, TypeError, psycopg.ProgrammingError])
def test_programming_errors_propagate(error):
with pytest.raises(error):
run_cycles(Mock(evaluate_current=Mock(side_effect=error("bad contract"))), Mock(), 1)
@pytest.mark.parametrize("interval", [0, -1, float("nan"), float("inf"), True])
def test_invalid_interval(interval):
with pytest.raises(CalculationConfigError):
MaterialEfficiencyRunner(Mock(), Mock(), poll_interval_seconds=interval)
@pytest.fixture
def runtime_files(tmp_path):
config = tmp_path / "config.yaml"
config.write_text(
yaml.safe_dump(
dict(
workplace="example",
machine_id="machine",
calculation_id="calc",
production_order_format="ORDER/{production_order}",
erp_timezone="Europe/Berlin",
)
)
)
postgres = tmp_path / "postgres.env"
postgres.write_text(
"POSTGRES_HOST=localhost\nPOSTGRES_PORT=5432\nPOSTGRES_DB=analytics\n"
"POSTGRES_USER=writer\nPOSTGRES_PASSWORD=secret\n"
)
erp = tmp_path / "erp.env"
erp.write_text(
"ERP_DB_HOST=localhost\nERP_DB_PORT=1433\nERP_DB_NAME=erp\n"
"ERP_DB_USER=reader\nERP_DB_PASSWORD=secret\n"
)
return config, postgres, erp
def test_configured_runtime_naive_feedback_to_persistence(runtime_files, storage):
config, postgres, erp = runtime_files
runner = build_material_efficiency_runner(config, secrets_file=postgres, erp_secrets_file=erp)
assert runner.interval == 60
runner.service.erp = Mock(
get_current_workplace_status=Mock(
return_value=CurrentWorkplaceStatus(
"example",
"00123",
None,
None,
datetime(2026, 9, 4, 12),
None,
800,
None,
None,
None,
)
)
)
runner.service.materials = Mock(
latest_at_or_before=Mock(
return_value=MaterialConsumptionSnapshot(
NOW,
"calc",
"machine",
"ORDER/00123",
"run",
1000,
)
)
)
runner.sleep = Mock(side_effect=KeyboardInterrupt)
runner.run()
row = storage[0].execute("SELECT * FROM material_efficiency_snapshots").fetchone()
assert row[:7] == (
NOW.isoformat(),
NOW.isoformat(),
"calc",
"machine",
"example",
"00123",
"ORDER/00123",
)
@pytest.mark.parametrize(
"changes",
[
{"erp_timezone": "invalid"},
{"production_order_format": "missing"},
{"workplace": ""},
{"machine_id": None},
{"poll_interval_seconds": 0},
{"unknown": 1},
],
)
def test_startup_config_errors(runtime_files, changes, capsys):
config, postgres, erp = runtime_files
config.write_text(yaml.safe_dump(yaml.safe_load(config.read_text()) | changes))
assert (
main(
[
"run",
"material-efficiency",
"--config",
str(config),
"--secrets-file",
str(postgres),
"--erp-secrets-file",
str(erp),
]
)
== 2
)
assert "Config/startup error" in capsys.readouterr().err
def test_k7_config_and_cli(runtime_files, monkeypatch):
_, postgres, erp = runtime_files
runner = build_material_efficiency_runner(
Path("config/k7-material-efficiency.yaml"), secrets_file=postgres, erp_secrets_file=erp
)
assert runner.service.workplace == "K7"
assert runner.service.machine_id == "c220f95c-a65e-4cb7-99b7-0626d6c7508c"
assert runner.service.calculation_id == "k7-fiber-consumption"
assert runner.service.format_template == "K 7-{production_order}"
assert runner.service.erp_timezone.key == "Europe/Berlin"
run = Mock()
monkeypatch.setattr(MaterialEfficiencyRunner, "run", run)
assert (
main(
[
"run",
"material-efficiency",
"--config",
"config/k7-material-efficiency.yaml",
"--secrets-file",
str(postgres),
"--erp-secrets-file",
str(erp),
]
)
== 0
)
run.assert_called_once()
def test_failed_insert_rolls_back_and_later_cycle_reconnects(storage):
database, connect, connection = storage
execute = connection.__enter__.return_value.execute
original = execute.side_effect
attempts = 0
def fail_once(sql, params):
nonlocal attempts
attempts += 1
if attempts == 1:
raise psycopg.OperationalError("secret SQL")
return original(sql, params)
execute.side_effect = fail_once
errors = run_cycles(
Mock(evaluate_current=Mock(return_value=SNAPSHOT)),
PostgresMaterialEfficiencyWriter(SETTINGS),
2,
)
assert "secret" not in errors
assert connect.call_count == 2
assert connection.__exit__.call_args_list[0].args[0] is psycopg.OperationalError
assert database.execute("SELECT count(*) FROM material_efficiency_snapshots").fetchone() == (1,)
def test_configured_interval_and_missing_settings(runtime_files, monkeypatch):
config, postgres, erp = runtime_files
config.write_text(
yaml.safe_dump(yaml.safe_load(config.read_text()) | {"poll_interval_seconds": 91})
)
runner = build_material_efficiency_runner(config, secrets_file=postgres, erp_secrets_file=erp)
assert runner.interval == 91
postgres.write_text("")
for key in ("HOST", "PORT", "DB", "USER", "PASSWORD"):
monkeypatch.delenv(f"POSTGRES_{key}", raising=False)
from production_analytics.enlyze.exploration import ConfigurationError
with pytest.raises(ConfigurationError, match="POSTGRES_HOST"):
build_material_efficiency_runner(config, secrets_file=postgres, erp_secrets_file=erp)
@pytest.mark.parametrize("result", [None, False, True])
def test_runner_success_output(result):
snapshot = replace(SNAPSHOT, enlyze_production_order="ORDER/00123")
service = Mock(evaluate_current=Mock(return_value=None if result is None else snapshot))
writer = Mock(write=Mock(return_value=result))
output = Mock(wraps=io.StringIO())
errors = run_cycles(service, writer, 1, stdout=output)
assert errors == ""
if result is None:
writer.write.assert_not_called()
else:
writer.write.assert_called_once_with(snapshot)
if result is True:
assert output.getvalue() == (
"2026-09-04T10:00:00+00:00 workplace='example' production_order='ORDER/00123' "
"good_m2=800.000 material_kg=1000.000 g_per_m2=1250.000\n"
)
output.flush.assert_called_once_with()
else:
assert output.getvalue() == ""
output.flush.assert_not_called()
def test_repeated_duplicates_are_silent(storage):
writer = PostgresMaterialEfficiencyWriter(SETTINGS)
assert writer.write(SNAPSHOT) is True
output = io.StringIO()
run_cycles(Mock(evaluate_current=Mock(return_value=SNAPSHOT)), writer, 3, stdout=output)
assert output.getvalue() == ""
def test_commit_failure_does_not_log_success(storage):
_, _, connection = storage
connection.__exit__.side_effect = [psycopg.OperationalError("secret driver SQL"), None]
output = io.StringIO()
errors = run_cycles(
Mock(evaluate_current=Mock(return_value=SNAPSHOT)),
PostgresMaterialEfficiencyWriter(SETTINGS),
1,
stdout=output,
)
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)
+101
View File
@@ -0,0 +1,101 @@
import sqlite3
import sys
from datetime import UTC, datetime
from unittest.mock import MagicMock, patch
import pytest
from production_analytics.service.postgres_material import (
PostgresMaterialSnapshotRepository,
PostgresSettings,
)
START = datetime(2026, 9, 4, 10, tzinfo=UTC)
ORDER = " K 7-123'; -- "
SETTINGS = PostgresSettings('localhost', 5432, 'analytics', 'reader', 'secret')
@pytest.fixture
def database():
# Execute the actual portable SELECT in SQLite; only adapt driver placeholders
# and datetime transport. This tests SQL semantics without a live PostgreSQL server.
with sqlite3.connect(':memory:') as database:
database.execute('''CREATE TABLE material_consumption_snapshots (
timestamp TEXT, calculation_id TEXT, machine_id TEXT, production_order TEXT,
run_id TEXT, consumption_kg REAL
)''')
rows = [
(START.isoformat(), 'calc', 'machine', ORDER, 'run1', 1000),
(START.replace(minute=10).isoformat(), 'calc', 'machine', ORDER, 'run2', 1100),
(START.replace(minute=4).isoformat(), 'other', 'machine', ORDER, 'run', 9999),
(START.replace(minute=4).isoformat(), 'calc', 'other', ORDER, 'run', 9999),
(START.replace(minute=4).isoformat(), 'calc', 'machine', ORDER.strip(), 'run', 9999),
(START.replace(minute=4).isoformat(), 'calc', 'machine', ORDER + '-other', 'run', 9999),
]
database.executemany(
'INSERT INTO material_consumption_snapshots VALUES (?,?,?,?,?,?)', rows,
)
yield database
@pytest.mark.parametrize('minute, expected_minute, consumption', [
(-1, None, None), (0, 0, 1000), (5, 0, 1000), (10, 10, 1100), (15, 10, 1100),
])
def test_aligned_lookup_sql(database, minute, expected_minute, consumption):
cutoff = START.replace(minute=minute) if minute >= 0 else START.replace(hour=9, minute=59)
driver = MagicMock()
connection = driver.connect.return_value.__enter__.return_value
def execute(sql, parameters):
assert parameters == ('calc', 'machine', ORDER, cutoff)
assert sql.count('%s') == 4
assert ORDER not in sql
assert 'ORDER BY timestamp DESC' in sql
assert 'LIMIT 1' in sql
row = database.execute(
sql.replace('%s', '?'), (*parameters[:3], parameters[3].isoformat()),
).fetchone()
if row is not None:
row = (datetime.fromisoformat(row[0]), *row[1:])
return MagicMock(fetchone=MagicMock(return_value=row))
connection.execute.side_effect = execute
with patch.dict(sys.modules, psycopg=driver):
result = PostgresMaterialSnapshotRepository(SETTINGS).latest_at_or_before(
calculation_id='calc', machine_id='machine', production_order=ORDER, timestamp=cutoff,
)
if expected_minute is None:
assert result is None
else:
assert result.timestamp == START.replace(minute=expected_minute)
assert result.timestamp <= cutoff
assert result.consumption_kg == consumption
assert (result.calculation_id, result.machine_id, result.production_order) == (
'calc', 'machine', ORDER,
)
assert result.run_id == ('run1' if expected_minute == 0 else 'run2')
connection.execute.assert_called_once()
driver.connect.assert_called_once_with(
host='localhost', port=5432, dbname='analytics', user='reader', password='secret',
connect_timeout=10, options='-c statement_timeout=10000',
)
driver.connect.return_value.__exit__.assert_called_once_with(None, None, None)
def test_naive_cutoff_rejected_before_connection():
driver = MagicMock()
with patch.dict(sys.modules, psycopg=driver), pytest.raises(ValueError, match='timezone-aware'):
PostgresMaterialSnapshotRepository(SETTINGS).latest_at_or_before(
calculation_id='calc', machine_id='machine', production_order=ORDER,
timestamp=START.replace(tzinfo=None),
)
driver.connect.assert_not_called()
def test_database_failure_propagates():
driver = MagicMock()
driver.connect.side_effect = RuntimeError('connection unavailable')
with patch.dict(sys.modules, psycopg=driver), pytest.raises(RuntimeError):
PostgresMaterialSnapshotRepository(SETTINGS).latest_at_or_before(
calculation_id='calc', machine_id='machine', production_order=ORDER, timestamp=START,
)
@@ -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')