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" def test_bento_uses_generic_efficiency_and_distinct_output_id(runtime_files, storage): _, postgres, erp = runtime_files runner = build_material_efficiency_runner( Path('config/bento1-material-efficiency.yaml'), secrets_file=postgres, erp_secrets_file=erp, ) from production_analytics.service.material_efficiency import MaterialEfficiencyService assert type(runner.service) is MaterialEfficiencyService service = runner.service status = CurrentWorkplaceStatus( 'Bento 1', '00123', None, None, NOW, None, 800, None, None, None, ) service.erp = Mock(get_current_workplace_status=Mock(return_value=status)) service.materials = Mock(latest_at_or_before=Mock(return_value=MaterialConsumptionSnapshot( NOW, 'bento1-fresh-bentonite-consumption', service.machine_id, 'Bento 1-00123', 'run', 2400, ))) snapshot = service.evaluate_current() assert snapshot.calculation_id == 'bento1-fresh-bentonite-efficiency' assert snapshot.material_consumption_g_per_m2 == 3000 assert snapshot.nominal_width_m is None service.materials.latest_at_or_before.assert_called_once_with( calculation_id='bento1-fresh-bentonite-consumption', machine_id=service.machine_id, production_order='Bento 1-00123', timestamp=NOW, ) runner.writer.write(snapshot) row = storage[0].execute( 'SELECT calculation_id, material_consumption_g_per_m2 FROM material_efficiency_snapshots' ).fetchone() assert row == ('bento1-fresh-bentonite-efficiency', 3000) @pytest.mark.parametrize('value', ['', None, 12, 'bad\x00id']) def test_invalid_output_calculation_id(runtime_files, value): config, postgres, erp = runtime_files document = yaml.safe_load(config.read_text()) | {'output_calculation_id': value} config.write_text(yaml.safe_dump(document)) with pytest.raises(CalculationConfigError, match='output_calculation_id'): build_material_efficiency_runner(config, secrets_file=postgres, erp_secrets_file=erp)