diff --git a/README.md b/README.md index 2f80351..a9aa2c5 100644 --- a/README.md +++ b/README.md @@ -318,7 +318,7 @@ across the order's runs. The service neither reintegrates nor sums snapshots and does not bridge run boundaries. The result includes workplace/machine/calculation identifiers, both order identifiers, -article context, nominal width, both original source timestamps, good area, aligned +article context, nominal width, the UTC ERP feedback instant and selected material timestamp, good area, aligned consumption, `material_consumption_kg_per_m2 = consumption_kg / good_quantity_m2`, and `material_consumption_g_per_m2 = material_consumption_kg_per_m2 * 1000`. Nominal width comes from generic ERP description parsing and is context only; @@ -352,9 +352,9 @@ result = service.evaluate_current() Aware ERP timestamps are compared as UTC instants. Naive ERP timestamps require explicit `erp_timezone`; there is no inferred system/database timezone. Ambiguous -or nonexistent local times during DST transitions are rejected. The original ERP -timestamp (including naivety) and selected material timestamp are preserved in the -result; retain the source timezone configuration alongside it for interpretation. +or nonexistent local times during DST transitions are rejected. The result carries the +ERP feedback instant normalized to UTC and the selected aware material timestamp. +The original ERP status object remains unchanged. This alignment avoids knowingly including future material consumption, but does not eliminate ERP roll-feedback timing uncertainty (feedback may be roughly one @@ -366,11 +366,10 @@ stored timestamps satisfy the cutoff; a later cumulative total cannot reconstruc an earlier value. Reliable final evaluation requires retaining final ERP feedback; the current-workplace adapter alone does not provide historical completed orders. -Call on changed ERP feedback; no new runner or polling loop is provided. Results -are returned in memory without modifying earlier results or material snapshots. -No schema, Grafana, or timeseries changes are needed. Daily per-machine 24h reporting -and aggregation remain future work. Repository SQL selection tests use an in-memory -SQLite fixture with driver transport adaptation; they require no live PostgreSQL. +The service itself returns results in memory without modifying material snapshots. +The standalone persistence runner is described below. Daily per-machine 24h reporting +and aggregation remain future work. Repository SQL tests use an in-memory SQLite +fixture with driver transport adaptation; they require no live PostgreSQL. ## Peak-cycle detection @@ -412,3 +411,56 @@ 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. + +### Feedback-driven material-efficiency persistence + +Apply the updated `db/schema.sql` using the PostgreSQL setup command above, then run +this independent foreground process: + +```bash +production-analytics run material-efficiency --config config/k7-material-efficiency.yaml +``` + +`poll_interval_seconds` in YAML defaults to 60 seconds (fixed delay after each +cycle). Workplace, machine, calculation, order format and `erp_timezone` are also +configured in YAML; K7 uses `Europe/Berlin`. PostgreSQL settings use the existing +exported `POSTGRES_*` variables and `--secrets-file` (default `secrets/enlyze.env`); +ERP uses `--erp-secrets-file` (default `secrets/erp.env`). No additional secrets are +needed, and `.env` is not loaded automatically. + +`PostgresMaterialEfficiencyWriter.write(snapshot) -> bool` validates aware +timestamps and finite numbers, then inserts all KPI fields in a short transaction. +`material_efficiency_snapshots` retains one immutable point per +`(calculation_id, machine_id, enlyze_production_order, erp_feedback_timestamp)`. +The writer uses `RETURNING 1` to return `True` for a new insert and `False` for a +duplicate. Only new inserts produce a flushed stdout line with feedback time, +workplace, order, good m², material kg and g/m². Unavailable evaluations and +duplicates remain silent. +Repeated evaluations attempt `ON CONFLICT DO NOTHING`; they never update history, +including after a restart. Missing/invalid KPI inputs produce no row and can be +retried on a later cycle. Nullable article context remains SQL NULL. + +Persistence frequency follows ERP feedback changes, not ENLYZE sample frequency. +This runner is independent of the 10-second material polling process. Live values +are plausibility indicators; final production-order (FA) values are more meaningful. +Grafana can read PostgreSQL alone: + +```sql +SELECT erp_feedback_timestamp AS "time", material_consumption_g_per_m2 +FROM material_efficiency_snapshots +WHERE $__timeFilter(erp_feedback_timestamp) + AND calculation_id = 'k7-fiber-consumption' + AND machine_id = 'c220f95c-a65e-4cb7-99b7-0626d6c7508c' +ORDER BY erp_feedback_timestamp; +``` + +ERP read errors and transient PostgreSQL connection/operational errors are reported +using error classes only and retried after the configured delay. Each database +operation opens a fresh connection. Configuration, schema/programming and invalid +timestamp contract errors terminate with a nonzero CLI exit; ambiguous/nonexistent +DST feedback remains rejected. Ctrl-C stops the foreground runner cleanly. + +The current-status ERP source cannot backfill feedback events missed between polls +or during outages, nor guarantee observation of final FA feedback before the order +changes. Daily 24h reports and final order summary tables remain future work. +No dashboard, aggregation or lag correction is added here. diff --git a/config/k7-material-efficiency.yaml b/config/k7-material-efficiency.yaml new file mode 100644 index 0000000..6d15d6e --- /dev/null +++ b/config/k7-material-efficiency.yaml @@ -0,0 +1,6 @@ +workplace: K7 +machine_id: c220f95c-a65e-4cb7-99b7-0626d6c7508c +calculation_id: k7-fiber-consumption +production_order_format: "K 7-{production_order}" +erp_timezone: Europe/Berlin +poll_interval_seconds: 60 diff --git a/db/schema.sql b/db/schema.sql index 0f3d5a2..53f99b7 100644 --- a/db/schema.sql +++ b/db/schema.sql @@ -11,3 +11,24 @@ CREATE TABLE IF NOT EXISTS material_consumption_snapshots ( -- 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); + +CREATE TABLE IF NOT EXISTS material_efficiency_snapshots ( + erp_feedback_timestamp timestamptz NOT NULL, + material_snapshot_timestamp timestamptz NOT NULL, + calculation_id text NOT NULL, + machine_id text NOT NULL, + workplace text NOT NULL, + erp_production_order text NOT NULL, + enlyze_production_order text NOT NULL, + article_number text NULL, + article_description text NULL, + nominal_width_m double precision NULL, + good_quantity_m2 double precision NOT NULL, + material_consumption_kg double precision NOT NULL, + material_consumption_kg_per_m2 double precision NOT NULL, + material_consumption_g_per_m2 double precision NOT NULL, + PRIMARY KEY (calculation_id, machine_id, enlyze_production_order, erp_feedback_timestamp) +); + +CREATE INDEX IF NOT EXISTS material_efficiency_snapshots_machine_time_idx + ON material_efficiency_snapshots (calculation_id, machine_id, erp_feedback_timestamp); diff --git a/src/production_analytics/cli/__main__.py b/src/production_analytics/cli/__main__.py index 379c425..4ce2aa9 100644 --- a/src/production_analytics/cli/__main__.py +++ b/src/production_analytics/cli/__main__.py @@ -37,10 +37,18 @@ def _parser() -> argparse.ArgumentParser: raw.add_argument("--query", action="append", type=_query_item, default=[], metavar="KEY=VALUE") raw.add_argument("--pretty", action="store_true", help="Pretty-print sanitized JSON") raw.add_argument("--verbose", action="store_true", help="Show safe request progress on stderr") - raw.add_argument("--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/") - raw.add_argument("--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/") - raw.add_argument("--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS) - raw.add_argument("--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS) + raw.add_argument( + "--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/" + ) + raw.add_argument( + "--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/" + ) + raw.add_argument( + "--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS + ) + raw.add_argument( + "--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS + ) raw.add_argument( "--secrets-file", type=Path, @@ -52,18 +60,44 @@ def _parser() -> argparse.ArgumentParser: timeseries.add_argument("--start", required=True, help="ISO 8601 datetime with timezone") timeseries.add_argument("--end", required=True, help="ISO 8601 datetime with timezone") timeseries.add_argument("--variable", required=True, help="Variable UUID") - timeseries.add_argument("--resampling-interval", type=int, help="Seconds; schema range is 10..604800") + timeseries.add_argument( + "--resampling-interval", type=int, help="Seconds; schema range is 10..604800" + ) timeseries.add_argument( "--resampling-method", - choices=["first", "last", "max", "min", "count", "sum", "avg", "median", "std", "q5", "q25", "q75", "q95"], + choices=[ + "first", + "last", + "max", + "min", + "count", + "sum", + "avg", + "median", + "std", + "q5", + "q25", + "q75", + "q95", + ], help="Optional schema-defined method for this variable", ) timeseries.add_argument("--pretty", action="store_true", help="Pretty-print sanitized JSON") - timeseries.add_argument("--verbose", action="store_true", help="Show safe request progress on stderr") - timeseries.add_argument("--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/") - timeseries.add_argument("--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/") - timeseries.add_argument("--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS) - timeseries.add_argument("--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS) + timeseries.add_argument( + "--verbose", action="store_true", help="Show safe request progress on stderr" + ) + timeseries.add_argument( + "--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/" + ) + timeseries.add_argument( + "--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/" + ) + timeseries.add_argument( + "--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS + ) + timeseries.add_argument( + "--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS + ) timeseries.add_argument( "--secrets-file", type=Path, @@ -78,6 +112,10 @@ def _parser() -> argparse.ArgumentParser: 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")) + efficiency = runners.add_parser("material-efficiency", help="Persist ERP feedback KPI points") + efficiency.add_argument("--config", type=Path, required=True) + efficiency.add_argument("--secrets-file", type=Path, default=Path("secrets/enlyze.env")) + efficiency.add_argument("--erp-secrets-file", type=Path, default=Path("secrets/erp.env")) return parser @@ -102,8 +140,10 @@ def _run_material(args: argparse.Namespace) -> int: 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, + 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: @@ -117,9 +157,35 @@ def _run_material(args: argparse.Namespace) -> int: return 0 +def _run_material_efficiency(args: argparse.Namespace) -> int: + from production_analytics.calculations.config import CalculationConfigError + from production_analytics.enlyze.exploration import ConfigurationError + from production_analytics.service.material_efficiency_runtime import ( + build_material_efficiency_runner, + ) + + try: + build_material_efficiency_runner( + args.config, + secrets_file=args.secrets_file, + erp_secrets_file=args.erp_secrets_file, + ).run() + except (CalculationConfigError, ConfigurationError) as exc: + print(f"Config/startup error: {exc}", file=sys.stderr) + return 2 + except Exception as exc: + print(f"Startup/runtime error: {type(exc).__name__}", file=sys.stderr) + return 1 + except KeyboardInterrupt: + return 0 + return 0 + + def main(argv: Sequence[str] | None = None) -> int: args = _parser().parse_args(argv) if args.namespace == "run": + if args.command == "material-efficiency": + return _run_material_efficiency(args) return _run_material(args) if args.namespace != "enlyze" or args.command not in {"raw", "timeseries"}: return 2 @@ -133,7 +199,8 @@ def main(argv: Sequence[str] | None = None) -> int: operation = "GET" if args.command == "raw" else "POST" target = args.path if args.command == "raw" else "/v2/timeseries" print( - f"Requesting {operation} {target} (timeout={settings.timeout_seconds:g}s; authentication {authentication}).", + f"Requesting {operation} {target} (timeout={settings.timeout_seconds:g}s; " + f"authentication {authentication}).", file=sys.stderr, ) client = ExplorationClient(settings) @@ -141,8 +208,13 @@ def main(argv: Sequence[str] | None = None) -> int: response = client.get(args.path, dict(args.query)) request = {"method": "GET", "path": response.path} else: - if args.resampling_interval is not None and not 10 <= args.resampling_interval <= 604800: - raise ExplorationError("--resampling-interval must be between 10 and 604800 seconds.") + if ( + args.resampling_interval is not None + and not 10 <= args.resampling_interval <= 604800 + ): + raise ExplorationError( + "--resampling-interval must be between 10 and 604800 seconds." + ) variable: dict[str, str] = {"uuid": args.variable} if args.resampling_method: variable["resampling_method"] = args.resampling_method @@ -170,9 +242,15 @@ def main(argv: Sequence[str] | None = None) -> int: indent = 2 if args.pretty or args.save_fixture else None print(json.dumps(payload, indent=indent, sort_keys=True)) if args.save_fixture: - print(f"Sanitized fixture written to {_write_fixture(args.fixture_dir, args.save_fixture, payload)}") + print( + "Sanitized fixture written to " + f"{_write_fixture(args.fixture_dir, args.save_fixture, payload)}" + ) if args.save_raw: - print(f"Raw response written to {_write_fixture(args.raw_dir, args.save_raw, raw_payload)}") + print( + "Raw response written to " + f"{_write_fixture(args.raw_dir, args.save_raw, raw_payload)}" + ) except ExplorationError as error: print(f"ENLYZE exploration error: {error}") return 1 diff --git a/src/production_analytics/service/material_efficiency.py b/src/production_analytics/service/material_efficiency.py index 70a032b..e5b4aa1 100644 --- a/src/production_analytics/service/material_efficiency.py +++ b/src/production_analytics/service/material_efficiency.py @@ -92,7 +92,7 @@ class MaterialEfficiencyService: enlyze_production_order=order, article_number=status.article_number, article_description=status.article_description, nominal_width_m=extract_nominal_width_m(status.article_description), - erp_feedback_timestamp=status.feedback_timestamp, + erp_feedback_timestamp=cutoff, material_snapshot_timestamp=material.timestamp, good_quantity_m2=quantity, material_consumption_kg=consumption, material_consumption_kg_per_m2=kg_per_m2, diff --git a/src/production_analytics/service/material_efficiency_runner.py b/src/production_analytics/service/material_efficiency_runner.py new file mode 100644 index 0000000..4129f8e --- /dev/null +++ b/src/production_analytics/service/material_efficiency_runner.py @@ -0,0 +1,60 @@ +"""Independent foreground ERP feedback polling; PostgreSQL owns deduplication.""" + +import sys +import time +from collections.abc import Callable +from typing import TextIO + +import psycopg + +from production_analytics.calculations.config import finite_number +from production_analytics.erp import ErpReadError +from production_analytics.service.material_efficiency import MaterialEfficiencyService +from production_analytics.service.postgres_material_efficiency import MaterialEfficiencyWriter + + +class MaterialEfficiencyRunner: + def __init__( + self, + service: MaterialEfficiencyService, + writer: MaterialEfficiencyWriter, + *, + poll_interval_seconds: float = 60, + sleep: Callable[[float], None] = time.sleep, + stderr: TextIO | None = None, + stdout: TextIO | None = None, + ) -> None: + self.service = service + self.writer = writer + self.interval = finite_number(poll_interval_seconds, "poll interval", positive=True) + self.sleep = sleep + self.stderr = stderr if stderr is not None else sys.stderr + self.stdout = stdout if stdout is not None else sys.stdout + + def run(self) -> None: + try: + while True: + try: + snapshot = self.service.evaluate_current() + if snapshot is not None and self.writer.write(snapshot): + print( + f"{snapshot.erp_feedback_timestamp.isoformat()} " + f"workplace={snapshot.workplace!r} " + f"production_order={snapshot.enlyze_production_order!r} " + f"good_m2={snapshot.good_quantity_m2:.3f} " + f"material_kg={snapshot.material_consumption_kg:.3f} " + f"g_per_m2={snapshot.material_consumption_g_per_m2:.3f}", + file=self.stdout, + flush=True, + ) + except (ErpReadError, psycopg.OperationalError, psycopg.InterfaceError) as exc: + # Never include driver messages, SQL, credentials or response bodies. + print( + f"Material efficiency cycle failed: {type(exc).__name__}", + file=self.stderr, + flush=True, + ) + # Fixed delay even after unavailable evaluations or transient failures. + self.sleep(self.interval) + except KeyboardInterrupt: + return diff --git a/src/production_analytics/service/material_efficiency_runtime.py b/src/production_analytics/service/material_efficiency_runtime.py new file mode 100644 index 0000000..0c4bc0c --- /dev/null +++ b/src/production_analytics/service/material_efficiency_runtime.py @@ -0,0 +1,82 @@ +"""Small strict YAML configuration and composition for KPI persistence.""" + +import os +from pathlib import Path +from zoneinfo import ZoneInfo, ZoneInfoNotFoundError + +import yaml + +from production_analytics.calculations.config import ( + CalculationConfigError, + _UniqueLoader, + finite_number, +) +from production_analytics.context import build_enlyze_production_order +from production_analytics.enlyze.exploration import ConfigurationError, load_secret_file +from production_analytics.erp import ErpSettings, ErpWorkplaceStatusGateway +from production_analytics.service.material_efficiency import MaterialEfficiencyService +from production_analytics.service.material_efficiency_runner import MaterialEfficiencyRunner +from production_analytics.service.postgres_material import ( + PostgresMaterialSnapshotRepository, + PostgresSettings, +) +from production_analytics.service.postgres_material_efficiency import ( + PostgresMaterialEfficiencyWriter, +) + + +def build_material_efficiency_runner( + config: Path, + *, + secrets_file: Path = Path("secrets/enlyze.env"), + erp_secrets_file: Path = Path("secrets/erp.env"), +) -> MaterialEfficiencyRunner: + try: + document = yaml.load(config.read_text(encoding="utf-8"), Loader=_UniqueLoader) + except (OSError, UnicodeError, yaml.YAMLError): + raise CalculationConfigError("Cannot read material efficiency YAML configuration") from None + required = { + "workplace", + "machine_id", + "calculation_id", + "production_order_format", + "erp_timezone", + } + if ( + not isinstance(document, dict) + or not required <= document.keys() + or document.keys() - required - {"poll_interval_seconds"} + ): + raise CalculationConfigError("Invalid material efficiency configuration fields") + for name in required: + value = document[name] + if not isinstance(value, str) or not value.strip() or "\x00" in value: + raise CalculationConfigError(f"{name} must be a non-empty string without NUL bytes") + try: + timezone = ZoneInfo(document["erp_timezone"]) + build_enlyze_production_order("0", document["production_order_format"]) + except (ValueError, ZoneInfoNotFoundError): + raise CalculationConfigError("Invalid ERP timezone or production order format") from None + interval = finite_number( + document.get("poll_interval_seconds", 60), "poll interval", positive=True + ) + environment = dict(os.environ) + try: + environment.update(load_secret_file(secrets_file)) + except Exception: + raise ConfigurationError("PostgreSQL settings file could not be loaded") from None + postgres = PostgresSettings.from_environment(environment) + erp = ErpSettings.from_secret_file(erp_secrets_file) + return MaterialEfficiencyRunner( + MaterialEfficiencyService( + ErpWorkplaceStatusGateway(erp), + PostgresMaterialSnapshotRepository(postgres), + workplace=document["workplace"], + machine_id=document["machine_id"], + calculation_id=document["calculation_id"], + format_template=document["production_order_format"], + erp_timezone=timezone, + ), + PostgresMaterialEfficiencyWriter(postgres), + poll_interval_seconds=interval, + ) diff --git a/src/production_analytics/service/postgres_material_efficiency.py b/src/production_analytics/service/postgres_material_efficiency.py new file mode 100644 index 0000000..37e9863 --- /dev/null +++ b/src/production_analytics/service/postgres_material_efficiency.py @@ -0,0 +1,71 @@ +"""Immutable PostgreSQL history of feedback-aligned KPI points.""" + +from typing import Protocol + +import psycopg + +from production_analytics.calculations.config import finite_number +from production_analytics.service.material_efficiency import MaterialEfficiencySnapshot +from production_analytics.service.postgres_material import PostgresSettings + + +class MaterialEfficiencyWriter(Protocol): + def write(self, snapshot: MaterialEfficiencySnapshot) -> bool: ... + + +class PostgresMaterialEfficiencyWriter: + def __init__(self, settings: PostgresSettings) -> None: + self.settings = settings + + def write(self, snapshot: MaterialEfficiencySnapshot) -> bool: + for name in ("erp_feedback_timestamp", "material_snapshot_timestamp"): + if getattr(snapshot, name).utcoffset() is None: + raise ValueError(f"{name} must be timezone-aware") + for name in ( + "good_quantity_m2", + "material_consumption_kg", + "material_consumption_kg_per_m2", + "material_consumption_g_per_m2", + ): + finite_number(getattr(snapshot, name), name) + if snapshot.nominal_width_m is not None: + finite_number(snapshot.nominal_width_m, "nominal_width_m") + # Context exit commits/rolls back and closes; later cycles reconnect after outages. + with psycopg.connect( + host=self.settings.host, + port=self.settings.port, + dbname=self.settings.dbname, + user=self.settings.user, + password=self.settings.password, + connect_timeout=10, + options="-c statement_timeout=10000", + ) as connection: + row = connection.execute( + """INSERT INTO material_efficiency_snapshots ( + erp_feedback_timestamp, material_snapshot_timestamp, calculation_id, + machine_id, workplace, erp_production_order, enlyze_production_order, + article_number, article_description, nominal_width_m, good_quantity_m2, + material_consumption_kg, material_consumption_kg_per_m2, + material_consumption_g_per_m2 + ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) + ON CONFLICT (calculation_id, machine_id, enlyze_production_order, + erp_feedback_timestamp) DO NOTHING + RETURNING 1""", + ( + snapshot.erp_feedback_timestamp, + snapshot.material_snapshot_timestamp, + snapshot.calculation_id, + snapshot.machine_id, + snapshot.workplace, + snapshot.erp_production_order, + snapshot.enlyze_production_order, + snapshot.article_number, + snapshot.article_description, + snapshot.nominal_width_m, + snapshot.good_quantity_m2, + snapshot.material_consumption_kg, + snapshot.material_consumption_kg_per_m2, + snapshot.material_consumption_g_per_m2, + ), + ).fetchone() + return row is not None diff --git a/tests/test_material_efficiency.py b/tests/test_material_efficiency.py index c193533..dcd8cfa 100644 --- a/tests/test_material_efficiency.py +++ b/tests/test_material_efficiency.py @@ -152,14 +152,15 @@ def test_temporal_regression_and_successive_feedback(): assert first.good_quantity_m2 == 800 -def test_naive_feedback_requires_explicit_timezone_and_preserves_source(): +def test_naive_feedback_requires_explicit_timezone_and_returns_instant(): naive = datetime(2026, 9, 4, 12, 5) status = replace(STATUS, feedback_timestamp=naive) with pytest.raises(ValueError, match='explicit erp_timezone'): service().evaluate(status) subject = service(erp_timezone=ZoneInfo('Europe/Berlin')) result = subject.evaluate(status) - assert result.erp_feedback_timestamp == naive + assert result.erp_feedback_timestamp == FEEDBACK + assert status.feedback_timestamp == naive assert subject.materials.latest_at_or_before.call_args.kwargs['timestamp'] == FEEDBACK @@ -171,10 +172,11 @@ def test_dst_gap_and_overlap_rejected(timestamp): ) -def test_aware_feedback_preserves_offset(): +def test_aware_feedback_returns_utc(): timestamp = FEEDBACK.astimezone(ZoneInfo('Europe/Berlin')) result = service().evaluate(replace(STATUS, feedback_timestamp=timestamp)) - assert result.erp_feedback_timestamp is timestamp + assert result.erp_feedback_timestamp == timestamp + assert result.erp_feedback_timestamp.tzinfo is UTC def test_database_failure_propagates(): diff --git a/tests/test_material_efficiency_persistence.py b/tests/test_material_efficiency_persistence.py new file mode 100644 index 0000000..fd8a36b --- /dev/null +++ b/tests/test_material_efficiency_persistence.py @@ -0,0 +1,447 @@ +import io +import sqlite3 +from dataclasses import replace +from datetime import UTC, datetime +from pathlib import Path +from unittest.mock import MagicMock, Mock + +import psycopg +import pytest +import yaml + +from production_analytics.calculations.config import CalculationConfigError +from production_analytics.cli.__main__ import main +from production_analytics.erp import CurrentWorkplaceStatus, ErpReadError +from production_analytics.service.material_efficiency import ( + MaterialEfficiencySnapshot, +) +from production_analytics.service.material_efficiency_runner import MaterialEfficiencyRunner +from production_analytics.service.material_efficiency_runtime import ( + build_material_efficiency_runner, +) +from production_analytics.service.postgres_material import ( + MaterialConsumptionSnapshot, + PostgresSettings, +) +from production_analytics.service.postgres_material_efficiency import ( + PostgresMaterialEfficiencyWriter, +) + +NOW = datetime(2026, 9, 4, 10, tzinfo=UTC) +SNAPSHOT = MaterialEfficiencySnapshot( + "example", + "machine", + "calc", + "00123", + " ORDER/00123'; -- ", + "article", + "fabric 5 x 50 m", + 5, + NOW, + NOW, + 800, + 1000, + 1.25, + 1250, +) +SETTINGS = PostgresSettings("localhost", 5432, "analytics", "writer", "secret") + + +@pytest.fixture +def storage(monkeypatch): + # Real schema and INSERT semantics, with only driver transport adapted for SQLite. + database = sqlite3.connect(":memory:") + database.executescript(Path("db/schema.sql").read_text()) + connection = MagicMock() + + def execute(sql, params): + assert sql.lstrip().startswith("INSERT INTO") + assert "SELECT" not in sql.upper() + assert "RETURNING 1" in sql + assert sql.count("%s") == 14 + assert SNAPSHOT.enlyze_production_order not in sql + return database.execute( + sql.replace("%s", "?"), + tuple(value.isoformat() if isinstance(value, datetime) else value for value in params), + ) + + connection.__enter__.return_value.execute.side_effect = execute + connect = Mock(return_value=connection) + monkeypatch.setattr(psycopg, "connect", connect) + yield database, connect, connection + database.close() + + +@pytest.mark.parametrize("nullable", [False, True]) +def test_all_fields_and_immutable_event(storage, nullable): + database, connect, connection = storage + snapshot = ( + replace(SNAPSHOT, article_number=None, article_description=None, nominal_width_m=None) + if nullable + else SNAPSHOT + ) + writer = PostgresMaterialEfficiencyWriter(SETTINGS) + assert writer.write(snapshot) is True + inserted = writer.write( + replace( + snapshot, + material_consumption_kg=9999, + material_snapshot_timestamp=NOW.replace(minute=1), + ) + ) + assert inserted is False + assert connection.__enter__.return_value.execute.call_count == 2 + assert database.execute("SELECT * FROM material_efficiency_snapshots").fetchall() == [ + ( + NOW.isoformat(), + NOW.isoformat(), + "calc", + "machine", + "example", + "00123", + SNAPSHOT.enlyze_production_order, + snapshot.article_number, + snapshot.article_description, + snapshot.nominal_width_m, + 800, + 1000, + 1.25, + 1250, + ) + ] + assert connect.call_count == 2 + assert connect.call_args.kwargs["connect_timeout"] == 10 + connection.__exit__.assert_called_with(None, None, None) + + +@pytest.mark.parametrize( + "field,value", + [ + ("calculation_id", "other"), + ("machine_id", "other"), + ("enlyze_production_order", SNAPSHOT.enlyze_production_order.strip()), + ("enlyze_production_order", SNAPSHOT.enlyze_production_order + "-combined"), + ("erp_feedback_timestamp", NOW.replace(minute=1)), + ], +) +def test_exact_identity(storage, field, value): + writer = PostgresMaterialEfficiencyWriter(SETTINGS) + writer.write(SNAPSHOT) + writer.write(replace(SNAPSHOT, **{field: value})) + assert storage[0].execute("SELECT count(*) FROM material_efficiency_snapshots").fetchone() == ( + 2, + ) + + +@pytest.mark.parametrize("field", ["erp_feedback_timestamp", "material_snapshot_timestamp"]) +def test_naive_rejected(storage, field): + with pytest.raises(ValueError, match="timezone-aware"): + PostgresMaterialEfficiencyWriter(SETTINGS).write( + replace(SNAPSHOT, **{field: NOW.replace(tzinfo=None)}), + ) + storage[1].assert_not_called() + + +@pytest.mark.parametrize( + "field", + [ + "nominal_width_m", + "good_quantity_m2", + "material_consumption_kg", + "material_consumption_kg_per_m2", + "material_consumption_g_per_m2", + ], +) +@pytest.mark.parametrize("value", [float("nan"), float("inf"), -float("inf")]) +def test_nonfinite_rejected(storage, field, value): + with pytest.raises(ValueError, match="finite"): + PostgresMaterialEfficiencyWriter(SETTINGS).write(replace(SNAPSHOT, **{field: value})) + storage[1].assert_not_called() + + +def run_cycles(service, writer, cycles, interval=60, stdout=None): + sleep = Mock(side_effect=[None] * (cycles - 1) + [KeyboardInterrupt]) + errors = io.StringIO() + MaterialEfficiencyRunner( + service, writer, poll_interval_seconds=interval, sleep=sleep, stderr=errors, stdout=stdout + ).run() + assert sleep.call_args_list == [((interval,),)] * cycles + return errors.getvalue() + + +def test_runner_history_and_none(storage): + service = Mock( + evaluate_current=Mock( + side_effect=[ + None, + SNAPSHOT, + SNAPSHOT, + replace(SNAPSHOT, erp_feedback_timestamp=NOW.replace(minute=15)), + ] + ) + ) + output = io.StringIO() + run_cycles(service, PostgresMaterialEfficiencyWriter(SETTINGS), 4, 73, stdout=output) + lines = output.getvalue().splitlines() + assert len(lines) == 2 + assert lines[0].startswith(NOW.isoformat()) + assert lines[1].startswith(NOW.replace(minute=15).isoformat()) + assert storage[1].call_count == 3 + assert storage[0].execute("SELECT count(*) FROM material_efficiency_snapshots").fetchone() == ( + 2, + ) + + +@pytest.mark.parametrize("error", [ErpReadError, psycopg.OperationalError, psycopg.InterfaceError]) +def test_read_recovers(error): + service = Mock(evaluate_current=Mock(side_effect=[error("secret"), SNAPSHOT])) + writer = Mock() + errors = run_cycles(service, writer, 2) + writer.write.assert_called_once_with(SNAPSHOT) + assert error.__name__ in errors and "secret" not in errors + + +def test_write_recovers_with_fresh_transaction(storage): + database, connect, connection = storage + connect.side_effect = [psycopg.OperationalError("secret"), connection] + errors = run_cycles( + Mock(evaluate_current=Mock(return_value=SNAPSHOT)), + PostgresMaterialEfficiencyWriter(SETTINGS), + 2, + ) + assert "OperationalError" in errors and "secret" not in errors + assert connect.call_count == 2 + assert database.execute("SELECT count(*) FROM material_efficiency_snapshots").fetchone() == (1,) + + +@pytest.mark.parametrize("error", [ValueError, TypeError, psycopg.ProgrammingError]) +def test_programming_errors_propagate(error): + with pytest.raises(error): + run_cycles(Mock(evaluate_current=Mock(side_effect=error("bad contract"))), Mock(), 1) + + +@pytest.mark.parametrize("interval", [0, -1, float("nan"), float("inf"), True]) +def test_invalid_interval(interval): + with pytest.raises(CalculationConfigError): + MaterialEfficiencyRunner(Mock(), Mock(), poll_interval_seconds=interval) + + +@pytest.fixture +def runtime_files(tmp_path): + config = tmp_path / "config.yaml" + config.write_text( + yaml.safe_dump( + dict( + workplace="example", + machine_id="machine", + calculation_id="calc", + production_order_format="ORDER/{production_order}", + erp_timezone="Europe/Berlin", + ) + ) + ) + postgres = tmp_path / "postgres.env" + postgres.write_text( + "POSTGRES_HOST=localhost\nPOSTGRES_PORT=5432\nPOSTGRES_DB=analytics\n" + "POSTGRES_USER=writer\nPOSTGRES_PASSWORD=secret\n" + ) + erp = tmp_path / "erp.env" + erp.write_text( + "ERP_DB_HOST=localhost\nERP_DB_PORT=1433\nERP_DB_NAME=erp\n" + "ERP_DB_USER=reader\nERP_DB_PASSWORD=secret\n" + ) + return config, postgres, erp + + +def test_configured_runtime_naive_feedback_to_persistence(runtime_files, storage): + config, postgres, erp = runtime_files + runner = build_material_efficiency_runner(config, secrets_file=postgres, erp_secrets_file=erp) + assert runner.interval == 60 + runner.service.erp = Mock( + get_current_workplace_status=Mock( + return_value=CurrentWorkplaceStatus( + "example", + "00123", + None, + None, + datetime(2026, 9, 4, 12), + None, + 800, + None, + None, + None, + ) + ) + ) + runner.service.materials = Mock( + latest_at_or_before=Mock( + return_value=MaterialConsumptionSnapshot( + NOW, + "calc", + "machine", + "ORDER/00123", + "run", + 1000, + ) + ) + ) + runner.sleep = Mock(side_effect=KeyboardInterrupt) + runner.run() + row = storage[0].execute("SELECT * FROM material_efficiency_snapshots").fetchone() + assert row[:7] == ( + NOW.isoformat(), + NOW.isoformat(), + "calc", + "machine", + "example", + "00123", + "ORDER/00123", + ) + + +@pytest.mark.parametrize( + "changes", + [ + {"erp_timezone": "invalid"}, + {"production_order_format": "missing"}, + {"workplace": ""}, + {"machine_id": None}, + {"poll_interval_seconds": 0}, + {"unknown": 1}, + ], +) +def test_startup_config_errors(runtime_files, changes, capsys): + config, postgres, erp = runtime_files + config.write_text(yaml.safe_dump(yaml.safe_load(config.read_text()) | changes)) + assert ( + main( + [ + "run", + "material-efficiency", + "--config", + str(config), + "--secrets-file", + str(postgres), + "--erp-secrets-file", + str(erp), + ] + ) + == 2 + ) + assert "Config/startup error" in capsys.readouterr().err + + +def test_k7_config_and_cli(runtime_files, monkeypatch): + _, postgres, erp = runtime_files + runner = build_material_efficiency_runner( + Path("config/k7-material-efficiency.yaml"), secrets_file=postgres, erp_secrets_file=erp + ) + assert runner.service.workplace == "K7" + assert runner.service.machine_id == "c220f95c-a65e-4cb7-99b7-0626d6c7508c" + assert runner.service.calculation_id == "k7-fiber-consumption" + assert runner.service.format_template == "K 7-{production_order}" + assert runner.service.erp_timezone.key == "Europe/Berlin" + run = Mock() + monkeypatch.setattr(MaterialEfficiencyRunner, "run", run) + assert ( + main( + [ + "run", + "material-efficiency", + "--config", + "config/k7-material-efficiency.yaml", + "--secrets-file", + str(postgres), + "--erp-secrets-file", + str(erp), + ] + ) + == 0 + ) + run.assert_called_once() + + +def test_failed_insert_rolls_back_and_later_cycle_reconnects(storage): + database, connect, connection = storage + execute = connection.__enter__.return_value.execute + original = execute.side_effect + attempts = 0 + + def fail_once(sql, params): + nonlocal attempts + attempts += 1 + if attempts == 1: + raise psycopg.OperationalError("secret SQL") + return original(sql, params) + + execute.side_effect = fail_once + errors = run_cycles( + Mock(evaluate_current=Mock(return_value=SNAPSHOT)), + PostgresMaterialEfficiencyWriter(SETTINGS), + 2, + ) + assert "secret" not in errors + assert connect.call_count == 2 + assert connection.__exit__.call_args_list[0].args[0] is psycopg.OperationalError + assert database.execute("SELECT count(*) FROM material_efficiency_snapshots").fetchone() == (1,) + + +def test_configured_interval_and_missing_settings(runtime_files, monkeypatch): + config, postgres, erp = runtime_files + config.write_text( + yaml.safe_dump(yaml.safe_load(config.read_text()) | {"poll_interval_seconds": 91}) + ) + runner = build_material_efficiency_runner(config, secrets_file=postgres, erp_secrets_file=erp) + assert runner.interval == 91 + postgres.write_text("") + for key in ("HOST", "PORT", "DB", "USER", "PASSWORD"): + monkeypatch.delenv(f"POSTGRES_{key}", raising=False) + from production_analytics.enlyze.exploration import ConfigurationError + + with pytest.raises(ConfigurationError, match="POSTGRES_HOST"): + build_material_efficiency_runner(config, secrets_file=postgres, erp_secrets_file=erp) + + +@pytest.mark.parametrize("result", [None, False, True]) +def test_runner_success_output(result): + snapshot = replace(SNAPSHOT, enlyze_production_order="ORDER/00123") + service = Mock(evaluate_current=Mock(return_value=None if result is None else snapshot)) + writer = Mock(write=Mock(return_value=result)) + output = Mock(wraps=io.StringIO()) + errors = run_cycles(service, writer, 1, stdout=output) + assert errors == "" + if result is None: + writer.write.assert_not_called() + else: + writer.write.assert_called_once_with(snapshot) + if result is True: + assert output.getvalue() == ( + "2026-09-04T10:00:00+00:00 workplace='example' production_order='ORDER/00123' " + "good_m2=800.000 material_kg=1000.000 g_per_m2=1250.000\n" + ) + output.flush.assert_called_once_with() + else: + assert output.getvalue() == "" + output.flush.assert_not_called() + + +def test_repeated_duplicates_are_silent(storage): + writer = PostgresMaterialEfficiencyWriter(SETTINGS) + assert writer.write(SNAPSHOT) is True + output = io.StringIO() + run_cycles(Mock(evaluate_current=Mock(return_value=SNAPSHOT)), writer, 3, stdout=output) + assert output.getvalue() == "" + + +def test_commit_failure_does_not_log_success(storage): + _, _, connection = storage + connection.__exit__.side_effect = [psycopg.OperationalError("secret driver SQL"), None] + output = io.StringIO() + errors = run_cycles( + Mock(evaluate_current=Mock(return_value=SNAPSHOT)), + PostgresMaterialEfficiencyWriter(SETTINGS), + 1, + stdout=output, + ) + assert output.getvalue() == "" + assert errors == "Material efficiency cycle failed: OperationalError\n"