Compare commits

...
7 Commits
23 changed files with 2521 additions and 61 deletions
+3
View File
@@ -29,3 +29,6 @@ secrets/*
.DS_Store
activate.sh
# Live polling checkpoints
data/state/
+164 -8
View File
@@ -7,8 +7,11 @@ production-order context, and calculation state.
## Status
This repository includes a read-only ENLYZE API exploration CLI. It contains
no verified ENLYZE operation wrappers, database migrations, or HTTP endpoints.
This repository includes an ENLYZE exploration CLI, a production-run/timeseries
gateway, and a configured continuous material polling runner with atomic JSON
checkpoints and PostgreSQL cumulative material-consumption snapshots for Grafana.
The live material-consumption MVP remains in progress.
There are no HTTP endpoints.
## Intended flow
@@ -76,7 +79,7 @@ Call `process(sample)` or `process_many(samples)` on the same instance for live
input, replay, or successive chunks. Both use the same calculation. The immutable
`state` snapshot exposes cumulative kg, integrated running seconds, last timestamp
(UTC), last rate, last gate, and whether the latest gate is active. State restoration
and persistence are not implemented yet. Naive timestamps and non-finite values
and JSON persistence are supported. Naive timestamps and non-finite values
are rejected; backwards timestamps raise without changing state. Duplicate
timestamps add no consumption but replace the baseline in arrival order.
Finite negative material rates are currently accepted and decrease cumulative
@@ -88,13 +91,167 @@ Run `python scripts/validate_k7_material_consumption.py` from the repository roo
with the ignored `k7-00842-throughput.raw.json` and `k7-00842-speed.raw.json`
captures in `data/raw/enlyze/`. The utility joins common non-null timestamps
without filling values and rejects unordered or duplicate capture records.
K7-specific signal UUIDs are confined to the utility: `Stundenleistung Anlage`
is gated by `Geschwindigkeit Gesamtanlage > 0.5 m/min`. With the explicit
K7-specific signal UUIDs are declared in the runner configuration and validation
utility: `Stundenleistung Anlage` is gated by `Geschwindigkeit Gesamtanlage > 0.5 m/min`. With the explicit
20-second validation gap limit (`--max-sample-gap-seconds` to override),
2,672 common samples yield **5.205555556 h** and **5,180.127150811 kg** for run
00842, matching the previous manual calculation's rounded results. Synthetic
tests verify equivalence to manual interval integration without local captures
or live ENLYZE access. Production polling, Grafana, and Bento 1 are not implemented.
or live ENLYZE access. Grafana and Bento 1 remain pending.
`MaterialPollingService.poll_once(now=...)` starts at the open run's start or
resumes inclusively at its last processed timestamp. Duplicate boundary samples
add no time. A new run UUID for the same order retains cumulative kg and running
seconds but resets the temporal baseline; different orders and machines have
separate state. Empty responses add no consumption. Naive `now` is rejected;
no open run or a window beginning after `now` returns `None` without saving.
`JsonMaterialStateStore` uses deterministic SHA-256 filenames derived from the
machine/order pair, preventing path traversal and sanitized-name collisions.
It flushes and fsyncs a temporary file in the state directory before atomic
replacement; a failed write leaves the previous primary file intact. Missing
state returns `None`; malformed JSON or state raises `ValueError` with file
context. Use an ignored runtime directory and one polling writer per state key;
atomic replacement does not coordinate concurrent read-modify-write cycles.
Files from the earlier uncommitted sanitized-filename prototype are not loaded
under the new names. The gateway rejects ambiguous open runs, naive windows,
and missing columns or malformed records instead of silently skipping them.
## Continuous material polling
The committed [K7 configuration](config/k7-material-consumption.yaml) declares
`k7-fiber-consumption`, a `material_consumption` calculation at version `"1"`.
It selects `Stundenleistung Anlage` (kg/h) as the rate and `Geschwindigkeit
Gesamtanlage` (m/min) as the gate, integrating only when the gate is **strictly
above 0.5**, with a maximum sample gap of 20 seconds. `Anlage läuft` is not the
primary gate. The output is `material_consumption` in `kg`, grouped by
`production_order`; production orders remain opaque strings, including leading
zeros and whitespace.
Install the declared dependencies in the repository venv, then run from the
repository root:
```bash
.venv/bin/python -m pip install -e '.[dev]'
.venv/bin/python -m production_analytics.cli run material-poll \
--config config/k7-material-consumption.yaml \
--calculation-id k7-fiber-consumption \
--poll-interval-seconds 10 \
--state-directory data/state/material \
--secrets-file secrets/enlyze.env
```
The installed `production-analytics run material-poll` command is equivalent.
The interval, state directory, and secrets path above are defaults. Paths are
relative to the current working directory. Existing `ENLYZE_BASE_URL`,
`ENLYZE_API_KEY`, and `ENLYZE_HTTP_TIMEOUT_SECONDS` handling is reused; values
in the secrets file override environment values. Keys are never CLI arguments.
Keep calculation fields in YAML and technical runtime settings in CLI options.
The loader requires every field shown in the K7 file, rejects unknown fields,
duplicate keys/ids, unsupported types/versions/outputs, and invalid or non-finite
numbers. Numeric fields must be YAML numbers, and `version` must be quoted.
Multiple material calculations may share a file, but then `--calculation-id` is
required. All entries are validated before selection. The unrelated
`config/calculations.example.yaml` remains illustrative and is not executable
by this material-only runner.
The foreground runner polls once immediately using UTC, then sleeps for the
configured positive interval after each completed cycle, including failures.
A slow poll delays the next cycle; there is no overlap or catch-up scheduling.
Ctrl-C/SIGINT exits cleanly. Successful cycles print machine, opaque order,
cumulative kg, integrated running seconds, and cycle-start timestamp. No open
run (or no eligible window) is normal and visible. Errors on stderr identify
state loading/saving or ENLYZE/gateway/polling and the exception class; arbitrary
exception text and response bodies are withheld to avoid leaking credentials.
Startup/configuration failures return non-zero; cycle failures retry after the
delay without replacing checkpoints with empty state.
Checkpoints live in ignored `data/state/material/` by default, with one hashed
JSON filename per machine/order pair. Run **one polling writer per state key**;
there is no multi-process locking. Use a separate state directory when changing
calculation parameters or running a different calculation for the same machine
and order: checkpoints are not namespaced by calculation id/version. These files
are integration checkpoints, not a Grafana metric store.
Each successful non-empty poll writes one cumulative snapshot to PostgreSQL before
printing its result. JSON remains the restart checkpoint; PostgreSQL stores derived
time-series snapshots only. ENLYZE remains the raw-data source of truth. Database
write failures emit a concise error and polling continues with the next cycle;
missed snapshots are not retried or backfilled and JSON state is not rolled back.
The next successful snapshot includes the continuing cumulative total.
A new order's first snapshot is **not guaranteed to be zero**: the service processes
available samples from the run start before returning, so its first total may
already be non-zero. A returned zero is stored normally. Totals continue across
runs for the same order; existing checkpoints for a previously seen order resume.
### PostgreSQL setup
Set all five runtime settings: `POSTGRES_HOST`, `POSTGRES_PORT`, `POSTGRES_DB`,
`POSTGRES_USER`, and `POSTGRES_PASSWORD`, using exported environment variables or
the existing secrets file (file values override the environment). The runner does
not automatically load `.env`. Missing or invalid settings fail at startup;
connectivity and schema errors are reported during polling. For the existing
Compose instance, use host `localhost`, the published port (default `5432`), and
`production_analytics` for both database and user. Set the same password in
Compose's `.env` and the runner environment/secrets file.
Start the database and apply the repeatable schema from the repository root:
```bash
docker compose up -d timescaledb
docker compose exec -T timescaledb psql -U production_analytics -d production_analytics \
-v ON_ERROR_STOP=1 < db/schema.sql
```
Wait until PostgreSQL is ready before applying the schema. Configure Grafana's
PostgreSQL data source to query `material_consumption_snapshots`, selecting
`timestamp` as time and `consumption_kg` as the cumulative value, filtered by
`calculation_id`, `machine_id`, and `production_order`. These are ordinary
PostgreSQL tables/indexes; no Timescale-specific features or hypertables are used.
The primary key deduplicates calculation/machine/order/timestamp (first write wins).
A short transaction opens and closes one synchronous connection per snapshot,
with 10-second connection and statement timeouts. Snapshot timestamps are poll
start times, not source-sample timestamps. Schema application is manual.
Bento 1 bentonite source and gate selection remain intentionally undefined pending
process validation; the generic integration core is unchanged and reusable.
## ERP current workplace status
The read-only ERP adapter reads production-order/article context from MSSQL
`NV_DWH.dbo.GRAFANA_WORKPLACE_STATUS`, for later live material KPIs. Credentials
belong in ignored `secrets/erp.env`, using `ERP_DB_HOST`, `ERP_DB_PORT`,
`ERP_DB_NAME`, `ERP_DB_USER`, and `ERP_DB_PASSWORD`. The verified endpoint is
`192.168.111.24:49601`, database `NV_DWH`, using SQL Server authentication and
a read-only account. All settings are required; ports must be integers 1–65535.
```python
from production_analytics.erp import ErpSettings, ErpWorkplaceStatusGateway
settings = ErpSettings.from_secret_file() # File values override environment values.
status = ErpWorkplaceStatusGateway(settings).get_current_workplace_status("K7")
```
Alternatively, use `ErpSettings.from_environment(mapping)`. The existing simple
dotenv loader is reused; no shell content is executed. Each call opens and closes
one connection, with 10-second login/query timeouts, and performs a parameterized
SELECT filtered by workplace. Zero rows returns `None`, one row returns an
immutable `CurrentWorkplaceStatus`, and multiple rows raise `ErpReadError`.
Driver failures and malformed rows also raise `ErpReadError` with safe messages.
Only trailing whitespace is removed from workplace, production order, article
number, and description; identifiers otherwise remain unchanged. Nullable fields
stay `None`. Decimal/integer quantities are explicitly converted to finite floats
(with normal floating-point precision); non-finite/overflowing values are rejected.
Quantity fields use m² and remaining time uses hours, following the supplied source
semantics. `feedback_timestamp` is preserved exactly, including a naive timezone.
It represents the latest ERP feedback, and the ERP export is delayed relative to
live process data; the adapter makes no freshness inference.
This milestone adds no kg/m² or material-efficiency calculation, ERP-to-ENLYZE
order mapping, persistence, polling, CLI, or Grafana integration. Automated ERP
tests use mocks and require no live connectivity.
## Peak-cycle detection
@@ -132,8 +289,7 @@ completed peaks, including a legitimate approximately 1195 kg cycle. Run
available. Neither capture should be added to version control or automated
tests.
`docker compose up -d timescaledb` is an optional local database design for a
future persistence milestone. It is not required for the bootstrap tests.
The normal unit test suite mocks PostgreSQL and requires no live database.
See [PROJECT_KNOWLEDGE.md](PROJECT_KNOWLEDGE.md) for durable project context
and [docs/roadmap.md](docs/roadmap.md) for the implementation sequence.
+14
View File
@@ -0,0 +1,14 @@
# Validated K7: Stundenleistung Anlage (kg/h), gated by
# Geschwindigkeit Gesamtanlage (m/min) > 0.5. Not Anlage läuft.
calculations:
- id: k7-fiber-consumption
type: material_consumption
version: "1"
machine_ref: c220f95c-a65e-4cb7-99b7-0626d6c7508c
rate_signal_ref: c9d06af5-f6d6-4ede-b6c4-5a98bac77129
gate_signal_ref: 823867bb-f5d2-40eb-b875-657155addfd0
gate_threshold: 0.5
max_sample_gap_seconds: 20.0
output_metric: material_consumption
output_unit: kg
group_by: production_order
+13
View File
@@ -0,0 +1,13 @@
CREATE TABLE IF NOT EXISTS material_consumption_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,
consumption_kg double precision NOT NULL,
PRIMARY KEY (calculation_id, machine_id, production_order, timestamp)
);
-- Machine time series across orders; the primary key covers a specific order.
CREATE INDEX IF NOT EXISTS material_consumption_snapshots_machine_time_idx
ON material_consumption_snapshots (calculation_id, machine_id, timestamp);
+17 -1
View File
@@ -13,7 +13,23 @@ Exploration established a K7 candidate mass-flow signal and speed gate; the
evidence and remaining uncertainties are recorded in `docs/enlyze-api.md`.
Continue to record sanitized fixtures for newly verified semantics.
## 2. Live material-consumption MVP (next implementation milestone)
## 2. Live material-consumption MVP (in progress)
The integration engine, ENLYZE gateway, single-cycle polling service, atomic
JSON checkpoint store, and configured continuous foreground polling runner are
implemented and covered by synthetic tests. Configure the validated K7 instance
in `config/k7-material-consumption.yaml`; start it with
`production-analytics run material-poll --config config/k7-material-consumption.yaml`.
Runtime options select the polling interval, state directory (default: ignored
`data/state/material/`), and existing secrets-file handling. Only one writer may
poll each machine/order state key.
The runner now persists cumulative material-consumption snapshots for Grafana in
PostgreSQL using `db/schema.sql`. JSON remains the restart checkpoint; PostgreSQL
stores derived time-series snapshots only. No Timescale-specific features are used.
A new order may first appear with a non-zero total after initial sample integration.
Grafana dashboard validation remains pending. No HTTP endpoints, application
containers, daemonization, or schedulers have been added.
Detect the currently active Production Run and production order, continuously
ingest new ENLYZE samples, maintain persistent incremental integration state,
+1 -1
View File
@@ -10,7 +10,7 @@ readme = "README.md"
requires-python = ">=3.11"
license = { text = "Proprietary" }
authors = [{ name = "Production Analytics Team" }]
dependencies = []
dependencies = ["PyYAML>=6.0", "psycopg[binary]>=3.2,<4", "pymssql>=2.4.0,<3"]
[project.optional-dependencies]
dev = [
@@ -0,0 +1,110 @@
"""Strict configuration for executable material-consumption instances."""
from dataclasses import dataclass, fields
from math import isfinite
from pathlib import Path
import yaml
class CalculationConfigError(ValueError):
"""An operator-readable configuration failure, without input values."""
class _UniqueLoader(yaml.SafeLoader):
pass
def _mapping(loader, node):
result = {}
for key_node, value_node in node.value:
key = loader.construct_object(key_node)
if not isinstance(key, str) or key in result:
raise CalculationConfigError("YAML mapping keys must be unique strings")
result[key] = loader.construct_object(value_node)
return result
_UniqueLoader.add_constructor(yaml.resolver.BaseResolver.DEFAULT_MAPPING_TAG, _mapping)
def finite_number(value: object, name: str, *, positive: bool = False) -> float:
if isinstance(value, bool) or not isinstance(value, (int, float)):
raise CalculationConfigError(f"{name} must be a finite number")
try:
number = float(value)
except (OverflowError, ValueError) as exc:
raise CalculationConfigError(f"{name} must be a finite number") from exc
if not isfinite(number) or (positive and number <= 0):
requirement = "finite and greater than zero" if positive else "finite"
raise CalculationConfigError(f"{name} must be {requirement}")
return number
@dataclass(frozen=True, slots=True)
class MaterialCalculationConfig:
id: str
type: str
version: str
machine_ref: str
rate_signal_ref: str
gate_signal_ref: str
gate_threshold: float
max_sample_gap_seconds: float
output_metric: str
output_unit: str
group_by: str
def load_material_calculation(
path: str | Path, calculation_id: str | None = None,
) -> MaterialCalculationConfig:
try:
with Path(path).open(encoding="utf-8") as stream:
document = yaml.load(stream, Loader=_UniqueLoader)
except (OSError, UnicodeError) as exc:
raise CalculationConfigError("Cannot read calculation configuration file") from exc
except yaml.YAMLError as exc:
raise CalculationConfigError("Invalid YAML in calculation configuration") from exc
if not isinstance(document, dict) or set(document) != {"calculations"}:
raise CalculationConfigError("Configuration must contain only 'calculations'")
entries = document["calculations"]
if not isinstance(entries, list) or not entries:
raise CalculationConfigError("calculations must be a non-empty list")
required = {field.name for field in fields(MaterialCalculationConfig)}
instances = {}
for index, entry in enumerate(entries):
prefix = f"calculations[{index}]"
if not isinstance(entry, dict):
raise CalculationConfigError(f"{prefix} must be a mapping")
if "type" in entry and entry["type"] != "material_consumption":
raise CalculationConfigError(f"{prefix}.type: only 'material_consumption' is supported")
missing = required - entry.keys()
if missing:
raise CalculationConfigError(f"{prefix}: missing fields: {', '.join(sorted(missing))}")
if entry.keys() - required:
raise CalculationConfigError(f"{prefix}: unknown fields (check spelling)")
for name in sorted(required - {"gate_threshold", "max_sample_gap_seconds"}):
if not isinstance(entry[name], str) or not entry[name].strip():
raise CalculationConfigError(f"{prefix}.{name} must be a non-empty string")
for name, expected in (
("type", "material_consumption"), ("version", "1"),
("output_metric", "material_consumption"), ("output_unit", "kg"),
("group_by", "production_order"),
):
if entry[name] != expected:
raise CalculationConfigError(f"{prefix}.{name}: only {expected!r} is supported")
entry["gate_threshold"] = finite_number(entry["gate_threshold"], f"{prefix}.gate_threshold")
entry["max_sample_gap_seconds"] = finite_number(
entry["max_sample_gap_seconds"], f"{prefix}.max_sample_gap_seconds", positive=True,
)
if entry["id"] in instances:
raise CalculationConfigError("Calculation ids must be unique")
instances[entry["id"]] = MaterialCalculationConfig(**entry)
if calculation_id is not None:
if calculation_id not in instances:
raise CalculationConfigError("Requested calculation id was not found")
return instances[calculation_id]
if len(instances) != 1:
raise CalculationConfigError("Multiple calculations: specify --calculation-id")
return next(iter(instances.values()))
@@ -1,6 +1,7 @@
"""Serialization helpers for persistent material-integration state."""
from datetime import UTC, datetime
from math import isfinite
from production_analytics.calculations.material_consumption import (
MaterialIntegrationState,
@@ -9,6 +10,11 @@ from production_analytics.calculations.material_consumption import (
def material_state_to_dict(state: MaterialIntegrationState) -> dict[str, object]:
"""Convert integration state to a JSON-serializable dictionary."""
if state.last_processed_timestamp is not None and (
state.last_processed_timestamp.tzinfo is None
or state.last_processed_timestamp.utcoffset() is None
):
raise ValueError("last_processed_timestamp must be timezone-aware")
return {
"cumulative_consumption_kg": state.cumulative_consumption_kg,
"integrated_running_seconds": state.integrated_running_seconds,
@@ -25,7 +31,33 @@ def material_state_to_dict(state: MaterialIntegrationState) -> dict[str, object]
def material_state_from_dict(data: dict[str, object]) -> MaterialIntegrationState:
"""Restore integration state from a JSON-compatible dictionary."""
if not isinstance(data, dict):
raise ValueError("integration_state must be an object")
for key in (
"cumulative_consumption_kg", "integrated_running_seconds",
"last_material_rate_kg_per_hour", "last_gate_value",
):
value = data[key]
if value is None and key.startswith("last_"):
continue
if isinstance(value, bool) or not isinstance(value, (int, float)) or not isfinite(value):
raise ValueError(f"{key} must be a finite number")
if data["integrated_running_seconds"] < 0:
raise ValueError("integrated_running_seconds must not be negative")
if not isinstance(data["integration_active"], bool):
raise ValueError("integration_active must be a boolean")
raw_timestamp = data["last_processed_timestamp"]
if raw_timestamp is not None:
if not isinstance(raw_timestamp, str):
raise ValueError("last_processed_timestamp must be an ISO timestamp")
parsed = datetime.fromisoformat(raw_timestamp)
if parsed.tzinfo is None or parsed.utcoffset() is None:
raise ValueError("last_processed_timestamp must be timezone-aware")
if data["integration_active"] and (
raw_timestamp is None or data["last_material_rate_kg_per_hour"] is None
or data["last_gate_value"] is None
):
raise ValueError("active integration requires a complete temporal baseline")
timestamp = (
datetime.fromisoformat(str(raw_timestamp)).astimezone(UTC)
+37 -1
View File
@@ -1,4 +1,4 @@
"""Command-line entry point for safe, read-only API exploration."""
"""Command-line entry point for exploration and production polling."""
from __future__ import annotations
@@ -70,6 +70,14 @@ def _parser() -> argparse.ArgumentParser:
default=Path("secrets/enlyze.env"),
help="Local dotenv-style secret file (default: secrets/enlyze.env)",
)
run = namespaces.add_parser("run", help="Foreground production runners")
runners = run.add_subparsers(dest="command", required=True)
material = runners.add_parser("material-poll", help="Continuously poll material consumption")
material.add_argument("--config", type=Path, required=True)
material.add_argument("--calculation-id")
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("--secrets-file", type=Path, default=Path("secrets/enlyze.env"))
return parser
@@ -83,8 +91,36 @@ def _write_fixture(directory: Path, name: str, payload: object) -> Path:
return destination
def _run_material(args: argparse.Namespace) -> int:
from production_analytics.calculations.config import (
CalculationConfigError,
load_material_calculation,
)
from production_analytics.enlyze.exploration import ConfigurationError
from production_analytics.service.material_runtime import build_material_runner
try:
calculation = load_material_calculation(args.config, args.calculation_id)
runner = build_material_runner(
calculation, poll_interval_seconds=args.poll_interval_seconds,
state_directory=args.state_directory, secrets_file=args.secrets_file,
)
runner.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:
args = _parser().parse_args(argv)
if args.namespace == "run":
return _run_material(args)
if args.namespace != "enlyze" or args.command not in {"raw", "timeseries"}:
return 2
+129 -50
View File
@@ -2,6 +2,7 @@
from dataclasses import dataclass
from datetime import UTC, datetime
from math import isfinite
from production_analytics.calculations.material_consumption import MaterialSample
from production_analytics.enlyze.exploration import ExplorationClient
@@ -22,29 +23,63 @@ class EnlyzeApiGateway:
self._client = client
def get_open_production_run(self, machine_id: str) -> EnlyzeProductionRun | None:
response = self._client.get(
"/v2/production-runs",
{"machine": machine_id},
)
for item in response.body.get("data", []):
if item.get("end") is not None:
continue
return EnlyzeProductionRun(
uuid=str(item["uuid"]),
machine_id=str(item["machine"]),
product_id=(
str(item["product"])
if item.get("product") is not None
else None
),
production_order=str(item["production_order"]),
start=_parse_timestamp(item["start"]),
end=None,
runs = [run for run in self.get_production_runs(machine_id) if run.end is None]
if len(runs) > 1:
raise ValueError(
"Invalid production-run response: multiple open Production Runs returned"
)
return runs[0] if runs else None
return None
def get_production_runs(self, machine_id: str) -> list[EnlyzeProductionRun]:
"""Retrieve all machine runs, following validated cursor pagination."""
params = {"machine": machine_id}
runs = []
followed_cursors: set[str] = set()
page = 1
while True:
response = self._client.get("/v2/production-runs", params)
try:
if not isinstance(response.body, dict):
raise ValueError("body must be an object")
items = response.body["data"]
if not isinstance(items, list):
raise ValueError("data must be a list")
for item in items:
if not isinstance(item, dict) or "end" not in item:
raise ValueError("run must be an object with an end field")
for key in ("uuid", "machine", "production_order"):
if not isinstance(item[key], str) or not item[key]:
raise ValueError(f"{key} must be a non-empty string")
if item["machine"] != machine_id:
raise ValueError("run machine does not match requested machine")
start = _parse_timestamp(item["start"])
end = _parse_timestamp(item["end"]) if item["end"] is not None else None
if end is not None and end < start:
raise ValueError("run end must not precede start")
runs.append(EnlyzeProductionRun(
uuid=item["uuid"], machine_id=item["machine"],
product_id=item.get("product"), production_order=item["production_order"],
start=start, end=end,
))
next_cursor = None
if "metadata" in response.body:
metadata = response.body["metadata"]
if not isinstance(metadata, dict) or "next_cursor" not in metadata:
raise ValueError("metadata must be an object with a next_cursor field")
next_cursor = metadata["next_cursor"]
if next_cursor is not None and (
not isinstance(next_cursor, str) or not next_cursor
):
raise ValueError("next_cursor must be null or a non-empty string")
if next_cursor is None:
return runs
if next_cursor in followed_cursors:
raise ValueError("next_cursor has already been followed")
followed_cursors.add(next_cursor)
params = {**params, "cursor": next_cursor}
page += 1
except (KeyError, TypeError, ValueError) as exc:
raise ValueError(f"Invalid production-run response on page {page}: {exc}") from exc
def get_material_samples(
self,
@@ -55,35 +90,79 @@ class EnlyzeApiGateway:
start: datetime,
end: datetime,
) -> list[MaterialSample]:
response = self._client.post_json(
"/v2/timeseries",
{
"machine": machine_id,
"start": start.astimezone(UTC).isoformat(),
"end": end.astimezone(UTC).isoformat(),
"variables": [
{"uuid": rate_variable_id},
{"uuid": gate_variable_id},
],
},
)
data = response.body["data"]
columns = data["columns"]
time_index = columns.index("time")
rate_index = columns.index(rate_variable_id)
gate_index = columns.index(gate_variable_id)
return [
MaterialSample(
timestamp=_parse_timestamp(record[time_index]),
material_rate_kg_per_hour=float(record[rate_index]),
gate_value=float(record[gate_index]),
)
for record in data["records"]
]
for name, value in (("start", start), ("end", end)):
if value.tzinfo is None or value.utcoffset() is None:
raise ValueError(f"{name} must be timezone-aware")
if end < start:
raise ValueError("end must not precede start")
request_body = {
"machine": machine_id,
"start": start.astimezone(UTC).isoformat(),
"end": end.astimezone(UTC).isoformat(),
"variables": [
{"uuid": rate_variable_id},
{"uuid": gate_variable_id},
],
}
samples = []
followed_cursors: set[str] = set()
page = 1
while True:
response = self._client.post_json("/v2/timeseries", request_body)
try:
data = response.body["data"]
columns = data["columns"]
if not isinstance(columns, list):
raise ValueError("columns must be a list")
for required in ("time", rate_variable_id, gate_variable_id):
if columns.count(required) != 1:
raise ValueError(f"required column {required!r} must occur exactly once")
time_index = columns.index("time")
rate_index = columns.index(rate_variable_id)
gate_index = columns.index(gate_variable_id)
if not isinstance(data["records"], list):
raise ValueError("records must be a list")
for index, record in enumerate(data["records"]):
try:
if not isinstance(record, list) or len(record) != len(columns):
raise ValueError("record must match columns")
rate, gate = record[rate_index], record[gate_index]
for value in (rate, gate):
if isinstance(value, bool) or not isinstance(value, (int, float)):
raise ValueError("rate and gate must be numeric")
if not isfinite(value):
raise ValueError("rate and gate must be finite")
samples.append(MaterialSample(
timestamp=_parse_timestamp(record[time_index]),
material_rate_kg_per_hour=float(rate), gate_value=float(gate),
))
except (TypeError, ValueError) as exc:
raise ValueError(f"malformed record {index}: {exc}") from exc
next_cursor = None
if "metadata" in response.body:
metadata = response.body["metadata"]
if not isinstance(metadata, dict) or "next_cursor" not in metadata:
raise ValueError("metadata must be an object with a next_cursor field")
next_cursor = metadata["next_cursor"]
if next_cursor is not None and (
not isinstance(next_cursor, str) or not next_cursor
):
raise ValueError("next_cursor must be null or a non-empty string")
if next_cursor is None:
return samples
if next_cursor in followed_cursors:
raise ValueError("next_cursor has already been followed")
followed_cursors.add(next_cursor)
request_body = {**request_body, "cursor": next_cursor}
page += 1
except (KeyError, TypeError, ValueError) as exc:
raise ValueError(f"Invalid timeseries response on page {page}: {exc}") from exc
def _parse_timestamp(value: str) -> datetime:
return datetime.fromisoformat(value.replace("Z", "+00:00")).astimezone(UTC)
if not isinstance(value, str):
raise ValueError("timestamp must be an ISO string")
timestamp = datetime.fromisoformat(value.replace("Z", "+00:00"))
if timestamp.tzinfo is None or timestamp.utcoffset() is None:
raise ValueError("timestamp must be timezone-aware")
return timestamp.astimezone(UTC)
+10
View File
@@ -0,0 +1,10 @@
"""Read-only ERP current workplace context."""
from production_analytics.erp.gateway import (
CurrentWorkplaceStatus,
ErpReadError,
ErpSettings,
ErpWorkplaceStatusGateway,
)
__all__ = ["CurrentWorkplaceStatus", "ErpReadError", "ErpSettings", "ErpWorkplaceStatusGateway"]
+148
View File
@@ -0,0 +1,148 @@
"""Focused MSSQL boundary for dbo.GRAFANA_WORKPLACE_STATUS."""
import os
from collections.abc import Mapping
from dataclasses import dataclass, field
from datetime import datetime
from decimal import Decimal
from math import isfinite
import pymssql
from production_analytics.enlyze.exploration import ConfigurationError, load_secret_file
@dataclass(frozen=True, slots=True)
class ErpSettings:
host: str
port: int
database: str
user: str = field(repr=False)
password: str = field(repr=False)
def __post_init__(self) -> None:
for attribute, suffix in (
("host", "HOST"), ("database", "NAME"), ("user", "USER"), ("password", "PASSWORD"),
):
value = getattr(self, attribute)
if not isinstance(value, str) or not value.strip() or "\x00" in value:
raise ConfigurationError(f"ERP_DB_{suffix} must be non-empty without NUL bytes")
if type(self.port) is not int or not 1 <= self.port <= 65535:
raise ConfigurationError("ERP_DB_PORT must be an integer from 1 to 65535")
@classmethod
def from_environment(cls, environment: Mapping[str, str]) -> "ErpSettings":
raw_port = environment.get("ERP_DB_PORT", "")
try:
if not isinstance(raw_port, str):
raise ValueError
port = int(raw_port)
except ValueError:
raise ConfigurationError("ERP_DB_PORT must be an integer from 1 to 65535") from None
return cls(
host=environment.get("ERP_DB_HOST", ""), port=port,
database=environment.get("ERP_DB_NAME", ""),
user=environment.get("ERP_DB_USER", ""),
password=environment.get("ERP_DB_PASSWORD", ""),
)
@classmethod
def from_secret_file(
cls, path: str | os.PathLike[str] = "secrets/erp.env",
*, environment: Mapping[str, str] | None = None,
) -> "ErpSettings":
"""Load dotenv values over the environment, following the existing convention."""
try:
values = load_secret_file(path)
except Exception:
raise ConfigurationError("ERP secret file could not be loaded") from None
base = os.environ if environment is None else environment
return cls.from_environment({**base, **values})
@dataclass(frozen=True, slots=True)
class CurrentWorkplaceStatus:
"""Latest ERP feedback; timestamp and timezone are preserved, with no freshness claim."""
workplace: str
production_order: str
article_number: str | None
article_description: str | None
feedback_timestamp: datetime
order_quantity_m2: float | None
good_quantity_m2: float | None
remaining_quantity_m2: float | None
remaining_time_hours: float | None
remaining_rolls: float | None
class ErpReadError(RuntimeError):
"""Safe operator-facing ERP read or result error."""
_CURRENT_STATUS_SQL = """SELECT
[Arbeitsplatz], [Fertigungsauftragsnummer], [Artikelnummer], [Artikelbezeichnung],
[Zeitstempel], [Auftragsmenge], [Gutmenge], [Restmenge],
[Verbleibende Zeit], [Verbleibende Rollen]
FROM [dbo].[GRAFANA_WORKPLACE_STATUS]
WHERE [Arbeitsplatz] = %s"""
class ErpWorkplaceStatusGateway:
def __init__(self, settings: ErpSettings) -> None:
self._settings = settings
def get_current_workplace_status(self, workplace: str) -> CurrentWorkplaceStatus | None:
if not isinstance(workplace, str) or not workplace.strip() or "\x00" in workplace:
raise ValueError("workplace must be a non-empty string without NUL bytes")
try:
with pymssql.connect(
server=self._settings.host, port=self._settings.port,
database=self._settings.database, user=self._settings.user,
password=self._settings.password, login_timeout=10, timeout=10,
) as connection:
with connection.cursor() as cursor:
cursor.execute(_CURRENT_STATUS_SQL, (workplace,))
# Two rows suffice to detect a violation without selecting an arbitrary row.
rows = cursor.fetchmany(2)
except Exception:
# Driver messages may include credentials; suppress their traceback chain too.
raise ErpReadError("ERP current workplace status read failed") from None
if not rows:
return None
if len(rows) > 1:
raise ErpReadError("ERP current workplace status returned multiple rows")
try:
row = rows[0]
if len(row) != 10 or not isinstance(row[4], datetime):
raise ValueError
return CurrentWorkplaceStatus(
workplace=_text(row[0]), production_order=_text(row[1]),
article_number=_text(row[2], nullable=True),
article_description=_text(row[3], nullable=True),
feedback_timestamp=row[4],
order_quantity_m2=_number(row[5]), good_quantity_m2=_number(row[6]),
remaining_quantity_m2=_number(row[7]), remaining_time_hours=_number(row[8]),
remaining_rolls=_number(row[9]),
)
except Exception:
raise ErpReadError("ERP current workplace status returned an invalid row") from None
def _text(value: object, *, nullable: bool = False) -> str | None:
if value is None and nullable:
return None
if not isinstance(value, str):
raise ValueError
return value.rstrip()
def _number(value: object) -> float | None:
if value is None:
return None
if isinstance(value, bool) or not isinstance(value, (Decimal, int, float)):
raise ValueError
result = float(value)
if not isfinite(result):
raise ValueError
return result
@@ -0,0 +1,146 @@
"""Single-cycle live material polling orchestration."""
from dataclasses import dataclass
from datetime import datetime
from typing import Protocol
from production_analytics.calculations.material_consumption import (
MaterialConsumptionIntegrator,
MaterialIntegrationState,
MaterialIntegratorConfig,
)
from production_analytics.enlyze.gateway import (
EnlyzeApiGateway,
EnlyzeProductionRun,
)
@dataclass(frozen=True, slots=True)
class MaterialPollingState:
run_id: str
integration_state: MaterialIntegrationState
class MaterialStateStore(Protocol):
def load(
self,
machine_id: str,
production_order: str,
) -> MaterialPollingState | None: ...
def save(
self,
machine_id: str,
production_order: str,
state: MaterialPollingState,
) -> None: ...
@dataclass(frozen=True, slots=True)
class MaterialPollResult:
run: EnlyzeProductionRun
state: MaterialIntegrationState
class MaterialPollingService:
def __init__(
self,
*,
gateway: EnlyzeApiGateway,
state_store: MaterialStateStore,
machine_id: str,
rate_variable_id: str,
gate_variable_id: str,
gate_threshold: float,
max_sample_gap_seconds: float,
) -> None:
self._gateway = gateway
self._state_store = state_store
self._machine_id = machine_id
self._rate_variable_id = rate_variable_id
self._gate_variable_id = gate_variable_id
self._config = MaterialIntegratorConfig(
gate_threshold=gate_threshold,
max_sample_gap_seconds=max_sample_gap_seconds,
)
def poll_once(self, *, now: datetime) -> MaterialPollResult | None:
if now.tzinfo is None or now.utcoffset() is None:
raise ValueError("now must be timezone-aware")
run = self._gateway.get_open_production_run(self._machine_id)
if run is None or now < run.start:
return None
saved_polling_state = self._state_store.load(
self._machine_id,
run.production_order,
)
if saved_polling_state is None:
initial_state = MaterialIntegrationState()
historical_runs = sorted(
(
previous for previous in self._gateway.get_production_runs(self._machine_id)
if previous.production_order == run.production_order
and previous.uuid != run.uuid
and previous.start <= run.start
),
key=lambda previous: previous.start,
)
for previous in historical_runs:
if previous.end is None or previous.end > run.start:
raise ValueError("Historical production run overlaps the current open run")
integrated = self._integrate(
initial_state, start=previous.start, end=previous.end,
)
initial_state = MaterialIntegrationState(
cumulative_consumption_kg=integrated.cumulative_consumption_kg,
integrated_running_seconds=integrated.integrated_running_seconds,
)
start = run.start
elif saved_polling_state.run_id == run.uuid:
initial_state = saved_polling_state.integration_state
start = (
initial_state.last_processed_timestamp
if initial_state.last_processed_timestamp is not None
else run.start
)
else:
previous = saved_polling_state.integration_state
initial_state = MaterialIntegrationState(
cumulative_consumption_kg=previous.cumulative_consumption_kg,
integrated_running_seconds=previous.integrated_running_seconds,
)
start = run.start
if now < start:
return None
state = self._integrate(initial_state, start=start, end=now)
self._state_store.save(
self._machine_id,
run.production_order,
MaterialPollingState(
run_id=run.uuid,
integration_state=state,
),
)
return MaterialPollResult(
run=run,
state=state,
)
def _integrate(
self, initial_state: MaterialIntegrationState, *, start: datetime, end: datetime,
) -> MaterialIntegrationState:
integrator = MaterialConsumptionIntegrator(self._config, initial_state=initial_state)
samples = self._gateway.get_material_samples(
machine_id=self._machine_id,
rate_variable_id=self._rate_variable_id,
gate_variable_id=self._gate_variable_id,
start=start,
end=end,
)
return integrator.process_many(samples)
@@ -0,0 +1,88 @@
"""Sequential foreground polling with injectable time and output."""
import sys
import time
from collections.abc import Callable
from datetime import UTC, datetime
from typing import TextIO
from production_analytics.calculations.config import finite_number
from production_analytics.enlyze.exploration import ConfigurationError, ExplorationError
from production_analytics.service.material_polling import MaterialPollingService
from production_analytics.service.postgres_material import MaterialSnapshotWriter
class MaterialStateError(RuntimeError):
"""State operation failed; message contains only the operation and error class."""
def utc_now() -> datetime:
return datetime.now(UTC)
class MaterialPollingRunner:
def __init__(
self, service: MaterialPollingService, *, machine_id: str,
calculation_id: str, snapshot_writer: MaterialSnapshotWriter,
poll_interval_seconds: float, clock: Callable[[], datetime] = utc_now,
sleep: Callable[[float], None] = time.sleep,
stdout: TextIO | None = None, stderr: TextIO | None = None,
) -> None:
self.service = service
self.machine_id = machine_id
self.calculation_id = calculation_id
self.snapshot_writer = snapshot_writer
self.interval = finite_number(poll_interval_seconds, "poll interval", positive=True)
self.clock = clock
self.sleep = sleep
self.stdout = stdout if stdout is not None else sys.stdout
self.stderr = stderr if stderr is not None else sys.stderr
def run(self) -> None:
try:
while True:
now = self.clock()
if now.tzinfo is None or now.utcoffset() is None:
raise ConfigurationError("Runner clock must return timezone-aware timestamps")
now = now.astimezone(UTC)
prefix = f"{now.isoformat()} machine={self.machine_id!r}"
try:
result = self.service.poll_once(now=now)
except ConfigurationError:
raise
except Exception as exc:
if isinstance(exc, MaterialStateError):
detail = str(exc)
elif isinstance(exc, ExplorationError):
detail = f"ENLYZE/gateway: {type(exc).__name__}"
else:
detail = f"gateway/poll cycle: {type(exc).__name__}"
# Do not echo arbitrary exception messages or response bodies.
print(f"{prefix} error {detail}", file=self.stderr, flush=True)
else:
if result is None:
detail = "no open Production Run / no eligible polling window"
else:
try:
self.snapshot_writer.write(
timestamp=now, calculation_id=self.calculation_id,
machine_id=self.machine_id,
production_order=result.run.production_order,
run_id=result.run.uuid,
consumption_kg=result.state.cumulative_consumption_kg,
)
except Exception:
print(
f"{prefix} error PostgreSQL snapshot write failed",
file=self.stderr, flush=True,
)
detail = (
f"production_order={result.run.production_order!r} "
f"consumption_kg={result.state.cumulative_consumption_kg:.9f} "
f"running_seconds={result.state.integrated_running_seconds:.3f}"
)
print(f"{prefix} {detail}", file=self.stdout, flush=True)
# Fixed delay after completion, including errors; never catch up or overlap.
self.sleep(self.interval)
except KeyboardInterrupt:
print("Material polling stopped", file=self.stdout, flush=True)
@@ -0,0 +1,85 @@
"""Compose the configured live material polling application."""
import math
import os
import tempfile
from pathlib import Path
from urllib.parse import urlsplit
from production_analytics.calculations.config import MaterialCalculationConfig, finite_number
from production_analytics.enlyze.exploration import (
ConfigurationError,
ExplorationClient,
ExplorationSettings,
load_secret_file,
)
from production_analytics.enlyze.gateway import EnlyzeApiGateway
from production_analytics.service.material_polling import (
MaterialPollingService,
MaterialPollingState,
)
from production_analytics.service.material_runner import MaterialPollingRunner, MaterialStateError
from production_analytics.service.material_state_store import JsonMaterialStateStore
from production_analytics.service.postgres_material import (
PostgresMaterialSnapshotWriter,
PostgresSettings,
)
class _ReportingStateStore:
def __init__(self, store: JsonMaterialStateStore) -> None:
self.store = store
def load(self, machine_id: str, production_order: str) -> MaterialPollingState | None:
try:
return self.store.load(machine_id, production_order)
except Exception as exc:
raise MaterialStateError(f"state loading: {type(exc).__name__}") from exc
def save(
self, machine_id: str, production_order: str, state: MaterialPollingState,
) -> None:
try:
self.store.save(machine_id, production_order, state)
except Exception as exc:
raise MaterialStateError(f"state saving: {type(exc).__name__}") from exc
def build_material_runner(
calculation: MaterialCalculationConfig, *, poll_interval_seconds: float,
state_directory: Path, secrets_file: Path,
) -> MaterialPollingRunner:
finite_number(poll_interval_seconds, "poll interval", positive=True)
environment = dict(os.environ)
environment.update(load_secret_file(secrets_file))
settings = ExplorationSettings.from_environment(environment)
parsed = urlsplit(settings.base_url)
if parsed.scheme not in {"http", "https"} or not parsed.hostname or parsed.username:
raise ConfigurationError(
"ENLYZE_BASE_URL must be an HTTP(S) server URL without credentials"
)
if not math.isfinite(settings.timeout_seconds):
raise ConfigurationError("ENLYZE_HTTP_TIMEOUT_SECONDS must be finite")
try:
state_directory.mkdir(parents=True, exist_ok=True)
with tempfile.TemporaryFile(dir=state_directory):
pass
except OSError as exc:
raise ConfigurationError(
f"State directory startup check failed: {type(exc).__name__}"
) from exc
postgres_settings = PostgresSettings.from_environment(environment)
service = MaterialPollingService(
gateway=EnlyzeApiGateway(ExplorationClient(settings)),
state_store=_ReportingStateStore(JsonMaterialStateStore(state_directory)),
machine_id=calculation.machine_ref,
rate_variable_id=calculation.rate_signal_ref,
gate_variable_id=calculation.gate_signal_ref,
gate_threshold=calculation.gate_threshold,
max_sample_gap_seconds=calculation.max_sample_gap_seconds,
)
return MaterialPollingRunner(
service, calculation_id=calculation.id,
snapshot_writer=PostgresMaterialSnapshotWriter(postgres_settings),
machine_id=calculation.machine_ref, poll_interval_seconds=poll_interval_seconds,
)
@@ -0,0 +1,80 @@
"""JSON-backed persistence for material polling state."""
import hashlib
import json
import os
import tempfile
from pathlib import Path
from production_analytics.calculations.material_state import (
material_state_from_dict,
material_state_to_dict,
)
from production_analytics.service.material_polling import MaterialPollingState
class JsonMaterialStateStore:
def __init__(self, directory: str | Path) -> None:
self._directory = Path(directory)
def load(
self,
machine_id: str,
production_order: str,
) -> MaterialPollingState | None:
path = self._path(machine_id, production_order)
if not path.exists():
return None
try:
with path.open(encoding="utf-8") as f:
data = json.load(f)
if not isinstance(data, dict) or not isinstance(data.get("run_id"), str):
raise ValueError("run_id must be a string")
if not data["run_id"]:
raise ValueError("run_id must not be empty")
return MaterialPollingState(
run_id=data["run_id"],
integration_state=material_state_from_dict(data["integration_state"]),
)
except (ValueError, KeyError, TypeError) as exc:
raise ValueError(f"Invalid material polling state in {path}: {exc}") from exc
def save(
self,
machine_id: str,
production_order: str,
state: MaterialPollingState,
) -> None:
self._directory.mkdir(parents=True, exist_ok=True)
path = self._path(machine_id, production_order)
payload = {
"run_id": state.run_id,
"integration_state": material_state_to_dict(state.integration_state),
}
# Validate before touching the primary file, including non-finite values.
material_state_from_dict(payload["integration_state"])
if not isinstance(state.run_id, str) or not state.run_id:
raise ValueError("run_id must be a non-empty string")
temporary_path = None
try:
with tempfile.NamedTemporaryFile(
mode="w", encoding="utf-8", dir=self._directory,
prefix=".material-state-", suffix=".tmp", delete=False,
) as f:
temporary_path = Path(f.name)
json.dump(payload, f, indent=2, allow_nan=False)
f.flush()
os.fsync(f.fileno())
os.replace(temporary_path, path)
finally:
if temporary_path is not None:
temporary_path.unlink(missing_ok=True)
def _path(self, machine_id: str, production_order: str) -> Path:
# Hash the structured pair: sanitizing alone aliases distinct identifiers.
identity = json.dumps([machine_id, production_order], ensure_ascii=True)
digest = hashlib.sha256(identity.encode("utf-8")).hexdigest()
return self._directory / f"{digest}.json"
@@ -0,0 +1,68 @@
"""PostgreSQL settings and derived material snapshot inserts."""
from collections.abc import Mapping
from dataclasses import dataclass, field
from datetime import datetime
from typing import Protocol
from production_analytics.enlyze.exploration import ConfigurationError
@dataclass(frozen=True, slots=True)
class PostgresSettings:
host: str
port: int
dbname: str
user: str
password: str = field(repr=False)
@classmethod
def from_environment(cls, environment: Mapping[str, str]) -> "PostgresSettings":
values = {}
for key in ('HOST', 'PORT', 'DB', 'USER', 'PASSWORD'):
name = f'POSTGRES_{key}'
value = environment.get(name)
if not isinstance(value, str) or not value.strip() or '\x00' in value:
raise ConfigurationError(f'{name} must be a non-empty value without NUL bytes')
values[key] = value
try:
port = int(values['PORT'])
except ValueError:
raise ConfigurationError('POSTGRES_PORT must be an integer from 1 to 65535') from None
if not 1 <= port <= 65535:
raise ConfigurationError('POSTGRES_PORT must be an integer from 1 to 65535')
return cls(values['HOST'], port, values['DB'], values['USER'], values['PASSWORD'])
class MaterialSnapshotWriter(Protocol):
def write(
self, *, timestamp: datetime, calculation_id: str, machine_id: str,
production_order: str, run_id: str, consumption_kg: float,
) -> None: ...
class PostgresMaterialSnapshotWriter:
def __init__(self, settings: PostgresSettings) -> None:
self.settings = settings
def write(
self, *, timestamp: datetime, calculation_id: str, machine_id: str,
production_order: str, run_id: str, consumption_kg: float,
) -> None:
import psycopg
# One short transaction per cycle; context exit commits or rolls back and closes.
# A fresh connection lets the next cycle recover after a database outage.
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:
connection.execute(
"""INSERT INTO material_consumption_snapshots
(timestamp, calculation_id, machine_id, production_order, run_id, consumption_kg)
VALUES (%s, %s, %s, %s, %s, %s)
ON CONFLICT (calculation_id, machine_id, production_order, timestamp)
DO NOTHING""",
(timestamp, calculation_id, machine_id, production_order, run_id, consumption_kg),
)
+220
View File
@@ -1,6 +1,8 @@
from datetime import UTC, datetime
from unittest.mock import Mock
import pytest
from production_analytics.enlyze.gateway import EnlyzeApiGateway
@@ -131,3 +133,221 @@ def test_get_material_samples_uses_column_names_not_fixed_positions() -> None:
assert samples[0].material_rate_kg_per_hour == 1028.5
assert samples[0].gate_value == 2.1
@pytest.fixture
def timeseries():
client = Mock()
client.post_json.return_value.body = {"data": {
"columns": ["time", "rate", "gate"], "records": [],
}}
args = dict(machine_id="m", rate_variable_id="rate", gate_variable_id="gate",
start=datetime(2026, 9, 3, tzinfo=UTC), end=datetime(2026, 9, 4, tzinfo=UTC))
return EnlyzeApiGateway(client), client, args
@pytest.mark.parametrize("name", ["start", "end"])
def test_naive_window_rejected(timeseries, name) -> None:
gateway, client, args = timeseries
args[name] = args[name].replace(tzinfo=None)
with pytest.raises(ValueError, match="timezone-aware"):
gateway.get_material_samples(**args)
client.post_json.assert_not_called()
def test_reversed_window_rejected(timeseries) -> None:
gateway, client, args = timeseries
args["start"], args["end"] = args["end"], args["start"]
with pytest.raises(ValueError, match="precede"):
gateway.get_material_samples(**args)
client.post_json.assert_not_called()
def test_empty_timeseries(timeseries) -> None:
gateway, client, args = timeseries
assert gateway.get_material_samples(**args) == []
@pytest.mark.parametrize("column", ["time", "rate", "gate"])
def test_missing_required_column(timeseries, column) -> None:
gateway, client, args = timeseries
client.post_json.return_value.body["data"]["columns"].remove(column)
with pytest.raises(ValueError, match="required column"):
gateway.get_material_samples(**args)
@pytest.mark.parametrize("record", [[], {}, ["bad", 1, 1], [None, 1, 1],
["2026-09-03T00:00:00", 1, 1], ["2026-09-03T00:00:00Z", None, 1],
["2026-09-03T00:00:00Z", 1, float("inf")], ["2026-09-03T00:00:00Z", True, 1]])
def test_malformed_records_fail_clearly(timeseries, record) -> None:
gateway, client, args = timeseries
client.post_json.return_value.body["data"]["records"] = [record]
with pytest.raises(ValueError, match="malformed record 0"):
gateway.get_material_samples(**args)
def test_multiple_open_runs_rejected() -> None:
client = Mock()
run = dict(uuid="r", machine="m", production_order="o", start="2026-09-03T00:00:00Z",
end=None)
client.get.return_value.body = {"data": [run, {**run, "uuid": "r2"}]}
with pytest.raises(ValueError, match="multiple open Production Runs"):
EnlyzeApiGateway(client).get_open_production_run("m")
@pytest.mark.parametrize("item", [{}, None, {"end": None},
dict(uuid="r", machine="m", production_order="o", start="bad", end=None)])
def test_malformed_run_fails_clearly(item) -> None:
client = Mock()
client.get.return_value.body = {"data": [item]}
with pytest.raises(ValueError, match="Invalid production-run response"):
EnlyzeApiGateway(client).get_open_production_run("m")
def test_timeseries_with_null_cursor(timeseries) -> None:
gateway, client, args = timeseries
client.post_json.return_value.body["metadata"] = {"next_cursor": None}
client.post_json.return_value.body["data"]["records"] = [
["2026-09-03T00:00:00Z", 10, 2],
]
samples = gateway.get_material_samples(**args)
assert len(samples) == 1
assert samples[0].material_rate_kg_per_hour == 10
client.post_json.assert_called_once()
def test_timeseries_pagination_resolves_columns_per_page(timeseries) -> None:
gateway, client, args = timeseries
client.post_json.side_effect = [
Mock(body={
"data": {"columns": ["time", "rate", "gate"],
"records": [["2026-09-03T00:00:00Z", 10, 2]]},
"metadata": {"next_cursor": "continuation"},
}),
Mock(body={
"data": {"columns": ["gate", "time", "rate"],
"records": [[3, "2026-09-03T01:00:00Z", 20]]},
"metadata": {"next_cursor": None},
}),
]
samples = gateway.get_material_samples(**args)
assert [(s.timestamp, s.material_rate_kg_per_hour, s.gate_value) for s in samples] == [
(datetime(2026, 9, 3, tzinfo=UTC), 10, 2),
(datetime(2026, 9, 3, 1, tzinfo=UTC), 20, 3),
]
assert client.post_json.call_count == 2
first, second = client.post_json.call_args_list
assert first.args == ("/v2/timeseries", {
"machine": "m", "start": args["start"].isoformat(), "end": args["end"].isoformat(),
"variables": [{"uuid": "rate"}, {"uuid": "gate"}],
})
assert second.args == ("/v2/timeseries", {**first.args[1], "cursor": "continuation"})
@pytest.mark.parametrize("metadata", [None, [], "bad", {},
{"next_cursor": ""}, {"next_cursor": 1}, {"next_cursor": False},
{"next_cursor": []}, {"next_cursor": {}}])
def test_invalid_pagination_metadata(timeseries, metadata) -> None:
gateway, client, args = timeseries
client.post_json.return_value.body["metadata"] = metadata
with pytest.raises(ValueError, match="Invalid timeseries response.*next_cursor"):
gateway.get_material_samples(**args)
client.post_json.assert_called_once()
@pytest.mark.parametrize("cursors", [["a", "a"], ["a", "b", "a"]])
def test_repeated_pagination_cursor(timeseries, cursors) -> None:
gateway, client, args = timeseries
data = client.post_json.return_value.body["data"]
client.post_json.side_effect = [
Mock(body={"data": data, "metadata": {"next_cursor": cursor}}) for cursor in cursors
]
with pytest.raises(ValueError, match="next_cursor has already been followed"):
gateway.get_material_samples(**args)
assert client.post_json.call_count == len(cursors)
@pytest.mark.parametrize("data, error", [
({"columns": ["time", "rate", "gate"], "records": [[]]}, "malformed record 0"),
({"columns": ["time", "rate", "rate", "gate"], "records": []}, "required column"),
({"columns": ["time", "rate"], "records": []}, "required column"),
({"columns": None, "records": []}, "columns must be a list"),
({"columns": ["time", "rate", "gate"], "records": None}, "records must be a list"),
({"columns": ["time", "rate", "gate"],
"records": [["2026-09-03T01:00:00", 10, 2]]}, "timezone-aware"),
({"columns": ["time", "rate", "gate"],
"records": [["2026-09-03T01:00:00Z", float("nan"), 2]]}, "finite"),
({"columns": ["time", "rate", "gate"],
"records": [["2026-09-03T01:00:00Z", 10, True]]}, "numeric"),
])
def test_malformed_later_page(timeseries, data, error) -> None:
gateway, client, args = timeseries
first_body = client.post_json.return_value.body
first_body["metadata"] = {"next_cursor": "continuation"}
client.post_json.side_effect = [
Mock(body=first_body), Mock(body={"data": data, "metadata": {"next_cursor": None}}),
]
with pytest.raises(ValueError, match=f"Invalid timeseries response on page 2:.*{error}"):
gateway.get_material_samples(**args)
assert client.post_json.call_count == 2
def test_production_runs_follow_pages_and_find_open_run() -> None:
client = Mock()
closed = dict(uuid="closed", machine="m", production_order="opaque-order",
start="2026-09-01T00:00:00Z", end="2026-09-01T01:00:00Z")
opened = {**closed, "uuid": "open", "end": None}
pages = [
Mock(body={"data": [closed], "metadata": {"next_cursor": "a"}}),
Mock(body={"data": [opened], "metadata": {"next_cursor": None}}),
]
client.get.side_effect = pages
gateway = EnlyzeApiGateway(client)
runs = gateway.get_production_runs("m")
assert [run.uuid for run in runs] == ["closed", "open"]
assert runs[0].end == datetime(2026, 9, 1, 1, tzinfo=UTC)
assert [call.args for call in client.get.call_args_list] == [
("/v2/production-runs", {"machine": "m"}),
("/v2/production-runs", {"machine": "m", "cursor": "a"}),
]
client.get.side_effect = pages
assert gateway.get_open_production_run("m") == runs[1]
@pytest.mark.parametrize("metadata", [
None, [], "bad", {}, {"next_cursor": ""}, {"next_cursor": 1},
{"next_cursor": False}, {"next_cursor": []}, {"next_cursor": {}},
])
def test_production_run_invalid_pagination(metadata) -> None:
client = Mock()
client.get.return_value.body = {"data": [], "metadata": metadata}
with pytest.raises(ValueError, match="Invalid production-run response.*next_cursor"):
EnlyzeApiGateway(client).get_production_runs("m")
@pytest.mark.parametrize("cursors", [["a", "a"], ["a", "b", "a"]])
def test_production_run_cyclic_pagination(cursors) -> None:
client = Mock()
client.get.side_effect = [
Mock(body={"data": [], "metadata": {"next_cursor": cursor}})
for cursor in cursors
]
with pytest.raises(ValueError, match="next_cursor has already been followed"):
EnlyzeApiGateway(client).get_production_runs("m")
assert client.get.call_count == len(cursors)
@pytest.mark.parametrize("body", [
None, [], {}, {"data": {}}, {"data": [None]},
{"data": [dict(uuid="r", machine="wrong", production_order="o",
start="2026-09-03T00:00:00Z", end=None)]},
{"data": [dict(uuid="r", machine="m", production_order="o",
start="2026-09-03T00:00:00Z", end="2026-09-02T00:00:00Z")]},
])
def test_production_run_malformed_later_page(body) -> None:
client = Mock()
client.get.side_effect = [
Mock(body={"data": [], "metadata": {"next_cursor": "a"}}), Mock(body=body),
]
with pytest.raises(ValueError, match="Invalid production-run response on page 2"):
EnlyzeApiGateway(client).get_production_runs("m")
+187
View File
@@ -0,0 +1,187 @@
import re
import traceback
from dataclasses import FrozenInstanceError, replace
from datetime import UTC, datetime
from decimal import Decimal
from unittest.mock import MagicMock
import pytest
from production_analytics.enlyze.exploration import ConfigurationError
from production_analytics.erp import ErpReadError, ErpSettings, ErpWorkplaceStatusGateway
ENV = dict(ERP_DB_HOST='localhost', ERP_DB_PORT='49601', ERP_DB_NAME='NV_DWH',
ERP_DB_USER='private-user', ERP_DB_PASSWORD='private-password')
NOW = datetime(2026, 9, 4, 21, 20, 50)
ROW = ('K7 ', '12026000815 ', '212520 ',
'Stex R 1501 C (PR) 5,80 x 50 m ', NOW,
Decimal('80040.000'), Decimal('69281.000'), Decimal('10759.000'), 17, 38)
def test_valid_settings():
settings = ErpSettings.from_environment(ENV)
assert (settings.host, settings.port, settings.database) == ('localhost', 49601, 'NV_DWH')
assert settings.user == ENV['ERP_DB_USER']
assert settings.password == ENV['ERP_DB_PASSWORD']
assert settings.user not in repr(settings)
assert settings.password not in repr(settings)
@pytest.mark.parametrize('key', ENV)
@pytest.mark.parametrize('value', [None, '', ' ', '\x00'])
def test_required_settings(key, value):
environment = dict(ENV)
if value is None:
del environment[key]
else:
environment[key] = value
with pytest.raises(ConfigurationError, match=key):
ErpSettings.from_environment(environment)
@pytest.mark.parametrize('port', ['0', '-1', '65536', '1.5', 'private-password'])
def test_invalid_port(port):
with pytest.raises(ConfigurationError, match='ERP_DB_PORT') as error:
ErpSettings.from_environment(dict(ENV, ERP_DB_PORT=port))
assert 'private-password' not in ''.join(traceback.format_exception(error.value))
@pytest.mark.parametrize('port', [1, 65535])
def test_port_boundaries(port):
assert ErpSettings.from_environment(dict(ENV, ERP_DB_PORT=str(port))).port == port
@pytest.mark.parametrize('changes', [{'port': True}, {'port': 1.5}, {'host': ''},
{'database': ''}, {'user': ''}, {'password': ''}])
def test_direct_settings_validation(changes):
with pytest.raises(ConfigurationError):
replace(ErpSettings.from_environment(ENV), **changes)
def test_secret_file(tmp_path):
path = tmp_path / 'erp.env'
path.write_text("ERP_DB_PASSWORD='file password'\n")
assert ErpSettings.from_secret_file(path, environment=ENV).password == 'file password'
assert ErpSettings.from_secret_file(tmp_path / 'missing', environment=ENV).port == 49601
path.write_text("ERP_DB_PASSWORD='private-password\n")
with pytest.raises(ConfigurationError) as error:
ErpSettings.from_secret_file(path, environment=ENV)
assert 'private-password' not in ''.join(traceback.format_exception(error.value))
@pytest.fixture
def db(monkeypatch):
connect = MagicMock()
monkeypatch.setattr('production_analytics.erp.gateway.pymssql.connect', connect)
connection = connect.return_value.__enter__.return_value
cursor = connection.cursor.return_value.__enter__.return_value
cursor.fetchmany.return_value = [ROW]
return connect, connection, cursor
def read():
return ErpWorkplaceStatusGateway(ErpSettings.from_environment(ENV))
def test_valid_row_and_resource_lifecycle(db):
connect, connection, cursor = db
status = read().get_current_workplace_status('K7')
assert status.workplace == 'K7'
assert status.production_order == '12026000815'
assert status.article_number == '212520'
assert status.article_description == 'Stex R 1501 C (PR) 5,80 x 50 m'
assert status.feedback_timestamp is NOW
assert status.feedback_timestamp.tzinfo is None
assert (status.order_quantity_m2, status.good_quantity_m2, status.remaining_quantity_m2,
status.remaining_time_hours, status.remaining_rolls) == (
80040., 69281., 10759., 17., 38.,
)
assert type(status.order_quantity_m2) is float
with pytest.raises(FrozenInstanceError):
status.workplace = 'other'
cursor.fetchmany.assert_called_once_with(2)
connection.cursor.return_value.__exit__.assert_called_once()
connect.return_value.__exit__.assert_called_once()
connect.assert_called_once_with(server='localhost', port=49601, database='NV_DWH',
user='private-user', password='private-password',
login_timeout=10, timeout=10)
connection.commit.assert_not_called()
def test_leading_whitespace_and_timezone_preserved(db):
row = list(ROW)
row[:4] = [' K7 ', ' K 7-001 ', ' 001 ', ' description ']
row[4] = NOW.replace(tzinfo=UTC)
db[2].fetchmany.return_value = [row]
status = read().get_current_workplace_status(' K7')
assert (status.workplace, status.production_order, status.article_number,
status.article_description) == (' K7', ' K 7-001', ' 001', ' description')
assert status.feedback_timestamp is row[4]
def test_nullable_fields(db):
db[2].fetchmany.return_value = [('K7', '001', None, None, NOW, *([None] * 5))]
status = read().get_current_workplace_status('K7')
assert all(getattr(status, key) is None for key in (
'article_number', 'article_description', 'order_quantity_m2', 'good_quantity_m2',
'remaining_quantity_m2', 'remaining_time_hours', 'remaining_rolls',
))
@pytest.mark.parametrize('value', [Decimal('NaN'), Decimal('Infinity'), Decimal('1e999'),
float('nan'), True, 'private-password'])
def test_invalid_numbers(db, value):
row = list(ROW)
row[5] = value
db[2].fetchmany.return_value = [row]
with pytest.raises(ErpReadError, match='invalid row') as error:
read().get_current_workplace_status('K7')
assert 'private-password' not in ''.join(traceback.format_exception(error.value))
def test_fractional_decimal(db):
row = list(ROW)
row[5:] = [Decimal('12.125')] * 5
db[2].fetchmany.return_value = [row]
status = read().get_current_workplace_status('K7')
assert status.order_quantity_m2 == status.remaining_rolls == 12.125
def test_zero_rows(db):
db[2].fetchmany.return_value = []
assert read().get_current_workplace_status('K7') is None
def test_multiple_rows(db):
db[2].fetchmany.return_value = [ROW, ROW]
with pytest.raises(ErpReadError, match='multiple rows'):
read().get_current_workplace_status('K7')
def test_parameterized_read_only_query(db):
workplace = "K7'; DELETE FROM anything; --"
read().get_current_workplace_status(workplace)
db[2].execute.assert_called_once()
sql, params = db[2].execute.call_args.args
assert params == (workplace,)
assert workplace not in sql
assert 'WHERE [Arbeitsplatz] = %s' in sql
assert re.findall(r'FROM\s+(\S+)', sql) == ['[dbo].[GRAFANA_WORKPLACE_STATUS]']
assert not re.search(r'\b(INSERT|UPDATE|DELETE|MERGE|EXEC|TOP|ORDER BY)\b', sql, re.I)
assert sql.lstrip().startswith('SELECT')
@pytest.mark.parametrize('stage', ['connect', 'execute', 'fetch', 'cursor_close', 'close'])
def test_driver_errors_are_safe(db, stage):
connect, connection, cursor = db
target = {'connect': connect, 'execute': cursor.execute, 'fetch': cursor.fetchmany,
'cursor_close': connection.cursor.return_value.__exit__,
'close': connect.return_value.__exit__}[stage]
target.side_effect = RuntimeError('private-user private-password')
with pytest.raises(ErpReadError, match='read failed') as error:
read().get_current_workplace_status('K7')
rendered = ''.join(traceback.format_exception(error.value))
assert 'private-user' not in rendered
assert 'private-password' not in rendered
if stage != 'connect':
connect.return_value.__exit__.assert_called_once()
+409
View File
@@ -0,0 +1,409 @@
from datetime import UTC, datetime
from unittest.mock import Mock
import pytest
from production_analytics.calculations.material_consumption import (
MaterialIntegrationState,
MaterialSample,
)
from production_analytics.enlyze.gateway import EnlyzeProductionRun
from production_analytics.service.material_polling import (
MaterialPollingService,
MaterialPollingState,
)
def test_poll_once_starts_at_run_start_without_saved_state() -> None:
gateway = Mock()
gateway.get_production_runs.return_value = []
state_store = Mock()
state_store.load.return_value = None
run = EnlyzeProductionRun(
uuid="run-1",
machine_id="machine-1",
product_id="product-1",
production_order="ORDER-1",
start=datetime(2026, 9, 3, 5, 7, 47, tzinfo=UTC),
end=None,
)
gateway.get_open_production_run.return_value = run
gateway.get_production_runs.return_value = [run]
samples = [
MaterialSample(
timestamp=datetime(2026, 9, 3, 5, 7, 50, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
MaterialSample(
timestamp=datetime(2026, 9, 3, 5, 8, 0, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
]
gateway.get_material_samples.return_value = samples
service = MaterialPollingService(
gateway=gateway,
state_store=state_store,
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
gate_threshold=0.5,
max_sample_gap_seconds=20.0,
)
result = service.poll_once(
now=datetime(2026, 9, 3, 5, 8, 10, tzinfo=UTC),
)
assert result is not None
assert result.run == run
assert result.state.cumulative_consumption_kg == 10.0
assert result.state.integrated_running_seconds == 10.0
gateway.get_material_samples.assert_called_once_with(
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
start=run.start,
end=datetime(2026, 9, 3, 5, 8, 10, tzinfo=UTC),
)
state_store.save.assert_called_once_with(
"machine-1",
run.production_order,
MaterialPollingState(
run_id=run.uuid,
integration_state=result.state,
),
)
def test_poll_once_resumes_from_saved_timestamp_without_double_counting() -> None:
gateway = Mock()
gateway.get_production_runs.return_value = []
state_store = Mock()
run = EnlyzeProductionRun(
uuid="run-1",
machine_id="machine-1",
product_id="product-1",
production_order="ORDER-1",
start=datetime(2026, 9, 3, 5, 7, 47, tzinfo=UTC),
end=None,
)
gateway.get_open_production_run.return_value = run
saved_state = MaterialIntegrationState(
cumulative_consumption_kg=10.0,
integrated_running_seconds=10.0,
last_processed_timestamp=datetime(2026, 9, 3, 5, 8, 0, tzinfo=UTC),
last_material_rate_kg_per_hour=3600.0,
last_gate_value=1.0,
integration_active=True,
)
state_store.load.return_value = MaterialPollingState(
run_id=run.uuid,
integration_state=saved_state,
)
gateway.get_material_samples.return_value = [
MaterialSample(
timestamp=datetime(2026, 9, 3, 5, 8, 0, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
MaterialSample(
timestamp=datetime(2026, 9, 3, 5, 8, 10, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
]
service = MaterialPollingService(
gateway=gateway,
state_store=state_store,
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
gate_threshold=0.5,
max_sample_gap_seconds=20.0,
)
result = service.poll_once(
now=datetime(2026, 9, 3, 5, 8, 20, tzinfo=UTC),
)
assert result is not None
gateway.get_production_runs.assert_not_called()
assert result.state.cumulative_consumption_kg == 20.0
assert result.state.integrated_running_seconds == 20.0
gateway.get_material_samples.assert_called_once_with(
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
start=saved_state.last_processed_timestamp,
end=datetime(2026, 9, 3, 5, 8, 20, tzinfo=UTC),
)
def test_poll_once_without_open_run_does_nothing() -> None:
gateway = Mock()
gateway.get_production_runs.return_value = []
state_store = Mock()
gateway.get_open_production_run.return_value = None
service = MaterialPollingService(
gateway=gateway,
state_store=state_store,
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
gate_threshold=0.5,
max_sample_gap_seconds=20.0,
)
result = service.poll_once(
now=datetime(2026, 9, 3, 5, 8, 20, tzinfo=UTC),
)
assert result is None
gateway.get_material_samples.assert_not_called()
state_store.load.assert_not_called()
state_store.save.assert_not_called()
def test_new_run_of_same_order_keeps_total_but_restarts_integration_at_run_start() -> None:
gateway = Mock()
gateway.get_production_runs.return_value = []
state_store = Mock()
run = EnlyzeProductionRun(
uuid="run-2",
machine_id="machine-1",
product_id="product-1",
production_order="ORDER-1",
start=datetime(2026, 9, 3, 8, 0, 0, tzinfo=UTC),
end=None,
)
gateway.get_open_production_run.return_value = run
previous_integration_state = MaterialIntegrationState(
cumulative_consumption_kg=100.0,
integrated_running_seconds=360.0,
last_processed_timestamp=datetime(2026, 9, 3, 6, 0, 0, tzinfo=UTC),
last_material_rate_kg_per_hour=3600.0,
last_gate_value=1.0,
integration_active=True,
)
state_store.load.return_value = MaterialPollingState(
run_id="run-1",
integration_state=previous_integration_state,
)
gateway.get_material_samples.return_value = [
MaterialSample(
timestamp=datetime(2026, 9, 3, 8, 0, 10, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
MaterialSample(
timestamp=datetime(2026, 9, 3, 8, 0, 20, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
]
service = MaterialPollingService(
gateway=gateway,
state_store=state_store,
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
gate_threshold=0.5,
max_sample_gap_seconds=20.0,
)
result = service.poll_once(
now=datetime(2026, 9, 3, 8, 0, 30, tzinfo=UTC),
)
assert result is not None
gateway.get_production_runs.assert_not_called()
assert result.state.cumulative_consumption_kg == 110.0
assert result.state.integrated_running_seconds == 370.0
state_store.load.assert_called_once_with(
"machine-1",
"ORDER-1",
)
gateway.get_material_samples.assert_called_once_with(
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
start=run.start,
end=datetime(2026, 9, 3, 8, 0, 30, tzinfo=UTC),
)
@pytest.fixture
def polling(tmp_path):
from production_analytics.service.material_state_store import JsonMaterialStateStore
gateway = Mock()
gateway.get_production_runs.return_value = []
gateway.get_open_production_run.return_value = EnlyzeProductionRun(
"run", "machine", None, "order", datetime(2026, 9, 3, tzinfo=UTC), None,
)
gateway.get_material_samples.return_value = []
store = JsonMaterialStateStore(tmp_path)
service = MaterialPollingService(
gateway=gateway, state_store=store, machine_id="machine",
rate_variable_id="rate", gate_variable_id="gate", gate_threshold=0.5,
max_sample_gap_seconds=20,
)
return service, gateway, store
def test_empty_response_persists_state_without_inferred_consumption(polling) -> None:
service, gateway, store = polling
now = datetime(2026, 9, 3, 1, tzinfo=UTC)
assert service.poll_once(now=now).state == MaterialIntegrationState()
saved = MaterialPollingState("run", MaterialIntegrationState(
cumulative_consumption_kg=10, integrated_running_seconds=10,
last_processed_timestamp=now, last_material_rate_kg_per_hour=3600,
last_gate_value=1, integration_active=True,
))
store.save("machine", "order", saved)
assert service.poll_once(now=now).state == saved.integration_state
assert store.load("machine", "order") == saved
def test_different_order_does_not_reuse_state(polling) -> None:
service, gateway, store = polling
store.save("machine", "other-order", MaterialPollingState(
"old", MaterialIntegrationState(cumulative_consumption_kg=100),
))
result = service.poll_once(now=datetime(2026, 9, 3, 1, tzinfo=UTC))
assert result.state == MaterialIntegrationState()
assert gateway.get_material_samples.call_args.kwargs["start"] == result.run.start
assert store.load("machine", "other-order").integration_state.cumulative_consumption_kg == 100
def test_now_before_run_start_does_not_query_or_save(polling) -> None:
service, gateway, store = polling
assert service.poll_once(now=datetime(2026, 9, 2, tzinfo=UTC)) is None
gateway.get_production_runs.assert_not_called()
gateway.get_material_samples.assert_not_called()
assert store.load("machine", "order") is None
def test_naive_now_is_rejected_before_io(polling) -> None:
service, gateway, store = polling
with pytest.raises(ValueError, match="now must be timezone-aware"):
service.poll_once(now=datetime(2026, 9, 3))
gateway.get_open_production_run.assert_not_called()
def test_new_run_empty_response_resets_baseline_and_retains_totals(polling) -> None:
service, gateway, store = polling
store.save("machine", "order", MaterialPollingState("previous", MaterialIntegrationState(
cumulative_consumption_kg=100, integrated_running_seconds=50,
last_processed_timestamp=datetime(2026, 9, 2, 23, 59, 59, tzinfo=UTC),
last_material_rate_kg_per_hour=3600, last_gate_value=1, integration_active=True,
)))
result = service.poll_once(now=datetime(2026, 9, 3, 1, tzinfo=UTC))
assert result.state == MaterialIntegrationState(
cumulative_consumption_kg=100, integrated_running_seconds=50,
)
assert store.load("machine", "order") == MaterialPollingState("run", result.state)
def test_now_before_saved_timestamp_leaves_state_unchanged(polling) -> None:
service, gateway, store = polling
saved = MaterialPollingState("run", MaterialIntegrationState(
last_processed_timestamp=datetime(2026, 9, 3, 2, tzinfo=UTC),
))
store.save("machine", "order", saved)
assert service.poll_once(now=datetime(2026, 9, 3, 1, tzinfo=UTC)) is None
gateway.get_material_samples.assert_not_called()
assert store.load("machine", "order") == saved
@pytest.mark.parametrize("order", [
"K 7-12026000815", "K 7-12025000074-K 7-12025000075",
])
@pytest.mark.parametrize("gap_seconds", [5, 86400])
def test_bootstrap_replays_exact_order_segments_and_then_resumes(
polling, order, gap_seconds,
) -> None:
from dataclasses import replace
from datetime import timedelta
service, gateway, store = polling
base = datetime(2026, 9, 3, tzinfo=UTC)
runs = [
EnlyzeProductionRun(
f"run-{i}", "machine", None, order,
base + timedelta(seconds=i * (10 + gap_seconds)),
base + timedelta(seconds=i * (10 + gap_seconds) + 10) if i < 4 else None,
)
for i in range(5)
]
current = runs[-1]
now = current.start + timedelta(seconds=10)
gateway.get_open_production_run.return_value = current
gateway.get_production_runs.return_value = [
current, runs[2], runs[0], runs[3], runs[1],
replace(runs[0], uuid="unrelated", production_order=order + "-suffix"),
replace(runs[0], uuid="component", production_order="K 7-12025000074"),
replace(runs[0], uuid="whitespace", production_order=order + " "),
replace(current, uuid="future", start=now + timedelta(days=1)),
]
def samples(**kwargs):
return [
MaterialSample(kwargs["start"], 3600, 1),
MaterialSample(kwargs["end"], 3600, 1),
]
gateway.get_material_samples.side_effect = samples
result = service.poll_once(now=now)
assert result.state.cumulative_consumption_kg == 50
assert result.state.integrated_running_seconds == 50
assert [
(call.kwargs["start"], call.kwargs["end"])
for call in gateway.get_material_samples.call_args_list
] == [(run.start, run.end or now) for run in runs]
assert store.load("machine", order) == MaterialPollingState(current.uuid, result.state)
gateway.get_production_runs.reset_mock()
gateway.get_material_samples.reset_mock()
result = service.poll_once(now=now + timedelta(seconds=10))
gateway.get_production_runs.assert_not_called()
assert gateway.get_material_samples.call_args.kwargs["start"] == now
assert result.state.cumulative_consumption_kg == 60
assert result.state.integrated_running_seconds == 60
def test_bootstrap_failure_does_not_persist_partial_totals(polling) -> None:
from dataclasses import replace
from datetime import timedelta
service, gateway, store = polling
current = gateway.get_open_production_run.return_value
gateway.get_production_runs.return_value = [
replace(current, uuid="old", start=current.start - timedelta(hours=1),
end=current.start - timedelta(minutes=30)),
current,
]
gateway.get_material_samples.side_effect = [[], ValueError("failed page")]
with pytest.raises(ValueError, match="failed page"):
service.poll_once(now=current.start + timedelta(seconds=10))
assert store.load("machine", "order") is None
+291
View File
@@ -0,0 +1,291 @@
import io
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 MaterialIntegrationState
from production_analytics.cli.__main__ import main
from production_analytics.enlyze.exploration import ConfigurationError
from production_analytics.enlyze.gateway import EnlyzeProductionRun
from production_analytics.service.material_polling import MaterialPollResult
from production_analytics.service.material_runner import MaterialPollingRunner, MaterialStateError
from production_analytics.service.material_runtime import (
_ReportingStateStore,
build_material_runner,
)
EXAMPLE = Path(__file__).resolve().parents[1] / 'config/k7-material-consumption.yaml'
NOW = datetime(2026, 1, 1, tzinfo=UTC)
def write_config(tmp_path, entries):
path = tmp_path / 'calculations.yaml'
path.write_text(yaml.safe_dump({'calculations': entries}))
return path
def entry():
return yaml.safe_load(EXAMPLE.read_text())['calculations'][0]
def test_valid_k7_config():
config = load_material_calculation(EXAMPLE)
assert config.id == 'k7-fiber-consumption'
assert config.machine_ref == 'c220f95c-a65e-4cb7-99b7-0626d6c7508c'
assert config.rate_signal_ref == 'c9d06af5-f6d6-4ede-b6c4-5a98bac77129'
assert config.gate_signal_ref == '823867bb-f5d2-40eb-b875-657155addfd0'
assert config.gate_threshold == 0.5
assert config.max_sample_gap_seconds == 20.0
@pytest.mark.parametrize('field', list(entry()))
def test_missing_fields(tmp_path, field):
data = entry()
del data[field]
with pytest.raises(CalculationConfigError, match=f'missing fields: {field}'):
load_material_calculation(write_config(tmp_path, [data]))
@pytest.mark.parametrize(('field', 'value'), [
('type', 'integration'), ('version', 1), ('version', '2'),
('machine_ref', ''), ('rate_signal_ref', ' '), ('gate_signal_ref', None), ('id', ''),
('gate_threshold', True), ('gate_threshold', '0.5'), ('gate_threshold', float('nan')),
('max_sample_gap_seconds', float('inf')), ('max_sample_gap_seconds', 0),
('max_sample_gap_seconds', -1), ('output_unit', 'tonnes'), ('group_by', 'run'),
('gate_treshold', 1),
])
def test_invalid_fields(tmp_path, field, value):
data = entry()
data[field] = value
with pytest.raises(CalculationConfigError):
load_material_calculation(write_config(tmp_path, [data]))
@pytest.mark.parametrize('content', [
'calculations: [', 'calculations: []', 'calculations: {}', 'other: []',
'calculations: []\ncalculations: []', 'calculations:\n - id: x\n id: y',
])
def test_malformed_documents(tmp_path, content):
path = tmp_path / 'bad.yaml'
path.write_text(content)
with pytest.raises(CalculationConfigError):
load_material_calculation(path)
def test_selection(tmp_path):
first, second = entry(), dict(entry(), id='second')
path = write_config(tmp_path, [first, second])
with pytest.raises(CalculationConfigError, match='specify --calculation-id'):
load_material_calculation(path)
assert load_material_calculation(path, 'second').id == 'second'
with pytest.raises(CalculationConfigError, match='not found'):
load_material_calculation(path, 'absent')
with pytest.raises(CalculationConfigError, match='unique'):
load_material_calculation(write_config(tmp_path, [first, first]), first['id'])
def test_repeated_sequential_polls_and_long_cycle():
events = []
elapsed = 0
def poll(*, now):
nonlocal elapsed
assert now == NOW + timedelta(seconds=elapsed)
events.append('poll-start')
elapsed += 25 # Longer than the interval; no real sleep or concurrent work.
events.append('poll-end')
def sleep(seconds):
nonlocal elapsed
events.append(('sleep', seconds))
elapsed += seconds
if elapsed >= 70:
raise KeyboardInterrupt
output = io.StringIO()
service = Mock(poll_once=Mock(side_effect=poll))
MaterialPollingRunner(
service, calculation_id='test-calculation', snapshot_writer=Mock(),
machine_id='machine', poll_interval_seconds=10,
clock=lambda: NOW + timedelta(seconds=elapsed), sleep=sleep, stdout=output,
).run()
assert events == ['poll-start', 'poll-end', ('sleep', 10)] * 2
assert service.poll_once.call_count == 2
assert output.getvalue().count('no open Production Run') == 2
assert 'stopped' in output.getvalue()
def test_error_then_success_preserves_opaque_order_and_reports_totals():
run = EnlyzeProductionRun('run', 'machine', None, ' 00842/ABC ', NOW, None)
result = MaterialPollResult(run, MaterialIntegrationState(12.5, 20))
service = Mock(poll_once=Mock(side_effect=[ValueError('SECRET'), result]))
output, errors = io.StringIO(), io.StringIO()
sleep = Mock(side_effect=[None, KeyboardInterrupt])
MaterialPollingRunner(
service, calculation_id='test-calculation', snapshot_writer=Mock(),
machine_id='machine', poll_interval_seconds=5, clock=lambda: NOW,
sleep=sleep, stdout=output, stderr=errors,
).run()
assert service.poll_once.call_count == 2
assert sleep.call_count == 2
assert 'gateway/poll cycle: ValueError' in errors.getvalue()
assert 'SECRET' not in errors.getvalue()
assert "production_order=' 00842/ABC '" in output.getvalue()
assert 'consumption_kg=12.500000000 running_seconds=20.000' in output.getvalue()
assert NOW.isoformat() in output.getvalue()
def test_interrupt_during_poll():
sleep, output = Mock(), io.StringIO()
MaterialPollingRunner(
Mock(poll_once=Mock(side_effect=KeyboardInterrupt)),
calculation_id='test-calculation', snapshot_writer=Mock(), machine_id='m',
poll_interval_seconds=1, sleep=sleep, stdout=output,
).run()
sleep.assert_not_called()
assert 'stopped' in output.getvalue()
def test_cycle_configuration_failure_is_fatal():
sleep = Mock()
runner = MaterialPollingRunner(
Mock(poll_once=Mock(side_effect=ConfigurationError('bad settings'))),
calculation_id='test-calculation', snapshot_writer=Mock(),
machine_id='m', poll_interval_seconds=1, sleep=sleep,
)
with pytest.raises(ConfigurationError):
runner.run()
sleep.assert_not_called()
@pytest.mark.parametrize('interval', [0, -1, float('nan'), float('inf'), True])
def test_invalid_interval(interval):
with pytest.raises(CalculationConfigError):
MaterialPollingRunner(
Mock(), machine_id='m', calculation_id='test', snapshot_writer=Mock(),
poll_interval_seconds=interval,
)
@pytest.mark.parametrize('operation', ['load', 'save'])
def test_state_error_classification(operation):
store = Mock()
getattr(store, operation).side_effect = OSError('SECRET')
wrapped = _ReportingStateStore(store)
arguments = ('machine', ' 00842 ') if operation == 'load' else ('machine', ' 00842 ', Mock())
phase = 'loading' if operation == 'load' else 'saving'
with pytest.raises(MaterialStateError, match=f'state {phase}: OSError') as error:
getattr(wrapped, operation)(*arguments)
assert 'SECRET' not in str(error.value)
def test_runtime_wiring(tmp_path, monkeypatch):
monkeypatch.setenv('ENLYZE_BASE_URL', 'https://example.invalid/api/')
secrets = tmp_path / 'secret.env'
for key, value in dict(HOST='localhost', PORT='5432', DB='analytics',
USER='user', PASSWORD='environment-password').items():
monkeypatch.setenv(f'POSTGRES_{key}', value)
secrets.write_text('ENLYZE_API_KEY=secret-from-file\nPOSTGRES_PASSWORD=file-password\n')
with patch('production_analytics.service.material_runtime.MaterialPollingService') as service:
runner = build_material_runner(
load_material_calculation(EXAMPLE), poll_interval_seconds=7,
state_directory=tmp_path / 'state', secrets_file=secrets,
)
kwargs = service.call_args.kwargs
assert kwargs['machine_id'] == entry()['machine_ref']
assert kwargs['rate_variable_id'] == entry()['rate_signal_ref']
assert kwargs['gate_variable_id'] == entry()['gate_signal_ref']
assert kwargs['gate_threshold'] == 0.5
assert kwargs['max_sample_gap_seconds'] == 20.0
assert kwargs['gateway']._client._settings.api_key == 'secret-from-file'
assert kwargs['state_store'].store._directory == tmp_path / 'state'
assert runner.service is service.return_value
assert runner.interval == 7
assert runner.calculation_id == 'k7-fiber-consumption'
assert runner.snapshot_writer.settings.password == 'file-password'
assert runner.snapshot_writer.settings.dbname == 'analytics'
def test_cli_bad_config(tmp_path, capsys):
assert main(['run', 'material-poll', '--config', str(tmp_path / 'missing')]) == 2
assert 'Config/startup' in capsys.readouterr().err
def test_cli_startup_and_success(tmp_path, capsys):
args = ['run', 'material-poll', '--config', str(EXAMPLE),
'--calculation-id', 'k7-fiber-consumption']
with patch('production_analytics.service.material_runtime.build_material_runner') as build:
assert main(args) == 0
build.return_value.run.assert_called_once_with()
assert build.call_args.kwargs['state_directory'] == Path('data/state/material')
build.side_effect = ConfigurationError('bad settings')
assert main(args) == 2
assert 'bad settings' in capsys.readouterr().err
def test_failed_cycle_keeps_checkpoint_and_recovers(tmp_path):
from production_analytics.calculations.material_consumption import MaterialSample
from production_analytics.service.material_polling import (
MaterialPollingService,
MaterialPollingState,
)
from production_analytics.service.material_state_store import JsonMaterialStateStore
order = ' 00842 '
store = JsonMaterialStateStore(tmp_path)
store.save('machine', order, MaterialPollingState('run', MaterialIntegrationState(
cumulative_consumption_kg=1, integrated_running_seconds=10,
last_processed_timestamp=NOW, last_material_rate_kg_per_hour=3600,
last_gate_value=1, integration_active=True,
)))
checkpoint = next(tmp_path.glob('*.json'))
original = checkpoint.read_bytes()
gateway = Mock()
gateway.get_open_production_run.return_value = EnlyzeProductionRun(
'run', 'machine', None, order, NOW, None,
)
gateway.get_material_samples.side_effect = [
ValueError('bad response'), [MaterialSample(NOW + timedelta(seconds=10), 3600, 1)],
]
service = MaterialPollingService(
gateway=gateway, state_store=_ReportingStateStore(store), machine_id='machine',
rate_variable_id='rate', gate_variable_id='gate', gate_threshold=.5,
max_sample_gap_seconds=20,
)
sleeps = 0
def sleep(seconds):
nonlocal sleeps
sleeps += 1
if sleeps == 1:
assert checkpoint.read_bytes() == original
else:
raise KeyboardInterrupt
MaterialPollingRunner(
service, calculation_id='test-calculation', snapshot_writer=Mock(),
machine_id='machine', poll_interval_seconds=1,
clock=lambda: NOW + timedelta(seconds=10), sleep=sleep,
stdout=io.StringIO(), stderr=io.StringIO(),
).run()
restored = store.load('machine', order)
assert restored.integration_state.cumulative_consumption_kg == 11
assert restored.integration_state.integrated_running_seconds == 20
def test_unwritable_state_location_is_startup_failure(tmp_path, monkeypatch, capsys):
monkeypatch.setenv('ENLYZE_BASE_URL', 'https://example.invalid/api/')
blocked = tmp_path / 'file'
blocked.write_text('not a directory')
assert main([
'run', 'material-poll', '--config', str(EXAMPLE),
'--state-directory', str(blocked), '--secrets-file', str(tmp_path / 'absent'),
]) == 2
assert 'State directory startup check failed' in capsys.readouterr().err
+141
View File
@@ -0,0 +1,141 @@
from datetime import UTC, datetime
import pytest
from production_analytics.calculations.material_consumption import (
MaterialIntegrationState,
)
from production_analytics.service.material_polling import MaterialPollingState
from production_analytics.service.material_state_store import JsonMaterialStateStore
def test_json_state_store_returns_none_for_missing_state(tmp_path) -> None:
store = JsonMaterialStateStore(tmp_path)
assert store.load("machine-1", "ORDER-1") is None
def test_json_state_store_roundtrip(tmp_path) -> None:
store = JsonMaterialStateStore(tmp_path)
original = MaterialPollingState(
run_id="run-1",
integration_state=MaterialIntegrationState(
cumulative_consumption_kg=123.4,
integrated_running_seconds=456.0,
last_processed_timestamp=datetime(
2026, 9, 3, 5, 10, 50, tzinfo=UTC
),
last_material_rate_kg_per_hour=1028.5,
last_gate_value=2.1,
integration_active=True,
),
)
store.save("machine-1", "ORDER-1", original)
restored = store.load("machine-1", "ORDER-1")
assert restored == original
def test_json_state_store_separates_orders(tmp_path) -> None:
store = JsonMaterialStateStore(tmp_path)
state_a = MaterialPollingState(
run_id="run-a",
integration_state=MaterialIntegrationState(
cumulative_consumption_kg=10.0,
),
)
state_b = MaterialPollingState(
run_id="run-b",
integration_state=MaterialIntegrationState(
cumulative_consumption_kg=20.0,
),
)
store.save("machine-1", "ORDER-1", state_a)
store.save("machine-1", "ORDER-2", state_b)
assert store.load("machine-1", "ORDER-1") == state_a
assert store.load("machine-1", "ORDER-2") == state_b
def test_identifiers_are_isolated_safe_and_deterministic(tmp_path) -> None:
store = JsonMaterialStateStore(tmp_path)
pairs = [("a/b", "x"), ("a_b", "x"), ("a", "b__x"), ("a__b", "x"),
("../..", "/tmp/escape"), ("", ""), ("a/b", "y")]
for index, pair in enumerate(pairs):
store.save(*pair, MaterialPollingState(str(index), MaterialIntegrationState()))
assert len(list(tmp_path.iterdir())) == len(pairs)
for index, pair in enumerate(pairs):
assert JsonMaterialStateStore(tmp_path).load(*pair).run_id == str(index)
assert all(path.parent == tmp_path and len(path.name) == 69 for path in tmp_path.iterdir())
def test_failed_write_preserves_primary_and_cleans_temporary_file(tmp_path, monkeypatch) -> None:
import production_analytics.service.material_state_store as module
store = JsonMaterialStateStore(tmp_path)
state = MaterialPollingState("old", MaterialIntegrationState())
store.save("machine", "order", state)
path = next(tmp_path.iterdir())
original = path.read_bytes()
def fail_dump(payload, file, **kwargs):
file.write('{"partial":')
raise OSError("disk full")
monkeypatch.setattr(module.json, "dump", fail_dump)
with pytest.raises(OSError, match="disk full"):
store.save("machine", "order", MaterialPollingState("new", MaterialIntegrationState()))
assert path.read_bytes() == original
assert list(tmp_path.iterdir()) == [path]
def test_replace_failure_preserves_primary(tmp_path, monkeypatch) -> None:
import production_analytics.service.material_state_store as module
store = JsonMaterialStateStore(tmp_path)
state = MaterialPollingState("old", MaterialIntegrationState())
store.save("m", "o", state)
path = next(tmp_path.iterdir())
def fail_replace(source, target):
assert source.parent == target.parent == tmp_path
assert store.load("m", "o") == state
raise OSError("replace failed")
monkeypatch.setattr(module.os, "replace", fail_replace)
with pytest.raises(OSError, match="replace failed"):
store.save("m", "o", MaterialPollingState("new", MaterialIntegrationState()))
assert store.load("m", "o") == state
assert list(tmp_path.iterdir()) == [path]
@pytest.mark.parametrize("contents", ['{', '[]', '{}', '{"run_id": null}',
'{"run_id": "r", "integration_state": {}}'])
def test_corrupt_state_fails_clearly(tmp_path, contents) -> None:
store = JsonMaterialStateStore(tmp_path)
store.save("m", "o", MaterialPollingState("r", MaterialIntegrationState()))
next(tmp_path.iterdir()).write_text(contents)
with pytest.raises(ValueError, match="Invalid material polling state"):
store.load("m", "o")
@pytest.mark.parametrize("field,value", [
("integration_active", "false"), ("cumulative_consumption_kg", float("nan")),
("integrated_running_seconds", -1), ("last_processed_timestamp", "2026-09-03T05:00:00"),
("integration_active", True), ("last_gate_value", True),
])
def test_malformed_integration_state_fails(tmp_path, field, value) -> None:
import json
store = JsonMaterialStateStore(tmp_path)
store.save("m", "o", MaterialPollingState("r", MaterialIntegrationState()))
path = next(tmp_path.iterdir())
payload = json.loads(path.read_text())
payload["integration_state"][field] = value
path.write_text(json.dumps(payload))
with pytest.raises(ValueError, match="Invalid material polling state"):
store.load("m", "o")
+128
View File
@@ -0,0 +1,128 @@
import io
import sys
from datetime import UTC, datetime
from unittest.mock import MagicMock, Mock, patch
import pytest
from production_analytics.calculations.material_consumption import MaterialIntegrationState
from production_analytics.enlyze.exploration import ConfigurationError
from production_analytics.enlyze.gateway import EnlyzeProductionRun
from production_analytics.service.material_polling import MaterialPollResult
from production_analytics.service.material_runner import MaterialPollingRunner
from production_analytics.service.postgres_material import (
PostgresMaterialSnapshotWriter,
PostgresSettings,
)
ENV = dict(POSTGRES_HOST='localhost', POSTGRES_PORT='5432', POSTGRES_DB='analytics',
POSTGRES_USER='writer', POSTGRES_PASSWORD='secret-password')
NOW = datetime(2026, 1, 1, tzinfo=UTC)
@pytest.mark.parametrize('key', ENV)
@pytest.mark.parametrize('value', [None, '', ' ', '\x00'])
def test_required_settings(key, value):
environment = dict(ENV)
if value is None:
del environment[key]
else:
environment[key] = value
with pytest.raises(ConfigurationError, match=key) as error:
PostgresSettings.from_environment(environment)
assert 'secret-password' not in str(error.value)
@pytest.mark.parametrize('port', ['0', '-1', '65536', '1.5', 'secret-password'])
def test_invalid_port(port):
with pytest.raises(ConfigurationError, match='POSTGRES_PORT') as error:
PostgresSettings.from_environment(dict(ENV, POSTGRES_PORT=port))
assert str(error.value) == 'POSTGRES_PORT must be an integer from 1 to 65535'
def test_valid_settings_and_password_repr():
settings = PostgresSettings.from_environment(ENV)
assert (settings.host, settings.port, settings.dbname, settings.user) == (
'localhost', 5432, 'analytics', 'writer',
)
assert settings.password == 'secret-password'
assert settings.password not in repr(settings)
def test_parameterized_insert_and_connection_lifecycle():
driver = MagicMock()
connection = driver.connect.return_value.__enter__.return_value
writer = PostgresMaterialSnapshotWriter(PostgresSettings.from_environment(ENV))
fields = dict(timestamp=NOW, calculation_id='calc', machine_id='machine',
production_order=" 00842'; DROP TABLE x; -- ", run_id='run', consumption_kg=12.5)
with patch.dict(sys.modules, psycopg=driver):
writer.write(**fields)
sql, parameters = connection.execute.call_args.args
assert sql.count('%s') == 6
assert fields['production_order'] not in sql
assert parameters == tuple(fields.values())
assert 'ON CONFLICT' in sql
connection.execute.assert_called_once()
driver.connect.assert_called_once_with(
host='localhost', port=5432, dbname='analytics', user='writer',
password='secret-password', connect_timeout=10, options='-c statement_timeout=10000',
)
driver.connect.return_value.__exit__.assert_called_once_with(None, None, None)
def test_writer_failure_propagates_and_next_write_reconnects():
driver = MagicMock()
driver.connect.return_value.__enter__.return_value.execute.side_effect = [
RuntimeError('secret'), None,
]
writer = PostgresMaterialSnapshotWriter(PostgresSettings.from_environment(ENV))
fields = dict(timestamp=NOW, calculation_id='calc', machine_id='machine',
production_order='order', run_id='run', consumption_kg=0)
with patch.dict(sys.modules, psycopg=driver):
with pytest.raises(RuntimeError):
writer.write(**fields)
writer.write(**fields)
assert driver.connect.call_count == 2
assert driver.connect.return_value.__exit__.call_count == 2
@pytest.mark.parametrize('consumption', [0, 12.5])
def test_runner_snapshot_fields_and_none(consumption):
result = MaterialPollResult(
EnlyzeProductionRun('run', 'machine', None, ' 00842 ', NOW, None),
MaterialIntegrationState(consumption, 20),
)
writer = Mock()
output = io.StringIO()
MaterialPollingRunner(
Mock(poll_once=Mock(side_effect=[None, result])), machine_id='machine',
calculation_id='configured-calculation', snapshot_writer=writer,
poll_interval_seconds=10, clock=lambda: NOW,
sleep=Mock(side_effect=[None, KeyboardInterrupt]), stdout=output,
).run()
writer.write.assert_called_once_with(
timestamp=NOW, calculation_id='configured-calculation', machine_id='machine',
production_order=' 00842 ', run_id='run', consumption_kg=consumption,
)
assert 'no open Production Run' in output.getvalue()
assert f'consumption_kg={consumption:.9f}' in output.getvalue()
def test_snapshot_failure_reports_safely_and_continues():
result = MaterialPollResult(
EnlyzeProductionRun('run', 'machine', None, 'order', NOW, None),
MaterialIntegrationState(12.5, 20),
)
service = Mock(poll_once=Mock(return_value=result))
writer = Mock(write=Mock(side_effect=[RuntimeError('secret-password arbitrary SQL'), None]))
output, errors = io.StringIO(), io.StringIO()
MaterialPollingRunner(
service, machine_id='machine', calculation_id='calc', snapshot_writer=writer,
poll_interval_seconds=10, clock=lambda: NOW,
sleep=Mock(side_effect=[None, KeyboardInterrupt]), stdout=output, stderr=errors,
).run()
assert service.poll_once.call_count == writer.write.call_count == 2
assert errors.getvalue().count('PostgreSQL snapshot write failed') == 1
assert 'secret-password' not in errors.getvalue()
assert 'arbitrary SQL' not in errors.getvalue()
assert output.getvalue().count('consumption_kg=12.500000000') == 2