From 640ec53330be578c0bf987471b82b5f0be0c8fff Mon Sep 17 00:00:00 2001 From: Martin Tazl Date: Thu, 1 Oct 2026 12:27:47 +0200 Subject: [PATCH] Add channel material and downtime analytics --- config/bento1-downtime.yaml | 7 + config/downtime.example.yaml | 9 + ...-channel-material-consumption.example.yaml | 30 ++ docs/downtime-architecture.md | 36 ++ .../calculations/channel_config.py | 127 +++++++ .../calculations/channel_material.py | 99 +++++ src/production_analytics/enlyze/gateway.py | 356 +++++++++++++++++- .../service/channel_material.py | 246 ++++++++++++ .../service/channel_material_runner.py | 98 +++++ .../service/channel_material_runtime.py | 69 ++++ .../service/channel_material_state_store.py | 59 +++ src/production_analytics/service/downtime.py | 214 +++++++++++ .../service/downtime_runner.py | 50 +++ .../service/downtime_runtime.py | 67 ++++ .../service/postgres_channel_material.py | 63 ++++ .../service/postgres_downtime.py | 112 ++++++ tests/conftest.py | 21 ++ tests/test_channel_material.py | 158 ++++++++ tests/test_channel_material_runner.py | 141 +++++++ tests/test_downtime.py | 140 +++++++ tests/test_enlyze_gateway.py | 265 ++++++++++--- tests/test_material_application.py | 4 +- tests/test_material_efficiency_persistence.py | 6 +- tests/test_schema.py | 61 +++ 24 files changed, 2362 insertions(+), 76 deletions(-) create mode 100644 config/bento1-downtime.yaml create mode 100644 config/downtime.example.yaml create mode 100644 config/e1-channel-material-consumption.example.yaml create mode 100644 docs/downtime-architecture.md create mode 100644 src/production_analytics/calculations/channel_config.py create mode 100644 src/production_analytics/calculations/channel_material.py create mode 100644 src/production_analytics/service/channel_material.py create mode 100644 src/production_analytics/service/channel_material_runner.py create mode 100644 src/production_analytics/service/channel_material_runtime.py create mode 100644 src/production_analytics/service/channel_material_state_store.py create mode 100644 src/production_analytics/service/downtime.py create mode 100644 src/production_analytics/service/downtime_runner.py create mode 100644 src/production_analytics/service/downtime_runtime.py create mode 100644 src/production_analytics/service/postgres_channel_material.py create mode 100644 src/production_analytics/service/postgres_downtime.py create mode 100644 tests/conftest.py create mode 100644 tests/test_channel_material.py create mode 100644 tests/test_channel_material_runner.py create mode 100644 tests/test_downtime.py create mode 100644 tests/test_schema.py diff --git a/config/bento1-downtime.yaml b/config/bento1-downtime.yaml new file mode 100644 index 0000000..f569553 --- /dev/null +++ b/config/bento1-downtime.yaml @@ -0,0 +1,7 @@ +machine_id: 5f42a4f6-9ca0-4f6f-9786-40d50a35b230 +erp_workplace: Bento 1 +production_order_format: "Bento 1-{production_order}" +poll_interval_seconds: 300 +lookback_hours: 48 +full_reconciliation_hours: 24 +completion_tolerance_m2: 0.001 diff --git a/config/downtime.example.yaml b/config/downtime.example.yaml new file mode 100644 index 0000000..f1e222c --- /dev/null +++ b/config/downtime.example.yaml @@ -0,0 +1,9 @@ +# Generic ENLYZE downtime reconciliation configuration; one runner per machine. +machine_id: machine-uuid +erp_workplace: WORKPLACE +production_order_format: "Machine-{production_order}" +poll_interval_seconds: 300 +lookback_hours: 48 +# Full source scan keeps very old open/UNKNOWN events eligible for delayed edits. +full_reconciliation_hours: 24 +completion_tolerance_m2: 0.001 diff --git a/config/e1-channel-material-consumption.example.yaml b/config/e1-channel-material-consumption.example.yaml new file mode 100644 index 0000000..d437707 --- /dev/null +++ b/config/e1-channel-material-consumption.example.yaml @@ -0,0 +1,30 @@ +# Supply E1's ENLYZE machine UUID before operating. Signal references are PLC +# origins, resolved uniquely by the ENLYZE gateway at poll time. +calculations: + - id: e1-channel-material-consumption + type: channel_material_consumption + version: "1" + machine_ref: 8302e3d1-b1e5-42f1-8540-615eb6c73e08 + max_sample_gap_seconds: 20.0 + percentage_tolerance: 1.0 + extruders: + - id: Ex1 + total_rate_signal_ref: DB107:18 + channels: + - {id: C1, status_signal_ref: DB107:122.0, percentage_signal_ref: DB107:126, screw_speed_signal_ref: DB107:170} + - {id: C2, status_signal_ref: DB107:222.0, percentage_signal_ref: DB107:226, screw_speed_signal_ref: DB107:270} + - {id: C3, status_signal_ref: DB107:322.0, percentage_signal_ref: DB107:326, screw_speed_signal_ref: DB107:370} + - {id: C4, status_signal_ref: DB107:422.0, percentage_signal_ref: DB107:426, screw_speed_signal_ref: DB107:470} + - {id: C5, status_signal_ref: DB107:522.0, percentage_signal_ref: DB107:526, screw_speed_signal_ref: DB107:570} + - {id: C6, status_signal_ref: DB107:622.0, percentage_signal_ref: DB107:626, screw_speed_signal_ref: DB107:670} + - {id: C7, status_signal_ref: DB107:722.0, percentage_signal_ref: DB107:726, screw_speed_signal_ref: DB107:770} + - id: Ex2 + total_rate_signal_ref: DB207:18 + channels: + - {id: C1, status_signal_ref: DB207:122.0, percentage_signal_ref: DB207:126, screw_speed_signal_ref: DB207:170} + - {id: C2, status_signal_ref: DB207:222.0, percentage_signal_ref: DB207:226, screw_speed_signal_ref: DB207:270} + - {id: C3, status_signal_ref: DB207:322.0, percentage_signal_ref: DB207:326, screw_speed_signal_ref: DB207:370} + - {id: C4, status_signal_ref: DB207:422.0, percentage_signal_ref: DB207:426, screw_speed_signal_ref: DB207:470} + - {id: C5, status_signal_ref: DB207:522.0, percentage_signal_ref: DB207:526, screw_speed_signal_ref: DB207:570} + - {id: C6, status_signal_ref: DB207:622.0, percentage_signal_ref: DB207:626, screw_speed_signal_ref: DB207:670} + - {id: C7, status_signal_ref: DB207:722.0, percentage_signal_ref: DB207:726, screw_speed_signal_ref: DB207:770} diff --git a/docs/downtime-architecture.md b/docs/downtime-architecture.md new file mode 100644 index 0000000..ab0da49 --- /dev/null +++ b/docs/downtime-architecture.md @@ -0,0 +1,36 @@ +# ENLYZE downtime architecture + +`GET /v2/downtimes` is the source of downtime events. Its UUID is the durable external +identity. `end: null` means the source event is currently open. `reason: null` is valid +and is persisted as `UNKNOWN`; it is never treated as `UNPLANNED`. ENLYZE's +`updated.timestamp` describes source metadata/reason editing, not finalization. + +The reconciliation runner polls each configured machine with a configurable source-start +lookback (48 hours by default), follows pagination, and upserts by UUID. Every 24 hours by +default it also scans the full machine source so an old open or UNKNOWN event remains eligible +for a delayed classification update. It always replaces +the end time, comment, reason metadata, category, source update timestamp, and attributed +timing. This is intentionally not append-only because supervisors can classify an event much +later. Operators can enlarge `lookback_hours` where delayed classification exceeds the normal +window. + +`production_downtime_events` retains ENLYZE source timing separately from FA-attributed +timing. An event only receives an order while the persisted ERP order boundary is active. +When ERP reports a new production order, the preceding boundary ends at the new feedback +timestamp. A boundary also ends when `remaining_quantity_m2 <= completion_tolerance_m2`, or +when `good_quantity_m2 + tolerance >= order_quantity_m2`. These are the exact fields from +`CurrentWorkplaceStatus`; `feedback_timestamp` supplies the boundary instant. No equality on +floating values is required. If neither signal proves completion, attribution stays open. + +This protects a completed FA from time in an ENLYZE downtime that remains open after work has +finished. It does not infer completion from downtime. The durable +`production_order_attribution_state` makes this clipping survive runner restarts. + +There is no machine-readable schedule today. No clock-time, shift, overnight, weekend, or Excel +planning inference occurs. ENLYZE `reason.category` is the only planned/unplanned source; +anything absent or unsupported remains `UNKNOWN`. A future SPA schedule provider can refine +classification or boundaries as a separate input without changing stored source facts. + +The resulting fields support totals per order, category and reason, individual source events, +and open events. Source duration is ENLYZE's `source_start/source_end`; attributed duration is +only the interval inside the proven FA boundary. diff --git a/src/production_analytics/calculations/channel_config.py b/src/production_analytics/calculations/channel_config.py new file mode 100644 index 0000000..003ef30 --- /dev/null +++ b/src/production_analytics/calculations/channel_config.py @@ -0,0 +1,127 @@ +"""Configuration for reusable channel-level extruder consumption.""" + +from dataclasses import dataclass +from pathlib import Path + +import yaml + +from production_analytics.calculations.config import CalculationConfigError, finite_number +from production_analytics.service.channel_material import ChannelDefinition + + +@dataclass(frozen=True, slots=True) +class ChannelMaterialCalculationConfig: + id: str + machine_ref: str + max_sample_gap_seconds: float + percentage_tolerance: float + channels: tuple[ChannelDefinition, ...] + material_names: dict[str, str] + + +def load_channel_material_calculation(path: str | Path) -> ChannelMaterialCalculationConfig: + try: + document = yaml.safe_load(Path(path).read_text(encoding="utf-8")) + except (OSError, UnicodeError, yaml.YAMLError) as exc: + raise CalculationConfigError("Cannot read channel material configuration") from exc + if ( + not isinstance(document, dict) + or set(document) != {"calculations"} + or not isinstance(document["calculations"], list) + ): + raise CalculationConfigError("Configuration must contain only calculations") + if len(document["calculations"]) != 1: + raise CalculationConfigError("Channel configuration requires one calculation") + entry = document["calculations"][0] + required = { + "id", + "type", + "version", + "machine_ref", + "max_sample_gap_seconds", + "percentage_tolerance", + "extruders", + } + if ( + not isinstance(entry, dict) + or set(entry) - (required | {"material_names"}) + or required - set(entry) + ): + raise CalculationConfigError("Invalid channel calculation fields") + if entry["type"] != "channel_material_consumption" or entry["version"] != "1": + raise CalculationConfigError("Unsupported channel calculation type or version") + if not all(isinstance(entry[x], str) and entry[x].strip() for x in ("id", "machine_ref")): + raise CalculationConfigError("id and machine_ref must be non-empty strings") + channels = [] + extruders = entry["extruders"] + if not isinstance(extruders, list) or not extruders: + raise CalculationConfigError("extruders must be a non-empty list") + for extruder in extruders: + if not isinstance(extruder, dict) or set(extruder) != { + "id", + "total_rate_signal_ref", + "channels", + }: + raise CalculationConfigError("Invalid extruder fields") + if not all( + isinstance(extruder[x], str) and extruder[x].strip() + for x in ("id", "total_rate_signal_ref") + ): + raise CalculationConfigError("Extruder id and total signal must be strings") + rows = extruder["channels"] + if not isinstance(rows, list) or len(rows) != 7: + raise CalculationConfigError("Each extruder requires exactly seven channels") + for row in rows: + allowed = { + "id", + "status_signal_ref", + "percentage_signal_ref", + "screw_speed_signal_ref", + "material_number_signal_ref", + } + if ( + not isinstance(row, dict) + or set(row) - allowed + or {"id", "status_signal_ref", "percentage_signal_ref", "screw_speed_signal_ref"} + - set(row) + ): + raise CalculationConfigError("Invalid dosing channel fields") + if not all( + isinstance(row[x], str) and row[x].strip() + for x in ( + "id", + "status_signal_ref", + "percentage_signal_ref", + "screw_speed_signal_ref", + ) + ): + raise CalculationConfigError("Dosing signal references must be non-empty strings") + material = row.get("material_number_signal_ref") + if material is not None and (not isinstance(material, str) or not material.strip()): + raise CalculationConfigError("material signal must be a string or null") + channels.append( + ChannelDefinition( + extruder["id"], + row["id"], + extruder["total_rate_signal_ref"], + row["status_signal_ref"], + row["percentage_signal_ref"], + row["screw_speed_signal_ref"], + material, + ) + ) + if len({(c.extruder, c.channel) for c in channels}) != len(channels): + raise CalculationConfigError("Extruder/channel identities must be unique") + names = entry.get("material_names", {}) + if not isinstance(names, dict) or any( + not isinstance(k, str) or not isinstance(v, str) for k, v in names.items() + ): + raise CalculationConfigError("material_names must map strings to strings") + return ChannelMaterialCalculationConfig( + entry["id"], + entry["machine_ref"], + finite_number(entry["max_sample_gap_seconds"], "max_sample_gap_seconds", positive=True), + finite_number(entry["percentage_tolerance"], "percentage_tolerance"), + tuple(channels), + names, + ) diff --git a/src/production_analytics/calculations/channel_material.py b/src/production_analytics/calculations/channel_material.py new file mode 100644 index 0000000..59bc5ed --- /dev/null +++ b/src/production_analytics/calculations/channel_material.py @@ -0,0 +1,99 @@ +"""Pure channel-level dosing calculations, shared by extruder implementations.""" + +from collections.abc import Iterable +from dataclasses import dataclass +from datetime import datetime +from math import isfinite + +from production_analytics.calculations.material_consumption import MaterialSample + +MATERIAL_ASSIGNED = "ASSIGNED" +MATERIAL_UNAVAILABLE = "UNAVAILABLE" +MATERIAL_AMBIGUOUS = "AMBIGUOUS" +MATERIAL_UNMAPPED = "UNMAPPED" + + +@dataclass(frozen=True, slots=True) +class DosingChannelSample: + timestamp: datetime + total_rate_kg_per_hour: float + channel_status: bool + dosing_percentage: float + screw_speed: float + material_number: str | None = None + + +@dataclass(frozen=True, slots=True) +class ChannelRate: + sample: MaterialSample + material_number: str | None + material_name: str | None + material_mapping_status: str + + +def validate_active_percentage_sum( + samples: Iterable[DosingChannelSample], *, tolerance: float +) -> bool: + """Return whether active channel percentages close to 100; never normalize them.""" + values = tuple(samples) + if not isfinite(tolerance) or tolerance < 0: + raise ValueError("percentage tolerance must be finite and non-negative") + total = 0.0 + for value in values: + if not isfinite(value.dosing_percentage): + raise ValueError("dosing percentage must be finite") + if value.channel_status: + total += value.dosing_percentage + return abs(total - 100.0) <= tolerance + + +def channel_rates( + samples: Iterable[DosingChannelSample], + *, + material_names: dict[str, str] | None = None, + percentage_tolerance: float = 1.0, +) -> tuple[ChannelRate, ...]: + """Derive rates and conservative material identity for one aligned timestamp. + + A material value which occurs on more than one channel is deliberately not + attributed to either channel: PLC assignments are not reliable enough to + resolve that ambiguity from a previous sample. + """ + values = tuple(samples) + validate_active_percentage_sum(values, tolerance=percentage_tolerance) + counts = { + material_number: sum(sample.material_number == material_number for sample in values) + for material_number in (sample.material_number for sample in values) + if material_number is not None + } + names = material_names or {} + result = [] + for value in values: + if not isfinite(value.total_rate_kg_per_hour) or not isfinite(value.screw_speed): + raise ValueError("total rate and screw speed must be finite") + if value.material_number is None: + number, name, status = None, None, MATERIAL_UNAVAILABLE + elif counts[value.material_number] > 1: + number, name, status = None, None, MATERIAL_AMBIGUOUS + elif value.material_number not in names and names: + number, name, status = value.material_number, None, MATERIAL_UNMAPPED + else: + number, name, status = ( + value.material_number, + names.get(value.material_number), + MATERIAL_ASSIGNED, + ) + rate = ( + value.total_rate_kg_per_hour * value.dosing_percentage / 100 + if value.channel_status + else 0.0 + ) + result.append( + ChannelRate( + sample=MaterialSample(value.timestamp, rate, 1.0 if value.channel_status else 0.0), + material_number=number, + material_name=name, + material_mapping_status=status, + ) + ) + return tuple(result) diff --git a/src/production_analytics/enlyze/gateway.py b/src/production_analytics/enlyze/gateway.py index 6d8424a..111745e 100644 --- a/src/production_analytics/enlyze/gateway.py +++ b/src/production_analytics/enlyze/gateway.py @@ -3,7 +3,9 @@ from dataclasses import dataclass from datetime import UTC, datetime from math import isfinite +from typing import TYPE_CHECKING +from production_analytics.calculations.channel_material import DosingChannelSample from production_analytics.calculations.material_consumption import ( MaterialSample, area_application_rate_kg_per_hour, @@ -11,6 +13,9 @@ from production_analytics.calculations.material_consumption import ( ) from production_analytics.enlyze.exploration import ExplorationClient +if TYPE_CHECKING: + from production_analytics.service.channel_material import ChannelDefinition + @dataclass(frozen=True, slots=True) class EnlyzeProductionRun: @@ -22,6 +27,24 @@ class EnlyzeProductionRun: end: datetime | None +@dataclass(frozen=True, slots=True) +class EnlyzeDowntime: + """ENLYZE downtime source record. ``end`` and ``reason`` are deliberately nullable.""" + + uuid: str + machine_id: str + source_type: str + start: datetime + end: datetime | None + comment: str | None + reason_id: str | None + reason_name: str | None + reason_description: str | None + reason_group: str | None + reason_category: str | None + source_updated_at: datetime | None + + class EnlyzeApiGateway: def __init__(self, client: ExplorationClient) -> None: self._client = client @@ -60,11 +83,16 @@ class EnlyzeApiGateway: end = _parse_timestamp(item["end"]) if item["end"] is not None else None if end is not None and end < start: raise ValueError("run end must not precede start") - runs.append(EnlyzeProductionRun( - uuid=item["uuid"], machine_id=item["machine"], - product_id=item.get("product"), production_order=item["production_order"], - start=start, end=end, - )) + runs.append( + EnlyzeProductionRun( + uuid=item["uuid"], + machine_id=item["machine"], + product_id=item.get("product"), + production_order=item["production_order"], + start=start, + end=end, + ) + ) next_cursor = None if "metadata" in response.body: metadata = response.body["metadata"] @@ -85,6 +113,45 @@ class EnlyzeApiGateway: except (KeyError, TypeError, ValueError) as exc: raise ValueError(f"Invalid production-run response on page {page}: {exc}") from exc + def get_downtimes( + self, + machine_id: str, + *, + start: datetime | None = None, + ) -> list[EnlyzeDowntime]: + """Retrieve machine downtimes and follow ENLYZE cursor pagination. + + ``start`` is an inclusive source-time lookback; callers use it for reconciliation, + not as a claim that records before it cannot subsequently be edited. + """ + if start is not None and (start.tzinfo is None or start.utcoffset() is None): + raise ValueError("start must be timezone-aware") + params = {"machine": machine_id} + if start is not None: + params["start"] = start.astimezone(UTC).isoformat() + result: list[EnlyzeDowntime] = [] + cursors: set[str] = set() + page = 1 + while True: + response = self._client.get("/v2/downtimes", params) + try: + if not isinstance(response.body, dict) or not isinstance( + response.body["data"], list + ): + raise ValueError("body data must be a list") + for item in response.body["data"]: + result.append(_parse_downtime(item, machine_id)) + cursor = _next_cursor(response.body) + if cursor is None: + return result + if cursor in cursors: + raise ValueError("next_cursor has already been followed") + cursors.add(cursor) + params = {**params, "cursor": cursor} + page += 1 + except (KeyError, TypeError, ValueError) as exc: + raise ValueError(f"Invalid downtime response on page {page}: {exc}") from exc + def get_material_samples( self, *, @@ -152,7 +219,9 @@ class EnlyzeApiGateway: if nominal_width_m is None or specific_discharge_kg_per_rev_m is None: raise ValueError("rotational discharge requires width and factor") rate = rotational_discharge_rate_kg_per_hour( - sources, nominal_width_m, specific_discharge_kg_per_rev_m, + sources, + nominal_width_m, + specific_discharge_kg_per_rev_m, ) elif application_variable_ids: if nominal_width_m is None: @@ -167,15 +236,21 @@ class EnlyzeApiGateway: from production_analytics.calculations.material_application import ( rotational_application_g_m2, ) + application = rotational_application_g_m2( - sources, specific_discharge_kg_per_rev_m, gate, + sources, + specific_discharge_kg_per_rev_m, + gate, process_application_gate_threshold, ) - samples.append(MaterialSample( - timestamp=_parse_timestamp(record[time_index]), - material_rate_kg_per_hour=float(rate), gate_value=float(gate), - application_g_m2=application, - )) + samples.append( + MaterialSample( + timestamp=_parse_timestamp(record[time_index]), + material_rate_kg_per_hour=float(rate), + gate_value=float(gate), + application_g_m2=application, + ) + ) except (TypeError, ValueError) as exc: raise ValueError(f"malformed record {index}: {exc}") from exc next_cursor = None @@ -198,6 +273,199 @@ class EnlyzeApiGateway: except (KeyError, TypeError, ValueError) as exc: raise ValueError(f"Invalid timeseries response on page {page}: {exc}") from exc + def get_dosing_channel_samples( + self, + *, + machine_id: str, + channels: tuple["ChannelDefinition", ...], + start: datetime, + end: datetime, + ) -> dict[str, list[DosingChannelSample]]: + """Read configured dosing channels through the gateway only. + + Configurations name stable PLC origins; ENLYZE UUIDs are resolved at poll + time and must be unique, preventing a silently stale variable mapping. + """ + if ( + start.tzinfo is None + or start.utcoffset() is None + or end.tzinfo is None + or end.utcoffset() is None + ): + raise ValueError("dosing channel window must be timezone-aware") + if end < start: + raise ValueError("end must not precede start") + resolved = self.resolve_dosing_channel_signal_refs( + machine_id=machine_id, + channels=channels, + ) + variable_ids = list( + dict.fromkeys( + variable_id + for references in resolved.values() + for variable_id in references + if variable_id is not None + ) + ) + body = { + "machine": machine_id, + "start": start.astimezone(UTC).isoformat(), + "end": end.astimezone(UTC).isoformat(), + "variables": [{"uuid": x} for x in variable_ids], + } + result = {key: [] for key in resolved} + cursors: set[str] = set() + while True: + response = self._client.post_json("/v2/timeseries", body) + try: + data = response.body["data"] + columns = data["columns"] + if not isinstance(columns, list) or any( + columns.count(x) != 1 for x in ("time", *variable_ids) + ): + raise ValueError("required columns must occur exactly once") + if not isinstance(data["records"], list): + raise ValueError("records must be a list") + for row in data["records"]: + if not isinstance(row, list) or len(row) != len(columns): + raise ValueError("record must match columns") + timestamp = _parse_timestamp(row[columns.index("time")]) + for key, (total, status, percentage, screw, material) in resolved.items(): + sample = self._dosing_sample( + row, columns, timestamp, total, status, percentage, screw, material + ) + if sample is not None: + result[key].append(sample) + cursor = _next_cursor(response.body) + if cursor is None: + return result + if cursor in cursors: + raise ValueError("next_cursor has already been followed") + cursors.add(cursor) + body = {**body, "cursor": cursor} + except (KeyError, TypeError, ValueError) as exc: + raise ValueError(f"Invalid dosing timeseries response: {exc}") from exc + + def resolve_dosing_channel_signal_refs( + self, + *, + machine_id: str, + channels: tuple["ChannelDefinition", ...], + ) -> dict[str, tuple[str, str, str, str, str | None]]: + """Resolve configured PLC origins without reading samples or writing state.""" + origins = self._variables_by_origin(machine_id) + resolved: dict[str, tuple[str, str, str, str, str | None]] = {} + for channel in channels: + refs = ( + channel.total_rate_signal_ref, + channel.status_signal_ref, + channel.percentage_signal_ref, + channel.screw_speed_signal_ref, + ) + ids = tuple(self._unique_origin(origins, ref) for ref in refs) + material = ( + self._unique_origin(origins, channel.material_number_signal_ref) + if channel.material_number_signal_ref + else None + ) + key = channel.channel + "\0" + channel.extruder + resolved[key] = (*ids, material) + return resolved + + @staticmethod + def _dosing_sample( + row: list[object], + columns: list[object], + timestamp: datetime, + total: str, + status: str, + percentage: str, + screw: str, + material: str | None, + ) -> DosingChannelSample | None: + def numeric(variable: str) -> float: + value = row[columns.index(variable)] + if ( + isinstance(value, bool) + or not isinstance(value, (int, float)) + or not isfinite(value) + ): + raise ValueError("dosing values must be finite numeric") + return float(value) + + for variable in (total, percentage, screw): + if row[columns.index(variable)] is None: + return None + + status_value = row[columns.index(status)] + if status_value is None: + enabled = False + elif isinstance(status_value, bool): + enabled = status_value + elif ( + isinstance(status_value, (int, float)) + and not isinstance(status_value, bool) + and isfinite(status_value) + ): + enabled = status_value != 0 + else: + raise ValueError("dosing status must be boolean or numeric") + material_value = None if material is None else row[columns.index(material)] + if material_value is not None and not isinstance(material_value, (str, int, float)): + raise ValueError("material assignment must be scalar or null") + return DosingChannelSample( + timestamp, + numeric(total), + enabled, + numeric(percentage), + numeric(screw), + None if material_value is None else str(material_value), + ) + + def _variables_by_origin(self, machine_id: str) -> dict[str, list[str]]: + params = {"machine": machine_id} + result: dict[str, list[str]] = {} + cursors: set[str] = set() + while True: + response = self._client.get("/v2/variables", params) + try: + items = response.body["data"] + if not isinstance(items, list): + raise ValueError("data must be a list") + for item in items: + uuid = item["uuid"] + if not isinstance(uuid, str): + raise ValueError("variable uuid must be a string") + details = item.get("details") + if not isinstance(details, dict): + raise ValueError("variable details must be an object") + origin_identifier = details.get("origin_identifier") + if origin_identifier is None: + # Derived ENLYZE variables do not have a PLC origin. + continue + if not isinstance(origin_identifier, dict): + raise ValueError("variable origin_identifier must be an object") + code = origin_identifier.get("code") + if not isinstance(code, str): + raise ValueError("variable origin must be a string") + result.setdefault(code, []).append(uuid) + cursor = _next_cursor(response.body) + if cursor is None: + return result + if cursor in cursors: + raise ValueError("next_cursor has already been followed") + cursors.add(cursor) + params = {**params, "cursor": cursor} + except (KeyError, TypeError, ValueError) as exc: + raise ValueError(f"Invalid variable response: {exc}") from exc + + @staticmethod + def _unique_origin(origins: dict[str, list[str]], origin: str) -> str: + matches = origins.get(origin, []) + if len(matches) != 1: + raise ValueError(f"Configured PLC origin {origin!r} is unavailable or ambiguous") + return matches[0] + def _parse_timestamp(value: str) -> datetime: if not isinstance(value, str): @@ -206,3 +474,67 @@ def _parse_timestamp(value: str) -> datetime: if timestamp.tzinfo is None or timestamp.utcoffset() is None: raise ValueError("timestamp must be timezone-aware") return timestamp.astimezone(UTC) + + +def _next_cursor(body: dict[object, object]) -> str | None: + if "metadata" not in body: + return None + metadata = body["metadata"] + if not isinstance(metadata, dict) or "next_cursor" not in metadata: + raise ValueError("metadata must be an object with a next_cursor field") + cursor = metadata["next_cursor"] + if cursor is not None and (not isinstance(cursor, str) or not cursor): + raise ValueError("next_cursor must be null or a non-empty string") + return cursor + + +def _nullable_text(value: object, name: str) -> str | None: + if value is not None and not isinstance(value, str): + raise ValueError(f"{name} must be a string or null") + return value + + +def _parse_downtime(item: object, machine_id: str) -> EnlyzeDowntime: + if not isinstance(item, dict): + raise ValueError("downtime must be an object") + for key in ("uuid", "machine", "type", "start"): + if not isinstance(item.get(key), str) or not item[key]: + raise ValueError(f"{key} must be a non-empty string") + if item["machine"] != machine_id: + raise ValueError("downtime machine does not match requested machine") + start = _parse_timestamp(item["start"]) + end_value = item.get("end") + end = _parse_timestamp(end_value) if end_value is not None else None + if end is not None and end < start: + raise ValueError("downtime end must not precede start") + reason = item.get("reason") + if reason is not None and not isinstance(reason, dict): + raise ValueError("reason must be an object or null") + if reason is not None: + for key in ("uuid", "name", "category"): + if not isinstance(reason.get(key), str) or not reason[key]: + raise ValueError(f"reason {key} must be a non-empty string") + updated = item.get("updated") + if updated is not None and not isinstance(updated, dict): + raise ValueError("updated must be an object or null") + updated_at = None + if updated is not None: + updated_at = _parse_timestamp(updated.get("timestamp")) + return EnlyzeDowntime( + uuid=item["uuid"], + machine_id=machine_id, + source_type=item["type"], + start=start, + end=end, + comment=_nullable_text(item.get("comment"), "comment"), + reason_id=None if reason is None else reason["uuid"], + reason_name=None if reason is None else reason["name"], + reason_description=None + if reason is None + else _nullable_text(reason.get("description"), "reason description"), + reason_group=None + if reason is None + else _nullable_text(reason.get("group"), "reason group"), + reason_category=None if reason is None else reason["category"], + source_updated_at=updated_at, + ) diff --git a/src/production_analytics/service/channel_material.py b/src/production_analytics/service/channel_material.py new file mode 100644 index 0000000..1869782 --- /dev/null +++ b/src/production_analytics/service/channel_material.py @@ -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 diff --git a/src/production_analytics/service/channel_material_runner.py b/src/production_analytics/service/channel_material_runner.py new file mode 100644 index 0000000..207776d --- /dev/null +++ b/src/production_analytics/service/channel_material_runner.py @@ -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) diff --git a/src/production_analytics/service/channel_material_runtime.py b/src/production_analytics/service/channel_material_runtime.py new file mode 100644 index 0000000..1cc2076 --- /dev/null +++ b/src/production_analytics/service/channel_material_runtime.py @@ -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, + ) diff --git a/src/production_analytics/service/channel_material_state_store.py b/src/production_analytics/service/channel_material_state_store.py new file mode 100644 index 0000000..64f8d90 --- /dev/null +++ b/src/production_analytics/service/channel_material_state_store.py @@ -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) diff --git a/src/production_analytics/service/downtime.py b/src/production_analytics/service/downtime.py new file mode 100644 index 0000000..77ec570 --- /dev/null +++ b/src/production_analytics/service/downtime.py @@ -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 diff --git a/src/production_analytics/service/downtime_runner.py b/src/production_analytics/service/downtime_runner.py new file mode 100644 index 0000000..91a92bd --- /dev/null +++ b/src/production_analytics/service/downtime_runner.py @@ -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 diff --git a/src/production_analytics/service/downtime_runtime.py b/src/production_analytics/service/downtime_runtime.py new file mode 100644 index 0000000..73544db --- /dev/null +++ b/src/production_analytics/service/downtime_runtime.py @@ -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 + ) diff --git a/src/production_analytics/service/postgres_channel_material.py b/src/production_analytics/service/postgres_channel_material.py new file mode 100644 index 0000000..e5d2e37 --- /dev/null +++ b/src/production_analytics/service/postgres_channel_material.py @@ -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, + ), + ) diff --git a/src/production_analytics/service/postgres_downtime.py b/src/production_analytics/service/postgres_downtime.py new file mode 100644 index 0000000..b7eb9f8 --- /dev/null +++ b/src/production_analytics/service/postgres_downtime.py @@ -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, + ), + ) diff --git a/tests/conftest.py b/tests/conftest.py new file mode 100644 index 0000000..61c04c0 --- /dev/null +++ b/tests/conftest.py @@ -0,0 +1,21 @@ +"""Shared test adaptations for the production PostgreSQL schema.""" + +import re +from pathlib import Path + +import pytest + + +@pytest.fixture(scope="session") +def sqlite_schema() -> str: + """Keep PostgreSQL DDL authoritative; adapt JSON storage and clock for SQLite. + + SQLite stores serialized JSON as text. These tests exercise schema and INSERT + semantics, not PostgreSQL JSONB operators or TimescaleDB behavior. + """ + schema = (Path(__file__).resolve().parents[1] / "db" / "schema.sql").read_text( + encoding="utf-8" + ) + schema = re.sub(r"::jsonb\b", "", schema) + schema = re.sub(r"\bjsonb\b", "text", schema) + return schema.replace("DEFAULT now()", "DEFAULT CURRENT_TIMESTAMP") diff --git a/tests/test_channel_material.py b/tests/test_channel_material.py new file mode 100644 index 0000000..601aa3d --- /dev/null +++ b/tests/test_channel_material.py @@ -0,0 +1,158 @@ +from datetime import UTC, datetime, timedelta +from pathlib import Path +from unittest.mock import Mock + +import pytest + +from production_analytics.calculations.channel_config import load_channel_material_calculation +from production_analytics.calculations.channel_material import ( + MATERIAL_AMBIGUOUS, + MATERIAL_UNAVAILABLE, + DosingChannelSample, + channel_rates, + validate_active_percentage_sum, +) +from production_analytics.enlyze.gateway import EnlyzeApiGateway, EnlyzeProductionRun +from production_analytics.service.channel_material import ( + ChannelDefinition, + ChannelMaterialPollingService, +) +from production_analytics.service.channel_material_state_store import JsonChannelMaterialStateStore + +NOW = datetime(2026, 9, 1, tzinfo=UTC) +E1_CONFIG = ( + Path(__file__).resolve().parents[1] / "config/e1-channel-material-consumption.example.yaml" +) + + +def sample(*, on=True, percent=100, material=None, timestamp=NOW): + return DosingChannelSample(timestamp, 120, on, percent, 10, material) + + +def test_off_channel_with_nonzero_displayed_percentage_contributes_zero(): + assert channel_rates((sample(on=False, percent=40),))[0].sample.material_rate_kg_per_hour == 0 + + +def test_e1_config_dry_run_resolves_every_signal_through_gateway_without_sampling(): + config = load_channel_material_calculation(E1_CONFIG) + origins = { + reference + for channel in config.channels + for reference in ( + channel.total_rate_signal_ref, + channel.status_signal_ref, + channel.percentage_signal_ref, + channel.screw_speed_signal_ref, + ) + } + client = Mock() + client.get.return_value.body = { + "data": [ + { + "uuid": f"variable-{index}", + "details": {"origin_identifier": {"code": origin}}, + } + for index, origin in enumerate(sorted(origins)) + ] + } + + resolved = EnlyzeApiGateway(client).resolve_dosing_channel_signal_refs( + machine_id=config.machine_ref, + channels=config.channels, + ) + + assert config.machine_ref == "8302e3d1-b1e5-42f1-8540-615eb6c73e08" + assert len(config.channels) == 14 + assert len(origins) == 44 + assert len(resolved) == 14 + assert all(len(refs) == 5 and refs[-1] is None for refs in resolved.values()) + client.get.assert_called_once_with("/v2/variables", {"machine": config.machine_ref}) + client.post_json.assert_not_called() + + +def test_active_percentages_sum_to_100_and_deviation_is_reported_not_normalized(): + active = (sample(percent=40), sample(percent=60)) + assert validate_active_percentage_sum(active, tolerance=0.1) + assert not validate_active_percentage_sum( + (sample(percent=40), sample(percent=50)), tolerance=0.1 + ) + assert [x.sample.material_rate_kg_per_hour for x in channel_rates(active)] == [48, 72] + + +def test_null_and_duplicate_material_assignments_are_not_inferred(): + assert ( + channel_rates((sample(material=None),))[0].material_mapping_status == MATERIAL_UNAVAILABLE + ) + rates = channel_rates((sample(material="123"), sample(material="123"))) + assert all(rate.material_mapping_status == MATERIAL_AMBIGUOUS for rate in rates) + assert all(rate.material_number is None and rate.material_name is None for rate in rates) + + +def test_run_change_and_json_restart_preserve_channel_total(tmp_path): + channels = (ChannelDefinition("Ex1", "C1", "rate", "on", "pct", "rpm"),) + store, gateway = JsonChannelMaterialStateStore(tmp_path), Mock() + gateway.get_production_runs.return_value = [] + gateway.get_open_production_run.return_value = EnlyzeProductionRun( + "r1", "m", None, "order", NOW, None + ) + key = "C1\0Ex1" + gateway.get_dosing_channel_samples.return_value = { + key: [sample(timestamp=NOW), sample(timestamp=NOW + timedelta(seconds=10))] + } + first = ChannelMaterialPollingService( + gateway=gateway, + state_store=store, + machine_id="m", + channels=channels, + max_sample_gap_seconds=20, + ) + assert first.poll_once(now=NOW + timedelta(seconds=10))[ + 0 + ].state.cumulative_consumption_kg == pytest.approx(10 / 30) + gateway.get_open_production_run.return_value = EnlyzeProductionRun( + "r2", "m", None, "order", NOW + timedelta(seconds=20), None + ) + gateway.get_dosing_channel_samples.return_value = { + key: [ + sample(timestamp=NOW + timedelta(seconds=20)), + sample(timestamp=NOW + timedelta(seconds=30)), + ] + } + restarted = ChannelMaterialPollingService( + gateway=gateway, + state_store=store, + machine_id="m", + channels=channels, + max_sample_gap_seconds=20, + ) + assert restarted.poll_once(now=NOW + timedelta(seconds=30))[ + 0 + ].state.cumulative_consumption_kg == pytest.approx(20 / 30) + + +def test_bootstrap_replays_closed_runs_without_integrating_the_inter_run_gap(tmp_path): + channels = (ChannelDefinition("Ex1", "C1", "rate", "on", "pct", "rpm"),) + closed = EnlyzeProductionRun("r1", "m", None, "order", NOW, NOW + timedelta(seconds=10)) + current = EnlyzeProductionRun("r2", "m", None, "order", NOW + timedelta(seconds=20), None) + gateway = Mock() + gateway.get_open_production_run.return_value = current + gateway.get_production_runs.return_value = [closed, current] + key = "C1\0Ex1" + gateway.get_dosing_channel_samples.side_effect = [ + {key: [sample(timestamp=NOW), sample(timestamp=NOW + timedelta(seconds=10))]}, + { + key: [ + sample(timestamp=NOW + timedelta(seconds=20)), + sample(timestamp=NOW + timedelta(seconds=30)), + ] + }, + ] + result = ChannelMaterialPollingService( + gateway=gateway, + state_store=JsonChannelMaterialStateStore(tmp_path), + machine_id="m", + channels=channels, + max_sample_gap_seconds=20, + ).poll_once(now=NOW + timedelta(seconds=30)) + assert result is not None + assert result[0].state.cumulative_consumption_kg == pytest.approx(20 / 30) diff --git a/tests/test_channel_material_runner.py b/tests/test_channel_material_runner.py new file mode 100644 index 0000000..ecbb828 --- /dev/null +++ b/tests/test_channel_material_runner.py @@ -0,0 +1,141 @@ +import io +from datetime import UTC, datetime +from pathlib import Path +from unittest.mock import Mock, patch + +from production_analytics.calculations.channel_config import load_channel_material_calculation +from production_analytics.calculations.material_consumption import MaterialIntegrationState +from production_analytics.cli.__main__ import _parser, main +from production_analytics.enlyze.gateway import EnlyzeProductionRun +from production_analytics.service.channel_material import ChannelPollResult +from production_analytics.service.channel_material_runner import ChannelMaterialPollingRunner +from production_analytics.service.channel_material_runtime import build_channel_material_runner + +EXAMPLE = ( + Path(__file__).resolve().parents[1] / "config/e1-channel-material-consumption.example.yaml" +) +NOW = datetime(2026, 9, 1, tzinfo=UTC) + + +def test_cli_channel_material_arguments(): + args = _parser().parse_args( + [ + "run", + "channel-material-poll", + "--config", + str(EXAMPLE), + "--poll-interval-seconds", + "15", + "--state-directory", + "state", + "--secrets-file", + "secrets.env", + "--once", + ] + ) + assert args.command == "channel-material-poll" + assert args.config == EXAMPLE + assert args.poll_interval_seconds == 15 + assert args.state_directory == Path("state") + assert args.secrets_file == Path("secrets.env") + assert args.once + + +def test_runtime_constructs_channel_service(tmp_path, monkeypatch): + monkeypatch.setenv("ENLYZE_BASE_URL", "https://example.invalid/api/") + for key, value in { + "HOST": "localhost", + "PORT": "5432", + "DB": "analytics", + "USER": "user", + "PASSWORD": "environment-password", + }.items(): + monkeypatch.setenv(f"POSTGRES_{key}", value) + secrets = tmp_path / "secrets.env" + secrets.write_text("ENLYZE_API_KEY=secret-from-file\nPOSTGRES_PASSWORD=file-password\n") + calculation = load_channel_material_calculation(EXAMPLE) + with patch( + "production_analytics.service.channel_material_runtime.ChannelMaterialPollingService" + ) as service: + runner = build_channel_material_runner( + calculation, + poll_interval_seconds=7, + state_directory=tmp_path / "state", + secrets_file=secrets, + ) + kwargs = service.call_args.kwargs + assert kwargs["machine_id"] == calculation.machine_ref + assert kwargs["channels"] == calculation.channels + assert kwargs["max_sample_gap_seconds"] == calculation.max_sample_gap_seconds + assert kwargs["percentage_tolerance"] == calculation.percentage_tolerance + assert kwargs["material_names"] == calculation.material_names + assert kwargs["gateway"]._client._settings.api_key == "secret-from-file" + assert kwargs["state_store"].directory == tmp_path / "state" + assert runner.service is service.return_value + assert runner.interval == 7 + assert runner.snapshot_writer.settings.password == "file-password" + + +def test_once_executes_one_cycle_and_exits(capsys): + with patch( + "production_analytics.service.channel_material_runtime.build_channel_material_runner" + ) as build: + assert main(["run", "channel-material-poll", "--config", str(EXAMPLE), "--once"]) == 0 + build.return_value.run_once.assert_called_once_with() + build.return_value.run.assert_not_called() + assert capsys.readouterr().err == "" + + +def test_no_snapshots_are_written_without_eligible_production_run(): + writer = Mock() + output = io.StringIO() + runner = ChannelMaterialPollingRunner( + Mock(poll_once=Mock(return_value=None)), + machine_id="machine", + calculation_id="calculation", + snapshot_writer=writer, + poll_interval_seconds=10, + clock=lambda: NOW, + stdout=output, + ) + assert not runner.run_once() + writer.write.assert_not_called() + assert "snapshots=0" in output.getvalue() + + +def test_one_cycle_writes_each_channel_snapshot(): + run = EnlyzeProductionRun("run", "machine", None, "order", NOW, None) + first = ChannelPollResult( + run, "Ex1", "C1", MaterialIntegrationState(1.5, 10), "42", "Material", "MAPPED", True + ) + second = ChannelPollResult( + run, "Ex1", "C2", MaterialIntegrationState(2.5, 10), None, None, "AMBIGUOUS", False + ) + writer = Mock() + output = io.StringIO() + runner = ChannelMaterialPollingRunner( + Mock(poll_once=Mock(return_value=(first, second))), + machine_id="machine", + calculation_id="calculation", + snapshot_writer=writer, + poll_interval_seconds=10, + clock=lambda: NOW, + stdout=output, + ) + assert runner.run_once() + assert writer.write.call_count == 2 + assert writer.write.call_args_list[0].kwargs == { + "timestamp": NOW, + "calculation_id": "calculation", + "machine_id": "machine", + "extruder": "Ex1", + "channel": "C1", + "production_order": "order", + "run_id": "run", + "cumulative_consumption_kg": 1.5, + "material_number": "42", + "material_name": "Material", + "material_mapping_status": "MAPPED", + "percentage_sum_valid": True, + } + assert "channel snapshots=2 written" in output.getvalue() diff --git a/tests/test_downtime.py b/tests/test_downtime.py new file mode 100644 index 0000000..292e1ac --- /dev/null +++ b/tests/test_downtime.py @@ -0,0 +1,140 @@ +from datetime import UTC, datetime + +from production_analytics.enlyze.gateway import EnlyzeDowntime +from production_analytics.erp import CurrentWorkplaceStatus +from production_analytics.service.downtime import ( + DowntimeCategory, + DowntimeReconciliationService, + ProductionOrderBoundary, +) + +NOW = datetime(2026, 9, 16, 12, tzinfo=UTC) + + +def source(*, identifier="d", end=NOW, category=None): + return EnlyzeDowntime( + identifier, + "machine", + "THRESHOLD", + datetime(2026, 9, 16, 10, tzinfo=UTC), + end, + None, + "reason" if category else None, + "reason" if category else None, + None, + None, + category, + None, + ) + + +def status(order="1", *, remaining=3, good=7, target=10, feedback=NOW): + return CurrentWorkplaceStatus( + "WP", order, None, None, feedback, target, good, remaining, None, None + ) + + +class Gateway: + def __init__(self, events): + self.events = events + + def get_downtimes(self, machine_id, *, start=None): + return self.events + + +class Erp: + def __init__(self, current): + self.current = current + + def get_current_workplace_status(self, workplace): + return self.current + + +class Repo: + def __init__(self): + self.boundary = None + self.events = {} + self.closures = [] + + def current_boundary(self, machine_id): + return self.boundary + + def save_boundary(self, boundary): + self.boundary = boundary + + def close_order_attribution(self, machine_id, production_order, ended_at): + self.closures.append((production_order, ended_at)) + for event in self.events.values(): + if event.production_order == production_order and event.attributed_end is None: + event # database-specific clipping is covered by SQL contract + + def upsert(self, event): + self.events[event.external_id] = event + + +def service(repo, events, current): + return DowntimeReconciliationService( + Gateway(events), + Erp(current), + repo, + machine_id="machine", + workplace="WP", + production_order_format="FA-{production_order}", + ) + + +def test_categories_open_and_reconciliation_upsert() -> None: + repo = Repo() + runner = service( + repo, + [ + source(identifier="p", category="PLANNED"), + source(identifier="u", category="UNPLANNED"), + source(identifier="x", end=None), + ], + status(), + ) + assert runner.reconcile_once(NOW) == 3 + assert {key: event.category for key, event in repo.events.items()} == { + "p": DowntimeCategory.PLANNED, + "u": DowntimeCategory.UNPLANNED, + "x": DowntimeCategory.UNKNOWN, + } + assert repo.events["x"].source_end is None + # Same UUID is overwritten, never duplicated; delayed classification is accepted. + runner.gateway.events = [source(identifier="x", end=NOW, category="PLANNED")] + runner.reconcile_once(NOW) + assert len(repo.events) == 3 + assert repo.events["x"].category is DowntimeCategory.PLANNED + assert repo.events["x"].source_end == NOW + + +def test_completion_clips_attribution_and_new_order_closes_previous() -> None: + repo = Repo() + repo.boundary = ProductionOrderBoundary( + "machine", + "FA-1", + datetime(2026, 9, 16, 9, tzinfo=UTC), + None, + ) + event = source(identifier="open", end=None) + first = service(repo, [event], status(remaining=0, feedback=NOW)) + first.reconcile_once(NOW) + assert repo.boundary is not None and repo.boundary.ended_at == NOW + assert repo.events["open"].production_order == "FA-1" + assert repo.events["open"].attributed_end == NOW + repo.boundary = ProductionOrderBoundary( + "machine", + "FA-1", + datetime(2026, 9, 16, 8, tzinfo=UTC), + None, + ) + service(repo, [], status(order="2", feedback=NOW)).reconcile_once(NOW) + assert repo.closures == [("FA-1", NOW), ("FA-1", NOW)] + assert repo.boundary.production_order == "FA-2" + + +def test_unknown_is_not_unplanned_and_no_schedule_is_used() -> None: + repo = Repo() + service(repo, [source(end=datetime(2026, 9, 20, tzinfo=UTC))], status()).reconcile_once(NOW) + assert repo.events["d"].category is DowntimeCategory.UNKNOWN diff --git a/tests/test_enlyze_gateway.py b/tests/test_enlyze_gateway.py index 7636f4d..2409238 100644 --- a/tests/test_enlyze_gateway.py +++ b/tests/test_enlyze_gateway.py @@ -6,6 +6,52 @@ import pytest from production_analytics.enlyze.gateway import EnlyzeApiGateway +def test_downtime_pagination_and_nullable_reason() -> None: + client = Mock() + base = { + "machine": "machine-1", + "type": "THRESHOLD", + "comment": None, + "start": "2026-09-03T04:00:00Z", + "end": None, + "reason": None, + "updated": None, + } + client.get.side_effect = [ + Mock(body={"data": [{**base, "uuid": "d1"}], "metadata": {"next_cursor": "next"}}), + Mock( + body={ + "data": [ + { + **base, + "uuid": "d2", + "end": "2026-09-03T05:00:00Z", + "reason": { + "uuid": "r", + "name": "Break", + "description": None, + "group": "General", + "category": "PLANNED", + }, + } + ], + "metadata": {"next_cursor": None}, + } + ), + ] + result = EnlyzeApiGateway(client).get_downtimes( + "machine-1", + start=datetime(2026, 9, 3, tzinfo=UTC), + ) + assert [item.uuid for item in result] == ["d1", "d2"] + assert result[0].end is None and result[0].reason_category is None + assert result[1].reason_category == "PLANNED" + assert client.get.call_args_list[1].args == ( + "/v2/downtimes", + {"machine": "machine-1", "start": "2026-09-03T00:00:00+00:00", "cursor": "next"}, + ) + + def test_get_open_production_run_returns_current_run() -> None: client = Mock() client.get.return_value.body = { @@ -138,11 +184,19 @@ def test_get_material_samples_uses_column_names_not_fixed_positions() -> None: @pytest.fixture def timeseries(): client = Mock() - client.post_json.return_value.body = {"data": { - "columns": ["time", "rate", "gate"], "records": [], - }} - args = dict(machine_id="m", rate_variable_id="rate", gate_variable_id="gate", - start=datetime(2026, 9, 3, tzinfo=UTC), end=datetime(2026, 9, 4, tzinfo=UTC)) + client.post_json.return_value.body = { + "data": { + "columns": ["time", "rate", "gate"], + "records": [], + } + } + args = dict( + machine_id="m", + rate_variable_id="rate", + gate_variable_id="gate", + start=datetime(2026, 9, 3, tzinfo=UTC), + end=datetime(2026, 9, 4, tzinfo=UTC), + ) return EnlyzeApiGateway(client), client, args @@ -176,9 +230,19 @@ def test_missing_required_column(timeseries, column) -> None: gateway.get_material_samples(**args) -@pytest.mark.parametrize("record", [[], {}, ["bad", 1, 1], [None, 1, 1], - ["2026-09-03T00:00:00", 1, 1], ["2026-09-03T00:00:00Z", None, 1], - ["2026-09-03T00:00:00Z", 1, float("inf")], ["2026-09-03T00:00:00Z", True, 1]]) +@pytest.mark.parametrize( + "record", + [ + [], + {}, + ["bad", 1, 1], + [None, 1, 1], + ["2026-09-03T00:00:00", 1, 1], + ["2026-09-03T00:00:00Z", None, 1], + ["2026-09-03T00:00:00Z", 1, float("inf")], + ["2026-09-03T00:00:00Z", True, 1], + ], +) def test_malformed_records_fail_clearly(timeseries, record) -> None: gateway, client, args = timeseries client.post_json.return_value.body["data"]["records"] = [record] @@ -188,15 +252,21 @@ def test_malformed_records_fail_clearly(timeseries, record) -> None: def test_multiple_open_runs_rejected() -> None: client = Mock() - run = dict(uuid="r", machine="m", production_order="o", start="2026-09-03T00:00:00Z", - end=None) + run = dict(uuid="r", machine="m", production_order="o", start="2026-09-03T00:00:00Z", end=None) client.get.return_value.body = {"data": [run, {**run, "uuid": "r2"}]} with pytest.raises(ValueError, match="multiple open Production Runs"): EnlyzeApiGateway(client).get_open_production_run("m") -@pytest.mark.parametrize("item", [{}, None, {"end": None}, - dict(uuid="r", machine="m", production_order="o", start="bad", end=None)]) +@pytest.mark.parametrize( + "item", + [ + {}, + None, + {"end": None}, + dict(uuid="r", machine="m", production_order="o", start="bad", end=None), + ], +) def test_malformed_run_fails_clearly(item) -> None: client = Mock() client.get.return_value.body = {"data": [item]} @@ -219,16 +289,24 @@ def test_timeseries_with_null_cursor(timeseries) -> None: def test_timeseries_pagination_resolves_columns_per_page(timeseries) -> None: gateway, client, args = timeseries client.post_json.side_effect = [ - Mock(body={ - "data": {"columns": ["time", "rate", "gate"], - "records": [["2026-09-03T00:00:00Z", 10, 2]]}, - "metadata": {"next_cursor": "continuation"}, - }), - Mock(body={ - "data": {"columns": ["gate", "time", "rate"], - "records": [[3, "2026-09-03T01:00:00Z", 20]]}, - "metadata": {"next_cursor": None}, - }), + Mock( + body={ + "data": { + "columns": ["time", "rate", "gate"], + "records": [["2026-09-03T00:00:00Z", 10, 2]], + }, + "metadata": {"next_cursor": "continuation"}, + } + ), + Mock( + body={ + "data": { + "columns": ["gate", "time", "rate"], + "records": [[3, "2026-09-03T01:00:00Z", 20]], + }, + "metadata": {"next_cursor": None}, + } + ), ] samples = gateway.get_material_samples(**args) assert [(s.timestamp, s.material_rate_kg_per_hour, s.gate_value) for s in samples] == [ @@ -237,16 +315,32 @@ def test_timeseries_pagination_resolves_columns_per_page(timeseries) -> None: ] assert client.post_json.call_count == 2 first, second = client.post_json.call_args_list - assert first.args == ("/v2/timeseries", { - "machine": "m", "start": args["start"].isoformat(), "end": args["end"].isoformat(), - "variables": [{"uuid": "rate"}, {"uuid": "gate"}], - }) + assert first.args == ( + "/v2/timeseries", + { + "machine": "m", + "start": args["start"].isoformat(), + "end": args["end"].isoformat(), + "variables": [{"uuid": "rate"}, {"uuid": "gate"}], + }, + ) assert second.args == ("/v2/timeseries", {**first.args[1], "cursor": "continuation"}) -@pytest.mark.parametrize("metadata", [None, [], "bad", {}, - {"next_cursor": ""}, {"next_cursor": 1}, {"next_cursor": False}, - {"next_cursor": []}, {"next_cursor": {}}]) +@pytest.mark.parametrize( + "metadata", + [ + None, + [], + "bad", + {}, + {"next_cursor": ""}, + {"next_cursor": 1}, + {"next_cursor": False}, + {"next_cursor": []}, + {"next_cursor": {}}, + ], +) def test_invalid_pagination_metadata(timeseries, metadata) -> None: gateway, client, args = timeseries client.post_json.return_value.body["metadata"] = metadata @@ -267,25 +361,38 @@ def test_repeated_pagination_cursor(timeseries, cursors) -> None: assert client.post_json.call_count == len(cursors) -@pytest.mark.parametrize("data, error", [ - ({"columns": ["time", "rate", "gate"], "records": [[]]}, "malformed record 0"), - ({"columns": ["time", "rate", "rate", "gate"], "records": []}, "required column"), - ({"columns": ["time", "rate"], "records": []}, "required column"), - ({"columns": None, "records": []}, "columns must be a list"), - ({"columns": ["time", "rate", "gate"], "records": None}, "records must be a list"), - ({"columns": ["time", "rate", "gate"], - "records": [["2026-09-03T01:00:00", 10, 2]]}, "timezone-aware"), - ({"columns": ["time", "rate", "gate"], - "records": [["2026-09-03T01:00:00Z", float("nan"), 2]]}, "finite"), - ({"columns": ["time", "rate", "gate"], - "records": [["2026-09-03T01:00:00Z", 10, True]]}, "numeric"), -]) +@pytest.mark.parametrize( + "data, error", + [ + ({"columns": ["time", "rate", "gate"], "records": [[]]}, "malformed record 0"), + ({"columns": ["time", "rate", "rate", "gate"], "records": []}, "required column"), + ({"columns": ["time", "rate"], "records": []}, "required column"), + ({"columns": None, "records": []}, "columns must be a list"), + ({"columns": ["time", "rate", "gate"], "records": None}, "records must be a list"), + ( + {"columns": ["time", "rate", "gate"], "records": [["2026-09-03T01:00:00", 10, 2]]}, + "timezone-aware", + ), + ( + { + "columns": ["time", "rate", "gate"], + "records": [["2026-09-03T01:00:00Z", float("nan"), 2]], + }, + "finite", + ), + ( + {"columns": ["time", "rate", "gate"], "records": [["2026-09-03T01:00:00Z", 10, True]]}, + "numeric", + ), + ], +) def test_malformed_later_page(timeseries, data, error) -> None: gateway, client, args = timeseries first_body = client.post_json.return_value.body first_body["metadata"] = {"next_cursor": "continuation"} client.post_json.side_effect = [ - Mock(body=first_body), Mock(body={"data": data, "metadata": {"next_cursor": None}}), + Mock(body=first_body), + Mock(body={"data": data, "metadata": {"next_cursor": None}}), ] with pytest.raises(ValueError, match=f"Invalid timeseries response on page 2:.*{error}"): gateway.get_material_samples(**args) @@ -294,8 +401,13 @@ def test_malformed_later_page(timeseries, data, error) -> None: def test_production_runs_follow_pages_and_find_open_run() -> None: client = Mock() - closed = dict(uuid="closed", machine="m", production_order="opaque-order", - start="2026-09-01T00:00:00Z", end="2026-09-01T01:00:00Z") + closed = dict( + uuid="closed", + machine="m", + production_order="opaque-order", + start="2026-09-01T00:00:00Z", + end="2026-09-01T01:00:00Z", + ) opened = {**closed, "uuid": "open", "end": None} pages = [ Mock(body={"data": [closed], "metadata": {"next_cursor": "a"}}), @@ -314,10 +426,20 @@ def test_production_runs_follow_pages_and_find_open_run() -> None: assert gateway.get_open_production_run("m") == runs[1] -@pytest.mark.parametrize("metadata", [ - None, [], "bad", {}, {"next_cursor": ""}, {"next_cursor": 1}, - {"next_cursor": False}, {"next_cursor": []}, {"next_cursor": {}}, -]) +@pytest.mark.parametrize( + "metadata", + [ + None, + [], + "bad", + {}, + {"next_cursor": ""}, + {"next_cursor": 1}, + {"next_cursor": False}, + {"next_cursor": []}, + {"next_cursor": {}}, + ], +) def test_production_run_invalid_pagination(metadata) -> None: client = Mock() client.get.return_value.body = {"data": [], "metadata": metadata} @@ -329,25 +451,50 @@ def test_production_run_invalid_pagination(metadata) -> None: def test_production_run_cyclic_pagination(cursors) -> None: client = Mock() client.get.side_effect = [ - Mock(body={"data": [], "metadata": {"next_cursor": cursor}}) - for cursor in cursors + Mock(body={"data": [], "metadata": {"next_cursor": cursor}}) for cursor in cursors ] with pytest.raises(ValueError, match="next_cursor has already been followed"): EnlyzeApiGateway(client).get_production_runs("m") assert client.get.call_count == len(cursors) -@pytest.mark.parametrize("body", [ - None, [], {}, {"data": {}}, {"data": [None]}, - {"data": [dict(uuid="r", machine="wrong", production_order="o", - start="2026-09-03T00:00:00Z", end=None)]}, - {"data": [dict(uuid="r", machine="m", production_order="o", - start="2026-09-03T00:00:00Z", end="2026-09-02T00:00:00Z")]}, -]) +@pytest.mark.parametrize( + "body", + [ + None, + [], + {}, + {"data": {}}, + {"data": [None]}, + { + "data": [ + dict( + uuid="r", + machine="wrong", + production_order="o", + start="2026-09-03T00:00:00Z", + end=None, + ) + ] + }, + { + "data": [ + dict( + uuid="r", + machine="m", + production_order="o", + start="2026-09-03T00:00:00Z", + end="2026-09-02T00:00:00Z", + ) + ] + }, + ], +) def test_production_run_malformed_later_page(body) -> None: client = Mock() client.get.side_effect = [ - Mock(body={"data": [], "metadata": {"next_cursor": "a"}}), Mock(body=body), + Mock(body={"data": [], "metadata": {"next_cursor": "a"}}), + Mock(body=body), ] with pytest.raises(ValueError, match="Invalid production-run response on page 2"): EnlyzeApiGateway(client).get_production_runs("m") diff --git a/tests/test_material_application.py b/tests/test_material_application.py index 5cf9a7b..441fae2 100644 --- a/tests/test_material_application.py +++ b/tests/test_material_application.py @@ -117,11 +117,11 @@ def test_persistence_before_checkpoint_retries_and_disjoint_runs(): store.save.assert_called_once() -def test_sql_persistence_idempotent_and_inactive_omitted(monkeypatch): +def test_sql_persistence_idempotent_and_inactive_omitted(monkeypatch, sqlite_schema): import sqlite3 database = sqlite3.connect(':memory:') - database.executescript(Path('db/schema.sql').read_text()) + database.executescript(sqlite_schema) connection = MagicMock() cursor = connection.__enter__.return_value.cursor.return_value.__enter__.return_value cursor.executemany.side_effect = lambda sql, rows: database.executemany( diff --git a/tests/test_material_efficiency_persistence.py b/tests/test_material_efficiency_persistence.py index 99ef2a8..bbdf3e7 100644 --- a/tests/test_material_efficiency_persistence.py +++ b/tests/test_material_efficiency_persistence.py @@ -48,10 +48,10 @@ SETTINGS = PostgresSettings("localhost", 5432, "analytics", "writer", "secret") @pytest.fixture -def storage(monkeypatch): - # Real schema and INSERT semantics, with only driver transport adapted for SQLite. +def storage(monkeypatch, sqlite_schema): + # Production schema and INSERT semantics, adapted for the SQLite dialect/transport. database = sqlite3.connect(":memory:") - database.executescript(Path("db/schema.sql").read_text()) + database.executescript(sqlite_schema) connection = MagicMock() def execute(sql, params): diff --git a/tests/test_schema.py b/tests/test_schema.py new file mode 100644 index 0000000..a67e8bb --- /dev/null +++ b/tests/test_schema.py @@ -0,0 +1,61 @@ +"""Regression coverage for SQLite initialization from production DDL.""" + +import json +import sqlite3 + +import pytest + + +def test_schema_initialization_and_power_meter_json(sqlite_schema): + with sqlite3.connect(":memory:") as database: + database.executescript(sqlite_schema) + database.executescript(sqlite_schema) + tables = { + row[0] for row in database.execute("SELECT name FROM sqlite_master WHERE type='table'") + } + assert { + "material_consumption_snapshots", + "material_efficiency_snapshots", + "material_application_snapshots", + "channel_material_consumption_snapshots", + "production_downtime_events", + "production_order_attribution_state", + "power_meter_readings", + "power_meter_monthly_reports", + } <= tables + + insert = """INSERT INTO power_meter_readings + (meter, timestamp, energy_total_kwh, energy_delta_kwh, + active_power_total_kw, grid_frequency_hz) + VALUES (?, ?, 100, 1, 4, 50)""" + database.execute(insert, ("B2", "2026-09-01T00:00:00Z")) + phases = {"L1": 1.25, "L2": 2.75} + database.execute(insert, ("B2", "2026-09-01T00:01:00Z")) + database.execute( + "UPDATE power_meter_readings SET phase_active_power_kw=? WHERE timestamp=?", + (json.dumps(phases), "2026-09-01T00:01:00Z"), + ) + rows = database.execute( + "SELECT phase_active_power_kw FROM power_meter_readings ORDER BY timestamp" + ).fetchall() + assert [json.loads(row[0]) for row in rows] == [{}, phases] + with pytest.raises(sqlite3.IntegrityError): + database.execute(insert, ("B2", "2026-09-01T00:00:00Z")) + + summary = {"consumption_kwh": 42.5, "first_timestamp": None, "outages": 0} + database.execute( + """INSERT INTO power_meter_monthly_reports (meter, month, report_path, summary) + VALUES (?, ?, ?, ?)""", + ("B2", "2026-09", "/reports/B2.pdf", json.dumps(summary)), + ) + path, stored_summary, generated_at = database.execute( + "SELECT report_path, summary, generated_at FROM power_meter_monthly_reports" + ).fetchone() + assert path == "/reports/B2.pdf" + assert json.loads(stored_summary) == summary + assert generated_at is not None + with pytest.raises(sqlite3.IntegrityError): + database.execute( + """INSERT INTO power_meter_monthly_reports (meter, month, report_path, summary) + VALUES ('B2', '2026-10', '/reports/B2.pdf', NULL)""" + )