Add channel material and downtime analytics

This commit is contained in:
2026-10-01 12:27:47 +02:00
parent 5452dc035b
commit 640ec53330
24 changed files with 2362 additions and 76 deletions
@@ -0,0 +1,246 @@
"""Stateful polling for configured extruder dosing channels."""
from collections import defaultdict
from dataclasses import dataclass
from datetime import datetime
from typing import Protocol
from production_analytics.calculations.channel_material import (
ChannelRate,
DosingChannelSample,
channel_rates,
validate_active_percentage_sum,
)
from production_analytics.calculations.material_consumption import (
MaterialConsumptionIntegrator,
MaterialIntegrationState,
MaterialIntegratorConfig,
)
from production_analytics.enlyze.gateway import EnlyzeProductionRun
@dataclass(frozen=True, slots=True)
class ChannelDefinition:
extruder: str
channel: str
total_rate_signal_ref: str
status_signal_ref: str
percentage_signal_ref: str
screw_speed_signal_ref: str
material_number_signal_ref: str | None = None
@dataclass(frozen=True, slots=True)
class ChannelPollingState:
run_id: str
integration_state: MaterialIntegrationState
class ChannelStateStore(Protocol):
def load(
self, machine_id: str, extruder: str, channel: str, production_order: str
) -> ChannelPollingState | None: ...
def save(
self,
machine_id: str,
extruder: str,
channel: str,
production_order: str,
state: ChannelPollingState,
) -> None: ...
class ChannelGateway(Protocol):
def get_open_production_run(self, machine_id: str) -> EnlyzeProductionRun | None: ...
def get_production_runs(self, machine_id: str) -> list[EnlyzeProductionRun]: ...
def get_dosing_channel_samples(
self,
*,
machine_id: str,
channels: tuple[ChannelDefinition, ...],
start: datetime,
end: datetime,
) -> dict[str, list[DosingChannelSample]]: ...
@dataclass(frozen=True, slots=True)
class ChannelPollResult:
run: EnlyzeProductionRun
extruder: str
channel: str
state: MaterialIntegrationState
material_number: str | None
material_name: str | None
material_mapping_status: str
percentage_sum_valid: bool
class ChannelMaterialPollingService:
def __init__(
self,
*,
gateway: ChannelGateway,
state_store: ChannelStateStore,
machine_id: str,
channels: tuple[ChannelDefinition, ...],
max_sample_gap_seconds: float,
percentage_tolerance: float = 1.0,
material_names: dict[str, str] | None = None,
) -> None:
if not channels or len({(c.extruder, c.channel) for c in channels}) != len(channels):
raise ValueError("channels must be non-empty and uniquely identified")
self.gateway, self.state_store, self.machine_id, self.channels = (
gateway,
state_store,
machine_id,
channels,
)
self.config = MaterialIntegratorConfig(0.5, max_sample_gap_seconds)
self.percentage_tolerance, self.material_names = percentage_tolerance, material_names or {}
def poll_once(self, *, now: datetime) -> tuple[ChannelPollResult, ...] | None:
if now.tzinfo is None or now.utcoffset() is None:
raise ValueError("now must be timezone-aware")
run = self.gateway.get_open_production_run(self.machine_id)
if run is None or now < run.start:
return None
# Each channel is separately checkpointed so an E4-style extra extruder
# can be introduced without changing state identity.
states: dict[str, ChannelPollingState | None] = {
c.channel + "\0" + c.extruder: self.state_store.load(
self.machine_id, c.extruder, c.channel, run.production_order
)
for c in self.channels
}
# A missing checkpoint is bootstrapped over each closed run separately.
# Resetting the temporal baseline at every run boundary is essential:
# a short wall-clock gap between ENLYZE runs must not become material.
if all(state is None for state in states.values()):
totals = {key: MaterialIntegrationState() for key in states}
previous_runs = sorted(
(
candidate
for candidate in self.gateway.get_production_runs(self.machine_id)
if candidate.production_order == run.production_order
and candidate.uuid != run.uuid
and candidate.end is not None
and candidate.end <= run.start
),
key=lambda candidate: candidate.start,
)
for previous in previous_runs:
historical = self._derive(
self.gateway.get_dosing_channel_samples(
machine_id=self.machine_id,
channels=self.channels,
start=previous.start,
end=previous.end,
)
)
for key, samples in historical.items():
integrator = MaterialConsumptionIntegrator(self.config, totals[key])
integrator.process_many(rate.sample for rate, _ in samples)
state = integrator.state
totals[key] = MaterialIntegrationState(
cumulative_consumption_kg=state.cumulative_consumption_kg,
integrated_running_seconds=state.integrated_running_seconds,
)
states = {key: ChannelPollingState("bootstrap", state) for key, state in totals.items()}
start = min(
(
s.integration_state.last_processed_timestamp
for s in states.values()
if s is not None
and s.run_id == run.uuid
and s.integration_state.last_processed_timestamp is not None
),
default=run.start,
)
if now < start:
return None
raw = self.gateway.get_dosing_channel_samples(
machine_id=self.machine_id, channels=self.channels, start=start, end=now
)
derived = self._derive(raw)
results = []
for definition in self.channels:
key = definition.channel + "\0" + definition.extruder
saved = states[key]
initial = (
saved.integration_state
if saved is not None and saved.run_id == run.uuid
else MaterialIntegrationState(
cumulative_consumption_kg=0.0
if saved is None
else saved.integration_state.cumulative_consumption_kg,
integrated_running_seconds=0.0
if saved is None
else saved.integration_state.integrated_running_seconds,
)
)
integrator = MaterialConsumptionIntegrator(self.config, initial)
samples = derived[key]
latest: ChannelRate | None = None
quality = True
for rate, sample_quality in samples:
if (
initial.last_processed_timestamp is not None
and rate.sample.timestamp < initial.last_processed_timestamp
):
continue
integrator.process(rate.sample)
latest = rate
quality = quality and sample_quality
state = integrator.state
self.state_store.save(
self.machine_id,
definition.extruder,
definition.channel,
run.production_order,
ChannelPollingState(run.uuid, state),
)
results.append(
ChannelPollResult(
run,
definition.extruder,
definition.channel,
state,
None if latest is None else latest.material_number,
None if latest is None else latest.material_name,
"UNAVAILABLE" if latest is None else latest.material_mapping_status,
quality,
)
)
return tuple(results)
def _derive(
self, raw: dict[str, list[DosingChannelSample]]
) -> dict[str, list[tuple[ChannelRate, bool]]]:
"""Evaluate aligned channel records with extruder-scoped quality rules."""
by_timestamp: dict[datetime, dict[str, DosingChannelSample]] = defaultdict(dict)
for definition in self.channels:
key = definition.channel + "\0" + definition.extruder
for sample in raw[key]:
by_timestamp[sample.timestamp][key] = sample
derived: dict[str, list[tuple[ChannelRate, bool]]] = defaultdict(list)
for timestamp, values in sorted(by_timestamp.items()):
if set(values) != set(raw):
raise ValueError(f"unaligned channel samples at {timestamp.isoformat()}")
grouped: dict[str, list[tuple[str, DosingChannelSample]]] = defaultdict(list)
for key, sample in values.items():
grouped[key.split("\0", 1)[1]].append((key, sample))
for siblings in grouped.values():
quality = validate_active_percentage_sum(
(sample for _, sample in siblings),
tolerance=self.percentage_tolerance,
)
for (key, _), rate in zip(
siblings,
channel_rates(
(sample for _, sample in siblings),
material_names=self.material_names,
percentage_tolerance=self.percentage_tolerance,
),
):
derived[key].append((rate, quality))
return derived
@@ -0,0 +1,98 @@
"""Sequential foreground runner for channel-level material consumption."""
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.channel_material import ChannelMaterialPollingService
from production_analytics.service.postgres_channel_material import (
PostgresChannelMaterialSnapshotWriter,
)
def utc_now() -> datetime:
return datetime.now(UTC)
class ChannelMaterialPollingRunner:
def __init__(
self,
service: ChannelMaterialPollingService,
*,
machine_id: str,
calculation_id: str,
snapshot_writer: PostgresChannelMaterialSnapshotWriter,
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
self.stdout = stdout if stdout is not None else sys.stdout
self.stderr = stderr if stderr is not None else sys.stderr
def run_once(self) -> bool:
"""Run and report one cycle; return whether channel snapshots were written."""
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}"
results = self.service.poll_once(now=now)
if results is None:
print(
f"{prefix} no open Production Run / no eligible polling window; snapshots=0",
file=self.stdout,
flush=True,
)
return False
for result in results:
self.snapshot_writer.write(
timestamp=now,
calculation_id=self.calculation_id,
machine_id=self.machine_id,
extruder=result.extruder,
channel=result.channel,
production_order=result.run.production_order,
run_id=result.run.uuid,
cumulative_consumption_kg=result.state.cumulative_consumption_kg,
material_number=result.material_number,
material_name=result.material_name,
material_mapping_status=result.material_mapping_status,
percentage_sum_valid=result.percentage_sum_valid,
)
print(
f"{prefix} production_order={results[0].run.production_order!r} "
f"channel snapshots={len(results)} written",
file=self.stdout,
flush=True,
)
return True
def run(self) -> None:
try:
while True:
try:
self.run_once()
except ConfigurationError:
raise
except Exception as exc:
if isinstance(exc, ExplorationError):
detail = f"ENLYZE/gateway: {type(exc).__name__}"
else:
detail = f"poll cycle: {type(exc).__name__}"
print(f"channel material error {detail}", file=self.stderr, flush=True)
self.sleep(self.interval)
except KeyboardInterrupt:
print("Channel material polling stopped", file=self.stdout, flush=True)
@@ -0,0 +1,69 @@
"""Compose the live channel-level material polling application."""
import math
import os
import tempfile
from pathlib import Path
from urllib.parse import urlsplit
from production_analytics.calculations.channel_config import ChannelMaterialCalculationConfig
from production_analytics.calculations.config import 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.channel_material import ChannelMaterialPollingService
from production_analytics.service.channel_material_runner import ChannelMaterialPollingRunner
from production_analytics.service.channel_material_state_store import JsonChannelMaterialStateStore
from production_analytics.service.postgres_channel_material import (
PostgresChannelMaterialSnapshotWriter,
)
from production_analytics.service.postgres_material import PostgresSettings
def build_channel_material_runner(
calculation: ChannelMaterialCalculationConfig,
*,
poll_interval_seconds: float,
state_directory: Path,
secrets_file: Path,
) -> ChannelMaterialPollingRunner:
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
postgres_settings = PostgresSettings.from_environment(environment)
service = ChannelMaterialPollingService(
gateway=EnlyzeApiGateway(ExplorationClient(settings)),
state_store=JsonChannelMaterialStateStore(state_directory),
machine_id=calculation.machine_ref,
channels=calculation.channels,
max_sample_gap_seconds=calculation.max_sample_gap_seconds,
percentage_tolerance=calculation.percentage_tolerance,
material_names=calculation.material_names,
)
return ChannelMaterialPollingRunner(
service,
machine_id=calculation.machine_ref,
calculation_id=calculation.id,
snapshot_writer=PostgresChannelMaterialSnapshotWriter(postgres_settings),
poll_interval_seconds=poll_interval_seconds,
)
@@ -0,0 +1,59 @@
"""Restart-safe JSON state storage keyed by channel identity."""
import hashlib
import json
import os
import tempfile
from pathlib import Path
from production_analytics.calculations.material_state import (
material_state_from_dict,
material_state_to_dict,
)
from production_analytics.service.channel_material import ChannelPollingState
class JsonChannelMaterialStateStore:
def __init__(self, directory: str | Path) -> None:
self.directory = Path(directory)
def _path(self, *identity: str) -> Path:
digest = hashlib.sha256(json.dumps(identity).encode()).hexdigest()
return self.directory / f"channel-material-{digest}.json"
def load(
self, machine_id: str, extruder: str, channel: str, production_order: str
) -> ChannelPollingState | None:
path = self._path(machine_id, extruder, channel, production_order)
if not path.exists():
return None
data = json.loads(path.read_text(encoding="utf-8"))
if not isinstance(data, dict) or not isinstance(data.get("run_id"), str):
raise ValueError("Invalid channel material state")
return ChannelPollingState(
data["run_id"], material_state_from_dict(data["integration_state"])
)
def save(
self,
machine_id: str,
extruder: str,
channel: str,
production_order: str,
state: ChannelPollingState,
) -> None:
self.directory.mkdir(parents=True, exist_ok=True)
payload = {
"run_id": state.run_id,
"integration_state": material_state_to_dict(state.integration_state),
}
material_state_from_dict(payload["integration_state"])
destination = self._path(machine_id, extruder, channel, production_order)
with tempfile.NamedTemporaryFile(
mode="w", encoding="utf-8", dir=self.directory, delete=False
) as handle:
temporary = Path(handle.name)
json.dump(payload, handle, allow_nan=False)
handle.flush()
os.fsync(handle.fileno())
os.replace(temporary, destination)
@@ -0,0 +1,214 @@
"""Generic ENLYZE downtime reconciliation and ERP order attribution."""
from dataclasses import dataclass
from datetime import UTC, datetime, timedelta
from enum import StrEnum
from math import isfinite
from typing import Protocol
from production_analytics.context import build_enlyze_production_order
from production_analytics.enlyze.gateway import EnlyzeDowntime
from production_analytics.erp import CurrentWorkplaceStatus
class DowntimeCategory(StrEnum):
PLANNED = "PLANNED"
UNPLANNED = "UNPLANNED"
UNKNOWN = "UNKNOWN"
@dataclass(frozen=True, slots=True)
class ProductionOrderBoundary:
machine_id: str
production_order: str
started_at: datetime
ended_at: datetime | None
@dataclass(frozen=True, slots=True)
class ProductionDowntimeEvent:
external_id: str
machine_id: str
production_order: str | None
source_type: str
source_start: datetime
source_end: datetime | None
attributed_start: datetime | None
attributed_end: datetime | None
attributed_duration_seconds: float | None
reason_id: str | None
reason_name: str | None
reason_description: str | None
reason_group: str | None
category: DowntimeCategory
comment: str | None
source_updated_at: datetime | None
last_reconciled_at: datetime
class DowntimeGateway(Protocol):
def get_downtimes(
self, machine_id: str, *, start: datetime | None = None
) -> list[EnlyzeDowntime]: ...
class WorkplaceStatusReader(Protocol):
def get_current_workplace_status(self, workplace: str) -> CurrentWorkplaceStatus | None: ...
class DowntimeRepository(Protocol):
def current_boundary(self, machine_id: str) -> ProductionOrderBoundary | None: ...
def save_boundary(self, boundary: ProductionOrderBoundary) -> None: ...
def close_order_attribution(
self, machine_id: str, production_order: str, ended_at: datetime
) -> None: ...
def upsert(self, event: ProductionDowntimeEvent) -> None: ...
class DowntimeReconciliationService:
"""Reconciles mutable source events without deriving a shift schedule."""
def __init__(
self,
gateway: DowntimeGateway,
erp: WorkplaceStatusReader,
repository: DowntimeRepository,
*,
machine_id: str,
workplace: str,
production_order_format: str,
lookback_hours: float = 48,
full_reconciliation_hours: float = 24,
completion_tolerance_m2: float = 0.001,
) -> None:
if not isfinite(lookback_hours) or lookback_hours <= 0:
raise ValueError("lookback_hours must be positive")
if not isfinite(full_reconciliation_hours) or full_reconciliation_hours <= 0:
raise ValueError("full_reconciliation_hours must be positive")
if not isfinite(completion_tolerance_m2) or completion_tolerance_m2 < 0:
raise ValueError("completion_tolerance_m2 must be non-negative")
self.gateway, self.erp, self.repository = gateway, erp, repository
self.machine_id, self.workplace = machine_id, workplace
self.production_order_format = production_order_format
self.lookback = timedelta(hours=lookback_hours)
self.full_reconciliation_interval = timedelta(hours=full_reconciliation_hours)
self._last_full_reconciliation: datetime | None = None
self.tolerance = completion_tolerance_m2
def reconcile_once(self, now: datetime) -> int:
if now.tzinfo is None or now.utcoffset() is None:
raise ValueError("now must be timezone-aware")
now = now.astimezone(UTC)
boundary = self._refresh_boundary(now)
events = self.gateway.get_downtimes(self.machine_id, start=now - self.lookback)
if (
self._last_full_reconciliation is None
or now - self._last_full_reconciliation >= self.full_reconciliation_interval
):
# ENLYZE does not guarantee that a delayed reason edit remains in a source-start
# lookback window. A periodic full scan makes old open/UNKNOWN UUIDs mutable too.
events = list(
{
event.uuid: event
for event in [
*events,
*self.gateway.get_downtimes(self.machine_id),
]
}.values()
)
self._last_full_reconciliation = now
for source in events:
self.repository.upsert(self._event(source, boundary, now))
return len(events)
def _refresh_boundary(self, now: datetime) -> ProductionOrderBoundary | None:
status = self.erp.get_current_workplace_status(self.workplace)
old = self.repository.current_boundary(self.machine_id)
if status is None:
return old
order = build_enlyze_production_order(status.production_order, self.production_order_format)
feedback_at = _aware_feedback(status.feedback_timestamp)
if old is not None and old.production_order != order and old.ended_at is None:
self.repository.close_order_attribution(
self.machine_id, old.production_order, feedback_at
)
if old is None or old.production_order != order:
current = ProductionOrderBoundary(self.machine_id, order, feedback_at, None)
else:
current = old
if _is_complete(status, self.tolerance) and current.ended_at is None:
current = ProductionOrderBoundary(
current.machine_id,
current.production_order,
current.started_at,
feedback_at,
)
self.repository.close_order_attribution(
self.machine_id, current.production_order, feedback_at
)
self.repository.save_boundary(current)
return current
def _event(
self,
source: EnlyzeDowntime,
boundary: ProductionOrderBoundary | None,
now: datetime,
) -> ProductionDowntimeEvent:
category = (
DowntimeCategory(source.reason_category)
if source.reason_category
in {
"PLANNED",
"UNPLANNED",
}
else DowntimeCategory.UNKNOWN
)
attributed_start = attributed_end = None
order = None
if boundary is not None:
lower = max(source.start, boundary.started_at)
upper = source.end
if boundary.ended_at is not None:
upper = boundary.ended_at if upper is None else min(upper, boundary.ended_at)
if upper is None or lower < upper:
order, attributed_start, attributed_end = boundary.production_order, lower, upper
duration = (
None
if attributed_start is None or attributed_end is None
else (attributed_end - attributed_start).total_seconds()
)
return ProductionDowntimeEvent(
source.uuid,
source.machine_id,
order,
source.source_type,
source.start,
source.end,
attributed_start,
attributed_end,
duration,
source.reason_id,
source.reason_name,
source.reason_description,
source.reason_group,
category,
source.comment,
source.source_updated_at,
now,
)
def _aware_feedback(value: datetime) -> datetime:
if value.tzinfo is None or value.utcoffset() is None:
raise ValueError("ERP feedback timestamp must be timezone-aware for downtime attribution")
return value.astimezone(UTC)
def _is_complete(status: CurrentWorkplaceStatus, tolerance: float) -> bool:
remaining = status.remaining_quantity_m2
if remaining is not None and remaining <= tolerance:
return True
if status.order_quantity_m2 is None or status.good_quantity_m2 is None:
return False
return status.good_quantity_m2 + tolerance >= status.order_quantity_m2
@@ -0,0 +1,50 @@
"""Foreground reconciliation runner suitable for a systemd service."""
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.service.downtime import DowntimeReconciliationService
class DowntimeReconciliationRunner:
def __init__(
self,
service: DowntimeReconciliationService,
*,
machine_id: str,
poll_interval_seconds: float = 300,
clock: Callable[[], datetime] | None = None,
sleep: Callable[[float], None] = time.sleep,
stdout: TextIO | None = None,
stderr: TextIO | None = None,
) -> None:
self.service, self.machine_id = service, machine_id
self.interval = finite_number(poll_interval_seconds, "poll interval", positive=True)
self.clock, self.sleep = clock or (lambda: datetime.now(UTC)), sleep
self.stdout, self.stderr = stdout or sys.stdout, stderr or sys.stderr
def run(self) -> None:
try:
while True:
now = self.clock()
try:
count = self.service.reconcile_once(now)
print(
f"{now.isoformat()} machine={self.machine_id!r} reconciled={count}",
file=self.stdout,
flush=True,
)
except Exception as exc:
print(
f"{now.isoformat()} machine={self.machine_id!r} reconciliation failed: "
f"{type(exc).__name__}",
file=self.stderr,
flush=True,
)
self.sleep(self.interval)
except KeyboardInterrupt:
return
@@ -0,0 +1,67 @@
"""Build the generic downtime runner from a small, explicit YAML configuration."""
import os
from dataclasses import dataclass
from pathlib import Path
import yaml
from production_analytics.enlyze.exploration import (
ExplorationClient,
ExplorationSettings,
load_secret_file,
)
from production_analytics.enlyze.gateway import EnlyzeApiGateway
from production_analytics.erp import ErpSettings, ErpWorkplaceStatusGateway
from production_analytics.service.downtime import DowntimeReconciliationService
from production_analytics.service.downtime_runner import DowntimeReconciliationRunner
from production_analytics.service.postgres_downtime import PostgresDowntimeRepository
from production_analytics.service.postgres_material import PostgresSettings
@dataclass(frozen=True, slots=True)
class DowntimeConfig:
machine_id: str
erp_workplace: str
production_order_format: str
poll_interval_seconds: float = 300
lookback_hours: float = 48
full_reconciliation_hours: float = 24
completion_tolerance_m2: float = 0.001
def load_downtime_config(path: Path) -> DowntimeConfig:
raw = yaml.safe_load(path.read_text(encoding="utf-8"))
if not isinstance(raw, dict):
raise ValueError("Downtime config must be a mapping")
allowed = {field.name for field in DowntimeConfig.__dataclass_fields__.values()}
if set(raw) - allowed or not {"machine_id", "erp_workplace", "production_order_format"} <= set(
raw
):
raise ValueError("Downtime config has unknown or missing fields")
for name in ("machine_id", "erp_workplace", "production_order_format"):
if not isinstance(raw[name], str) or not raw[name].strip():
raise ValueError(f"{name} must be a non-empty string")
return DowntimeConfig(**raw)
def build_downtime_runner(
config_path: Path, *, secrets_file: Path, erp_secrets_file: Path
) -> DowntimeReconciliationRunner:
config = load_downtime_config(config_path)
environment = {**os.environ, **load_secret_file(secrets_file)}
erp_environment = {**os.environ, **load_secret_file(erp_secrets_file)}
service = DowntimeReconciliationService(
EnlyzeApiGateway(ExplorationClient(ExplorationSettings.from_environment(environment))),
ErpWorkplaceStatusGateway(ErpSettings.from_environment(erp_environment)),
PostgresDowntimeRepository(PostgresSettings.from_environment(os.environ)),
machine_id=config.machine_id,
workplace=config.erp_workplace,
production_order_format=config.production_order_format,
lookback_hours=config.lookback_hours,
full_reconciliation_hours=config.full_reconciliation_hours,
completion_tolerance_m2=config.completion_tolerance_m2,
)
return DowntimeReconciliationRunner(
service, machine_id=config.machine_id, poll_interval_seconds=config.poll_interval_seconds
)
@@ -0,0 +1,63 @@
"""PostgreSQL writer for channel-level cumulative consumption snapshots."""
from datetime import datetime
from production_analytics.service.postgres_material import PostgresSettings
class PostgresChannelMaterialSnapshotWriter:
def __init__(self, settings: PostgresSettings) -> None:
self.settings = settings
def write(
self,
*,
timestamp: datetime,
calculation_id: str,
machine_id: str,
extruder: str,
channel: str,
production_order: str,
run_id: str,
cumulative_consumption_kg: float,
material_number: str | None,
material_name: str | None,
material_mapping_status: str,
percentage_sum_valid: bool,
) -> None:
import psycopg
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 channel_material_consumption_snapshots
(timestamp, calculation_id, machine_id, extruder, doser_channel,
production_order, run_id, cumulative_consumption_kg,
material_number, material_name, material_mapping_status,
percentage_sum_valid)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (calculation_id, machine_id, extruder, doser_channel,
production_order, timestamp)
DO NOTHING""",
(
timestamp,
calculation_id,
machine_id,
extruder,
channel,
production_order,
run_id,
cumulative_consumption_kg,
material_number,
material_name,
material_mapping_status,
percentage_sum_valid,
),
)
@@ -0,0 +1,112 @@
"""PostgreSQL persistence for reconciled ENLYZE downtime events."""
from production_analytics.service.downtime import ProductionDowntimeEvent, ProductionOrderBoundary
from production_analytics.service.postgres_material import PostgresSettings
class PostgresDowntimeRepository:
def __init__(self, settings: PostgresSettings) -> None:
self.settings = settings
def _connect(self):
import psycopg
return 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",
)
def current_boundary(self, machine_id: str) -> ProductionOrderBoundary | None:
with self._connect() as connection:
row = connection.execute(
"""SELECT machine_id, production_order, started_at, ended_at
FROM production_order_attribution_state WHERE machine_id = %s""",
(machine_id,),
).fetchone()
return None if row is None else ProductionOrderBoundary(*row)
def save_boundary(self, boundary: ProductionOrderBoundary) -> None:
with self._connect() as connection:
connection.execute(
"""INSERT INTO production_order_attribution_state
(machine_id, production_order, started_at, ended_at)
VALUES (%s, %s, %s, %s)
ON CONFLICT (machine_id) DO UPDATE SET
production_order = EXCLUDED.production_order,
started_at = EXCLUDED.started_at, ended_at = EXCLUDED.ended_at""",
(
boundary.machine_id,
boundary.production_order,
boundary.started_at,
boundary.ended_at,
),
)
def close_order_attribution(self, machine_id: str, production_order: str, ended_at) -> None:
with self._connect() as connection:
connection.execute(
"""UPDATE production_downtime_events
SET attributed_end = CASE WHEN source_end IS NULL OR source_end > %s THEN %s
ELSE source_end END,
attributed_duration_seconds = EXTRACT(EPOCH FROM
(CASE WHEN source_end IS NULL OR source_end > %s
THEN %s ELSE source_end END
- attributed_start))
WHERE machine_id = %s AND production_order = %s
AND attributed_start IS NOT NULL AND attributed_end IS NULL""",
(ended_at, ended_at, ended_at, ended_at, machine_id, production_order),
)
def upsert(self, event: ProductionDowntimeEvent) -> None:
with self._connect() as connection:
connection.execute(
"""INSERT INTO production_downtime_events
(external_id, machine_id, production_order, source_type, source_start,
source_end, attributed_start, attributed_end, attributed_duration_seconds,
reason_id, reason_name, reason_description, reason_group, category, comment,
source_updated_at, last_reconciled_at)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (external_id) DO UPDATE SET
machine_id = EXCLUDED.machine_id,
production_order = COALESCE(
EXCLUDED.production_order, production_downtime_events.production_order),
source_type = EXCLUDED.source_type, source_start = EXCLUDED.source_start,
source_end = EXCLUDED.source_end,
attributed_start = COALESCE(
EXCLUDED.attributed_start, production_downtime_events.attributed_start),
attributed_end = COALESCE(
EXCLUDED.attributed_end, production_downtime_events.attributed_end),
attributed_duration_seconds = COALESCE(
EXCLUDED.attributed_duration_seconds,
production_downtime_events.attributed_duration_seconds),
reason_id = EXCLUDED.reason_id, reason_name = EXCLUDED.reason_name,
reason_description = EXCLUDED.reason_description,
reason_group = EXCLUDED.reason_group,
category = EXCLUDED.category, comment = EXCLUDED.comment,
source_updated_at = EXCLUDED.source_updated_at,
last_reconciled_at = EXCLUDED.last_reconciled_at""",
(
event.external_id,
event.machine_id,
event.production_order,
event.source_type,
event.source_start,
event.source_end,
event.attributed_start,
event.attributed_end,
event.attributed_duration_seconds,
event.reason_id,
event.reason_name,
event.reason_description,
event.reason_group,
event.category.value,
event.comment,
event.source_updated_at,
event.last_reconciled_at,
),
)