From 6ddb76fd0352dfe0e7134061170d01fd6c4b556a Mon Sep 17 00:00:00 2001 From: Martin Tazl Date: Sat, 5 Sep 2026 07:55:40 +0200 Subject: [PATCH] Add continuous material polling runner --- .gitignore | 3 + README.md | 72 ++++- config/k7-material-consumption.yaml | 14 + docs/roadmap.md | 17 +- pyproject.toml | 2 +- .../calculations/config.py | 110 +++++++ src/production_analytics/cli/__main__.py | 38 ++- .../service/material_runner.py | 71 +++++ .../service/material_runtime.py | 78 +++++ tests/test_material_runner.py | 277 ++++++++++++++++++ 10 files changed, 672 insertions(+), 10 deletions(-) create mode 100644 config/k7-material-consumption.yaml create mode 100644 src/production_analytics/calculations/config.py create mode 100644 src/production_analytics/service/material_runner.py create mode 100644 src/production_analytics/service/material_runtime.py create mode 100644 tests/test_material_runner.py diff --git a/.gitignore b/.gitignore index 81e4e89..ff02673 100644 --- a/.gitignore +++ b/.gitignore @@ -29,3 +29,6 @@ secrets/* .DS_Store activate.sh + +# Live polling checkpoints +data/state/ diff --git a/README.md b/README.md index 9005610..e197580 100644 --- a/README.md +++ b/README.md @@ -8,8 +8,10 @@ production-order context, and calculation state. ## Status This repository includes an ENLYZE exploration CLI, a production-run/timeseries -gateway, and single-cycle live material polling with JSON state persistence. -Database migrations, a continuous polling runner, and HTTP endpoints remain pending. +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. +There are no HTTP endpoints. ## Intended flow @@ -89,8 +91,8 @@ Run `python scripts/validate_k7_material_consumption.py` from the repository roo with the ignored `k7-00842-throughput.raw.json` and `k7-00842-speed.raw.json` captures in `data/raw/enlyze/`. The utility joins common non-null timestamps without filling values and rejects unordered or duplicate capture records. -K7-specific signal UUIDs are confined to the utility: `Stundenleistung Anlage` -is gated by `Geschwindigkeit Gesamtanlage > 0.5 m/min`. With the explicit +K7-specific signal UUIDs are declared in the runner configuration and validation +utility: `Stundenleistung Anlage` is gated by `Geschwindigkeit Gesamtanlage > 0.5 m/min`. With the explicit 20-second validation gap limit (`--max-sample-gap-seconds` to override), 2,672 common samples yield **5.205555556 h** and **5,180.127150811 kg** for run 00842, matching the previous manual calculation's rounded results. Synthetic @@ -115,6 +117,68 @@ Files from the earlier uncommitted sanitized-filename prototype are not loaded under the new names. The gateway rejects ambiguous open runs, naive windows, and missing columns or malformed records instead of silently skipping them. +## Continuous material polling + +The committed [K7 configuration](config/k7-material-consumption.yaml) declares +`k7-fiber-consumption`, a `material_consumption` calculation at version `"1"`. +It selects `Stundenleistung Anlage` (kg/h) as the rate and `Geschwindigkeit +Gesamtanlage` (m/min) as the gate, integrating only when the gate is **strictly +above 0.5**, with a maximum sample gap of 20 seconds. `Anlage läuft` is not the +primary gate. The output is `material_consumption` in `kg`, grouped by +`production_order`; production orders remain opaque strings, including leading +zeros and whitespace. + +Install the declared dependencies in the repository venv, then run from the +repository root: + +```bash +.venv/bin/python -m pip install -e '.[dev]' +.venv/bin/python -m production_analytics.cli run material-poll \ + --config config/k7-material-consumption.yaml \ + --calculation-id k7-fiber-consumption \ + --poll-interval-seconds 10 \ + --state-directory data/state/material \ + --secrets-file secrets/enlyze.env +``` + +The installed `production-analytics run material-poll` command is equivalent. +The interval, state directory, and secrets path above are defaults. Paths are +relative to the current working directory. Existing `ENLYZE_BASE_URL`, +`ENLYZE_API_KEY`, and `ENLYZE_HTTP_TIMEOUT_SECONDS` handling is reused; values +in the secrets file override environment values. Keys are never CLI arguments. + +Keep calculation fields in YAML and technical runtime settings in CLI options. +The loader requires every field shown in the K7 file, rejects unknown fields, +duplicate keys/ids, unsupported types/versions/outputs, and invalid or non-finite +numbers. Numeric fields must be YAML numbers, and `version` must be quoted. +Multiple material calculations may share a file, but then `--calculation-id` is +required. All entries are validated before selection. The unrelated +`config/calculations.example.yaml` remains illustrative and is not executable +by this material-only runner. + +The foreground runner polls once immediately using UTC, then sleeps for the +configured positive interval after each completed cycle, including failures. +A slow poll delays the next cycle; there is no overlap or catch-up scheduling. +Ctrl-C/SIGINT exits cleanly. Successful cycles print machine, opaque order, +cumulative kg, integrated running seconds, and cycle-start timestamp. No open +run (or no eligible window) is normal and visible. Errors on stderr identify +state loading/saving or ENLYZE/gateway/polling and the exception class; arbitrary +exception text and response bodies are withheld to avoid leaking credentials. +Startup/configuration failures return non-zero; cycle failures retry after the +delay without replacing checkpoints with empty state. + +Checkpoints live in ignored `data/state/material/` by default, with one hashed +JSON filename per machine/order pair. Run **one polling writer per state key**; +there is no multi-process locking. Use a separate state directory when changing +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. +Bento 1 bentonite source and gate selection remain intentionally undefined pending +process validation; the generic integration core is unchanged and reusable. + ## Peak-cycle detection `PeakCycleDetector` is a pure calculation-domain component for roll length, diff --git a/config/k7-material-consumption.yaml b/config/k7-material-consumption.yaml new file mode 100644 index 0000000..1a0bcb7 --- /dev/null +++ b/config/k7-material-consumption.yaml @@ -0,0 +1,14 @@ +# Validated K7: Stundenleistung Anlage (kg/h), gated by +# Geschwindigkeit Gesamtanlage (m/min) > 0.5. Not Anlage läuft. +calculations: + - id: k7-fiber-consumption + type: material_consumption + version: "1" + machine_ref: c220f95c-a65e-4cb7-99b7-0626d6c7508c + rate_signal_ref: c9d06af5-f6d6-4ede-b6c4-5a98bac77129 + gate_signal_ref: 823867bb-f5d2-40eb-b875-657155addfd0 + gate_threshold: 0.5 + max_sample_gap_seconds: 20.0 + output_metric: material_consumption + output_unit: kg + group_by: production_order diff --git a/docs/roadmap.md b/docs/roadmap.md index cfdd0c3..f4e085c 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -15,10 +15,19 @@ Continue to record sanitized fixtures for newly verified semantics. ## 2. Live material-consumption MVP (in progress) -The integration engine, ENLYZE gateway, single-cycle polling service, and atomic -JSON checkpoint store are implemented and covered by synthetic tests. Next: -wire a configured continuous polling runner and expose/store derived totals -for Grafana, followed by the TimescaleDB persistence foundation below. +The integration engine, ENLYZE gateway, single-cycle polling service, atomic +JSON checkpoint store, and configured continuous foreground polling runner are +implemented and covered by synthetic tests. Configure the validated K7 instance +in `config/k7-material-consumption.yaml`; start it with +`production-analytics run material-poll --config config/k7-material-consumption.yaml`. +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 +containers, daemonization, or schedulers have been added. Detect the currently active Production Run and production order, continuously ingest new ENLYZE samples, maintain persistent incremental integration state, diff --git a/pyproject.toml b/pyproject.toml index 57fabe9..62a5a1a 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 = [] +dependencies = ["PyYAML>=6.0"] [project.optional-dependencies] dev = [ diff --git a/src/production_analytics/calculations/config.py b/src/production_analytics/calculations/config.py new file mode 100644 index 0000000..e209f89 --- /dev/null +++ b/src/production_analytics/calculations/config.py @@ -0,0 +1,110 @@ +"""Strict configuration for executable material-consumption instances.""" + +from dataclasses import dataclass, fields +from math import isfinite +from pathlib import Path + +import yaml + + +class CalculationConfigError(ValueError): + """An operator-readable configuration failure, without input values.""" + + +class _UniqueLoader(yaml.SafeLoader): + pass + + +def _mapping(loader, node): + result = {} + for key_node, value_node in node.value: + key = loader.construct_object(key_node) + if not isinstance(key, str) or key in result: + raise CalculationConfigError("YAML mapping keys must be unique strings") + result[key] = loader.construct_object(value_node) + return result + + +_UniqueLoader.add_constructor(yaml.resolver.BaseResolver.DEFAULT_MAPPING_TAG, _mapping) + + +def finite_number(value: object, name: str, *, positive: bool = False) -> float: + if isinstance(value, bool) or not isinstance(value, (int, float)): + raise CalculationConfigError(f"{name} must be a finite number") + try: + number = float(value) + except (OverflowError, ValueError) as exc: + raise CalculationConfigError(f"{name} must be a finite number") from exc + if not isfinite(number) or (positive and number <= 0): + requirement = "finite and greater than zero" if positive else "finite" + raise CalculationConfigError(f"{name} must be {requirement}") + return number + + +@dataclass(frozen=True, slots=True) +class MaterialCalculationConfig: + id: str + type: str + version: str + machine_ref: str + rate_signal_ref: str + gate_signal_ref: str + gate_threshold: float + max_sample_gap_seconds: float + output_metric: str + output_unit: str + group_by: str + + +def load_material_calculation( + path: str | Path, calculation_id: str | None = None, +) -> MaterialCalculationConfig: + try: + with Path(path).open(encoding="utf-8") as stream: + document = yaml.load(stream, Loader=_UniqueLoader) + except (OSError, UnicodeError) as exc: + raise CalculationConfigError("Cannot read calculation configuration file") from exc + except yaml.YAMLError as exc: + raise CalculationConfigError("Invalid YAML in calculation configuration") from exc + if not isinstance(document, dict) or set(document) != {"calculations"}: + raise CalculationConfigError("Configuration must contain only 'calculations'") + entries = document["calculations"] + if not isinstance(entries, list) or not entries: + raise CalculationConfigError("calculations must be a non-empty list") + required = {field.name for field in fields(MaterialCalculationConfig)} + instances = {} + for index, entry in enumerate(entries): + prefix = f"calculations[{index}]" + if not isinstance(entry, dict): + raise CalculationConfigError(f"{prefix} must be a mapping") + if "type" in entry and entry["type"] != "material_consumption": + raise CalculationConfigError(f"{prefix}.type: only 'material_consumption' is supported") + missing = required - entry.keys() + if missing: + raise CalculationConfigError(f"{prefix}: missing fields: {', '.join(sorted(missing))}") + if entry.keys() - required: + raise CalculationConfigError(f"{prefix}: unknown fields (check spelling)") + for name in sorted(required - {"gate_threshold", "max_sample_gap_seconds"}): + if not isinstance(entry[name], str) or not entry[name].strip(): + raise CalculationConfigError(f"{prefix}.{name} must be a non-empty string") + for name, expected in ( + ("type", "material_consumption"), ("version", "1"), + ("output_metric", "material_consumption"), ("output_unit", "kg"), + ("group_by", "production_order"), + ): + if entry[name] != expected: + raise CalculationConfigError(f"{prefix}.{name}: only {expected!r} is supported") + entry["gate_threshold"] = finite_number(entry["gate_threshold"], f"{prefix}.gate_threshold") + entry["max_sample_gap_seconds"] = finite_number( + entry["max_sample_gap_seconds"], f"{prefix}.max_sample_gap_seconds", positive=True, + ) + if entry["id"] in instances: + raise CalculationConfigError("Calculation ids must be unique") + instances[entry["id"]] = MaterialCalculationConfig(**entry) + if calculation_id is not None: + if calculation_id not in instances: + raise CalculationConfigError("Requested calculation id was not found") + return instances[calculation_id] + if len(instances) != 1: + raise CalculationConfigError("Multiple calculations: specify --calculation-id") + return next(iter(instances.values())) diff --git a/src/production_analytics/cli/__main__.py b/src/production_analytics/cli/__main__.py index 8c01de4..379c425 100644 --- a/src/production_analytics/cli/__main__.py +++ b/src/production_analytics/cli/__main__.py @@ -1,4 +1,4 @@ -"""Command-line entry point for safe, read-only API exploration.""" +"""Command-line entry point for exploration and production polling.""" from __future__ import annotations @@ -70,6 +70,14 @@ def _parser() -> argparse.ArgumentParser: default=Path("secrets/enlyze.env"), help="Local dotenv-style secret file (default: secrets/enlyze.env)", ) + run = namespaces.add_parser("run", help="Foreground production runners") + runners = run.add_subparsers(dest="command", required=True) + material = runners.add_parser("material-poll", help="Continuously poll material consumption") + material.add_argument("--config", type=Path, required=True) + material.add_argument("--calculation-id") + material.add_argument("--poll-interval-seconds", type=float, default=10.0) + material.add_argument("--state-directory", type=Path, default=Path("data/state/material")) + material.add_argument("--secrets-file", type=Path, default=Path("secrets/enlyze.env")) return parser @@ -83,8 +91,36 @@ def _write_fixture(directory: Path, name: str, payload: object) -> Path: return destination +def _run_material(args: argparse.Namespace) -> int: + from production_analytics.calculations.config import ( + CalculationConfigError, + load_material_calculation, + ) + from production_analytics.enlyze.exploration import ConfigurationError + from production_analytics.service.material_runtime import build_material_runner + + try: + calculation = load_material_calculation(args.config, args.calculation_id) + runner = build_material_runner( + calculation, poll_interval_seconds=args.poll_interval_seconds, + state_directory=args.state_directory, secrets_file=args.secrets_file, + ) + runner.run() + except (CalculationConfigError, ConfigurationError) as exc: + print(f"Config/startup error: {exc}", file=sys.stderr) + return 2 + except Exception as exc: + print(f"Startup/runtime error: {type(exc).__name__}", file=sys.stderr) + return 1 + except KeyboardInterrupt: + return 0 + return 0 + + def main(argv: Sequence[str] | None = None) -> int: args = _parser().parse_args(argv) + if args.namespace == "run": + return _run_material(args) if args.namespace != "enlyze" or args.command not in {"raw", "timeseries"}: return 2 diff --git a/src/production_analytics/service/material_runner.py b/src/production_analytics/service/material_runner.py new file mode 100644 index 0000000..c13cf55 --- /dev/null +++ b/src/production_analytics/service/material_runner.py @@ -0,0 +1,71 @@ +"""Sequential foreground polling with injectable time and output.""" + +import sys +import time +from collections.abc import Callable +from datetime import UTC, datetime +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 + + +class MaterialStateError(RuntimeError): + """State operation failed; message contains only the operation and error class.""" + + +def utc_now() -> datetime: + return datetime.now(UTC) + + +class MaterialPollingRunner: + def __init__( + self, service: MaterialPollingService, *, machine_id: str, + 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.interval = finite_number(poll_interval_seconds, "poll interval", positive=True) + self.clock = clock + self.sleep = sleep + self.stdout = stdout if stdout is not None else sys.stdout + self.stderr = stderr if stderr is not None else sys.stderr + + def run(self) -> None: + try: + while True: + now = self.clock() + if now.tzinfo is None or now.utcoffset() is None: + raise ConfigurationError("Runner clock must return timezone-aware timestamps") + now = now.astimezone(UTC) + prefix = f"{now.isoformat()} machine={self.machine_id!r}" + try: + result = self.service.poll_once(now=now) + except ConfigurationError: + raise + except Exception as exc: + if isinstance(exc, MaterialStateError): + detail = str(exc) + elif isinstance(exc, ExplorationError): + detail = f"ENLYZE/gateway: {type(exc).__name__}" + else: + detail = f"gateway/poll cycle: {type(exc).__name__}" + # Do not echo arbitrary exception messages or response bodies. + print(f"{prefix} error {detail}", file=self.stderr, flush=True) + else: + if result is None: + detail = "no open Production Run / no eligible polling window" + else: + detail = ( + f"production_order={result.run.production_order!r} " + f"consumption_kg={result.state.cumulative_consumption_kg:.9f} " + f"running_seconds={result.state.integrated_running_seconds:.3f}" + ) + print(f"{prefix} {detail}", file=self.stdout, flush=True) + # Fixed delay after completion, including errors; never catch up or overlap. + self.sleep(self.interval) + except KeyboardInterrupt: + print("Material polling stopped", file=self.stdout, flush=True) diff --git a/src/production_analytics/service/material_runtime.py b/src/production_analytics/service/material_runtime.py new file mode 100644 index 0000000..10b0b7f --- /dev/null +++ b/src/production_analytics/service/material_runtime.py @@ -0,0 +1,78 @@ +"""Compose the configured live material polling application.""" + +import math +import os +import tempfile +from pathlib import Path +from urllib.parse import urlsplit + +from production_analytics.calculations.config import MaterialCalculationConfig, finite_number +from production_analytics.enlyze.exploration import ( + ConfigurationError, + ExplorationClient, + ExplorationSettings, + load_secret_file, +) +from production_analytics.enlyze.gateway import EnlyzeApiGateway +from production_analytics.service.material_polling import ( + MaterialPollingService, + MaterialPollingState, +) +from production_analytics.service.material_runner import MaterialPollingRunner, MaterialStateError +from production_analytics.service.material_state_store import JsonMaterialStateStore + + +class _ReportingStateStore: + def __init__(self, store: JsonMaterialStateStore) -> None: + self.store = store + + def load(self, machine_id: str, production_order: str) -> MaterialPollingState | None: + try: + return self.store.load(machine_id, production_order) + except Exception as exc: + raise MaterialStateError(f"state loading: {type(exc).__name__}") from exc + + def save( + self, machine_id: str, production_order: str, state: MaterialPollingState, + ) -> None: + try: + self.store.save(machine_id, production_order, state) + except Exception as exc: + raise MaterialStateError(f"state saving: {type(exc).__name__}") from exc + + +def build_material_runner( + calculation: MaterialCalculationConfig, *, poll_interval_seconds: float, + state_directory: Path, secrets_file: Path, +) -> MaterialPollingRunner: + finite_number(poll_interval_seconds, "poll interval", positive=True) + environment = dict(os.environ) + environment.update(load_secret_file(secrets_file)) + settings = ExplorationSettings.from_environment(environment) + parsed = urlsplit(settings.base_url) + if parsed.scheme not in {"http", "https"} or not parsed.hostname or parsed.username: + raise ConfigurationError( + "ENLYZE_BASE_URL must be an HTTP(S) server URL without credentials" + ) + if not math.isfinite(settings.timeout_seconds): + raise ConfigurationError("ENLYZE_HTTP_TIMEOUT_SECONDS must be finite") + try: + state_directory.mkdir(parents=True, exist_ok=True) + with tempfile.TemporaryFile(dir=state_directory): + pass + except OSError as exc: + raise ConfigurationError( + f"State directory startup check failed: {type(exc).__name__}" + ) from exc + service = MaterialPollingService( + gateway=EnlyzeApiGateway(ExplorationClient(settings)), + state_store=_ReportingStateStore(JsonMaterialStateStore(state_directory)), + machine_id=calculation.machine_ref, + rate_variable_id=calculation.rate_signal_ref, + gate_variable_id=calculation.gate_signal_ref, + gate_threshold=calculation.gate_threshold, + max_sample_gap_seconds=calculation.max_sample_gap_seconds, + ) + return MaterialPollingRunner( + service, machine_id=calculation.machine_ref, poll_interval_seconds=poll_interval_seconds, + ) diff --git a/tests/test_material_runner.py b/tests/test_material_runner.py new file mode 100644 index 0000000..65085cb --- /dev/null +++ b/tests/test_material_runner.py @@ -0,0 +1,277 @@ +import io +from datetime import UTC, datetime, timedelta +from pathlib import Path +from unittest.mock import Mock, patch + +import pytest +import yaml + +from production_analytics.calculations.config import ( + CalculationConfigError, + load_material_calculation, +) +from production_analytics.calculations.material_consumption import MaterialIntegrationState +from production_analytics.cli.__main__ import main +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, MaterialStateError +from production_analytics.service.material_runtime import ( + _ReportingStateStore, + build_material_runner, +) + +EXAMPLE = Path(__file__).resolve().parents[1] / 'config/k7-material-consumption.yaml' +NOW = datetime(2026, 1, 1, tzinfo=UTC) + + +def write_config(tmp_path, entries): + path = tmp_path / 'calculations.yaml' + path.write_text(yaml.safe_dump({'calculations': entries})) + return path + + +def entry(): + return yaml.safe_load(EXAMPLE.read_text())['calculations'][0] + + +def test_valid_k7_config(): + config = load_material_calculation(EXAMPLE) + assert config.id == 'k7-fiber-consumption' + assert config.machine_ref == 'c220f95c-a65e-4cb7-99b7-0626d6c7508c' + assert config.rate_signal_ref == 'c9d06af5-f6d6-4ede-b6c4-5a98bac77129' + assert config.gate_signal_ref == '823867bb-f5d2-40eb-b875-657155addfd0' + assert config.gate_threshold == 0.5 + assert config.max_sample_gap_seconds == 20.0 + + +@pytest.mark.parametrize('field', list(entry())) +def test_missing_fields(tmp_path, field): + data = entry() + del data[field] + with pytest.raises(CalculationConfigError, match=f'missing fields: {field}'): + load_material_calculation(write_config(tmp_path, [data])) + + +@pytest.mark.parametrize(('field', 'value'), [ + ('type', 'integration'), ('version', 1), ('version', '2'), + ('machine_ref', ''), ('rate_signal_ref', ' '), ('gate_signal_ref', None), ('id', ''), + ('gate_threshold', True), ('gate_threshold', '0.5'), ('gate_threshold', float('nan')), + ('max_sample_gap_seconds', float('inf')), ('max_sample_gap_seconds', 0), + ('max_sample_gap_seconds', -1), ('output_unit', 'tonnes'), ('group_by', 'run'), + ('gate_treshold', 1), +]) +def test_invalid_fields(tmp_path, field, value): + data = entry() + data[field] = value + with pytest.raises(CalculationConfigError): + load_material_calculation(write_config(tmp_path, [data])) + + +@pytest.mark.parametrize('content', [ + 'calculations: [', 'calculations: []', 'calculations: {}', 'other: []', + 'calculations: []\ncalculations: []', 'calculations:\n - id: x\n id: y', +]) +def test_malformed_documents(tmp_path, content): + path = tmp_path / 'bad.yaml' + path.write_text(content) + with pytest.raises(CalculationConfigError): + load_material_calculation(path) + + +def test_selection(tmp_path): + first, second = entry(), dict(entry(), id='second') + path = write_config(tmp_path, [first, second]) + with pytest.raises(CalculationConfigError, match='specify --calculation-id'): + load_material_calculation(path) + assert load_material_calculation(path, 'second').id == 'second' + with pytest.raises(CalculationConfigError, match='not found'): + load_material_calculation(path, 'absent') + with pytest.raises(CalculationConfigError, match='unique'): + load_material_calculation(write_config(tmp_path, [first, first]), first['id']) + + +def test_repeated_sequential_polls_and_long_cycle(): + events = [] + elapsed = 0 + + def poll(*, now): + nonlocal elapsed + assert now == NOW + timedelta(seconds=elapsed) + events.append('poll-start') + elapsed += 25 # Longer than the interval; no real sleep or concurrent work. + events.append('poll-end') + + def sleep(seconds): + nonlocal elapsed + events.append(('sleep', seconds)) + elapsed += seconds + if elapsed >= 70: + raise KeyboardInterrupt + + output = io.StringIO() + service = Mock(poll_once=Mock(side_effect=poll)) + MaterialPollingRunner( + service, 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 + assert service.poll_once.call_count == 2 + assert output.getvalue().count('no open Production Run') == 2 + assert 'stopped' in output.getvalue() + + +def test_error_then_success_preserves_opaque_order_and_reports_totals(): + run = EnlyzeProductionRun('run', 'machine', None, ' 00842/ABC ', NOW, None) + result = MaterialPollResult(run, MaterialIntegrationState(12.5, 20)) + service = Mock(poll_once=Mock(side_effect=[ValueError('SECRET'), result])) + output, errors = io.StringIO(), io.StringIO() + sleep = Mock(side_effect=[None, KeyboardInterrupt]) + MaterialPollingRunner( + service, machine_id='machine', poll_interval_seconds=5, clock=lambda: NOW, + sleep=sleep, stdout=output, stderr=errors, + ).run() + assert service.poll_once.call_count == 2 + assert sleep.call_count == 2 + assert 'gateway/poll cycle: ValueError' in errors.getvalue() + assert 'SECRET' not in errors.getvalue() + assert "production_order=' 00842/ABC '" in output.getvalue() + assert 'consumption_kg=12.500000000 running_seconds=20.000' in output.getvalue() + assert NOW.isoformat() in output.getvalue() + + +def test_interrupt_during_poll(): + sleep, output = Mock(), io.StringIO() + MaterialPollingRunner( + Mock(poll_once=Mock(side_effect=KeyboardInterrupt)), machine_id='m', + poll_interval_seconds=1, sleep=sleep, stdout=output, + ).run() + sleep.assert_not_called() + assert 'stopped' in output.getvalue() + + +def test_cycle_configuration_failure_is_fatal(): + sleep = Mock() + runner = MaterialPollingRunner( + Mock(poll_once=Mock(side_effect=ConfigurationError('bad settings'))), + machine_id='m', poll_interval_seconds=1, sleep=sleep, + ) + with pytest.raises(ConfigurationError): + runner.run() + sleep.assert_not_called() + + +@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) + + +@pytest.mark.parametrize('operation', ['load', 'save']) +def test_state_error_classification(operation): + store = Mock() + getattr(store, operation).side_effect = OSError('SECRET') + wrapped = _ReportingStateStore(store) + arguments = ('machine', ' 00842 ') if operation == 'load' else ('machine', ' 00842 ', Mock()) + phase = 'loading' if operation == 'load' else 'saving' + with pytest.raises(MaterialStateError, match=f'state {phase}: OSError') as error: + getattr(wrapped, operation)(*arguments) + assert 'SECRET' not in str(error.value) + + +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') + with patch('production_analytics.service.material_runtime.MaterialPollingService') as service: + runner = build_material_runner( + load_material_calculation(EXAMPLE), poll_interval_seconds=7, + state_directory=tmp_path / 'state', secrets_file=secrets, + ) + kwargs = service.call_args.kwargs + assert kwargs['machine_id'] == entry()['machine_ref'] + assert kwargs['rate_variable_id'] == entry()['rate_signal_ref'] + assert kwargs['gate_variable_id'] == entry()['gate_signal_ref'] + assert kwargs['gate_threshold'] == 0.5 + assert kwargs['max_sample_gap_seconds'] == 20.0 + assert kwargs['gateway']._client._settings.api_key == 'secret-from-file' + assert kwargs['state_store'].store._directory == tmp_path / 'state' + assert runner.service is service.return_value + assert runner.interval == 7 + + +def test_cli_bad_config(tmp_path, capsys): + assert main(['run', 'material-poll', '--config', str(tmp_path / 'missing')]) == 2 + assert 'Config/startup' in capsys.readouterr().err + + +def test_cli_startup_and_success(tmp_path, capsys): + args = ['run', 'material-poll', '--config', str(EXAMPLE), + '--calculation-id', 'k7-fiber-consumption'] + with patch('production_analytics.service.material_runtime.build_material_runner') as build: + assert main(args) == 0 + build.return_value.run.assert_called_once_with() + assert build.call_args.kwargs['state_directory'] == Path('data/state/material') + build.side_effect = ConfigurationError('bad settings') + assert main(args) == 2 + assert 'bad settings' in capsys.readouterr().err + + +def test_failed_cycle_keeps_checkpoint_and_recovers(tmp_path): + from production_analytics.calculations.material_consumption import MaterialSample + from production_analytics.service.material_polling import ( + MaterialPollingService, + MaterialPollingState, + ) + from production_analytics.service.material_state_store import JsonMaterialStateStore + + order = ' 00842 ' + store = JsonMaterialStateStore(tmp_path) + store.save('machine', order, MaterialPollingState('run', MaterialIntegrationState( + cumulative_consumption_kg=1, integrated_running_seconds=10, + last_processed_timestamp=NOW, last_material_rate_kg_per_hour=3600, + last_gate_value=1, integration_active=True, + ))) + checkpoint = next(tmp_path.glob('*.json')) + original = checkpoint.read_bytes() + gateway = Mock() + gateway.get_open_production_run.return_value = EnlyzeProductionRun( + 'run', 'machine', None, order, NOW, None, + ) + gateway.get_material_samples.side_effect = [ + ValueError('bad response'), [MaterialSample(NOW + timedelta(seconds=10), 3600, 1)], + ] + service = MaterialPollingService( + gateway=gateway, state_store=_ReportingStateStore(store), machine_id='machine', + rate_variable_id='rate', gate_variable_id='gate', gate_threshold=.5, + max_sample_gap_seconds=20, + ) + sleeps = 0 + + def sleep(seconds): + nonlocal sleeps + sleeps += 1 + if sleeps == 1: + assert checkpoint.read_bytes() == original + else: + raise KeyboardInterrupt + + MaterialPollingRunner( + service, machine_id='machine', poll_interval_seconds=1, + clock=lambda: NOW + timedelta(seconds=10), sleep=sleep, + stdout=io.StringIO(), stderr=io.StringIO(), + ).run() + restored = store.load('machine', order) + assert restored.integration_state.cumulative_consumption_kg == 11 + assert restored.integration_state.integrated_running_seconds == 20 + + +def test_unwritable_state_location_is_startup_failure(tmp_path, monkeypatch, capsys): + monkeypatch.setenv('ENLYZE_BASE_URL', 'https://example.invalid/api/') + blocked = tmp_path / 'file' + blocked.write_text('not a directory') + assert main([ + 'run', 'material-poll', '--config', str(EXAMPLE), + '--state-directory', str(blocked), '--secrets-file', str(tmp_path / 'absent'), + ]) == 2 + assert 'State directory startup check failed' in capsys.readouterr().err