Persist material consumption snapshots

This commit is contained in:
2026-09-05 11:35:25 +02:00
parent 7c029759f5
commit ae13c3c35a
9 changed files with 303 additions and 17 deletions
+43 -6
View File
@@ -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.
+13
View File
@@ -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);
+5 -3
View File
@@ -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
+1 -1
View File
@@ -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 = [
@@ -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} "
@@ -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,
)
@@ -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),
)
+20 -6
View File
@@ -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()
+128
View File
@@ -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