Persist material efficiency snapshots
This commit is contained in:
@@ -37,10 +37,18 @@ def _parser() -> argparse.ArgumentParser:
|
||||
raw.add_argument("--query", action="append", type=_query_item, default=[], metavar="KEY=VALUE")
|
||||
raw.add_argument("--pretty", action="store_true", help="Pretty-print sanitized JSON")
|
||||
raw.add_argument("--verbose", action="store_true", help="Show safe request progress on stderr")
|
||||
raw.add_argument("--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/")
|
||||
raw.add_argument("--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/")
|
||||
raw.add_argument("--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS)
|
||||
raw.add_argument("--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS)
|
||||
raw.add_argument(
|
||||
"--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/"
|
||||
)
|
||||
raw.add_argument(
|
||||
"--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/"
|
||||
)
|
||||
raw.add_argument(
|
||||
"--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS
|
||||
)
|
||||
raw.add_argument(
|
||||
"--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS
|
||||
)
|
||||
raw.add_argument(
|
||||
"--secrets-file",
|
||||
type=Path,
|
||||
@@ -52,18 +60,44 @@ def _parser() -> argparse.ArgumentParser:
|
||||
timeseries.add_argument("--start", required=True, help="ISO 8601 datetime with timezone")
|
||||
timeseries.add_argument("--end", required=True, help="ISO 8601 datetime with timezone")
|
||||
timeseries.add_argument("--variable", required=True, help="Variable UUID")
|
||||
timeseries.add_argument("--resampling-interval", type=int, help="Seconds; schema range is 10..604800")
|
||||
timeseries.add_argument(
|
||||
"--resampling-interval", type=int, help="Seconds; schema range is 10..604800"
|
||||
)
|
||||
timeseries.add_argument(
|
||||
"--resampling-method",
|
||||
choices=["first", "last", "max", "min", "count", "sum", "avg", "median", "std", "q5", "q25", "q75", "q95"],
|
||||
choices=[
|
||||
"first",
|
||||
"last",
|
||||
"max",
|
||||
"min",
|
||||
"count",
|
||||
"sum",
|
||||
"avg",
|
||||
"median",
|
||||
"std",
|
||||
"q5",
|
||||
"q25",
|
||||
"q75",
|
||||
"q95",
|
||||
],
|
||||
help="Optional schema-defined method for this variable",
|
||||
)
|
||||
timeseries.add_argument("--pretty", action="store_true", help="Pretty-print sanitized JSON")
|
||||
timeseries.add_argument("--verbose", action="store_true", help="Show safe request progress on stderr")
|
||||
timeseries.add_argument("--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/")
|
||||
timeseries.add_argument("--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/")
|
||||
timeseries.add_argument("--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS)
|
||||
timeseries.add_argument("--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS)
|
||||
timeseries.add_argument(
|
||||
"--verbose", action="store_true", help="Show safe request progress on stderr"
|
||||
)
|
||||
timeseries.add_argument(
|
||||
"--save-fixture", metavar="NAME", help="Save sanitized JSON under fixtures/enlyze/"
|
||||
)
|
||||
timeseries.add_argument(
|
||||
"--save-raw", metavar="NAME", help="Save unmodified response under ignored data/raw/enlyze/"
|
||||
)
|
||||
timeseries.add_argument(
|
||||
"--fixture-dir", type=Path, default=Path("fixtures/enlyze"), help=argparse.SUPPRESS
|
||||
)
|
||||
timeseries.add_argument(
|
||||
"--raw-dir", type=Path, default=Path("data/raw/enlyze"), help=argparse.SUPPRESS
|
||||
)
|
||||
timeseries.add_argument(
|
||||
"--secrets-file",
|
||||
type=Path,
|
||||
@@ -78,6 +112,10 @@ def _parser() -> argparse.ArgumentParser:
|
||||
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"))
|
||||
efficiency = runners.add_parser("material-efficiency", help="Persist ERP feedback KPI points")
|
||||
efficiency.add_argument("--config", type=Path, required=True)
|
||||
efficiency.add_argument("--secrets-file", type=Path, default=Path("secrets/enlyze.env"))
|
||||
efficiency.add_argument("--erp-secrets-file", type=Path, default=Path("secrets/erp.env"))
|
||||
return parser
|
||||
|
||||
|
||||
@@ -102,8 +140,10 @@ def _run_material(args: argparse.Namespace) -> int:
|
||||
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,
|
||||
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:
|
||||
@@ -117,9 +157,35 @@ def _run_material(args: argparse.Namespace) -> int:
|
||||
return 0
|
||||
|
||||
|
||||
def _run_material_efficiency(args: argparse.Namespace) -> int:
|
||||
from production_analytics.calculations.config import CalculationConfigError
|
||||
from production_analytics.enlyze.exploration import ConfigurationError
|
||||
from production_analytics.service.material_efficiency_runtime import (
|
||||
build_material_efficiency_runner,
|
||||
)
|
||||
|
||||
try:
|
||||
build_material_efficiency_runner(
|
||||
args.config,
|
||||
secrets_file=args.secrets_file,
|
||||
erp_secrets_file=args.erp_secrets_file,
|
||||
).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":
|
||||
if args.command == "material-efficiency":
|
||||
return _run_material_efficiency(args)
|
||||
return _run_material(args)
|
||||
if args.namespace != "enlyze" or args.command not in {"raw", "timeseries"}:
|
||||
return 2
|
||||
@@ -133,7 +199,8 @@ def main(argv: Sequence[str] | None = None) -> int:
|
||||
operation = "GET" if args.command == "raw" else "POST"
|
||||
target = args.path if args.command == "raw" else "/v2/timeseries"
|
||||
print(
|
||||
f"Requesting {operation} {target} (timeout={settings.timeout_seconds:g}s; authentication {authentication}).",
|
||||
f"Requesting {operation} {target} (timeout={settings.timeout_seconds:g}s; "
|
||||
f"authentication {authentication}).",
|
||||
file=sys.stderr,
|
||||
)
|
||||
client = ExplorationClient(settings)
|
||||
@@ -141,8 +208,13 @@ def main(argv: Sequence[str] | None = None) -> int:
|
||||
response = client.get(args.path, dict(args.query))
|
||||
request = {"method": "GET", "path": response.path}
|
||||
else:
|
||||
if args.resampling_interval is not None and not 10 <= args.resampling_interval <= 604800:
|
||||
raise ExplorationError("--resampling-interval must be between 10 and 604800 seconds.")
|
||||
if (
|
||||
args.resampling_interval is not None
|
||||
and not 10 <= args.resampling_interval <= 604800
|
||||
):
|
||||
raise ExplorationError(
|
||||
"--resampling-interval must be between 10 and 604800 seconds."
|
||||
)
|
||||
variable: dict[str, str] = {"uuid": args.variable}
|
||||
if args.resampling_method:
|
||||
variable["resampling_method"] = args.resampling_method
|
||||
@@ -170,9 +242,15 @@ def main(argv: Sequence[str] | None = None) -> int:
|
||||
indent = 2 if args.pretty or args.save_fixture else None
|
||||
print(json.dumps(payload, indent=indent, sort_keys=True))
|
||||
if args.save_fixture:
|
||||
print(f"Sanitized fixture written to {_write_fixture(args.fixture_dir, args.save_fixture, payload)}")
|
||||
print(
|
||||
"Sanitized fixture written to "
|
||||
f"{_write_fixture(args.fixture_dir, args.save_fixture, payload)}"
|
||||
)
|
||||
if args.save_raw:
|
||||
print(f"Raw response written to {_write_fixture(args.raw_dir, args.save_raw, raw_payload)}")
|
||||
print(
|
||||
"Raw response written to "
|
||||
f"{_write_fixture(args.raw_dir, args.save_raw, raw_payload)}"
|
||||
)
|
||||
except ExplorationError as error:
|
||||
print(f"ENLYZE exploration error: {error}")
|
||||
return 1
|
||||
|
||||
@@ -92,7 +92,7 @@ class MaterialEfficiencyService:
|
||||
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,
|
||||
erp_feedback_timestamp=cutoff,
|
||||
material_snapshot_timestamp=material.timestamp,
|
||||
good_quantity_m2=quantity, material_consumption_kg=consumption,
|
||||
material_consumption_kg_per_m2=kg_per_m2,
|
||||
|
||||
@@ -0,0 +1,60 @@
|
||||
"""Independent foreground ERP feedback polling; PostgreSQL owns deduplication."""
|
||||
|
||||
import sys
|
||||
import time
|
||||
from collections.abc import Callable
|
||||
from typing import TextIO
|
||||
|
||||
import psycopg
|
||||
|
||||
from production_analytics.calculations.config import finite_number
|
||||
from production_analytics.erp import ErpReadError
|
||||
from production_analytics.service.material_efficiency import MaterialEfficiencyService
|
||||
from production_analytics.service.postgres_material_efficiency import MaterialEfficiencyWriter
|
||||
|
||||
|
||||
class MaterialEfficiencyRunner:
|
||||
def __init__(
|
||||
self,
|
||||
service: MaterialEfficiencyService,
|
||||
writer: MaterialEfficiencyWriter,
|
||||
*,
|
||||
poll_interval_seconds: float = 60,
|
||||
sleep: Callable[[float], None] = time.sleep,
|
||||
stderr: TextIO | None = None,
|
||||
stdout: TextIO | None = None,
|
||||
) -> None:
|
||||
self.service = service
|
||||
self.writer = writer
|
||||
self.interval = finite_number(poll_interval_seconds, "poll interval", positive=True)
|
||||
self.sleep = sleep
|
||||
self.stderr = stderr if stderr is not None else sys.stderr
|
||||
self.stdout = stdout if stdout is not None else sys.stdout
|
||||
|
||||
def run(self) -> None:
|
||||
try:
|
||||
while True:
|
||||
try:
|
||||
snapshot = self.service.evaluate_current()
|
||||
if snapshot is not None and self.writer.write(snapshot):
|
||||
print(
|
||||
f"{snapshot.erp_feedback_timestamp.isoformat()} "
|
||||
f"workplace={snapshot.workplace!r} "
|
||||
f"production_order={snapshot.enlyze_production_order!r} "
|
||||
f"good_m2={snapshot.good_quantity_m2:.3f} "
|
||||
f"material_kg={snapshot.material_consumption_kg:.3f} "
|
||||
f"g_per_m2={snapshot.material_consumption_g_per_m2:.3f}",
|
||||
file=self.stdout,
|
||||
flush=True,
|
||||
)
|
||||
except (ErpReadError, psycopg.OperationalError, psycopg.InterfaceError) as exc:
|
||||
# Never include driver messages, SQL, credentials or response bodies.
|
||||
print(
|
||||
f"Material efficiency cycle failed: {type(exc).__name__}",
|
||||
file=self.stderr,
|
||||
flush=True,
|
||||
)
|
||||
# Fixed delay even after unavailable evaluations or transient failures.
|
||||
self.sleep(self.interval)
|
||||
except KeyboardInterrupt:
|
||||
return
|
||||
@@ -0,0 +1,82 @@
|
||||
"""Small strict YAML configuration and composition for KPI persistence."""
|
||||
|
||||
import os
|
||||
from pathlib import Path
|
||||
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
|
||||
|
||||
import yaml
|
||||
|
||||
from production_analytics.calculations.config import (
|
||||
CalculationConfigError,
|
||||
_UniqueLoader,
|
||||
finite_number,
|
||||
)
|
||||
from production_analytics.context import build_enlyze_production_order
|
||||
from production_analytics.enlyze.exploration import ConfigurationError, load_secret_file
|
||||
from production_analytics.erp import ErpSettings, ErpWorkplaceStatusGateway
|
||||
from production_analytics.service.material_efficiency import MaterialEfficiencyService
|
||||
from production_analytics.service.material_efficiency_runner import MaterialEfficiencyRunner
|
||||
from production_analytics.service.postgres_material import (
|
||||
PostgresMaterialSnapshotRepository,
|
||||
PostgresSettings,
|
||||
)
|
||||
from production_analytics.service.postgres_material_efficiency import (
|
||||
PostgresMaterialEfficiencyWriter,
|
||||
)
|
||||
|
||||
|
||||
def build_material_efficiency_runner(
|
||||
config: Path,
|
||||
*,
|
||||
secrets_file: Path = Path("secrets/enlyze.env"),
|
||||
erp_secrets_file: Path = Path("secrets/erp.env"),
|
||||
) -> MaterialEfficiencyRunner:
|
||||
try:
|
||||
document = yaml.load(config.read_text(encoding="utf-8"), Loader=_UniqueLoader)
|
||||
except (OSError, UnicodeError, yaml.YAMLError):
|
||||
raise CalculationConfigError("Cannot read material efficiency YAML configuration") from None
|
||||
required = {
|
||||
"workplace",
|
||||
"machine_id",
|
||||
"calculation_id",
|
||||
"production_order_format",
|
||||
"erp_timezone",
|
||||
}
|
||||
if (
|
||||
not isinstance(document, dict)
|
||||
or not required <= document.keys()
|
||||
or document.keys() - required - {"poll_interval_seconds"}
|
||||
):
|
||||
raise CalculationConfigError("Invalid material efficiency configuration fields")
|
||||
for name in required:
|
||||
value = document[name]
|
||||
if not isinstance(value, str) or not value.strip() or "\x00" in value:
|
||||
raise CalculationConfigError(f"{name} must be a non-empty string without NUL bytes")
|
||||
try:
|
||||
timezone = ZoneInfo(document["erp_timezone"])
|
||||
build_enlyze_production_order("0", document["production_order_format"])
|
||||
except (ValueError, ZoneInfoNotFoundError):
|
||||
raise CalculationConfigError("Invalid ERP timezone or production order format") from None
|
||||
interval = finite_number(
|
||||
document.get("poll_interval_seconds", 60), "poll interval", positive=True
|
||||
)
|
||||
environment = dict(os.environ)
|
||||
try:
|
||||
environment.update(load_secret_file(secrets_file))
|
||||
except Exception:
|
||||
raise ConfigurationError("PostgreSQL settings file could not be loaded") from None
|
||||
postgres = PostgresSettings.from_environment(environment)
|
||||
erp = ErpSettings.from_secret_file(erp_secrets_file)
|
||||
return MaterialEfficiencyRunner(
|
||||
MaterialEfficiencyService(
|
||||
ErpWorkplaceStatusGateway(erp),
|
||||
PostgresMaterialSnapshotRepository(postgres),
|
||||
workplace=document["workplace"],
|
||||
machine_id=document["machine_id"],
|
||||
calculation_id=document["calculation_id"],
|
||||
format_template=document["production_order_format"],
|
||||
erp_timezone=timezone,
|
||||
),
|
||||
PostgresMaterialEfficiencyWriter(postgres),
|
||||
poll_interval_seconds=interval,
|
||||
)
|
||||
@@ -0,0 +1,71 @@
|
||||
"""Immutable PostgreSQL history of feedback-aligned KPI points."""
|
||||
|
||||
from typing import Protocol
|
||||
|
||||
import psycopg
|
||||
|
||||
from production_analytics.calculations.config import finite_number
|
||||
from production_analytics.service.material_efficiency import MaterialEfficiencySnapshot
|
||||
from production_analytics.service.postgres_material import PostgresSettings
|
||||
|
||||
|
||||
class MaterialEfficiencyWriter(Protocol):
|
||||
def write(self, snapshot: MaterialEfficiencySnapshot) -> bool: ...
|
||||
|
||||
|
||||
class PostgresMaterialEfficiencyWriter:
|
||||
def __init__(self, settings: PostgresSettings) -> None:
|
||||
self.settings = settings
|
||||
|
||||
def write(self, snapshot: MaterialEfficiencySnapshot) -> bool:
|
||||
for name in ("erp_feedback_timestamp", "material_snapshot_timestamp"):
|
||||
if getattr(snapshot, name).utcoffset() is None:
|
||||
raise ValueError(f"{name} must be timezone-aware")
|
||||
for name in (
|
||||
"good_quantity_m2",
|
||||
"material_consumption_kg",
|
||||
"material_consumption_kg_per_m2",
|
||||
"material_consumption_g_per_m2",
|
||||
):
|
||||
finite_number(getattr(snapshot, name), name)
|
||||
if snapshot.nominal_width_m is not None:
|
||||
finite_number(snapshot.nominal_width_m, "nominal_width_m")
|
||||
# Context exit commits/rolls back and closes; later cycles reconnect after outages.
|
||||
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(
|
||||
"""INSERT INTO material_efficiency_snapshots (
|
||||
erp_feedback_timestamp, material_snapshot_timestamp, calculation_id,
|
||||
machine_id, workplace, erp_production_order, enlyze_production_order,
|
||||
article_number, article_description, nominal_width_m, good_quantity_m2,
|
||||
material_consumption_kg, material_consumption_kg_per_m2,
|
||||
material_consumption_g_per_m2
|
||||
) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
||||
ON CONFLICT (calculation_id, machine_id, enlyze_production_order,
|
||||
erp_feedback_timestamp) DO NOTHING
|
||||
RETURNING 1""",
|
||||
(
|
||||
snapshot.erp_feedback_timestamp,
|
||||
snapshot.material_snapshot_timestamp,
|
||||
snapshot.calculation_id,
|
||||
snapshot.machine_id,
|
||||
snapshot.workplace,
|
||||
snapshot.erp_production_order,
|
||||
snapshot.enlyze_production_order,
|
||||
snapshot.article_number,
|
||||
snapshot.article_description,
|
||||
snapshot.nominal_width_m,
|
||||
snapshot.good_quantity_m2,
|
||||
snapshot.material_consumption_kg,
|
||||
snapshot.material_consumption_kg_per_m2,
|
||||
snapshot.material_consumption_g_per_m2,
|
||||
),
|
||||
).fetchone()
|
||||
return row is not None
|
||||
Reference in New Issue
Block a user