From c411581a1aa31c42d10f9156a37eb5bca4f48e1e Mon Sep 17 00:00:00 2001 From: Martin Tazl Date: Sun, 6 Sep 2026 05:45:11 +0200 Subject: [PATCH] Add feedback-aligned material efficiency KPI --- README.md | 77 +++++++- .../service/material_efficiency.py | 116 +++++++++++ .../service/postgres_material.py | 50 ++++- tests/test_material_efficiency.py | 184 ++++++++++++++++++ tests/test_material_snapshot_repository.py | 101 ++++++++++ 5 files changed, 524 insertions(+), 4 deletions(-) create mode 100644 src/production_analytics/service/material_efficiency.py create mode 100644 tests/test_material_efficiency.py create mode 100644 tests/test_material_snapshot_repository.py diff --git a/README.md b/README.md index 87db174..2f80351 100644 --- a/README.md +++ b/README.md @@ -297,9 +297,80 @@ finished-product width: it belongs to the MAHLO measurement system and observed values differ from nominal article widths. Structured ERP product master data would be preferred when available. -These helpers are independent of database access and are not yet wired into live -processing. No normalized context model, kg/m² or material-efficiency calculation, -persistence, polling, runner, timeseries, or Grafana changes are introduced here. +These helpers are independent of database access. The service below composes them +for material efficiency; no background processing is attached. + +## Feedback-aligned material efficiency + +`MaterialEfficiencyService` in `service.material_efficiency` combines ERP good area +with cumulative material consumption for one exact production order. Its +`evaluate_current()` reads the existing ERP workplace adapter once; `evaluate(status)` +evaluates an already supplied `CurrentWorkplaceStatus`. Both return an immutable +`MaterialEfficiencySnapshot` or `None` when inputs are unavailable. + +`PostgresMaterialSnapshotRepository(settings).latest_at_or_before(...)` in +`service.postgres_material` takes keyword arguments `calculation_id`, `machine_id`, +`production_order`, and an aware `timestamp`. A parameterized query selects the +latest snapshot **at or before ERP feedback time**, matching calculation, machine, +and mapped order exactly. It returns `MaterialConsumptionSnapshot` or `None`. +Run IDs do not restrict this lookup: stored consumption is already cumulative +across the order's runs. The service neither reintegrates nor sums snapshots and +does not bridge run boundaries. + +The result includes workplace/machine/calculation identifiers, both order identifiers, +article context, nominal width, both original source timestamps, good area, aligned +consumption, `material_consumption_kg_per_m2 = consumption_kg / good_quantity_m2`, +and `material_consumption_g_per_m2 = material_consumption_kg_per_m2 * 1000`. +Nominal width comes from generic ERP description parsing and is context only; +ERP good square metres remain authoritative even when width cannot be parsed. + +Missing ERP status, missing aligned material, missing/non-positive/non-finite good +area, invalid consumption (negative or non-finite), or non-finite computed ratios +produce `None`. Zero consumption is valid. Configuration errors, invalid timestamp +alignment, and database failures raise rather than masquerading as missing data. + +All machine/workplace settings are supplied explicitly. Example using the existing +K7 material calculation (with previously constructed ERP gateway and PG settings): + +```python +from zoneinfo import ZoneInfo + +from production_analytics.service.material_efficiency import MaterialEfficiencyService +from production_analytics.service.postgres_material import PostgresMaterialSnapshotRepository + +service = MaterialEfficiencyService( + erp_gateway, PostgresMaterialSnapshotRepository(postgres_settings), + workplace="K7", + machine_id="c220f95c-a65e-4cb7-99b7-0626d6c7508c", + calculation_id="k7-fiber-consumption", + format_template="K 7-{production_order}", + # Supply only after verifying the ERP timestamp's source timezone: + erp_timezone=ZoneInfo("Europe/Berlin"), +) +result = service.evaluate_current() +``` + +Aware ERP timestamps are compared as UTC instants. Naive ERP timestamps require +explicit `erp_timezone`; there is no inferred system/database timezone. Ambiguous +or nonexistent local times during DST transitions are rejected. The original ERP +timestamp (including naivety) and selected material timestamp are preserved in the +result; retain the source timezone configuration alongside it for interpretation. + +This alignment avoids knowingly including future material consumption, but does +not eliminate ERP roll-feedback timing uncertainty (feedback may be roughly one +roll ahead or otherwise offset). Live values are plausibility indicators; final +production-order values are more meaningful as relative timing error diminishes. +There is no interpolation, lag correction, smoothing, freshness threshold, or +estimated timestamp. Historical bootstrap snapshots are usable only when their +stored timestamps satisfy the cutoff; a later cumulative total cannot reconstruct +an earlier value. Reliable final evaluation requires retaining final ERP feedback; +the current-workplace adapter alone does not provide historical completed orders. + +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. ## Peak-cycle detection diff --git a/src/production_analytics/service/material_efficiency.py b/src/production_analytics/service/material_efficiency.py new file mode 100644 index 0000000..70a032b --- /dev/null +++ b/src/production_analytics/service/material_efficiency.py @@ -0,0 +1,116 @@ +"""Feedback-aligned, machine-independent production-order material efficiency.""" + +from dataclasses import dataclass +from datetime import UTC, datetime +from math import isfinite +from typing import Protocol +from zoneinfo import ZoneInfo + +from production_analytics.context import build_enlyze_production_order, extract_nominal_width_m +from production_analytics.erp import CurrentWorkplaceStatus +from production_analytics.service.postgres_material import MaterialSnapshotRepository + + +class WorkplaceStatusReader(Protocol): + def get_current_workplace_status(self, workplace: str) -> CurrentWorkplaceStatus | None: ... + + +@dataclass(frozen=True, slots=True) +class MaterialEfficiencySnapshot: + workplace: str + machine_id: str + calculation_id: str + erp_production_order: str + enlyze_production_order: str + article_number: str | None + article_description: str | None + nominal_width_m: float | None + erp_feedback_timestamp: datetime + material_snapshot_timestamp: datetime + good_quantity_m2: float + material_consumption_kg: float + material_consumption_kg_per_m2: float + material_consumption_g_per_m2: float + + +class MaterialEfficiencyService: + def __init__( + self, erp: WorkplaceStatusReader, materials: MaterialSnapshotRepository, *, + workplace: str, machine_id: str, calculation_id: str, format_template: str, + erp_timezone: ZoneInfo | None = None, + ) -> None: + self.erp = erp + self.materials = materials + self.workplace = workplace + self.machine_id = machine_id + self.calculation_id = calculation_id + self.format_template = format_template + self.erp_timezone = erp_timezone + + def evaluate_current(self) -> MaterialEfficiencySnapshot | None: + """Read current ERP feedback once; no scheduling or mutable KPI history.""" + status = self.erp.get_current_workplace_status(self.workplace) + return None if status is None else self.evaluate(status) + + def evaluate(self, status: CurrentWorkplaceStatus) -> MaterialEfficiencySnapshot | None: + """Evaluate supplied feedback; missing/invalid numeric inputs yield None. + + Configuration, timestamp and adapter contract errors raise ValueError. + Database failures propagate, rather than being treated as missing data. + """ + if status.workplace != self.workplace: + raise ValueError("ERP workplace does not match configured workplace") + order = build_enlyze_production_order(status.production_order, self.format_template) + quantity = status.good_quantity_m2 + if quantity is None or not isfinite(quantity) or quantity <= 0: + return None + cutoff = _feedback_instant(status.feedback_timestamp, self.erp_timezone) + material = self.materials.latest_at_or_before( + calculation_id=self.calculation_id, machine_id=self.machine_id, + production_order=order, timestamp=cutoff, + ) + if material is None: + return None + if ( + material.calculation_id != self.calculation_id + or material.machine_id != self.machine_id + or material.production_order != order + or material.timestamp.utcoffset() is None + or material.timestamp > cutoff + ): + raise ValueError("Material repository returned an unaligned snapshot") + consumption = material.consumption_kg + if not isfinite(consumption) or consumption < 0: + return None + kg_per_m2 = consumption / quantity + g_per_m2 = kg_per_m2 * 1000 + if not isfinite(kg_per_m2) or not isfinite(g_per_m2): + return None + return MaterialEfficiencySnapshot( + workplace=status.workplace, machine_id=self.machine_id, + calculation_id=self.calculation_id, erp_production_order=status.production_order, + enlyze_production_order=order, article_number=status.article_number, + article_description=status.article_description, + nominal_width_m=extract_nominal_width_m(status.article_description), + erp_feedback_timestamp=status.feedback_timestamp, + material_snapshot_timestamp=material.timestamp, + good_quantity_m2=quantity, material_consumption_kg=consumption, + material_consumption_kg_per_m2=kg_per_m2, + material_consumption_g_per_m2=g_per_m2, + ) + + +def _feedback_instant(timestamp: datetime, timezone: ZoneInfo | None) -> datetime: + if timestamp.utcoffset() is not None: + return timestamp.astimezone(UTC) + if timezone is None: + raise ValueError("Naive ERP feedback timestamp requires explicit erp_timezone") + # Round trips reject DST gaps; two distinct instants indicate a DST overlap. + instants = set() + for fold in (0, 1): + instant = timestamp.replace(tzinfo=timezone, fold=fold).astimezone(UTC) + if instant.astimezone(timezone).replace(tzinfo=None) == timestamp: + instants.add(instant) + if len(instants) != 1: + raise ValueError("ERP feedback timestamp is ambiguous or nonexistent in erp_timezone") + return instants.pop() diff --git a/src/production_analytics/service/postgres_material.py b/src/production_analytics/service/postgres_material.py index 64b6568..4574dcf 100644 --- a/src/production_analytics/service/postgres_material.py +++ b/src/production_analytics/service/postgres_material.py @@ -1,4 +1,4 @@ -"""PostgreSQL settings and derived material snapshot inserts.""" +"""PostgreSQL settings and derived material snapshot reads/writes.""" from collections.abc import Mapping from dataclasses import dataclass, field @@ -41,6 +41,54 @@ class MaterialSnapshotWriter(Protocol): ) -> None: ... +@dataclass(frozen=True, slots=True) +class MaterialConsumptionSnapshot: + timestamp: datetime + calculation_id: str + machine_id: str + production_order: str + run_id: str + consumption_kg: float + + +class MaterialSnapshotRepository(Protocol): + def latest_at_or_before( + self, *, calculation_id: str, machine_id: str, production_order: str, + timestamp: datetime, + ) -> MaterialConsumptionSnapshot | None: ... + + +class PostgresMaterialSnapshotRepository: + def __init__(self, settings: PostgresSettings) -> None: + self.settings = settings + + def latest_at_or_before( + self, *, calculation_id: str, machine_id: str, production_order: str, + timestamp: datetime, + ) -> MaterialConsumptionSnapshot | None: + """Read one exact-order cumulative snapshot, never later than the aware cutoff.""" + if timestamp.utcoffset() is None: + raise ValueError("Material snapshot cutoff must be timezone-aware") + import psycopg + + with psycopg.connect( + host=self.settings.host, port=self.settings.port, dbname=self.settings.dbname, + user=self.settings.user, password=self.settings.password, connect_timeout=10, + options='-c statement_timeout=10000', + ) as connection: + row = connection.execute( + """SELECT timestamp, calculation_id, machine_id, production_order, + run_id, consumption_kg + FROM material_consumption_snapshots + WHERE calculation_id = %s AND machine_id = %s + AND production_order = %s AND timestamp <= %s + ORDER BY timestamp DESC + LIMIT 1""", + (calculation_id, machine_id, production_order, timestamp), + ).fetchone() + return None if row is None else MaterialConsumptionSnapshot(*row) + + class PostgresMaterialSnapshotWriter: def __init__(self, settings: PostgresSettings) -> None: self.settings = settings diff --git a/tests/test_material_efficiency.py b/tests/test_material_efficiency.py new file mode 100644 index 0000000..c193533 --- /dev/null +++ b/tests/test_material_efficiency.py @@ -0,0 +1,184 @@ +from dataclasses import FrozenInstanceError, replace +from datetime import UTC, datetime +from unittest.mock import Mock +from zoneinfo import ZoneInfo + +import pytest + +from production_analytics.erp import CurrentWorkplaceStatus +from production_analytics.service.material_efficiency import MaterialEfficiencyService +from production_analytics.service.postgres_material import MaterialConsumptionSnapshot + +CALC = 'k7-fiber-consumption' +MACHINE = 'c220f95c-a65e-4cb7-99b7-0626d6c7508c' +START = datetime(2026, 9, 4, 10, tzinfo=UTC) +FEEDBACK = START.replace(minute=5) +STATUS = CurrentWorkplaceStatus( + 'K7', '12026000815', '212520', 'Stex R 1501 C (PR) 5,80 x 50 m', + FEEDBACK, None, 800, None, None, None, +) +MATERIAL = MaterialConsumptionSnapshot(START, CALC, MACHINE, 'K 7-12026000815', 'run', 1000) + + +def service(repository=None, **config): + return MaterialEfficiencyService( + Mock(get_current_workplace_status=Mock(return_value=STATUS)), + repository if repository is not None else Mock(latest_at_or_before=Mock( + return_value=MATERIAL, + )), + **dict(workplace='K7', machine_id=MACHINE, calculation_id=CALC, + format_template='K 7-{production_order}', **config), + ) + + +def test_valid_current_feedback(): + subject = service() + result = subject.evaluate_current() + assert result.material_consumption_kg_per_m2 == 1.25 + assert result.material_consumption_g_per_m2 == 1250 + assert result.nominal_width_m == 5.8 + assert result.good_quantity_m2 == 800 + assert result.material_consumption_kg == 1000 + assert result.erp_feedback_timestamp == FEEDBACK + assert result.material_snapshot_timestamp == START + assert (result.workplace, result.machine_id, result.calculation_id) == ('K7', MACHINE, CALC) + assert (result.erp_production_order, result.enlyze_production_order) == ( + '12026000815', 'K 7-12026000815', + ) + assert (result.article_number, result.article_description) == ( + STATUS.article_number, STATUS.article_description, + ) + subject.erp.get_current_workplace_status.assert_called_once_with('K7') + subject.materials.latest_at_or_before.assert_called_once_with( + calculation_id=CALC, machine_id=MACHINE, production_order='K 7-12026000815', + timestamp=FEEDBACK, + ) + with pytest.raises(FrozenInstanceError): + result.good_quantity_m2 = 5 + + +@pytest.mark.parametrize('quantity', [None, 0, -1, float('nan'), float('inf'), -float('inf')]) +def test_invalid_quantity(quantity): + subject = service() + assert subject.evaluate(replace(STATUS, good_quantity_m2=quantity)) is None + subject.materials.latest_at_or_before.assert_not_called() + + +@pytest.mark.parametrize('description', [None, 'unstructured article']) +def test_width_is_optional_context(description): + result = service().evaluate(replace(STATUS, article_description=description)) + assert result.nominal_width_m is None + assert result.material_consumption_kg_per_m2 == 1.25 + + +def test_unavailable_sources(): + subject = service() + subject.erp.get_current_workplace_status.return_value = None + assert subject.evaluate_current() is None + subject.materials.latest_at_or_before.return_value = None + assert subject.evaluate(STATUS) is None + + +@pytest.mark.parametrize('consumption, quantity', [ + (float('nan'), 800), (float('inf'), 800), (-1, 800), + (1e308, 1e-308), (1e308, 1), +]) +def test_invalid_consumption_and_overflow(consumption, quantity): + subject = service() + subject.materials.latest_at_or_before.return_value = replace( + MATERIAL, consumption_kg=consumption, + ) + assert subject.evaluate(replace(STATUS, good_quantity_m2=quantity)) is None + + +def test_zero_consumption_is_valid(): + subject = service() + subject.materials.latest_at_or_before.return_value = replace(MATERIAL, consumption_kg=0) + assert subject.evaluate(STATUS).material_consumption_kg_per_m2 == 0 + + +@pytest.mark.parametrize('changes', [ + dict(timestamp=START.replace(minute=10)), dict(timestamp=START.replace(tzinfo=None)), + dict(production_order='K 7-12026000815-K 7-12026000816'), + dict(machine_id='other'), dict(calculation_id='other'), +]) +def test_repository_contract_is_checked(changes): + subject = service() + subject.materials.latest_at_or_before.return_value = replace(MATERIAL, **changes) + with pytest.raises(ValueError, match='unaligned'): + subject.evaluate(STATUS) + + +def test_other_explicit_configuration(): + repository = Mock() + repository.latest_at_or_before.return_value = replace( + MATERIAL, machine_id='example-machine', production_order='ORDER/00123', + ) + subject = MaterialEfficiencyService( + Mock(), repository, workplace='example', machine_id='example-machine', + calculation_id=CALC, format_template='ORDER/{production_order}', + ) + result = subject.evaluate(replace(STATUS, workplace='example', production_order=' 00123 ')) + assert result.enlyze_production_order == 'ORDER/00123' + + +def test_wrong_workplace_and_combined_order_rejected(): + with pytest.raises(ValueError, match='workplace'): + service().evaluate(replace(STATUS, workplace='other')) + with pytest.raises(ValueError, match='ERP production order'): + service().evaluate(replace(STATUS, production_order='K 7-123-K 7-456')) + + +def test_temporal_regression_and_successive_feedback(): + later = replace( + MATERIAL, timestamp=START.replace(minute=10), consumption_kg=1100, run_id='run2', + ) + repository = Mock() + repository.latest_at_or_before.side_effect = lambda **kw: max( + (row for row in [MATERIAL, later] if row.timestamp <= kw['timestamp']), + key=lambda row: row.timestamp, default=None, + ) + subject = service(repository) + first = subject.evaluate(STATUS) + second = subject.evaluate(replace( + STATUS, feedback_timestamp=later.timestamp, good_quantity_m2=1000, + )) + assert first.material_snapshot_timestamp == START + assert first.material_consumption_kg == 1000 + assert first.material_consumption_kg_per_m2 == 1.25 + assert second.material_snapshot_timestamp == later.timestamp + assert second.material_consumption_kg == 1100 + assert second.material_consumption_kg_per_m2 == 1.1 + assert first.good_quantity_m2 == 800 + + +def test_naive_feedback_requires_explicit_timezone_and_preserves_source(): + 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 subject.materials.latest_at_or_before.call_args.kwargs['timestamp'] == FEEDBACK + + +@pytest.mark.parametrize('timestamp', [datetime(2026, 3, 29, 2, 30), datetime(2026, 10, 25, 2, 30)]) +def test_dst_gap_and_overlap_rejected(timestamp): + with pytest.raises(ValueError, match='ambiguous or nonexistent'): + service(erp_timezone=ZoneInfo('Europe/Berlin')).evaluate( + replace(STATUS, feedback_timestamp=timestamp), + ) + + +def test_aware_feedback_preserves_offset(): + timestamp = FEEDBACK.astimezone(ZoneInfo('Europe/Berlin')) + result = service().evaluate(replace(STATUS, feedback_timestamp=timestamp)) + assert result.erp_feedback_timestamp is timestamp + + +def test_database_failure_propagates(): + subject = service() + subject.materials.latest_at_or_before.side_effect = RuntimeError('unavailable') + with pytest.raises(RuntimeError, match='unavailable'): + subject.evaluate(STATUS) diff --git a/tests/test_material_snapshot_repository.py b/tests/test_material_snapshot_repository.py new file mode 100644 index 0000000..b11a836 --- /dev/null +++ b/tests/test_material_snapshot_repository.py @@ -0,0 +1,101 @@ +import sqlite3 +import sys +from datetime import UTC, datetime +from unittest.mock import MagicMock, patch + +import pytest + +from production_analytics.service.postgres_material import ( + PostgresMaterialSnapshotRepository, + PostgresSettings, +) + +START = datetime(2026, 9, 4, 10, tzinfo=UTC) +ORDER = " K 7-123'; -- " +SETTINGS = PostgresSettings('localhost', 5432, 'analytics', 'reader', 'secret') + + +@pytest.fixture +def database(): + # Execute the actual portable SELECT in SQLite; only adapt driver placeholders + # and datetime transport. This tests SQL semantics without a live PostgreSQL server. + with sqlite3.connect(':memory:') as database: + database.execute('''CREATE TABLE material_consumption_snapshots ( + timestamp TEXT, calculation_id TEXT, machine_id TEXT, production_order TEXT, + run_id TEXT, consumption_kg REAL + )''') + rows = [ + (START.isoformat(), 'calc', 'machine', ORDER, 'run1', 1000), + (START.replace(minute=10).isoformat(), 'calc', 'machine', ORDER, 'run2', 1100), + (START.replace(minute=4).isoformat(), 'other', 'machine', ORDER, 'run', 9999), + (START.replace(minute=4).isoformat(), 'calc', 'other', ORDER, 'run', 9999), + (START.replace(minute=4).isoformat(), 'calc', 'machine', ORDER.strip(), 'run', 9999), + (START.replace(minute=4).isoformat(), 'calc', 'machine', ORDER + '-other', 'run', 9999), + ] + database.executemany( + 'INSERT INTO material_consumption_snapshots VALUES (?,?,?,?,?,?)', rows, + ) + yield database + + +@pytest.mark.parametrize('minute, expected_minute, consumption', [ + (-1, None, None), (0, 0, 1000), (5, 0, 1000), (10, 10, 1100), (15, 10, 1100), +]) +def test_aligned_lookup_sql(database, minute, expected_minute, consumption): + cutoff = START.replace(minute=minute) if minute >= 0 else START.replace(hour=9, minute=59) + driver = MagicMock() + connection = driver.connect.return_value.__enter__.return_value + + def execute(sql, parameters): + assert parameters == ('calc', 'machine', ORDER, cutoff) + assert sql.count('%s') == 4 + assert ORDER not in sql + assert 'ORDER BY timestamp DESC' in sql + assert 'LIMIT 1' in sql + row = database.execute( + sql.replace('%s', '?'), (*parameters[:3], parameters[3].isoformat()), + ).fetchone() + if row is not None: + row = (datetime.fromisoformat(row[0]), *row[1:]) + return MagicMock(fetchone=MagicMock(return_value=row)) + + connection.execute.side_effect = execute + with patch.dict(sys.modules, psycopg=driver): + result = PostgresMaterialSnapshotRepository(SETTINGS).latest_at_or_before( + calculation_id='calc', machine_id='machine', production_order=ORDER, timestamp=cutoff, + ) + if expected_minute is None: + assert result is None + else: + assert result.timestamp == START.replace(minute=expected_minute) + assert result.timestamp <= cutoff + assert result.consumption_kg == consumption + assert (result.calculation_id, result.machine_id, result.production_order) == ( + 'calc', 'machine', ORDER, + ) + assert result.run_id == ('run1' if expected_minute == 0 else 'run2') + connection.execute.assert_called_once() + driver.connect.assert_called_once_with( + host='localhost', port=5432, dbname='analytics', user='reader', password='secret', + connect_timeout=10, options='-c statement_timeout=10000', + ) + driver.connect.return_value.__exit__.assert_called_once_with(None, None, None) + + +def test_naive_cutoff_rejected_before_connection(): + driver = MagicMock() + with patch.dict(sys.modules, psycopg=driver), pytest.raises(ValueError, match='timezone-aware'): + PostgresMaterialSnapshotRepository(SETTINGS).latest_at_or_before( + calculation_id='calc', machine_id='machine', production_order=ORDER, + timestamp=START.replace(tzinfo=None), + ) + driver.connect.assert_not_called() + + +def test_database_failure_propagates(): + driver = MagicMock() + driver.connect.side_effect = RuntimeError('connection unavailable') + with patch.dict(sys.modules, psycopg=driver), pytest.raises(RuntimeError): + PostgresMaterialSnapshotRepository(SETTINGS).latest_at_or_before( + calculation_id='calc', machine_id='machine', production_order=ORDER, timestamp=START, + )