From ae13c3c35a0b28d72e7a2cffc763b72d449b6ad0 Mon Sep 17 00:00:00 2001 From: Martin Tazl Date: Sat, 5 Sep 2026 11:35:25 +0200 Subject: [PATCH] Persist material consumption snapshots --- README.md | 49 ++++++- db/schema.sql | 13 ++ docs/roadmap.md | 8 +- pyproject.toml | 2 +- .../service/material_runner.py | 17 +++ .../service/material_runtime.py | 9 +- .../service/postgres_material.py | 68 ++++++++++ tests/test_material_runner.py | 26 +++- tests/test_postgres_material.py | 128 ++++++++++++++++++ 9 files changed, 303 insertions(+), 17 deletions(-) create mode 100644 db/schema.sql create mode 100644 src/production_analytics/service/postgres_material.py create mode 100644 tests/test_postgres_material.py diff --git a/README.md b/README.md index e197580..821fdc2 100644 --- a/README.md +++ b/README.md @@ -9,8 +9,8 @@ production-order context, and calculation state. This repository includes an ENLYZE exploration CLI, a production-run/timeseries gateway, and a configured continuous material polling runner with atomic JSON -checkpoints. The live material-consumption MVP remains in progress: Grafana-oriented -derived-total persistence/exposure and database migrations are still pending. +checkpoints and PostgreSQL cumulative material-consumption snapshots for Grafana. +The live material-consumption MVP remains in progress. There are no HTTP endpoints. ## Intended flow @@ -174,8 +174,46 @@ calculation parameters or running a different calculation for the same machine and order: checkpoints are not namespaced by calculation id/version. These files are integration checkpoints, not a Grafana metric store. -This runner exposes no Grafana/TimescaleDB metric yet. Grafana-oriented derived-total -persistence/exposure remains the next milestone; the MVP is not complete. +Each successful non-empty poll writes one cumulative snapshot to PostgreSQL before +printing its result. JSON remains the restart checkpoint; PostgreSQL stores derived +time-series snapshots only. ENLYZE remains the raw-data source of truth. Database +write failures emit a concise error and polling continues with the next cycle; +missed snapshots are not retried or backfilled and JSON state is not rolled back. +The next successful snapshot includes the continuing cumulative total. + +A new order's first snapshot is **not guaranteed to be zero**: the service processes +available samples from the run start before returning, so its first total may +already be non-zero. A returned zero is stored normally. Totals continue across +runs for the same order; existing checkpoints for a previously seen order resume. + +### PostgreSQL setup + +Set all five runtime settings: `POSTGRES_HOST`, `POSTGRES_PORT`, `POSTGRES_DB`, +`POSTGRES_USER`, and `POSTGRES_PASSWORD`, using exported environment variables or +the existing secrets file (file values override the environment). The runner does +not automatically load `.env`. Missing or invalid settings fail at startup; +connectivity and schema errors are reported during polling. For the existing +Compose instance, use host `localhost`, the published port (default `5432`), and +`production_analytics` for both database and user. Set the same password in +Compose's `.env` and the runner environment/secrets file. + +Start the database and apply the repeatable schema from the repository root: + +```bash +docker compose up -d timescaledb +docker compose exec -T timescaledb psql -U production_analytics -d production_analytics \ + -v ON_ERROR_STOP=1 < db/schema.sql +``` + +Wait until PostgreSQL is ready before applying the schema. Configure Grafana's +PostgreSQL data source to query `material_consumption_snapshots`, selecting +`timestamp` as time and `consumption_kg` as the cumulative value, filtered by +`calculation_id`, `machine_id`, and `production_order`. These are ordinary +PostgreSQL tables/indexes; no Timescale-specific features or hypertables are used. +The primary key deduplicates calculation/machine/order/timestamp (first write wins). +A short transaction opens and closes one synchronous connection per snapshot, +with 10-second connection and statement timeouts. Snapshot timestamps are poll +start times, not source-sample timestamps. Schema application is manual. Bento 1 bentonite source and gate selection remain intentionally undefined pending process validation; the generic integration core is unchanged and reusable. @@ -215,8 +253,7 @@ completed peaks, including a legitimate approximately 1195 kg cycle. Run available. Neither capture should be added to version control or automated tests. -`docker compose up -d timescaledb` is an optional local database design for a -future persistence milestone. It is not required for the bootstrap tests. +The normal unit test suite mocks PostgreSQL and requires no live database. See [PROJECT_KNOWLEDGE.md](PROJECT_KNOWLEDGE.md) for durable project context and [docs/roadmap.md](docs/roadmap.md) for the implementation sequence. diff --git a/db/schema.sql b/db/schema.sql new file mode 100644 index 0000000..0f3d5a2 --- /dev/null +++ b/db/schema.sql @@ -0,0 +1,13 @@ +CREATE TABLE IF NOT EXISTS material_consumption_snapshots ( + timestamp timestamptz NOT NULL, + calculation_id text NOT NULL, + machine_id text NOT NULL, + production_order text NOT NULL, + run_id text NOT NULL, + consumption_kg double precision NOT NULL, + PRIMARY KEY (calculation_id, machine_id, production_order, timestamp) +); + +-- Machine time series across orders; the primary key covers a specific order. +CREATE INDEX IF NOT EXISTS material_consumption_snapshots_machine_time_idx + ON material_consumption_snapshots (calculation_id, machine_id, timestamp); diff --git a/docs/roadmap.md b/docs/roadmap.md index f4e085c..ed89f05 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -24,9 +24,11 @@ Runtime options select the polling interval, state directory (default: ignored `data/state/material/`), and existing secrets-file handling. Only one writer may poll each machine/order state key. -The runner exposes no Grafana/TimescaleDB metric yet and does not complete the -MVP. Grafana-oriented derived-total persistence/exposure is the next milestone, -with the TimescaleDB persistence foundation below. No HTTP endpoints, application +The runner now persists cumulative material-consumption snapshots for Grafana in +PostgreSQL using `db/schema.sql`. JSON remains the restart checkpoint; PostgreSQL +stores derived time-series snapshots only. No Timescale-specific features are used. +A new order may first appear with a non-zero total after initial sample integration. +Grafana dashboard validation remains pending. No HTTP endpoints, application containers, daemonization, or schedulers have been added. Detect the currently active Production Run and production order, continuously diff --git a/pyproject.toml b/pyproject.toml index 62a5a1a..17ca6a9 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -10,7 +10,7 @@ readme = "README.md" requires-python = ">=3.11" license = { text = "Proprietary" } authors = [{ name = "Production Analytics Team" }] -dependencies = ["PyYAML>=6.0"] +dependencies = ["PyYAML>=6.0", "psycopg[binary]>=3.2,<4"] [project.optional-dependencies] dev = [ diff --git a/src/production_analytics/service/material_runner.py b/src/production_analytics/service/material_runner.py index c13cf55..6335102 100644 --- a/src/production_analytics/service/material_runner.py +++ b/src/production_analytics/service/material_runner.py @@ -9,6 +9,7 @@ from typing import TextIO from production_analytics.calculations.config import finite_number from production_analytics.enlyze.exploration import ConfigurationError, ExplorationError from production_analytics.service.material_polling import MaterialPollingService +from production_analytics.service.postgres_material import MaterialSnapshotWriter class MaterialStateError(RuntimeError): @@ -22,12 +23,15 @@ def utc_now() -> datetime: class MaterialPollingRunner: def __init__( self, service: MaterialPollingService, *, machine_id: str, + calculation_id: str, snapshot_writer: MaterialSnapshotWriter, poll_interval_seconds: float, clock: Callable[[], datetime] = utc_now, sleep: Callable[[float], None] = time.sleep, stdout: TextIO | None = None, stderr: TextIO | None = None, ) -> None: self.service = service self.machine_id = machine_id + self.calculation_id = calculation_id + self.snapshot_writer = snapshot_writer self.interval = finite_number(poll_interval_seconds, "poll interval", positive=True) self.clock = clock self.sleep = sleep @@ -59,6 +63,19 @@ class MaterialPollingRunner: if result is None: detail = "no open Production Run / no eligible polling window" else: + try: + self.snapshot_writer.write( + timestamp=now, calculation_id=self.calculation_id, + machine_id=self.machine_id, + production_order=result.run.production_order, + run_id=result.run.uuid, + consumption_kg=result.state.cumulative_consumption_kg, + ) + except Exception: + print( + f"{prefix} error PostgreSQL snapshot write failed", + file=self.stderr, flush=True, + ) detail = ( f"production_order={result.run.production_order!r} " f"consumption_kg={result.state.cumulative_consumption_kg:.9f} " diff --git a/src/production_analytics/service/material_runtime.py b/src/production_analytics/service/material_runtime.py index 10b0b7f..1fbd61a 100644 --- a/src/production_analytics/service/material_runtime.py +++ b/src/production_analytics/service/material_runtime.py @@ -20,6 +20,10 @@ from production_analytics.service.material_polling import ( ) from production_analytics.service.material_runner import MaterialPollingRunner, MaterialStateError from production_analytics.service.material_state_store import JsonMaterialStateStore +from production_analytics.service.postgres_material import ( + PostgresMaterialSnapshotWriter, + PostgresSettings, +) class _ReportingStateStore: @@ -64,6 +68,7 @@ def build_material_runner( raise ConfigurationError( f"State directory startup check failed: {type(exc).__name__}" ) from exc + postgres_settings = PostgresSettings.from_environment(environment) service = MaterialPollingService( gateway=EnlyzeApiGateway(ExplorationClient(settings)), state_store=_ReportingStateStore(JsonMaterialStateStore(state_directory)), @@ -74,5 +79,7 @@ def build_material_runner( max_sample_gap_seconds=calculation.max_sample_gap_seconds, ) return MaterialPollingRunner( - service, machine_id=calculation.machine_ref, poll_interval_seconds=poll_interval_seconds, + service, calculation_id=calculation.id, + snapshot_writer=PostgresMaterialSnapshotWriter(postgres_settings), + machine_id=calculation.machine_ref, poll_interval_seconds=poll_interval_seconds, ) diff --git a/src/production_analytics/service/postgres_material.py b/src/production_analytics/service/postgres_material.py new file mode 100644 index 0000000..64b6568 --- /dev/null +++ b/src/production_analytics/service/postgres_material.py @@ -0,0 +1,68 @@ +"""PostgreSQL settings and derived material snapshot inserts.""" + +from collections.abc import Mapping +from dataclasses import dataclass, field +from datetime import datetime +from typing import Protocol + +from production_analytics.enlyze.exploration import ConfigurationError + + +@dataclass(frozen=True, slots=True) +class PostgresSettings: + host: str + port: int + dbname: str + user: str + password: str = field(repr=False) + + @classmethod + def from_environment(cls, environment: Mapping[str, str]) -> "PostgresSettings": + values = {} + for key in ('HOST', 'PORT', 'DB', 'USER', 'PASSWORD'): + name = f'POSTGRES_{key}' + value = environment.get(name) + if not isinstance(value, str) or not value.strip() or '\x00' in value: + raise ConfigurationError(f'{name} must be a non-empty value without NUL bytes') + values[key] = value + try: + port = int(values['PORT']) + except ValueError: + raise ConfigurationError('POSTGRES_PORT must be an integer from 1 to 65535') from None + if not 1 <= port <= 65535: + raise ConfigurationError('POSTGRES_PORT must be an integer from 1 to 65535') + return cls(values['HOST'], port, values['DB'], values['USER'], values['PASSWORD']) + + +class MaterialSnapshotWriter(Protocol): + def write( + self, *, timestamp: datetime, calculation_id: str, machine_id: str, + production_order: str, run_id: str, consumption_kg: float, + ) -> None: ... + + +class PostgresMaterialSnapshotWriter: + def __init__(self, settings: PostgresSettings) -> None: + self.settings = settings + + def write( + self, *, timestamp: datetime, calculation_id: str, machine_id: str, + production_order: str, run_id: str, consumption_kg: float, + ) -> None: + import psycopg + + # One short transaction per cycle; context exit commits or rolls back and closes. + # A fresh connection lets the next cycle recover after a database outage. + with psycopg.connect( + host=self.settings.host, port=self.settings.port, dbname=self.settings.dbname, + user=self.settings.user, password=self.settings.password, connect_timeout=10, + options='-c statement_timeout=10000', + ) as connection: + connection.execute( + """INSERT INTO material_consumption_snapshots + (timestamp, calculation_id, machine_id, production_order, run_id, consumption_kg) + VALUES (%s, %s, %s, %s, %s, %s) + ON CONFLICT (calculation_id, machine_id, production_order, timestamp) + DO NOTHING""", + (timestamp, calculation_id, machine_id, production_order, run_id, consumption_kg), + ) diff --git a/tests/test_material_runner.py b/tests/test_material_runner.py index 65085cb..713b0e4 100644 --- a/tests/test_material_runner.py +++ b/tests/test_material_runner.py @@ -112,7 +112,8 @@ def test_repeated_sequential_polls_and_long_cycle(): output = io.StringIO() service = Mock(poll_once=Mock(side_effect=poll)) MaterialPollingRunner( - service, machine_id='machine', poll_interval_seconds=10, + service, calculation_id='test-calculation', snapshot_writer=Mock(), + machine_id='machine', poll_interval_seconds=10, clock=lambda: NOW + timedelta(seconds=elapsed), sleep=sleep, stdout=output, ).run() assert events == ['poll-start', 'poll-end', ('sleep', 10)] * 2 @@ -128,7 +129,8 @@ def test_error_then_success_preserves_opaque_order_and_reports_totals(): output, errors = io.StringIO(), io.StringIO() sleep = Mock(side_effect=[None, KeyboardInterrupt]) MaterialPollingRunner( - service, machine_id='machine', poll_interval_seconds=5, clock=lambda: NOW, + service, calculation_id='test-calculation', snapshot_writer=Mock(), + machine_id='machine', poll_interval_seconds=5, clock=lambda: NOW, sleep=sleep, stdout=output, stderr=errors, ).run() assert service.poll_once.call_count == 2 @@ -143,7 +145,8 @@ def test_error_then_success_preserves_opaque_order_and_reports_totals(): def test_interrupt_during_poll(): sleep, output = Mock(), io.StringIO() MaterialPollingRunner( - Mock(poll_once=Mock(side_effect=KeyboardInterrupt)), machine_id='m', + Mock(poll_once=Mock(side_effect=KeyboardInterrupt)), + calculation_id='test-calculation', snapshot_writer=Mock(), machine_id='m', poll_interval_seconds=1, sleep=sleep, stdout=output, ).run() sleep.assert_not_called() @@ -154,6 +157,7 @@ def test_cycle_configuration_failure_is_fatal(): sleep = Mock() runner = MaterialPollingRunner( Mock(poll_once=Mock(side_effect=ConfigurationError('bad settings'))), + calculation_id='test-calculation', snapshot_writer=Mock(), machine_id='m', poll_interval_seconds=1, sleep=sleep, ) with pytest.raises(ConfigurationError): @@ -164,7 +168,10 @@ def test_cycle_configuration_failure_is_fatal(): @pytest.mark.parametrize('interval', [0, -1, float('nan'), float('inf'), True]) def test_invalid_interval(interval): with pytest.raises(CalculationConfigError): - MaterialPollingRunner(Mock(), machine_id='m', poll_interval_seconds=interval) + MaterialPollingRunner( + Mock(), machine_id='m', calculation_id='test', snapshot_writer=Mock(), + poll_interval_seconds=interval, + ) @pytest.mark.parametrize('operation', ['load', 'save']) @@ -182,7 +189,10 @@ def test_state_error_classification(operation): def test_runtime_wiring(tmp_path, monkeypatch): monkeypatch.setenv('ENLYZE_BASE_URL', 'https://example.invalid/api/') secrets = tmp_path / 'secret.env' - secrets.write_text('ENLYZE_API_KEY=secret-from-file\n') + for key, value in dict(HOST='localhost', PORT='5432', DB='analytics', + USER='user', PASSWORD='environment-password').items(): + monkeypatch.setenv(f'POSTGRES_{key}', value) + secrets.write_text('ENLYZE_API_KEY=secret-from-file\nPOSTGRES_PASSWORD=file-password\n') with patch('production_analytics.service.material_runtime.MaterialPollingService') as service: runner = build_material_runner( load_material_calculation(EXAMPLE), poll_interval_seconds=7, @@ -198,6 +208,9 @@ def test_runtime_wiring(tmp_path, monkeypatch): assert kwargs['state_store'].store._directory == tmp_path / 'state' assert runner.service is service.return_value assert runner.interval == 7 + assert runner.calculation_id == 'k7-fiber-consumption' + assert runner.snapshot_writer.settings.password == 'file-password' + assert runner.snapshot_writer.settings.dbname == 'analytics' def test_cli_bad_config(tmp_path, capsys): @@ -257,7 +270,8 @@ def test_failed_cycle_keeps_checkpoint_and_recovers(tmp_path): raise KeyboardInterrupt MaterialPollingRunner( - service, machine_id='machine', poll_interval_seconds=1, + service, calculation_id='test-calculation', snapshot_writer=Mock(), + machine_id='machine', poll_interval_seconds=1, clock=lambda: NOW + timedelta(seconds=10), sleep=sleep, stdout=io.StringIO(), stderr=io.StringIO(), ).run() diff --git a/tests/test_postgres_material.py b/tests/test_postgres_material.py new file mode 100644 index 0000000..f8e5fd1 --- /dev/null +++ b/tests/test_postgres_material.py @@ -0,0 +1,128 @@ +import io +import sys +from datetime import UTC, datetime +from unittest.mock import MagicMock, Mock, patch + +import pytest + +from production_analytics.calculations.material_consumption import MaterialIntegrationState +from production_analytics.enlyze.exploration import ConfigurationError +from production_analytics.enlyze.gateway import EnlyzeProductionRun +from production_analytics.service.material_polling import MaterialPollResult +from production_analytics.service.material_runner import MaterialPollingRunner +from production_analytics.service.postgres_material import ( + PostgresMaterialSnapshotWriter, + PostgresSettings, +) + +ENV = dict(POSTGRES_HOST='localhost', POSTGRES_PORT='5432', POSTGRES_DB='analytics', + POSTGRES_USER='writer', POSTGRES_PASSWORD='secret-password') +NOW = datetime(2026, 1, 1, tzinfo=UTC) + + +@pytest.mark.parametrize('key', ENV) +@pytest.mark.parametrize('value', [None, '', ' ', '\x00']) +def test_required_settings(key, value): + environment = dict(ENV) + if value is None: + del environment[key] + else: + environment[key] = value + with pytest.raises(ConfigurationError, match=key) as error: + PostgresSettings.from_environment(environment) + assert 'secret-password' not in str(error.value) + + +@pytest.mark.parametrize('port', ['0', '-1', '65536', '1.5', 'secret-password']) +def test_invalid_port(port): + with pytest.raises(ConfigurationError, match='POSTGRES_PORT') as error: + PostgresSettings.from_environment(dict(ENV, POSTGRES_PORT=port)) + assert str(error.value) == 'POSTGRES_PORT must be an integer from 1 to 65535' + + +def test_valid_settings_and_password_repr(): + settings = PostgresSettings.from_environment(ENV) + assert (settings.host, settings.port, settings.dbname, settings.user) == ( + 'localhost', 5432, 'analytics', 'writer', + ) + assert settings.password == 'secret-password' + assert settings.password not in repr(settings) + + +def test_parameterized_insert_and_connection_lifecycle(): + driver = MagicMock() + connection = driver.connect.return_value.__enter__.return_value + writer = PostgresMaterialSnapshotWriter(PostgresSettings.from_environment(ENV)) + fields = dict(timestamp=NOW, calculation_id='calc', machine_id='machine', + production_order=" 00842'; DROP TABLE x; -- ", run_id='run', consumption_kg=12.5) + with patch.dict(sys.modules, psycopg=driver): + writer.write(**fields) + sql, parameters = connection.execute.call_args.args + assert sql.count('%s') == 6 + assert fields['production_order'] not in sql + assert parameters == tuple(fields.values()) + assert 'ON CONFLICT' in sql + connection.execute.assert_called_once() + driver.connect.assert_called_once_with( + host='localhost', port=5432, dbname='analytics', user='writer', + password='secret-password', connect_timeout=10, options='-c statement_timeout=10000', + ) + driver.connect.return_value.__exit__.assert_called_once_with(None, None, None) + + +def test_writer_failure_propagates_and_next_write_reconnects(): + driver = MagicMock() + driver.connect.return_value.__enter__.return_value.execute.side_effect = [ + RuntimeError('secret'), None, + ] + writer = PostgresMaterialSnapshotWriter(PostgresSettings.from_environment(ENV)) + fields = dict(timestamp=NOW, calculation_id='calc', machine_id='machine', + production_order='order', run_id='run', consumption_kg=0) + with patch.dict(sys.modules, psycopg=driver): + with pytest.raises(RuntimeError): + writer.write(**fields) + writer.write(**fields) + assert driver.connect.call_count == 2 + assert driver.connect.return_value.__exit__.call_count == 2 + + +@pytest.mark.parametrize('consumption', [0, 12.5]) +def test_runner_snapshot_fields_and_none(consumption): + result = MaterialPollResult( + EnlyzeProductionRun('run', 'machine', None, ' 00842 ', NOW, None), + MaterialIntegrationState(consumption, 20), + ) + writer = Mock() + output = io.StringIO() + MaterialPollingRunner( + Mock(poll_once=Mock(side_effect=[None, result])), machine_id='machine', + calculation_id='configured-calculation', snapshot_writer=writer, + poll_interval_seconds=10, clock=lambda: NOW, + sleep=Mock(side_effect=[None, KeyboardInterrupt]), stdout=output, + ).run() + writer.write.assert_called_once_with( + timestamp=NOW, calculation_id='configured-calculation', machine_id='machine', + production_order=' 00842 ', run_id='run', consumption_kg=consumption, + ) + assert 'no open Production Run' in output.getvalue() + assert f'consumption_kg={consumption:.9f}' in output.getvalue() + + +def test_snapshot_failure_reports_safely_and_continues(): + result = MaterialPollResult( + EnlyzeProductionRun('run', 'machine', None, 'order', NOW, None), + MaterialIntegrationState(12.5, 20), + ) + service = Mock(poll_once=Mock(return_value=result)) + writer = Mock(write=Mock(side_effect=[RuntimeError('secret-password arbitrary SQL'), None])) + output, errors = io.StringIO(), io.StringIO() + MaterialPollingRunner( + service, machine_id='machine', calculation_id='calc', snapshot_writer=writer, + poll_interval_seconds=10, clock=lambda: NOW, + sleep=Mock(side_effect=[None, KeyboardInterrupt]), stdout=output, stderr=errors, + ).run() + assert service.poll_once.call_count == writer.write.call_count == 2 + assert errors.getvalue().count('PostgreSQL snapshot write failed') == 1 + assert 'secret-password' not in errors.getvalue() + assert 'arbitrary SQL' not in errors.getvalue() + assert output.getvalue().count('consumption_kg=12.500000000') == 2