Add feedback-aligned material efficiency KPI
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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()
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
@@ -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,
|
||||
)
|
||||
Reference in New Issue
Block a user