Compare commits

...
2 Commits
Author SHA1 Message Date
admin be8d76507a Add Bento 1 fresh bentonite consumption 2026-09-06 13:29:06 +02:00
admin ed0e20a04b Persist material efficiency snapshots 2026-09-06 08:23:28 +02:00
19 changed files with 1281 additions and 46 deletions
+63 -9
View File
@@ -318,7 +318,7 @@ across the order's runs. The service neither reintegrates nor sums snapshots and
does not bridge run boundaries. does not bridge run boundaries.
The result includes workplace/machine/calculation identifiers, both order identifiers, The result includes workplace/machine/calculation identifiers, both order identifiers,
article context, nominal width, both original source timestamps, good area, aligned 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`, consumption, `material_consumption_kg_per_m2 = consumption_kg / good_quantity_m2`,
and `material_consumption_g_per_m2 = material_consumption_kg_per_m2 * 1000`. and `material_consumption_g_per_m2 = material_consumption_kg_per_m2 * 1000`.
Nominal width comes from generic ERP description parsing and is context only; Nominal width comes from generic ERP description parsing and is context only;
@@ -352,9 +352,9 @@ result = service.evaluate_current()
Aware ERP timestamps are compared as UTC instants. Naive ERP timestamps require Aware ERP timestamps are compared as UTC instants. Naive ERP timestamps require
explicit `erp_timezone`; there is no inferred system/database timezone. Ambiguous explicit `erp_timezone`; there is no inferred system/database timezone. Ambiguous
or nonexistent local times during DST transitions are rejected. The original ERP or nonexistent local times during DST transitions are rejected. The result carries the
timestamp (including naivety) and selected material timestamp are preserved in the ERP feedback instant normalized to UTC and the selected aware material timestamp.
result; retain the source timezone configuration alongside it for interpretation. The original ERP status object remains unchanged.
This alignment avoids knowingly including future material consumption, but does This alignment avoids knowingly including future material consumption, but does
not eliminate ERP roll-feedback timing uncertainty (feedback may be roughly one not eliminate ERP roll-feedback timing uncertainty (feedback may be roughly one
@@ -366,11 +366,10 @@ stored timestamps satisfy the cutoff; a later cumulative total cannot reconstruc
an earlier value. Reliable final evaluation requires retaining final ERP feedback; an earlier value. Reliable final evaluation requires retaining final ERP feedback;
the current-workplace adapter alone does not provide historical completed orders. the current-workplace adapter alone does not provide historical completed orders.
Call on changed ERP feedback; no new runner or polling loop is provided. Results The service itself returns results in memory without modifying material snapshots.
are returned in memory without modifying earlier results or material snapshots. The standalone persistence runner is described below. Daily per-machine 24h reporting
No schema, Grafana, or timeseries changes are needed. Daily per-machine 24h reporting and aggregation remain future work. Repository SQL tests use an in-memory SQLite
and aggregation remain future work. Repository SQL selection tests use an in-memory fixture with driver transport adaptation; they require no live PostgreSQL.
SQLite fixture with driver transport adaptation; they require no live PostgreSQL.
## Peak-cycle detection ## Peak-cycle detection
@@ -412,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).
+19
View File
@@ -0,0 +1,19 @@
# Fresh bentonite only: the recycled/recovered third spreader is excluded.
calculations:
- id: bento1-fresh-bentonite-consumption
type: material_consumption
version: "1"
machine_ref: 5f42a4f6-9ca0-4f6f-9786-40d50a35b230
source_mode: area_application
application_signal_refs:
- 19ea65d2-bd35-4d02-89de-4495583d9026
- d9f47615-69cf-4f97-88df-199585764491
# Transport Auszug Geschwindigkeit Istwert, m/min (also the production gate).
gate_signal_ref: fef41976-1103-4090-b780-eaeecc02fdfa
gate_threshold: 0.3
max_sample_gap_seconds: 20.0
erp_workplace: Bento 1
production_order_format: "Bento 1-{production_order}"
output_metric: material_consumption
output_unit: kg
group_by: production_order
+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
+21
View File
@@ -11,3 +11,24 @@ 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);
+60
View File
@@ -0,0 +1,60 @@
# Bento 1 fresh bentonite consumption
`config/bento1-material-consumption.yaml` defines
`bento1-fresh-bentonite-consumption`. This KPI estimates **fresh bentonite
consumption** from the two fresh-spreader setpoints in g/m². The third spreader
uses recycled/recovered bentonite and is deliberately excluded; it is neither
missing data nor estimated. This is not total bentonite deposited on the product,
and setpoints are not a measurement of actual mass flow.
The generic `area_application` source sums the configured application signals
and converts them to kg/h using `sum(g/m²) × nominal_width_m × speed_m/min × 60 / 1000`.
The existing previous-value integrator then accumulates kilograms only while
speed is strictly greater than 0.3 m/min. `direct_mass_rate` remains the default
for existing configurations, including K7. The persistent state schema and
production-order bootstrap/resume behavior are unchanged; each disjoint run
starts with a new sample baseline while retaining order totals.
Width comes from the existing ERP workplace status adapter and generic nominal
width parser. The ERP order must format to the ENLYZE order using
`Bento 1-{production_order}`. Missing/mismatched ERP context or an unparseable
width rejects the polling cycle without saving state. ENLYZE product metadata is
not used for width. The current ERP context supplies width for all runs of that
same order; this assumes the order's article width stays constant. Historical
orders no longer present in current ERP workplace status require a separate
historical context source for offline replay.
Run with:
```sh
production-analytics run material-poll \
--config config/bento1-material-consumption.yaml \
--erp-secrets-file secrets/erp.env
```
`erp_workplace: Bento 1` assumes that exact ERP workplace identifier; verify it
against the deployment's workplace view. The 20-second maximum sample gap is a
generic policy shared with K7, intended for a nominal 10-second sampling cadence.
One missing sample is tolerated: the resulting 20-second interval is integrated
using the preceding rate and gate. Gaps greater than 20 seconds are not
integrated, so prolonged or repeated gaps can undercount consumption.
Missing/nonfinite values from either configured fresh spreader reject the cycle.
Intervals exceeding the maximum gap and unobserved run tails are not inferred.
Before implementation, the calculation was independently validated manually
against the same three real ENLYZE signals over the two runs below, yielding
approximately **147021.9 kg** of fresh bentonite for `Bento 1-12026000814`.
The post-implementation live replay could not be repeated because ENLYZE
connectivity was temporarily unavailable (DNS resolution failure), not because
of a known implementation issue. Exact production-path replay remains a
follow-up verification rather than a blocker for this milestone.
The regression uses article 180305, description `Bfix NSP 5300, 5,00 x 40 m`
(width 5.00 m), and the supplied boundaries of runs
`b675818c-4636-4976-be6c-3e0e24da9e8a` and
`25d67424-8346-42d4-b844-75b4097382bd` for order `Bento 1-12026000814`.
It constructs synthetic samples from the supplied aggregates: 1256.0 active
minutes, 31294.6 m² and 4698.0 g/m², yielding approximately 147021.9 kg.
It verifies bootstrap across both runs and persisted resume without counting the
inter-run gap. It is an aggregate plausibility regression, not an independent
replay of recorded signals; raw reference samples were not supplied.
@@ -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
@@ -54,6 +54,10 @@ class MaterialCalculationConfig:
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, ...] = ()
erp_workplace: str = ""
production_order_format: str = ""
def load_material_calculation( def load_material_calculation(
@@ -71,7 +75,9 @@ 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 = {}
for index, entry in enumerate(entries): for index, entry in enumerate(entries):
prefix = f"calculations[{index}]" prefix = f"calculations[{index}]"
@@ -79,10 +85,13 @@ def load_material_calculation(
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"
if area:
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 not isinstance(entry[name], str) or not entry[name].strip(): if not isinstance(entry[name], str) or not entry[name].strip():
@@ -98,6 +107,32 @@ def load_material_calculation(
entry["max_sample_gap_seconds"] = finite_number( entry["max_sample_gap_seconds"] = 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"}:
raise CalculationConfigError("Unsupported material source_mode")
if area:
refs = entry.get("application_signal_refs")
if (not isinstance(refs, list) or not refs
or any(not isinstance(ref, str) or not ref.strip() for ref in refs)
or len(set(refs)) != len(refs)):
raise CalculationConfigError(
"application_signal_refs must be unique signal strings"
)
if entry["rate_signal_ref"] != "unused":
raise CalculationConfigError("area_application cannot specify rate_signal_ref")
for name in ("erp_workplace", "production_order_format"):
if not isinstance(entry.get(name), str) or not entry[name].strip():
raise CalculationConfigError(f"area_application requires {name}")
from production_analytics.context import build_enlyze_production_order
try:
build_enlyze_production_order("0", entry["production_order_format"])
except ValueError as exc:
raise CalculationConfigError("Invalid production_order_format") from exc
entry["application_signal_refs"] = tuple(refs)
elif any(name in entry for name in (
"application_signal_refs", "erp_workplace", "production_order_format",
)):
raise CalculationConfigError("Area context fields require area_application")
if entry["id"] in instances: 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)
@@ -102,3 +102,20 @@ 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
+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
+20 -9
View File
@@ -4,7 +4,10 @@ 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,
)
from production_analytics.enlyze.exploration import ExplorationClient from production_analytics.enlyze.exploration import ExplorationClient
@@ -89,20 +92,21 @@ class EnlyzeApiGateway:
gate_variable_id: str, gate_variable_id: str,
start: datetime, start: datetime,
end: datetime, end: datetime,
application_variable_ids: tuple[str, ...] = (),
nominal_width_m: 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")
source_ids = application_variable_ids or (rate_variable_id,)
variable_ids = list(dict.fromkeys((*source_ids, gate_variable_id)))
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,11 +118,11 @@ 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 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")
@@ -126,12 +130,19 @@ class EnlyzeApiGateway:
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]
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 application_variable_ids:
if nominal_width_m is None:
raise ValueError("area application requires nominal width")
rate = area_application_rate_kg_per_hour(sources, nominal_width_m, gate)
else:
rate = sources[0]
samples.append(MaterialSample( 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),
@@ -0,0 +1,24 @@
"""ERP nominal-width resolution for an explicitly matched production order."""
from production_analytics.context import build_enlyze_production_order, extract_nominal_width_m
from production_analytics.erp import ErpWorkplaceStatusGateway
class ErpNominalWidthProvider:
def __init__(
self, gateway: ErpWorkplaceStatusGateway, *, workplace: str, format_template: str,
) -> None:
self.gateway = gateway
self.workplace = workplace
self.format_template = format_template
def __call__(self, production_order: str) -> float:
status = self.gateway.get_current_workplace_status(self.workplace)
if (status is None or status.workplace != self.workplace
or build_enlyze_production_order(status.production_order, self.format_template)
!= production_order):
raise ValueError("ERP context does not match the production order")
width = extract_nominal_width_m(status.article_description)
if width is None:
raise ValueError("ERP article description has no unambiguous nominal width")
return width
@@ -92,7 +92,7 @@ class MaterialEfficiencyService:
enlyze_production_order=order, article_number=status.article_number, enlyze_production_order=order, article_number=status.article_number,
article_description=status.article_description, article_description=status.article_description,
nominal_width_m=extract_nominal_width_m(status.article_description), nominal_width_m=extract_nominal_width_m(status.article_description),
erp_feedback_timestamp=status.feedback_timestamp, erp_feedback_timestamp=cutoff,
material_snapshot_timestamp=material.timestamp, material_snapshot_timestamp=material.timestamp,
good_quantity_m2=quantity, material_consumption_kg=consumption, good_quantity_m2=quantity, material_consumption_kg=consumption,
material_consumption_kg_per_m2=kg_per_m2, material_consumption_kg_per_m2=kg_per_m2,
@@ -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,82 @@
"""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"}
):
raise CalculationConfigError("Invalid material efficiency configuration fields")
for name in required:
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,
),
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
@@ -53,7 +54,13 @@ class MaterialPollingService:
gate_variable_id: str, gate_variable_id: str,
gate_threshold: float, gate_threshold: float,
max_sample_gap_seconds: float, max_sample_gap_seconds: float,
application_variable_ids: tuple[str, ...] = (),
nominal_width_provider: Callable[[str], float] | None = None,
) -> None: ) -> None:
self._application_variable_ids = application_variable_ids
self._nominal_width_provider = nominal_width_provider
if application_variable_ids and nominal_width_provider is None:
raise ValueError("area application requires a nominal width provider")
self._gateway = gateway self._gateway = gateway
self._state_store = state_store self._state_store = state_store
self._machine_id = machine_id self._machine_id = machine_id
@@ -71,6 +78,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 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,
@@ -91,7 +100,7 @@ class MaterialPollingService:
if previous.end is None or previous.end > run.start: if previous.end is None or previous.end > run.start:
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,
) )
initial_state = MaterialIntegrationState( initial_state = MaterialIntegrationState(
cumulative_consumption_kg=integrated.cumulative_consumption_kg, cumulative_consumption_kg=integrated.cumulative_consumption_kg,
@@ -116,7 +125,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)
self._state_store.save( self._state_store.save(
self._machine_id, self._machine_id,
@@ -134,13 +143,20 @@ 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,
) -> 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,
)
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) return integrator.process_many(samples)
@@ -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,
@@ -48,6 +50,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 +72,16 @@ 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 == "area_application":
source_options = dict(
application_variable_ids=calculation.application_signal_refs,
nominal_width_provider=ErpNominalWidthProvider(
ErpWorkplaceStatusGateway(ErpSettings.from_secret_file(erp_secrets_file)),
workplace=calculation.erp_workplace,
format_template=calculation.production_order_format,
),
)
service = MaterialPollingService( service = MaterialPollingService(
gateway=EnlyzeApiGateway(ExplorationClient(settings)), gateway=EnlyzeApiGateway(ExplorationClient(settings)),
state_store=_ReportingStateStore(JsonMaterialStateStore(state_directory)), state_store=_ReportingStateStore(JsonMaterialStateStore(state_directory)),
@@ -77,6 +90,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,
@@ -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
+216
View File
@@ -0,0 +1,216 @@
from dataclasses import replace
from datetime import datetime, timedelta
from types import SimpleNamespace
from unittest.mock import Mock
import pytest
import yaml
from production_analytics.calculations.config import (
CalculationConfigError,
load_material_calculation,
)
from production_analytics.calculations.material_consumption import (
MaterialConsumptionIntegrator,
MaterialIntegratorConfig,
area_application_rate_kg_per_hour,
)
from production_analytics.enlyze.gateway import EnlyzeApiGateway, EnlyzeProductionRun
from production_analytics.service.material_context import ErpNominalWidthProvider
from production_analytics.service.material_polling import MaterialPollingService
from production_analytics.service.material_state_store import JsonMaterialStateStore
def timestamp(value):
return datetime.fromisoformat(value.replace('Z', '+00:00'))
def test_area_units_and_strict_gate():
client = Mock()
start = timestamp('2026-08-25T05:27:18Z')
client.post_json.return_value.body = {'data': {
'columns': ['speed', 's2', 'time', 's1'],
'records': [[speed, 2000, (start + timedelta(minutes=i)).isoformat(), 3000]
for i, speed in enumerate([2, 0.3, 0, 2, 2])],
}}
samples = EnlyzeApiGateway(client).get_material_samples(
machine_id='machine', rate_variable_id='unused', gate_variable_id='speed',
application_variable_ids=('s1', 's2'), nominal_width_m=5,
start=start, end=start + timedelta(minutes=4),
)
state = MaterialConsumptionIntegrator(MaterialIntegratorConfig(0.3, 60)).process_many(samples)
assert state.cumulative_consumption_kg == 100
assert state.integrated_running_seconds == 120
assert client.post_json.call_args.args[1]['variables'] == [
{'uuid': 's1'}, {'uuid': 's2'}, {'uuid': 'speed'},
]
@pytest.mark.parametrize('sources,width,speed', [
([], 5, 2), ([float('nan')], 5, 2), ([1], 0, 2), ([1], 5, float('inf')),
])
def test_invalid_area_inputs(sources, width, speed):
with pytest.raises(ValueError):
area_application_rate_kg_per_hour(sources, width, speed)
@pytest.mark.parametrize('bad', [None, True, float('nan')])
def test_missing_or_invalid_spreader_fails(bad):
client = Mock()
now = timestamp('2026-08-25T05:27:18Z')
client.post_json.return_value.body = {'data': {
'columns': ['time', 's1', 's2', 'speed'],
'records': [[now.isoformat(), 3000, bad, 2]],
}}
with pytest.raises(ValueError):
EnlyzeApiGateway(client).get_material_samples(
machine_id='machine', rate_variable_id='unused', gate_variable_id='speed',
application_variable_ids=('s1', 's2'), nominal_width_m=5, start=now, end=now,
)
def width_provider():
erp = Mock()
erp.get_current_workplace_status.return_value = SimpleNamespace(
workplace='Bento 1', production_order='12026000814', article_number='180305',
article_description='Bfix NSP 5300, 5,00 x 40 m',
)
return ErpNominalWidthProvider(
erp, workplace='Bento 1', format_template='Bento 1-{production_order}',
)
@pytest.mark.parametrize('field,value', [
('production_order', '12026000815'), ('workplace', 'K7'),
('article_description', None), ('article_description', '180305'),
])
def test_width_context_fails_closed(field, value):
provider = width_provider()
setattr(provider.gateway.get_current_workplace_status.return_value, field, value)
with pytest.raises(ValueError):
provider('Bento 1-12026000814')
def test_bento_reference_aggregate_bootstrap_and_persistent_resume(tmp_path):
"""Synthetic aggregate fixture, not recorded ENLYZE samples.
Uses reported active time/area/application within the actual run boundaries.
"""
config = load_material_calculation('config/bento1-material-consumption.yaml')
order = 'Bento 1-12026000814'
runs = [EnlyzeProductionRun(
uuid, config.machine_ref, 'external-product', order, timestamp(start), timestamp(end),
) for uuid, start, end in [
('b675818c-4636-4976-be6c-3e0e24da9e8a',
'2026-08-25T05:27:18Z', '2026-08-26T04:40:58Z'),
('25d67424-8346-42d4-b844-75b4097382bd',
'2026-08-26T05:43:39Z', '2026-08-26T05:54:05Z'),
]]
speed = 31294.6 / (5 * 1256)
client = Mock()
def response(path, body):
start, end = timestamp(body['start']), timestamp(body['end'])
run = next(run for run in runs if run.start <= start <= run.end)
active_seconds = 1256 * 60 - 626 if run == runs[0] else 626
stop = run.start + timedelta(seconds=active_seconds)
times = {start, end}
t = start
while t < end:
times.add(t)
t += timedelta(seconds=10)
if start <= stop <= end:
times.add(stop)
return SimpleNamespace(body={'data': {
'columns': ['time', *config.application_signal_refs, config.gate_signal_ref],
'records': [[t.isoformat(), 2300, 2398, speed if t < stop else 0]
for t in sorted(times)],
}})
client.post_json.side_effect = response
gateway = EnlyzeApiGateway(client)
gateway.get_open_production_run = Mock(return_value=replace(runs[1], end=None))
gateway.get_production_runs = Mock(return_value=runs)
store = JsonMaterialStateStore(tmp_path)
def service():
return MaterialPollingService(
gateway=gateway, state_store=store, machine_id=config.machine_ref,
rate_variable_id=config.rate_signal_ref, gate_variable_id=config.gate_signal_ref,
gate_threshold=config.gate_threshold,
max_sample_gap_seconds=config.max_sample_gap_seconds,
application_variable_ids=config.application_signal_refs,
nominal_width_provider=width_provider(),
)
result = service().poll_once(now=runs[1].end)
assert result.state.integrated_running_seconds / 60 == pytest.approx(1256)
assert result.state.cumulative_consumption_kg == pytest.approx(147021.9, abs=0.2)
assert result.state.integrated_running_seconds / 60 * speed * 5 == pytest.approx(31294.6)
assert service().poll_once(now=runs[1].end).state == result.state
assert gateway.get_production_runs.call_count == 1
assert [(timestamp(call.args[1]['start']), timestamp(call.args[1]['end']))
for call in client.post_json.call_args_list[:2]] == [
(run.start, run.end) for run in runs]
def test_k7_config_defaults_remain_direct():
config = load_material_calculation('config/k7-material-consumption.yaml')
assert config.source_mode == 'direct_mass_rate'
assert config.application_signal_refs == ()
assert config.gate_threshold == 0.5
assert config.rate_signal_ref == 'c9d06af5-f6d6-4ede-b6c4-5a98bac77129'
@pytest.mark.parametrize('field,value', [
('source_mode', 'unknown'), ('application_signal_refs', []),
('application_signal_refs', ['same', 'same']), ('erp_workplace', ''),
('production_order_format', '{wrong}'), ('rate_signal_ref', 'direct'),
])
def test_invalid_area_configuration(tmp_path, field, value):
with open('config/bento1-material-consumption.yaml') as stream:
document = yaml.safe_load(stream)
document['calculations'][0][field] = value
path = tmp_path / 'config.yaml'
path.write_text(yaml.safe_dump(document))
with pytest.raises(CalculationConfigError):
load_material_calculation(path)
def test_area_context_failure_does_not_touch_persistent_state(tmp_path):
config = load_material_calculation('config/bento1-material-consumption.yaml')
gateway = Mock()
now = timestamp('2026-08-25T05:27:18Z')
gateway.get_open_production_run.return_value = EnlyzeProductionRun(
'run', config.machine_ref, None, 'Bento 1-12026000814', now, None,
)
provider = width_provider()
provider.gateway.get_current_workplace_status.return_value = None
store = Mock()
service = MaterialPollingService(
gateway=gateway, state_store=store, machine_id=config.machine_ref,
rate_variable_id=config.rate_signal_ref, gate_variable_id=config.gate_signal_ref,
gate_threshold=config.gate_threshold, max_sample_gap_seconds=20,
application_variable_ids=config.application_signal_refs, nominal_width_provider=provider,
)
with pytest.raises(ValueError, match='context'):
service.poll_once(now=now)
gateway.get_material_samples.assert_not_called()
store.save.assert_not_called()
def test_area_gateway_pagination():
client = Mock()
now = timestamp('2026-08-25T05:27:18Z')
client.post_json.side_effect = [SimpleNamespace(body={
'metadata': {'next_cursor': cursor},
'data': {'columns': ['time', 's1', 's2', 'speed'],
'records': [[(now + timedelta(seconds=i * 10)).isoformat(), 3000, 2000, 2]]},
}) for i, cursor in enumerate(['next', None])]
samples = EnlyzeApiGateway(client).get_material_samples(
machine_id='machine', rate_variable_id='unused', gate_variable_id='speed',
application_variable_ids=('s1', 's2'), nominal_width_m=5,
start=now, end=now + timedelta(seconds=10),
)
assert [sample.material_rate_kg_per_hour for sample in samples] == [3000, 3000]
assert client.post_json.call_args.args[1]['cursor'] == 'next'
+6 -4
View File
@@ -152,14 +152,15 @@ def test_temporal_regression_and_successive_feedback():
assert first.good_quantity_m2 == 800 assert first.good_quantity_m2 == 800
def test_naive_feedback_requires_explicit_timezone_and_preserves_source(): def test_naive_feedback_requires_explicit_timezone_and_returns_instant():
naive = datetime(2026, 9, 4, 12, 5) naive = datetime(2026, 9, 4, 12, 5)
status = replace(STATUS, feedback_timestamp=naive) status = replace(STATUS, feedback_timestamp=naive)
with pytest.raises(ValueError, match='explicit erp_timezone'): with pytest.raises(ValueError, match='explicit erp_timezone'):
service().evaluate(status) service().evaluate(status)
subject = service(erp_timezone=ZoneInfo('Europe/Berlin')) subject = service(erp_timezone=ZoneInfo('Europe/Berlin'))
result = subject.evaluate(status) result = subject.evaluate(status)
assert result.erp_feedback_timestamp == naive assert result.erp_feedback_timestamp == FEEDBACK
assert status.feedback_timestamp == naive
assert subject.materials.latest_at_or_before.call_args.kwargs['timestamp'] == FEEDBACK assert subject.materials.latest_at_or_before.call_args.kwargs['timestamp'] == FEEDBACK
@@ -171,10 +172,11 @@ def test_dst_gap_and_overlap_rejected(timestamp):
) )
def test_aware_feedback_preserves_offset(): def test_aware_feedback_returns_utc():
timestamp = FEEDBACK.astimezone(ZoneInfo('Europe/Berlin')) timestamp = FEEDBACK.astimezone(ZoneInfo('Europe/Berlin'))
result = service().evaluate(replace(STATUS, feedback_timestamp=timestamp)) result = service().evaluate(replace(STATUS, feedback_timestamp=timestamp))
assert result.erp_feedback_timestamp is timestamp assert result.erp_feedback_timestamp == timestamp
assert result.erp_feedback_timestamp.tzinfo is UTC
def test_database_failure_propagates(): def test_database_failure_propagates():
@@ -0,0 +1,447 @@
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"