From aa8d3dc2674c159aa110151205c99ac5b436dba1 Mon Sep 17 00:00:00 2001 From: Victor Fraile Garcia Date: Mon, 20 Jul 2026 12:42:08 +0200 Subject: [PATCH] feat(ucepsa): persist Shop Floor context incident episodes --- ...opfloor_context_incident_history_v0.3.7.md | 187 ++++ ...sa-shopfloor-context-health.dashboard.json | 362 +++++++- ...context_incident_history_dashboard_v037.sh | 240 ++++++ ...opfloor_context_incident_collector_v037.sh | 125 +++ ...hopfloor_context_incident_history_v037.sql | 263 ++++++ ...hopfloor_context_incident_history_v037.sql | 738 ++++++++++++++++ .../tools/shopfloor_context_incident_admin.py | 369 ++++++++ .../shopfloor_context_incident_collector.py | 806 ++++++++++++++++++ 8 files changed, 3088 insertions(+), 2 deletions(-) create mode 100644 ucepsa/edge-oee-demo/docs/runbooks/105_shopfloor_context_incident_history_v0.3.7.md create mode 100755 ucepsa/edge-oee-demo/ops/deploy_shopfloor_context_incident_history_dashboard_v037.sh create mode 100755 ucepsa/edge-oee-demo/ops/manage_shopfloor_context_incident_collector_v037.sh create mode 100644 ucepsa/edge-oee-demo/ops/validate_shopfloor_context_incident_history_v037.sql create mode 100644 ucepsa/edge-oee-demo/sql/versions/109_ucepsa_shopfloor_context_incident_history_v037.sql create mode 100755 ucepsa/edge-oee-demo/tools/shopfloor_context_incident_admin.py create mode 100755 ucepsa/edge-oee-demo/tools/shopfloor_context_incident_collector.py diff --git a/ucepsa/edge-oee-demo/docs/runbooks/105_shopfloor_context_incident_history_v0.3.7.md b/ucepsa/edge-oee-demo/docs/runbooks/105_shopfloor_context_incident_history_v0.3.7.md new file mode 100644 index 0000000..575773a --- /dev/null +++ b/ucepsa/edge-oee-demo/docs/runbooks/105_shopfloor_context_incident_history_v0.3.7.md @@ -0,0 +1,187 @@ +# UCEPSA — Historial persistente de incidencias Shop Floor v0.3.7 + +## Objetivo + +Convertir incidencias vivas en episodios persistentes: + +```text +incidencia aparece +→ abre un episodio + +incidencia continúa +→ actualiza last_seen_at + +incidencia desaparece durante 45 s +→ cierra el episodio + +incidencia reaparece más tarde +→ abre un episodio nuevo +``` + +No genera una fila cada 15 segundos. + +## Alcance + +La versión: + +- lee la vista viva `li_shopfloor_context_incidents_v1`; +- guarda episodios SHADOW; +- permite revisión humana o técnica; +- amplía el dashboard existente; +- no crea sesiones `LEGACY_ERP`; +- no modifica Odoo; +- no atribuye pérdidas; +- no alimenta el Ledger oficial. + +`first_seen_at` es la primera observación realizada por MESAVAULT después de +activar el colector. No debe interpretarse como el inicio físico exacto de una +situación que ya estuviera activa. + +## Objetos + +```text +shopfloor_context_incident_policies +shopfloor_context_incident_episodes +shopfloor_context_incident_reviews +shopfloor_context_incident_collector_state + +li_shopfloor_context_incident_history_v1 +li_shopfloor_context_incident_backlog_v1 +li_shopfloor_context_incident_daily_summary_v1 +li_shopfloor_context_incident_collector_health_v1 +``` + +## Clasificaciones de revisión + +```text +AUTHORIZED_LEGACY_PRODUCTION +TEST_OR_SETUP +RESIDUAL_MATERIAL +ODOO_START_OMITTED +DATA_ISSUE +NOT_RELEVANT +OTHER +``` + +Una revisión solo documenta lo sucedido. No abre un contexto ni convierte el +episodio en pérdida. + +## Despliegue + +### Migración + +```bash +FILE=/srv/mesavault/40-clients/ucepsa/edge-oee-demo/sql/versions/109_ucepsa_shopfloor_context_incident_history_v037.sql + +docker exec -i mv_ucepsa_postgres_hot sh -lc \ + 'psql -v ON_ERROR_STOP=1 -U "$POSTGRES_USER" -d "$POSTGRES_DB"' \ + < "$FILE" +``` + +### Colector + +Copiar las herramientas al runtime: + +```bash +install -m 0644 \ + /srv/mesavault/40-clients/ucepsa/edge-oee-demo/tools/shopfloor_context_incident_collector.py \ + /srv/mesavault/edge-oee-ucepsa/tools/ + +install -m 0644 \ + /srv/mesavault/40-clients/ucepsa/edge-oee-demo/tools/shopfloor_context_incident_admin.py \ + /srv/mesavault/edge-oee-ucepsa/tools/ +``` + +Arranque: + +```bash +MANAGER=/srv/mesavault/40-clients/ucepsa/edge-oee-demo/ops/manage_shopfloor_context_incident_collector_v037.sh + +"$MANAGER" dry-run +"$MANAGER" once +"$MANAGER" start +"$MANAGER" status +``` + +### Validación + +Esperar al menos 20 segundos: + +```bash +sleep 20 + +VALIDATION=/srv/mesavault/40-clients/ucepsa/edge-oee-demo/ops/validate_shopfloor_context_incident_history_v037.sql + +docker exec -i mv_ucepsa_postgres_hot sh -lc \ + 'psql -v ON_ERROR_STOP=1 -U "$POSTGRES_USER" -d "$POSTGRES_DB"' \ + < "$VALIDATION" +``` + +Valores obligatorios: + +```text +last_seen_before_start = 0 +end_before_start = 0 +open_with_end = 0 +closed_without_end = 0 +official_episode_rows = 0 +duplicate_open_fingerprints = 0 +live_incidents_without_open_episode = 0 +expired_open_episodes_not_live = 0 +review_trigger_mismatches = 0 +official_operator_queue_rows = 0 +``` + +## Revisión de un episodio + +Listar backlog: + +```bash +"$MANAGER" admin list +``` + +Ver detalle: + +```bash +"$MANAGER" admin show --episode-id 1 +``` + +Ejemplo de clasificación, solo después de confirmación humana: + +```bash +"$MANAGER" admin review \ + --episode-id 1 \ + --classification AUTHORIZED_LEGACY_PRODUCTION \ + --reviewed-by "Juan Pablo" \ + --notes "Producción confirmada con el sistema anterior." +``` + +No ejecutar esa clasificación automáticamente para CORT-00. + +## Dashboard + +El despliegue mantiene el mismo UID y añade tres paneles: + +```text +Incidencias pendientes de revisión +Historial de las últimas 24 horas +Resumen diario de los últimos 7 días +``` + +```bash +read -rsp "Token temporal de Grafana: " GRAFANA_API_TOKEN +echo +export GRAFANA_API_TOKEN + +DASH=/srv/mesavault/40-clients/ucepsa/edge-oee-demo/ops/deploy_shopfloor_context_incident_history_dashboard_v037.sh + +"$DASH" deploy +"$DASH" validate + +unset GRAFANA_API_TOKEN +``` + +## Workorder 340 + +Esta versión puede desplegarse con el workorder 340 todavía abierto. Su cierre +sigue siendo una validación independiente del ciclo de sesión Shop Floor. diff --git a/ucepsa/edge-oee-demo/grafana/dashboards/ucepsa-shopfloor-context-health.dashboard.json b/ucepsa/edge-oee-demo/grafana/dashboards/ucepsa-shopfloor-context-health.dashboard.json index 9415ec1..e97645d 100644 --- a/ucepsa/edge-oee-demo/grafana/dashboards/ucepsa-shopfloor-context-health.dashboard.json +++ b/ucepsa/edge-oee-demo/grafana/dashboards/ucepsa-shopfloor-context-health.dashboard.json @@ -15,7 +15,7 @@ } ] }, - "description": "Supervisión SHADOW de máquina física, orden Odoo, sesión Shop Floor, sincronizador y calidad temporal. Presentación operativa v0.3.6.2. No alimenta el Ledger oficial.", + "description": "Supervisión SHADOW de máquina física, orden Odoo, sesión Shop Floor, sincronizador, calidad temporal e historial persistente de incidencias. No alimenta el Ledger oficial.", "editable": true, "fiscalYearStartMonth": 0, "graphTooltip": 1, @@ -1258,6 +1258,362 @@ ], "title": "Calidad de sesiones", "type": "table" + }, + { + "datasource": { + "type": "postgres", + "uid": "bfnbcasbm6hhca" + }, + "fieldConfig": { + "defaults": { + "custom": { + "align": "auto", + "cellOptions": { + "type": "auto" + }, + "filterable": true, + "inspect": false + }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": null + } + ] + } + }, + "overrides": [ + { + "matcher": { + "id": "byName", + "options": "Severidad" + }, + "properties": [ + { + "id": "mappings", + "value": [ + { + "type": "value", + "options": { + "WARNING": { + "text": "WARNING", + "color": "orange", + "index": 0 + }, + "CRITICAL": { + "text": "CRITICAL", + "color": "red", + "index": 1 + } + } + } + ] + }, + { + "id": "custom.cellOptions", + "value": { + "type": "color-background" + } + }, + { + "id": "custom.width", + "value": 110 + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "Situación" + }, + "properties": [ + { + "id": "custom.width", + "value": 420 + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "Acción" + }, + "properties": [ + { + "id": "custom.width", + "value": 520 + } + ] + } + ] + }, + "gridPos": { + "h": 10, + "w": 24, + "x": 0, + "y": 42 + }, + "id": 11, + "options": { + "cellHeight": "sm", + "footer": { + "countRows": false, + "fields": "", + "reducer": [ + "sum" + ], + "show": false + }, + "showHeader": true, + "sortBy": [] + }, + "pluginVersion": "11.0.0", + "targets": [ + { + "datasource": { + "type": "postgres", + "uid": "bfnbcasbm6hhca" + }, + "editorMode": "code", + "format": "table", + "rawQuery": true, + "rawSql": "SELECT\n backlog_rank AS \"Prioridad\",\n episode_id AS \"Episodio\",\n CASE episode_status\n WHEN 'OPEN' THEN 'Abierto'\n ELSE 'Cerrado'\n END AS \"Estado\",\n max_severity AS \"Severidad\",\n first_seen_at AS \"Primera observación\",\n episode_duration_min AS \"Duración min\",\n COALESCE(machine_id, '—') AS \"Máquina\",\n incident_title AS \"Situación\",\n COALESCE(production_order, '—') AS \"Orden\",\n CASE review_owner\n WHEN 'PRODUCTION' THEN 'Producción'\n ELSE 'Técnico'\n END AS \"Responsable\",\n latest_recommended_action AS \"Acción\"\nFROM\n mv_reports_ucepsa_prod\n .li_shopfloor_context_incident_backlog_v1\nORDER BY backlog_rank\nLIMIT 50", + "refId": "A" + } + ], + "title": "Incidencias pendientes de revisión", + "type": "table" + }, + { + "datasource": { + "type": "postgres", + "uid": "bfnbcasbm6hhca" + }, + "fieldConfig": { + "defaults": { + "custom": { + "align": "auto", + "cellOptions": { + "type": "auto" + }, + "filterable": true, + "inspect": false + }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": null + } + ] + } + }, + "overrides": [ + { + "matcher": { + "id": "byName", + "options": "Severidad" + }, + "properties": [ + { + "id": "mappings", + "value": [ + { + "type": "value", + "options": { + "WARNING": { + "text": "WARNING", + "color": "orange", + "index": 0 + }, + "CRITICAL": { + "text": "CRITICAL", + "color": "red", + "index": 1 + } + } + } + ] + }, + { + "id": "custom.cellOptions", + "value": { + "type": "color-background" + } + }, + { + "id": "custom.width", + "value": 110 + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "Situación" + }, + "properties": [ + { + "id": "custom.width", + "value": 420 + } + ] + } + ] + }, + "gridPos": { + "h": 10, + "w": 24, + "x": 0, + "y": 52 + }, + "id": 12, + "options": { + "cellHeight": "sm", + "footer": { + "countRows": false, + "fields": "", + "reducer": [ + "sum" + ], + "show": false + }, + "showHeader": true, + "sortBy": [] + }, + "pluginVersion": "11.0.0", + "targets": [ + { + "datasource": { + "type": "postgres", + "uid": "bfnbcasbm6hhca" + }, + "editorMode": "code", + "format": "table", + "rawQuery": true, + "rawSql": "SELECT\n episode_id AS \"Episodio\",\n CASE episode_status\n WHEN 'OPEN' THEN 'Abierto'\n ELSE 'Cerrado'\n END AS \"Estado\",\n max_severity AS \"Severidad\",\n first_seen_at AS \"Inicio observado\",\n ended_at AS \"Fin observado\",\n episode_duration_min AS \"Duración min\",\n COALESCE(machine_id, '—') AS \"Máquina\",\n incident_title AS \"Situación\",\n COALESCE(production_order, '—') AS \"Orden\",\n CASE review_status\n WHEN 'PENDING' THEN 'Pendiente'\n WHEN 'REVIEWED' THEN 'Revisado'\n WHEN 'DISMISSED' THEN 'Descartado'\n ELSE 'No requerido'\n END AS \"Revisión\",\n COALESCE(review_classification, '—')\n AS \"Clasificación\",\n COALESCE(reviewed_by, '—') AS \"Revisado por\"\nFROM\n mv_reports_ucepsa_prod\n .li_shopfloor_context_incident_history_v1\nWHERE first_seen_at >= now() - interval '24 hours'\nORDER BY first_seen_at DESC\nLIMIT 100", + "refId": "A" + } + ], + "title": "Historial de incidencias · últimas 24 horas", + "type": "table" + }, + { + "datasource": { + "type": "postgres", + "uid": "bfnbcasbm6hhca" + }, + "fieldConfig": { + "defaults": { + "custom": { + "align": "auto", + "cellOptions": { + "type": "auto" + }, + "filterable": true, + "inspect": false + }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": null + } + ] + } + }, + "overrides": [ + { + "matcher": { + "id": "byName", + "options": "Tipo" + }, + "properties": [ + { + "id": "custom.width", + "value": 460 + } + ] + }, + { + "matcher": { + "id": "byName", + "options": "Severidad máxima" + }, + "properties": [ + { + "id": "mappings", + "value": [ + { + "type": "value", + "options": { + "WARNING": { + "text": "WARNING", + "color": "orange", + "index": 0 + }, + "CRITICAL": { + "text": "CRITICAL", + "color": "red", + "index": 1 + } + } + } + ] + }, + { + "id": "custom.cellOptions", + "value": { + "type": "color-text" + } + } + ] + } + ] + }, + "gridPos": { + "h": 9, + "w": 24, + "x": 0, + "y": 62 + }, + "id": 13, + "options": { + "cellHeight": "sm", + "footer": { + "countRows": false, + "fields": "", + "reducer": [ + "sum" + ], + "show": false + }, + "showHeader": true, + "sortBy": [] + }, + "pluginVersion": "11.0.0", + "targets": [ + { + "datasource": { + "type": "postgres", + "uid": "bfnbcasbm6hhca" + }, + "editorMode": "code", + "format": "table", + "rawQuery": true, + "rawSql": "SELECT\n local_day AS \"Día\",\n COALESCE(machine_id, '—') AS \"Máquina\",\n incident_title AS \"Tipo\",\n episode_count AS \"Episodios\",\n open_episode_count AS \"Abiertos\",\n pending_review_count AS \"Pendientes\",\n observed_minutes AS \"Minutos observados\",\n average_episode_minutes AS \"Media min\",\n max_severity AS \"Severidad máxima\"\nFROM\n mv_reports_ucepsa_prod\n .li_shopfloor_context_incident_daily_summary_v1\nWHERE local_day >=\n (\n now()\n AT TIME ZONE 'Europe/Madrid'\n )::date - 6\nORDER BY\n local_day DESC,\n machine_id NULLS FIRST,\n incident_title", + "refId": "A" + } + ], + "title": "Resumen diario de incidencias · últimos 7 días", + "type": "table" } ], "refresh": "5s", @@ -1268,7 +1624,9 @@ "shop-floor", "context-health", "shadow", - "presentation-v0362" + "presentation-v0362", + "incident-history", + "v037" ], "templating": { "list": [] diff --git a/ucepsa/edge-oee-demo/ops/deploy_shopfloor_context_incident_history_dashboard_v037.sh b/ucepsa/edge-oee-demo/ops/deploy_shopfloor_context_incident_history_dashboard_v037.sh new file mode 100755 index 0000000..7cb0b26 --- /dev/null +++ b/ucepsa/edge-oee-demo/ops/deploy_shopfloor_context_incident_history_dashboard_v037.sh @@ -0,0 +1,240 @@ +#!/usr/bin/env bash +set -euo pipefail + +ACTION="${1:-deploy}" +GRAFANA_CONTAINER="${GRAFANA_CONTAINER:-mv_ucepsa_grafana}" +REPO_ROOT="${SHOPFLOOR_DASHBOARD_REPO_ROOT:-/srv/mesavault/40-clients/ucepsa/edge-oee-demo}" +DASHBOARD_FILE="${SHOPFLOOR_DASHBOARD_FILE:-$REPO_ROOT/grafana/dashboards/ucepsa-shopfloor-context-health.dashboard.json}" +FOLDER_UID="${SHOPFLOOR_DASHBOARD_FOLDER_UID:-ucepsa-mesavault}" +DASHBOARD_UID="ucepsa-shopfloor-context-health" +DATASOURCE_UID="bfnbcasbm6hhca" + +resolve_url() { + if [[ -n "${GRAFANA_URL:-}" ]]; then + printf '%s' "$GRAFANA_URL" + return + fi + + local port_line host_port container_ip + + port_line="$( + docker port "$GRAFANA_CONTAINER" 3000/tcp 2>/dev/null \ + | head -n 1 || true + )" + + if [[ -n "$port_line" ]]; then + host_port="${port_line##*:}" + printf 'http://127.0.0.1:%s' "$host_port" + return + fi + + container_ip="$( + docker inspect "$GRAFANA_CONTAINER" \ + --format '{{range .NetworkSettings.Networks}}{{.IPAddress}}{{end}}' + )" + + if [[ -z "$container_ip" ]]; then + echo "ERROR: no se pudo resolver Grafana." >&2 + exit 1 + fi + + printf 'http://%s:3000' "$container_ip" +} + +if [[ -z "${GRAFANA_API_TOKEN:-}" ]]; then + cat >&2 <<'MSG' +ERROR: falta GRAFANA_API_TOKEN. + + read -rsp "Token temporal de Grafana: " GRAFANA_API_TOKEN + echo + export GRAFANA_API_TOKEN +MSG + exit 1 +fi + +export GRAFANA_URL_RESOLVED +GRAFANA_URL_RESOLVED="$(resolve_url)" + +python3 - \ + "$ACTION" \ + "$DASHBOARD_FILE" \ + "$FOLDER_UID" \ + "$DASHBOARD_UID" \ + "$DATASOURCE_UID" <<'PY' +import json +import os +import sys +import urllib.error +import urllib.parse +import urllib.request + +( + action, + dashboard_file, + folder_uid, + dashboard_uid, + datasource_uid, +) = sys.argv[1:] + +base_url = os.environ["GRAFANA_URL_RESOLVED"].rstrip("/") +token = os.environ["GRAFANA_API_TOKEN"] + + +def request(method, path, payload=None): + body = None + headers = { + "Accept": "application/json", + "Authorization": f"Bearer {token}", + } + if payload is not None: + body = json.dumps( + payload, + ensure_ascii=False, + ).encode("utf-8") + headers["Content-Type"] = "application/json" + + req = urllib.request.Request( + base_url + path, + data=body, + headers=headers, + method=method, + ) + try: + with urllib.request.urlopen( + req, + timeout=20, + ) as response: + raw = response.read() + return ( + json.loads(raw.decode("utf-8")) + if raw + else None + ) + except urllib.error.HTTPError as exc: + detail = exc.read().decode( + "utf-8", + errors="replace", + ) + raise RuntimeError( + f"Grafana API {method} {path}: " + f"HTTP {exc.code}: {detail}" + ) from exc + + +health = request("GET", "/api/health") +if not health or health.get("database") != "ok": + raise RuntimeError( + f"Grafana no está saludable: {health!r}" + ) + +if action == "deploy": + with open( + dashboard_file, + encoding="utf-8", + ) as handle: + dashboard = json.load(handle) + + result = request( + "POST", + "/api/dashboards/db", + { + "dashboard": dashboard, + "folderUid": folder_uid, + "overwrite": True, + "message": ( + "MESAVAULT Shop Floor incident " + "history v0.3.7" + ), + }, + ) + print(json.dumps( + { + "status": "deployed", + "dashboard_uid": dashboard_uid, + "folder_uid": folder_uid, + "response": result, + }, + ensure_ascii=False, + indent=2, + )) + +elif action == "validate": + result = request( + "GET", + ( + "/api/dashboards/uid/" + + urllib.parse.quote(dashboard_uid) + ), + ) + dashboard = result["dashboard"] + panels = { + panel.get("id"): panel + for panel in dashboard.get("panels", []) + } + datasource_uids = sorted({ + panel.get("datasource", {}).get("uid") + for panel in dashboard.get("panels", []) + if isinstance( + panel.get("datasource"), + dict, + ) + and panel.get("datasource", {}).get("uid") + }) + errors = [] + + if len(dashboard.get("panels", [])) != 13: + errors.append("El dashboard no tiene 13 paneles") + if datasource_uids != [datasource_uid]: + errors.append( + f"Datasource inesperado: {datasource_uids!r}" + ) + expected = { + 11: "li_shopfloor_context_incident_backlog_v1", + 12: "li_shopfloor_context_incident_history_v1", + 13: "li_shopfloor_context_incident_daily_summary_v1", + } + for panel_id, fragment in expected.items(): + sql = ( + panels.get(panel_id, {}) + .get("targets", [{}])[0] + .get("rawSql", "") + ) + if fragment not in sql: + errors.append( + f"Panel {panel_id} no consulta {fragment}" + ) + + output = { + "status": "ok" if not errors else "error", + "dashboard_uid": dashboard.get("uid"), + "folder_uid": result.get("meta", {}).get( + "folderUid" + ), + "panel_count": len( + dashboard.get("panels", []) + ), + "datasource_uids": datasource_uids, + "version": dashboard.get("version"), + "url": result.get("meta", {}).get("url"), + "history_panels": { + "backlog": 11 in panels, + "history_24h": 12 in panels, + "daily_summary": 13 in panels, + }, + "errors": errors, + } + print(json.dumps( + output, + ensure_ascii=False, + indent=2, + )) + if errors: + raise SystemExit(1) + +else: + raise SystemExit( + "Uso: " + "deploy_shopfloor_context_incident_history_dashboard_v037.sh " + "{deploy|validate}" + ) +PY diff --git a/ucepsa/edge-oee-demo/ops/manage_shopfloor_context_incident_collector_v037.sh b/ucepsa/edge-oee-demo/ops/manage_shopfloor_context_incident_collector_v037.sh new file mode 100755 index 0000000..c9a56ca --- /dev/null +++ b/ucepsa/edge-oee-demo/ops/manage_shopfloor_context_incident_collector_v037.sh @@ -0,0 +1,125 @@ +#!/usr/bin/env bash +set -euo pipefail + +ACTION="${1:-status}" +ROOT="${SHOPFLOOR_RUNTIME_ROOT:-/srv/mesavault/edge-oee-ucepsa}" +CONTAINER="${SHOPFLOOR_INCIDENT_CONTAINER:-mv_ucepsa_shopfloor_context_incident_collector}" +NETWORK="${SHOPFLOOR_DOCKER_NETWORK:-mv_ucepsa_net}" +INTERVAL="${SHOPFLOOR_INCIDENT_INTERVAL_S:-15}" + +cd "$ROOT" + +image() { + docker inspect \ + mv_ucepsa_odoo_order_state_pg_sink \ + --format '{{.Config.Image}}' +} + +start_collector() { + local collector_image + collector_image="$(image)" + + docker rm -f "$CONTAINER" \ + >/dev/null 2>&1 || true + + docker run -d \ + --name "$CONTAINER" \ + --restart unless-stopped \ + --network "$NETWORK" \ + --env-file .env \ + -v "$ROOT:/app:ro" \ + -w /app \ + --entrypoint python \ + "$collector_image" \ + tools/shopfloor_context_incident_collector.py \ + --watch \ + --interval "$INTERVAL" + + docker ps \ + --format 'table {{.Names}}\t{{.Status}}' \ + | grep -E "$CONTAINER|NAMES" +} + +one_shot() { + local collector_image + collector_image="$(image)" + + docker run --rm -i \ + --network "$NETWORK" \ + --env-file .env \ + -v "$ROOT:/app:ro" \ + -w /app \ + --entrypoint python \ + "$collector_image" \ + tools/shopfloor_context_incident_collector.py \ + --once +} + +dry_run() { + local collector_image + collector_image="$(image)" + + docker run --rm -i \ + --network "$NETWORK" \ + --env-file .env \ + -v "$ROOT:/app:ro" \ + -w /app \ + --entrypoint python \ + "$collector_image" \ + tools/shopfloor_context_incident_collector.py \ + --once \ + --dry-run +} + +admin() { + shift + local collector_image + collector_image="$(image)" + + docker run --rm -i \ + --network "$NETWORK" \ + --env-file .env \ + -v "$ROOT:/app:ro" \ + -w /app \ + --entrypoint python \ + "$collector_image" \ + tools/shopfloor_context_incident_admin.py \ + "$@" +} + +case "$ACTION" in + start) + start_collector + ;; + stop) + docker stop "$CONTAINER" + ;; + restart) + docker rm -f "$CONTAINER" \ + >/dev/null 2>&1 || true + start_collector + ;; + status) + docker ps -a \ + --format 'table {{.Names}}\t{{.Status}}' \ + | grep -E "$CONTAINER|NAMES" + ;; + logs) + docker logs --tail 200 -f "$CONTAINER" + ;; + once) + one_shot + ;; + dry-run) + dry_run + ;; + admin) + admin "$@" + ;; + *) + echo \ + "Uso: $0 {start|stop|restart|status|logs|once|dry-run|admin ...}" \ + >&2 + exit 2 + ;; +esac diff --git a/ucepsa/edge-oee-demo/ops/validate_shopfloor_context_incident_history_v037.sql b/ucepsa/edge-oee-demo/ops/validate_shopfloor_context_incident_history_v037.sql new file mode 100644 index 0000000..be752a6 --- /dev/null +++ b/ucepsa/edge-oee-demo/ops/validate_shopfloor_context_incident_history_v037.sql @@ -0,0 +1,263 @@ +\pset pager off + +\echo '=== 1. Objetos v0.3.7 ===' + +SELECT + to_regclass( + 'mv_loss_intelligence.shopfloor_context_incident_policies' + ) AS policy_table, + to_regclass( + 'mv_loss_intelligence.shopfloor_context_incident_episodes' + ) AS episodes_table, + to_regclass( + 'mv_loss_intelligence.shopfloor_context_incident_reviews' + ) AS reviews_table, + to_regclass( + 'mv_loss_intelligence.shopfloor_context_incident_collector_state' + ) AS collector_state_table, + to_regclass( + 'mv_reports_ucepsa_prod.li_shopfloor_context_incident_history_v1' + ) AS history_view, + to_regclass( + 'mv_reports_ucepsa_prod.li_shopfloor_context_incident_backlog_v1' + ) AS backlog_view, + to_regclass( + 'mv_reports_ucepsa_prod.li_shopfloor_context_incident_daily_summary_v1' + ) AS daily_summary_view, + to_regclass( + 'mv_reports_ucepsa_prod.li_shopfloor_context_incident_collector_health_v1' + ) AS collector_health_view; + +\echo '=== 2. Política activa ===' + +SELECT * +FROM + mv_loss_intelligence + .shopfloor_context_incident_policies +WHERE tenant = 'ucepsa' + AND site = 'ucepsa_onpremise' + AND policy_code = + 'SHOPFLOOR_CONTEXT_INCIDENT_HISTORY_V1'; + +\echo '=== 3. Salud del colector ===' + +SELECT * +FROM + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_collector_health_v1; + +\echo '=== 4. Invariantes de episodios ===' + +SELECT + count(*) FILTER ( + WHERE last_seen_at < first_seen_at + ) AS last_seen_before_start, + count(*) FILTER ( + WHERE ended_at IS NOT NULL + AND ended_at < first_seen_at + ) AS end_before_start, + count(*) FILTER ( + WHERE episode_status = 'OPEN' + AND ended_at IS NOT NULL + ) AS open_with_end, + count(*) FILTER ( + WHERE episode_status = 'CLOSED' + AND ended_at IS NULL + ) AS closed_without_end, + count(*) FILTER ( + WHERE official_eligible + ) AS official_episode_rows +FROM + mv_loss_intelligence + .shopfloor_context_incident_episodes; + +\echo '=== 5. No hay dos episodios OPEN con la misma huella ===' + +WITH duplicated AS ( + SELECT + tenant, + site, + incident_fingerprint, + count(*) AS n + FROM + mv_loss_intelligence + .shopfloor_context_incident_episodes + WHERE episode_status = 'OPEN' + GROUP BY + tenant, + site, + incident_fingerprint + HAVING count(*) > 1 +) +SELECT count(*) AS duplicate_open_fingerprints +FROM duplicated; + +\echo '=== 6. Incidencias vivas con episodio OPEN ===' + +SELECT count(*) AS live_incidents_without_open_episode +FROM + mv_reports_ucepsa_prod + .li_shopfloor_context_incidents_v1 i +WHERE NOT EXISTS ( + SELECT 1 + FROM + mv_loss_intelligence + .shopfloor_context_incident_episodes e + WHERE e.tenant = 'ucepsa' + AND e.site = 'ucepsa_onpremise' + AND e.episode_status = 'OPEN' + AND e.incident_key = i.incident_key + AND e.incident_scope = i.incident_scope + AND e.incident_code = i.incident_code + AND e.machine_id + IS NOT DISTINCT FROM i.machine_id + AND e.context_key + IS NOT DISTINCT FROM i.context_key + AND e.production_order + IS NOT DISTINCT FROM i.production_order + AND e.odoo_workorder_id + IS NOT DISTINCT FROM i.odoo_workorder_id +); + +\echo '=== 7. Episodios OPEN vencidos que ya no están vivos ===' + +WITH policy AS ( + SELECT close_grace_s + FROM + mv_loss_intelligence + .shopfloor_context_incident_policies + WHERE tenant = 'ucepsa' + AND site = 'ucepsa_onpremise' + AND policy_code = + 'SHOPFLOOR_CONTEXT_INCIDENT_HISTORY_V1' + AND active + LIMIT 1 +) +SELECT count(*) AS expired_open_episodes_not_live +FROM + mv_loss_intelligence + .shopfloor_context_incident_episodes e +CROSS JOIN policy p +WHERE e.tenant = 'ucepsa' + AND e.site = 'ucepsa_onpremise' + AND e.episode_status = 'OPEN' + AND e.last_seen_at + < now() + - ( + p.close_grace_s + * interval '1 second' + ) + AND NOT EXISTS ( + SELECT 1 + FROM + mv_reports_ucepsa_prod + .li_shopfloor_context_incidents_v1 i + WHERE i.incident_key = e.incident_key + AND i.incident_scope = e.incident_scope + AND i.incident_code = e.incident_code + AND i.machine_id + IS NOT DISTINCT FROM e.machine_id + AND i.context_key + IS NOT DISTINCT FROM e.context_key + AND i.production_order + IS NOT DISTINCT FROM e.production_order + AND i.odoo_workorder_id + IS NOT DISTINCT FROM e.odoo_workorder_id + ); + +\echo '=== 8. Episodios actuales ===' + +SELECT + episode_id, + episode_status, + max_severity, + first_seen_at, + last_seen_at, + episode_duration_min, + machine_id, + incident_title, + production_order, + review_required, + review_owner, + review_status, + official_eligible +FROM + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_history_v1 +ORDER BY first_seen_at DESC +LIMIT 30; + +\echo '=== 9. Backlog de revisión ===' + +SELECT + backlog_rank, + episode_id, + episode_status, + max_severity, + machine_id, + incident_title, + production_order, + episode_duration_min, + review_owner, + review_status, + latest_recommended_action +FROM + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_backlog_v1 +ORDER BY backlog_rank +LIMIT 30; + +\echo '=== 10. Resumen diario ===' + +SELECT * +FROM + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_daily_summary_v1 +ORDER BY + local_day DESC, + machine_id NULLS FIRST, + incident_code +LIMIT 30; + +\echo '=== 11. Revisiones coherentes ===' + +SELECT count(*) AS review_trigger_mismatches +FROM + mv_loss_intelligence + .shopfloor_context_incident_reviews r +JOIN + mv_loss_intelligence + .shopfloor_context_incident_episodes e + ON e.episode_id = r.episode_id +WHERE r.review_id = ( + SELECT max(r2.review_id) + FROM + mv_loss_intelligence + .shopfloor_context_incident_reviews r2 + WHERE r2.episode_id = r.episode_id + ) + AND ( + e.review_status + IS DISTINCT FROM r.new_review_status + OR e.review_classification + IS DISTINCT FROM r.review_classification + OR e.reviewed_by + IS DISTINCT FROM r.reviewed_by + ); + +\echo '=== 12. Nada oficial ===' + +SELECT + ( + SELECT count(*) + FROM + mv_loss_intelligence + .shopfloor_context_incident_episodes + WHERE official_eligible + ) AS official_episode_rows, + ( + SELECT count(*) + FROM + mv_reports_ucepsa_prod + .li_operator_stop_queue_v2 + ) AS official_operator_queue_rows; diff --git a/ucepsa/edge-oee-demo/sql/versions/109_ucepsa_shopfloor_context_incident_history_v037.sql b/ucepsa/edge-oee-demo/sql/versions/109_ucepsa_shopfloor_context_incident_history_v037.sql new file mode 100644 index 0000000..a314ac5 --- /dev/null +++ b/ucepsa/edge-oee-demo/sql/versions/109_ucepsa_shopfloor_context_incident_history_v037.sql @@ -0,0 +1,738 @@ +BEGIN; + +CREATE TABLE IF NOT EXISTS +mv_loss_intelligence.shopfloor_context_incident_policies ( + tenant text NOT NULL, + site text NOT NULL, + policy_code text NOT NULL, + collector_interval_s integer NOT NULL DEFAULT 15 + CHECK (collector_interval_s > 0), + close_grace_s integer NOT NULL DEFAULT 45 + CHECK (close_grace_s >= collector_interval_s), + collector_delayed_after_s integer NOT NULL DEFAULT 45 + CHECK (collector_delayed_after_s > 0), + collector_stale_after_s integer NOT NULL DEFAULT 120 + CHECK ( + collector_stale_after_s + > collector_delayed_after_s + ), + active boolean NOT NULL DEFAULT true, + notes text, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now(), + PRIMARY KEY ( + tenant, + site, + policy_code + ) +); + +COMMENT ON TABLE +mv_loss_intelligence.shopfloor_context_incident_policies IS +'Política de captura persistente de incidencias vivas Shop Floor.'; + +INSERT INTO +mv_loss_intelligence.shopfloor_context_incident_policies ( + tenant, + site, + policy_code, + collector_interval_s, + close_grace_s, + collector_delayed_after_s, + collector_stale_after_s, + active, + notes +) +VALUES ( + 'ucepsa', + 'ucepsa_onpremise', + 'SHOPFLOOR_CONTEXT_INCIDENT_HISTORY_V1', + 15, + 45, + 45, + 120, + true, + 'Captura cada 15 s y cierra un episodio cuando la incidencia lleva 45 s sin aparecer. El inicio representa la primera observación de MESAVAULT, no necesariamente el inicio físico exacto.' +) +ON CONFLICT ( + tenant, + site, + policy_code +) +DO UPDATE SET + collector_interval_s = + EXCLUDED.collector_interval_s, + close_grace_s = + EXCLUDED.close_grace_s, + collector_delayed_after_s = + EXCLUDED.collector_delayed_after_s, + collector_stale_after_s = + EXCLUDED.collector_stale_after_s, + active = true, + notes = EXCLUDED.notes, + updated_at = now(); + +CREATE TABLE IF NOT EXISTS +mv_loss_intelligence.shopfloor_context_incident_episodes ( + episode_id bigserial PRIMARY KEY, + tenant text NOT NULL, + site text NOT NULL, + incident_fingerprint text NOT NULL, + incident_key text NOT NULL, + incident_scope text NOT NULL, + incident_code text NOT NULL, + machine_id text, + context_key text, + production_order text, + odoo_workorder_id bigint, + operator_name text, + initial_severity text NOT NULL, + current_severity text NOT NULL, + max_severity text NOT NULL, + episode_status text NOT NULL DEFAULT 'OPEN', + first_seen_at timestamptz NOT NULL, + last_seen_at timestamptz NOT NULL, + ended_at timestamptz, + closed_detected_at timestamptz, + close_reason text, + observation_count bigint NOT NULL DEFAULT 1 + CHECK (observation_count > 0), + initial_detail text, + latest_detail text, + initial_recommended_action text, + latest_recommended_action text, + review_required boolean NOT NULL DEFAULT true, + review_owner text NOT NULL, + review_status text NOT NULL DEFAULT 'PENDING', + review_classification text, + reviewed_by text, + reviewed_at timestamptz, + review_notes text, + source_payload_json jsonb NOT NULL + DEFAULT '{}'::jsonb, + official_eligible boolean NOT NULL DEFAULT false, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now(), + CONSTRAINT ck_shopfloor_context_episode_interval + CHECK ( + ended_at IS NULL + OR ended_at >= first_seen_at + ), + CONSTRAINT ck_shopfloor_context_episode_last_seen + CHECK ( + last_seen_at >= first_seen_at + ), + CONSTRAINT ck_shopfloor_context_episode_status + CHECK ( + episode_status IN ( + 'OPEN', + 'CLOSED' + ) + ), + CONSTRAINT ck_shopfloor_context_episode_severity + CHECK ( + initial_severity IN ( + 'WARNING', + 'CRITICAL' + ) + AND current_severity IN ( + 'WARNING', + 'CRITICAL' + ) + AND max_severity IN ( + 'WARNING', + 'CRITICAL' + ) + ), + CONSTRAINT ck_shopfloor_context_episode_review_owner + CHECK ( + review_owner IN ( + 'PRODUCTION', + 'TECHNICAL' + ) + ), + CONSTRAINT ck_shopfloor_context_episode_review_status + CHECK ( + review_status IN ( + 'PENDING', + 'REVIEWED', + 'DISMISSED', + 'NOT_REQUIRED' + ) + ), + CONSTRAINT ck_shopfloor_context_episode_review_classification + CHECK ( + review_classification IS NULL + OR review_classification IN ( + 'AUTHORIZED_LEGACY_PRODUCTION', + 'TEST_OR_SETUP', + 'RESIDUAL_MATERIAL', + 'ODOO_START_OMITTED', + 'DATA_ISSUE', + 'NOT_RELEVANT', + 'OTHER' + ) + ), + CONSTRAINT ck_shopfloor_context_episode_shadow_only + CHECK (official_eligible = false) +); + +CREATE UNIQUE INDEX IF NOT EXISTS +ux_shopfloor_context_incident_open_fingerprint +ON +mv_loss_intelligence.shopfloor_context_incident_episodes ( + tenant, + site, + incident_fingerprint +) +WHERE episode_status = 'OPEN'; + +CREATE INDEX IF NOT EXISTS +idx_shopfloor_context_incident_episode_time +ON +mv_loss_intelligence.shopfloor_context_incident_episodes ( + tenant, + site, + first_seen_at DESC +); + +CREATE INDEX IF NOT EXISTS +idx_shopfloor_context_incident_episode_machine +ON +mv_loss_intelligence.shopfloor_context_incident_episodes ( + tenant, + site, + machine_id, + first_seen_at DESC +); + +CREATE INDEX IF NOT EXISTS +idx_shopfloor_context_incident_episode_backlog +ON +mv_loss_intelligence.shopfloor_context_incident_episodes ( + tenant, + site, + review_required, + review_status, + episode_status, + max_severity, + first_seen_at +); + +CREATE TABLE IF NOT EXISTS +mv_loss_intelligence.shopfloor_context_incident_reviews ( + review_id bigserial PRIMARY KEY, + episode_id bigint NOT NULL + REFERENCES + mv_loss_intelligence + .shopfloor_context_incident_episodes ( + episode_id + ) + ON DELETE CASCADE, + previous_values jsonb NOT NULL, + new_review_status text NOT NULL, + review_classification text NOT NULL, + reviewed_by text NOT NULL, + reviewed_at timestamptz NOT NULL DEFAULT now(), + notes text, + evidence_json jsonb NOT NULL + DEFAULT '{}'::jsonb, + CONSTRAINT ck_shopfloor_context_review_status + CHECK ( + new_review_status IN ( + 'REVIEWED', + 'DISMISSED' + ) + ), + CONSTRAINT ck_shopfloor_context_review_classification + CHECK ( + review_classification IN ( + 'AUTHORIZED_LEGACY_PRODUCTION', + 'TEST_OR_SETUP', + 'RESIDUAL_MATERIAL', + 'ODOO_START_OMITTED', + 'DATA_ISSUE', + 'NOT_RELEVANT', + 'OTHER' + ) + ) +); + +CREATE INDEX IF NOT EXISTS +idx_shopfloor_context_incident_reviews_episode +ON +mv_loss_intelligence.shopfloor_context_incident_reviews ( + episode_id, + reviewed_at DESC +); + +CREATE OR REPLACE FUNCTION +mv_loss_intelligence.apply_shopfloor_context_incident_review() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +BEGIN + UPDATE + mv_loss_intelligence + .shopfloor_context_incident_episodes + SET + review_status = + NEW.new_review_status, + review_classification = + NEW.review_classification, + reviewed_by = + NEW.reviewed_by, + reviewed_at = + NEW.reviewed_at, + review_notes = + NEW.notes, + updated_at = + now() + WHERE episode_id = + NEW.episode_id; + + RETURN NEW; +END; +$$; + +DROP TRIGGER IF EXISTS +trg_apply_shopfloor_context_incident_review +ON +mv_loss_intelligence.shopfloor_context_incident_reviews; + +CREATE TRIGGER +trg_apply_shopfloor_context_incident_review +AFTER INSERT +ON +mv_loss_intelligence.shopfloor_context_incident_reviews +FOR EACH ROW +EXECUTE FUNCTION +mv_loss_intelligence.apply_shopfloor_context_incident_review(); + +CREATE TABLE IF NOT EXISTS +mv_loss_intelligence.shopfloor_context_incident_collector_state ( + tenant text NOT NULL, + site text NOT NULL, + collector_name text NOT NULL, + last_success_at timestamptz, + last_current_count integer, + last_opened_count integer, + last_updated_count integer, + last_closed_count integer, + last_error text, + updated_at timestamptz NOT NULL DEFAULT now(), + PRIMARY KEY ( + tenant, + site, + collector_name + ) +); + +CREATE OR REPLACE VIEW +mv_reports_ucepsa_prod.li_shopfloor_context_incident_history_v1 +AS +WITH review_counts AS ( + SELECT + episode_id, + COUNT(*) AS review_count + FROM + mv_loss_intelligence + .shopfloor_context_incident_reviews + GROUP BY episode_id +) +SELECT + e.episode_id, + ( + 'SHOPFLOOR_INCIDENT_EPISODE:' + || e.episode_id::text + ) AS episode_key, + e.tenant, + e.site, + e.incident_fingerprint, + e.incident_key, + e.incident_scope, + e.incident_code, + CASE e.incident_code + WHEN 'RUNNING_WITHOUT_SHOPFLOOR_SESSION' + THEN + 'Máquina en marcha sin orden ni sesión Shop Floor' + WHEN 'RUNNING_WITH_ORDER_NO_SHOPFLOOR_SESSION' + THEN + 'Máquina en marcha con orden, pero sin sesión Shop Floor' + WHEN 'ORDER_WITHOUT_SHOPFLOOR_SESSION' + THEN + 'Orden publicada sin sesión Shop Floor' + WHEN 'RUNNING_SESSION_WITHOUT_PUBLISHED_ORDER' + THEN + 'Sesión abierta y máquina en marcha, pero sin orden publicada' + WHEN 'OPEN_SESSION_WITHOUT_PUBLISHED_ORDER' + THEN + 'Sesión abierta sin orden publicada' + WHEN 'OPEN_SESSION_ORDER_MISMATCH' + THEN + 'La orden publicada no coincide con la sesión abierta' + WHEN 'MULTIPLE_OPEN_SESSIONS' + THEN + 'Varias sesiones abiertas en la misma máquina' + WHEN 'OPEN_SESSION_DATA_QUALITY_ERROR' + THEN + 'Error de calidad en una sesión abierta' + WHEN 'SHOPFLOOR_SYNC_DELAYED' + THEN + 'Sincronización Shop Floor retrasada' + WHEN 'SHOPFLOOR_SYNC_STALE' + THEN + 'Sincronización Shop Floor obsoleta' + WHEN 'SHOPFLOOR_SYNC_ERROR' + THEN + 'Error del sincronizador Shop Floor' + WHEN 'SHOPFLOOR_SYNC_NOT_INITIALIZED' + THEN + 'Sincronizador Shop Floor no inicializado' + WHEN 'PENDING_HAS_TIME' + THEN + 'Una sesión pendiente contiene tiempo productivo' + WHEN 'OPEN_MISSING_START' + THEN + 'Sesión abierta sin hora de inicio' + WHEN 'OPEN_HAS_END' + THEN + 'Sesión abierta con hora de fin' + WHEN 'CLOSED_MISSING_START' + THEN + 'Sesión cerrada sin hora de inicio' + WHEN 'CLOSED_MISSING_END' + THEN + 'Sesión cerrada sin hora de fin' + WHEN 'INVALID_INTERVAL' + THEN + 'Intervalo temporal inválido' + WHEN 'WORKORDER_SESSION_STATE_MISMATCH' + THEN + 'Estado Odoo y estado de sesión no coinciden' + WHEN 'STALE_OPEN_SESSION' + THEN + 'Sesión abierta durante más tiempo del esperado' + WHEN 'FUTURE_SESSION_START' + THEN + 'Sesión con inicio futuro' + ELSE e.incident_code + END AS incident_title, + e.machine_id, + e.context_key, + e.production_order, + e.odoo_workorder_id, + e.operator_name, + e.initial_severity, + e.current_severity, + e.max_severity, + e.episode_status, + e.first_seen_at, + e.last_seen_at, + e.ended_at, + e.closed_detected_at, + e.close_reason, + e.observation_count, + ROUND( + EXTRACT( + EPOCH FROM ( + COALESCE( + e.ended_at, + now() + ) + - e.first_seen_at + ) + )::numeric, + 2 + ) AS episode_duration_s, + ROUND( + EXTRACT( + EPOCH FROM ( + COALESCE( + e.ended_at, + now() + ) + - e.first_seen_at + ) + )::numeric / 60.0, + 2 + ) AS episode_duration_min, + ROUND( + EXTRACT( + EPOCH FROM ( + e.last_seen_at + - e.first_seen_at + ) + )::numeric, + 2 + ) AS confirmed_observed_duration_s, + e.initial_detail, + e.latest_detail, + e.initial_recommended_action, + e.latest_recommended_action, + e.review_required, + e.review_owner, + e.review_status, + e.review_classification, + e.reviewed_by, + e.reviewed_at, + e.review_notes, + COALESCE( + r.review_count, + 0 + ) AS review_count, + e.source_payload_json, + e.official_eligible, + e.created_at, + e.updated_at +FROM + mv_loss_intelligence + .shopfloor_context_incident_episodes e +LEFT JOIN review_counts r + ON r.episode_id = + e.episode_id; + +COMMENT ON VIEW +mv_reports_ucepsa_prod.li_shopfloor_context_incident_history_v1 IS +'Historial persistente de episodios derivados de incidencias vivas Shop Floor. first_seen_at es la primera observación MESAVAULT.'; + +CREATE OR REPLACE VIEW +mv_reports_ucepsa_prod.li_shopfloor_context_incident_backlog_v1 +AS +SELECT + ROW_NUMBER() OVER ( + ORDER BY + CASE max_severity + WHEN 'CRITICAL' THEN 1 + ELSE 2 + END, + CASE episode_status + WHEN 'OPEN' THEN 1 + ELSE 2 + END, + first_seen_at + ) AS backlog_rank, + h.* +FROM + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_history_v1 h +WHERE h.review_required + AND h.review_status = 'PENDING'; + +COMMENT ON VIEW +mv_reports_ucepsa_prod.li_shopfloor_context_incident_backlog_v1 IS +'Cola SHADOW de episodios pendientes de revisión humana o técnica. No crea contextos ni pérdidas oficiales.'; + +CREATE OR REPLACE VIEW +mv_reports_ucepsa_prod.li_shopfloor_context_incident_daily_summary_v1 +AS +SELECT + ( + first_seen_at + AT TIME ZONE 'Europe/Madrid' + )::date AS local_day, + machine_id, + incident_scope, + incident_code, + incident_title, + review_owner, + COUNT(*) AS episode_count, + COUNT(*) FILTER ( + WHERE episode_status = 'OPEN' + ) AS open_episode_count, + COUNT(*) FILTER ( + WHERE review_required + AND review_status = 'PENDING' + ) AS pending_review_count, + ROUND( + SUM(episode_duration_s) + / 60.0, + 2 + ) AS observed_minutes, + ROUND( + AVG(episode_duration_s) + / 60.0, + 2 + ) AS average_episode_minutes, + CASE + WHEN BOOL_OR( + max_severity = 'CRITICAL' + ) + THEN 'CRITICAL' + ELSE 'WARNING' + END AS max_severity, + false AS official_eligible +FROM + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_history_v1 +GROUP BY + ( + first_seen_at + AT TIME ZONE 'Europe/Madrid' + )::date, + machine_id, + incident_scope, + incident_code, + incident_title, + review_owner; + +COMMENT ON VIEW +mv_reports_ucepsa_prod.li_shopfloor_context_incident_daily_summary_v1 IS +'Resumen diario de episodios Shop Floor por máquina y tipo de incidencia.'; + +CREATE OR REPLACE VIEW +mv_reports_ucepsa_prod.li_shopfloor_context_incident_collector_health_v1 +AS +WITH policy AS ( + SELECT * + FROM + mv_loss_intelligence + .shopfloor_context_incident_policies + WHERE tenant = 'ucepsa' + AND site = 'ucepsa_onpremise' + AND policy_code = + 'SHOPFLOOR_CONTEXT_INCIDENT_HISTORY_V1' + AND active + ORDER BY updated_at DESC + LIMIT 1 +), +state AS ( + SELECT * + FROM + mv_loss_intelligence + .shopfloor_context_incident_collector_state + WHERE tenant = 'ucepsa' + AND site = 'ucepsa_onpremise' + AND collector_name = + 'shopfloor_context_incident_collector_v037' +) +SELECT + p.tenant, + p.site, + 'shopfloor_context_incident_collector_v037'::text + AS collector_name, + p.policy_code, + p.collector_interval_s, + p.close_grace_s, + p.collector_delayed_after_s, + p.collector_stale_after_s, + s.last_success_at, + s.last_current_count, + s.last_opened_count, + s.last_updated_count, + s.last_closed_count, + NULLIF( + btrim(s.last_error), + '' + ) AS last_error, + CASE + WHEN s.last_success_at IS NULL + THEN NULL::numeric + ELSE ROUND( + GREATEST( + EXTRACT( + EPOCH FROM ( + now() + - s.last_success_at + ) + ), + 0 + )::numeric, + 2 + ) + END AS seconds_since_success, + CASE + WHEN s.tenant IS NULL + THEN 'NOT_INITIALIZED' + WHEN NULLIF( + btrim(s.last_error), + '' + ) IS NOT NULL + THEN 'ERROR' + WHEN s.last_success_at IS NULL + THEN 'NOT_INITIALIZED' + WHEN now() - s.last_success_at + > make_interval( + secs => + p.collector_stale_after_s + ::double precision + ) + THEN 'STALE' + WHEN now() - s.last_success_at + > make_interval( + secs => + p.collector_delayed_after_s + ::double precision + ) + THEN 'DELAYED' + ELSE 'HEALTHY' + END AS collector_health_status, + CASE + WHEN s.tenant IS NULL + OR s.last_success_at IS NULL + OR NULLIF( + btrim(s.last_error), + '' + ) IS NOT NULL + OR now() - s.last_success_at + > make_interval( + secs => + p.collector_stale_after_s + ::double precision + ) + THEN 'CRITICAL' + WHEN now() - s.last_success_at + > make_interval( + secs => + p.collector_delayed_after_s + ::double precision + ) + THEN 'WARNING' + ELSE 'OK' + END AS severity, + false AS official_eligible, + now() AS observed_at +FROM policy p +LEFT JOIN state s + ON s.tenant = p.tenant + AND s.site = p.site; + +COMMENT ON VIEW +mv_reports_ucepsa_prod.li_shopfloor_context_incident_collector_health_v1 IS +'Salud del colector persistente de incidencias Shop Floor.'; + +DO $$ +BEGIN + IF EXISTS ( + SELECT 1 + FROM pg_roles + WHERE rolname = 'grafana_ucepsa_ro' + ) THEN + GRANT SELECT + ON + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_history_v1 + TO grafana_ucepsa_ro; + + GRANT SELECT + ON + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_backlog_v1 + TO grafana_ucepsa_ro; + + GRANT SELECT + ON + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_daily_summary_v1 + TO grafana_ucepsa_ro; + + GRANT SELECT + ON + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_collector_health_v1 + TO grafana_ucepsa_ro; + END IF; +END +$$; + +COMMIT; diff --git a/ucepsa/edge-oee-demo/tools/shopfloor_context_incident_admin.py b/ucepsa/edge-oee-demo/tools/shopfloor_context_incident_admin.py new file mode 100755 index 0000000..a4b93bf --- /dev/null +++ b/ucepsa/edge-oee-demo/tools/shopfloor_context_incident_admin.py @@ -0,0 +1,369 @@ +#!/usr/bin/env python3 +"""Administración de revisiones de episodios Shop Floor v0.3.7.""" + +from __future__ import annotations + +import argparse +import json +import os +from typing import Any + +import psycopg2 +from psycopg2.extras import Json, RealDictCursor + + +CLASSIFICATIONS = ( + "AUTHORIZED_LEGACY_PRODUCTION", + "TEST_OR_SETUP", + "RESIDUAL_MATERIAL", + "ODOO_START_OMITTED", + "DATA_ISSUE", + "NOT_RELEVANT", + "OTHER", +) + + +def env_first(*names: str, default: str | None = None) -> str | None: + for name in names: + value = os.getenv(name) + if value not in (None, ""): + return value + return default + + +def connect(): + required = { + "dbname": env_first( + "PGDATABASE", + "POSTGRES_DB", + ), + "user": env_first( + "PGUSER", + "POSTGRES_USER", + ), + "password": env_first( + "PGPASSWORD", + "POSTGRES_PASSWORD", + ), + } + missing = [ + key + for key, value in required.items() + if not value + ] + if missing: + raise RuntimeError( + "Faltan variables PostgreSQL: " + + ", ".join(missing) + ) + + return psycopg2.connect( + host=env_first( + "PGHOST", + default="mv_ucepsa_postgres_hot", + ), + port=int( + env_first("PGPORT", default="5432") + or "5432" + ), + dbname=required["dbname"], + user=required["user"], + password=required["password"], + connect_timeout=10, + application_name=( + "shopfloor_context_incident_admin_v037" + ), + ) + + +def print_json(value: Any) -> None: + print( + json.dumps( + value, + ensure_ascii=False, + indent=2, + default=str, + ) + ) + + +def list_backlog(limit: int) -> None: + with connect() as connection: + with connection.cursor( + cursor_factory=RealDictCursor, + ) as cursor: + cursor.execute( + """ + SELECT + backlog_rank, + episode_id, + episode_status, + max_severity, + first_seen_at, + last_seen_at, + episode_duration_min, + machine_id, + incident_title, + production_order, + operator_name, + review_owner, + review_status, + latest_recommended_action + FROM + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_backlog_v1 + ORDER BY backlog_rank + LIMIT %s + """, + (limit,), + ) + print_json( + [dict(row) for row in cursor.fetchall()] + ) + + +def show_episode(episode_id: int) -> None: + with connect() as connection: + with connection.cursor( + cursor_factory=RealDictCursor, + ) as cursor: + cursor.execute( + """ + SELECT * + FROM + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_history_v1 + WHERE episode_id = %s + """, + (episode_id,), + ) + episode = cursor.fetchone() + if not episode: + raise RuntimeError( + f"No existe el episodio {episode_id}" + ) + + cursor.execute( + """ + SELECT * + FROM + mv_loss_intelligence + .shopfloor_context_incident_reviews + WHERE episode_id = %s + ORDER BY reviewed_at, review_id + """, + (episode_id,), + ) + reviews = [ + dict(row) + for row in cursor.fetchall() + ] + print_json({ + "episode": dict(episode), + "reviews": reviews, + }) + + +def apply_review( + episode_id: int, + status: str, + classification: str, + reviewed_by: str, + notes: str | None, +) -> None: + with connect() as connection: + with connection.cursor( + cursor_factory=RealDictCursor, + ) as cursor: + cursor.execute( + """ + SELECT + episode_id, + review_status, + review_classification, + reviewed_by, + reviewed_at, + review_notes, + episode_status, + incident_code, + machine_id, + production_order + FROM + mv_loss_intelligence + .shopfloor_context_incident_episodes + WHERE episode_id = %s + FOR UPDATE + """, + (episode_id,), + ) + previous = cursor.fetchone() + if not previous: + raise RuntimeError( + f"No existe el episodio {episode_id}" + ) + + cursor.execute( + """ + INSERT INTO + mv_loss_intelligence + .shopfloor_context_incident_reviews ( + episode_id, + previous_values, + new_review_status, + review_classification, + reviewed_by, + notes, + evidence_json + ) + VALUES ( + %s, %s, %s, %s, %s, %s, %s + ) + RETURNING * + """, + ( + episode_id, + Json(dict(previous)), + status, + classification, + reviewed_by, + notes, + Json({ + "source": + "shopfloor_context_incident_admin_v037" + }), + ), + ) + review = dict(cursor.fetchone()) + + cursor.execute( + """ + SELECT * + FROM + mv_reports_ucepsa_prod + .li_shopfloor_context_incident_history_v1 + WHERE episode_id = %s + """, + (episode_id,), + ) + current = dict(cursor.fetchone()) + + print_json({ + "review": review, + "episode": current, + "warning": ( + "La revisión no crea una sesión de producción, " + "no atribuye pérdida y no modifica Odoo." + ), + }) + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser( + description=( + "Revisión SHADOW de episodios " + "de contexto Shop Floor." + ) + ) + sub = parser.add_subparsers( + dest="command", + required=True, + ) + + list_cmd = sub.add_parser( + "list", + help="Lista el backlog pendiente.", + ) + list_cmd.add_argument( + "--limit", + type=int, + default=50, + ) + + show_cmd = sub.add_parser( + "show", + help="Muestra un episodio y sus revisiones.", + ) + show_cmd.add_argument( + "--episode-id", + type=int, + required=True, + ) + + review_cmd = sub.add_parser( + "review", + help="Clasifica un episodio.", + ) + review_cmd.add_argument( + "--episode-id", + type=int, + required=True, + ) + review_cmd.add_argument( + "--classification", + choices=CLASSIFICATIONS, + required=True, + ) + review_cmd.add_argument( + "--reviewed-by", + required=True, + ) + review_cmd.add_argument( + "--notes", + ) + + dismiss_cmd = sub.add_parser( + "dismiss", + help="Descarta un episodio como no relevante.", + ) + dismiss_cmd.add_argument( + "--episode-id", + type=int, + required=True, + ) + dismiss_cmd.add_argument( + "--reviewed-by", + required=True, + ) + dismiss_cmd.add_argument( + "--notes", + ) + + return parser.parse_args() + + +def main() -> int: + args = parse_args() + + if args.command == "list": + list_backlog(args.limit) + return 0 + + if args.command == "show": + show_episode(args.episode_id) + return 0 + + if args.command == "review": + apply_review( + episode_id=args.episode_id, + status="REVIEWED", + classification=args.classification, + reviewed_by=args.reviewed_by, + notes=args.notes, + ) + return 0 + + if args.command == "dismiss": + apply_review( + episode_id=args.episode_id, + status="DISMISSED", + classification="NOT_RELEVANT", + reviewed_by=args.reviewed_by, + notes=args.notes, + ) + return 0 + + raise RuntimeError("Comando no soportado") + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/ucepsa/edge-oee-demo/tools/shopfloor_context_incident_collector.py b/ucepsa/edge-oee-demo/tools/shopfloor_context_incident_collector.py new file mode 100755 index 0000000..9717c80 --- /dev/null +++ b/ucepsa/edge-oee-demo/tools/shopfloor_context_incident_collector.py @@ -0,0 +1,806 @@ +#!/usr/bin/env python3 +"""Persistencia de incidencias vivas Shop Floor en episodios SHADOW. + +Lee exclusivamente: + mv_reports_ucepsa_prod.li_shopfloor_context_incidents_v1 + +Escribe exclusivamente en las tablas de historial v0.3.7. +No modifica Odoo, sesiones, paros, balizas ni el Ledger oficial. +""" + +from __future__ import annotations + +import argparse +import hashlib +import json +import os +import signal +import sys +import time +from dataclasses import dataclass +from datetime import datetime +from typing import Any + +import psycopg2 +from psycopg2.extras import Json, RealDictCursor + + +VERSION = "0.3.7" +COLLECTOR_NAME = "shopfloor_context_incident_collector_v037" +POLICY_CODE = "SHOPFLOOR_CONTEXT_INCIDENT_HISTORY_V1" +RUNNING = True + + +def stop_handler(signum: int, frame: Any) -> None: + del signum, frame + global RUNNING + RUNNING = False + + +signal.signal(signal.SIGTERM, stop_handler) +signal.signal(signal.SIGINT, stop_handler) + + +def env_first(*names: str, default: str | None = None) -> str | None: + for name in names: + value = os.getenv(name) + if value not in (None, ""): + return value + return default + + +@dataclass(frozen=True) +class Config: + tenant: str + site: str + pg_host: str + pg_port: int + pg_db: str + pg_user: str + pg_password: str + + +def load_config() -> Config: + required = { + "PGDATABASE": env_first("PGDATABASE", "POSTGRES_DB"), + "PGUSER": env_first("PGUSER", "POSTGRES_USER"), + "PGPASSWORD": env_first("PGPASSWORD", "POSTGRES_PASSWORD"), + } + missing = [name for name, value in required.items() if not value] + if missing: + raise RuntimeError( + "Faltan variables obligatorias: " + ", ".join(missing) + ) + + return Config( + tenant=env_first( + "TENANT", + "TENANT_OVERRIDE", + default="ucepsa", + ) + or "ucepsa", + site=env_first( + "SITE", + "SITE_OVERRIDE", + default="ucepsa_onpremise", + ) + or "ucepsa_onpremise", + pg_host=env_first( + "PGHOST", + default="mv_ucepsa_postgres_hot", + ) + or "mv_ucepsa_postgres_hot", + pg_port=int( + env_first("PGPORT", default="5432") + or "5432" + ), + pg_db=str(required["PGDATABASE"]), + pg_user=str(required["PGUSER"]), + pg_password=str(required["PGPASSWORD"]), + ) + + +def connect(config: Config): + return psycopg2.connect( + host=config.pg_host, + port=config.pg_port, + dbname=config.pg_db, + user=config.pg_user, + password=config.pg_password, + connect_timeout=10, + application_name=COLLECTOR_NAME, + ) + + +def canonical_fingerprint(row: dict[str, Any]) -> str: + identity = { + "incident_key": row.get("incident_key"), + "incident_scope": row.get("incident_scope"), + "incident_code": row.get("incident_code"), + "machine_id": row.get("machine_id"), + "context_key": row.get("context_key"), + "production_order": row.get("production_order"), + "odoo_workorder_id": row.get("odoo_workorder_id"), + } + raw = json.dumps( + identity, + ensure_ascii=False, + sort_keys=True, + separators=(",", ":"), + default=str, + ) + return hashlib.sha256(raw.encode("utf-8")).hexdigest() + + +def severity_rank(value: str | None) -> int: + return { + "WARNING": 1, + "CRITICAL": 2, + }.get(value or "", 0) + + +def review_policy(row: dict[str, Any]) -> tuple[bool, str]: + code = str(row.get("incident_code") or "") + scope = str(row.get("incident_scope") or "") + severity = str(row.get("severity") or "") + + production_codes = { + "RUNNING_WITHOUT_SHOPFLOOR_SESSION", + "ORDER_WITHOUT_SHOPFLOOR_SESSION", + } + + if code in production_codes: + return True, "PRODUCTION" + + if severity == "CRITICAL": + return True, "TECHNICAL" + + if scope in {"MACHINE", "SESSION"}: + return True, "TECHNICAL" + + if code in { + "SHOPFLOOR_SYNC_STALE", + "SHOPFLOOR_SYNC_ERROR", + "SHOPFLOOR_SYNC_NOT_INITIALIZED", + }: + return True, "TECHNICAL" + + return False, "TECHNICAL" + + +def source_payload(row: dict[str, Any]) -> dict[str, Any]: + return { + "incident_key": row.get("incident_key"), + "incident_scope": row.get("incident_scope"), + "incident_code": row.get("incident_code"), + "machine_id": row.get("machine_id"), + "context_key": row.get("context_key"), + "production_order": row.get("production_order"), + "odoo_workorder_id": row.get("odoo_workorder_id"), + "operator_name": row.get("operator_name"), + "severity": row.get("severity"), + "detail": row.get("detail"), + "recommended_action": row.get("recommended_action"), + "source_observed_at": ( + row.get("observed_at").isoformat() + if isinstance(row.get("observed_at"), datetime) + else row.get("observed_at") + ), + } + + +def fetch_policy(cursor, config: Config) -> dict[str, Any]: + cursor.execute( + """ + SELECT + collector_interval_s, + close_grace_s, + collector_delayed_after_s, + collector_stale_after_s + FROM + mv_loss_intelligence + .shopfloor_context_incident_policies + WHERE tenant = %s + AND site = %s + AND policy_code = %s + AND active + ORDER BY updated_at DESC + LIMIT 1 + """, + ( + config.tenant, + config.site, + POLICY_CODE, + ), + ) + row = cursor.fetchone() + if not row: + raise RuntimeError( + "No existe política activa " + f"{POLICY_CODE!r}" + ) + return dict(row) + + +def fetch_current_incidents( + cursor, + config: Config, +) -> list[dict[str, Any]]: + cursor.execute( + """ + SELECT + incident_key, + incident_scope, + observed_at, + severity, + incident_code, + machine_id, + context_key, + production_order, + odoo_workorder_id, + operator_name, + detail, + recommended_action + FROM + mv_reports_ucepsa_prod + .li_shopfloor_context_incidents_v1 + WHERE official_eligible = false + ORDER BY + incident_scope, + machine_id NULLS FIRST, + incident_code + """ + ) + return [dict(row) for row in cursor.fetchall()] + + +def current_db_time(cursor) -> datetime: + cursor.execute("SELECT now() AS db_now") + row = cursor.fetchone() + return row["db_now"] + + +def max_severity( + existing: str, + incoming: str, +) -> str: + if severity_rank(incoming) > severity_rank(existing): + return incoming + return existing + + +def process_cycle( + config: Config, + dry_run: bool = False, +) -> dict[str, Any]: + opened = 0 + updated = 0 + closed = 0 + + with connect(config) as connection: + with connection.cursor( + cursor_factory=RealDictCursor, + ) as cursor: + cursor.execute( + """ + SELECT pg_try_advisory_xact_lock( + hashtext(%s) + ) AS locked + """, + ( + ( + f"{config.tenant}:" + f"{config.site}:" + f"{COLLECTOR_NAME}" + ), + ), + ) + if not cursor.fetchone()["locked"]: + return { + "version": VERSION, + "status": "skipped_lock_busy", + } + + policy = fetch_policy(cursor, config) + db_now = current_db_time(cursor) + incidents = fetch_current_incidents( + cursor, + config, + ) + + deduplicated: dict[str, dict[str, Any]] = {} + for incident in incidents: + fingerprint = canonical_fingerprint(incident) + current = deduplicated.get(fingerprint) + if current is None: + deduplicated[fingerprint] = incident + continue + if severity_rank( + incident.get("severity") + ) >= severity_rank( + current.get("severity") + ): + deduplicated[fingerprint] = incident + + if dry_run: + return { + "version": VERSION, + "status": "dry_run", + "db_now": db_now.isoformat(), + "policy": policy, + "current_incident_count": len( + deduplicated + ), + "incidents": [ + { + **row, + "incident_fingerprint": + fingerprint, + } + for fingerprint, row + in deduplicated.items() + ], + } + + current_fingerprints: list[str] = [] + + for fingerprint, row in deduplicated.items(): + current_fingerprints.append(fingerprint) + review_required, review_owner = review_policy( + row + ) + observed_at = ( + row.get("observed_at") + or db_now + ) + payload = source_payload(row) + + cursor.execute( + """ + SELECT + episode_id, + max_severity, + review_status, + review_required + FROM + mv_loss_intelligence + .shopfloor_context_incident_episodes + WHERE tenant = %s + AND site = %s + AND incident_fingerprint = %s + AND episode_status = 'OPEN' + FOR UPDATE + """, + ( + config.tenant, + config.site, + fingerprint, + ), + ) + existing = cursor.fetchone() + + if existing is None: + initial_review_status = ( + "PENDING" + if review_required + else "NOT_REQUIRED" + ) + cursor.execute( + """ + INSERT INTO + mv_loss_intelligence + .shopfloor_context_incident_episodes ( + tenant, + site, + incident_fingerprint, + incident_key, + incident_scope, + incident_code, + machine_id, + context_key, + production_order, + odoo_workorder_id, + operator_name, + initial_severity, + current_severity, + max_severity, + episode_status, + first_seen_at, + last_seen_at, + observation_count, + initial_detail, + latest_detail, + initial_recommended_action, + latest_recommended_action, + review_required, + review_owner, + review_status, + source_payload_json, + official_eligible + ) + VALUES ( + %s, %s, %s, %s, %s, %s, + %s, %s, %s, %s, %s, + %s, %s, %s, + 'OPEN', + %s, %s, 1, + %s, %s, %s, %s, + %s, %s, %s, + %s, false + ) + """, + ( + config.tenant, + config.site, + fingerprint, + row.get("incident_key"), + row.get("incident_scope"), + row.get("incident_code"), + row.get("machine_id"), + row.get("context_key"), + row.get("production_order"), + row.get("odoo_workorder_id"), + row.get("operator_name"), + row.get("severity"), + row.get("severity"), + row.get("severity"), + observed_at, + observed_at, + row.get("detail"), + row.get("detail"), + row.get( + "recommended_action" + ), + row.get( + "recommended_action" + ), + review_required, + review_owner, + initial_review_status, + Json(payload), + ), + ) + opened += 1 + continue + + next_max = max_severity( + str(existing["max_severity"]), + str(row.get("severity")), + ) + next_required = ( + bool(existing["review_required"]) + or review_required + ) + next_status = str( + existing["review_status"] + ) + if ( + next_required + and next_status == "NOT_REQUIRED" + ): + next_status = "PENDING" + + cursor.execute( + """ + UPDATE + mv_loss_intelligence + .shopfloor_context_incident_episodes + SET + incident_key = %s, + incident_scope = %s, + incident_code = %s, + machine_id = %s, + context_key = %s, + production_order = %s, + odoo_workorder_id = %s, + operator_name = %s, + current_severity = %s, + max_severity = %s, + last_seen_at = GREATEST( + last_seen_at, + %s + ), + observation_count = + observation_count + 1, + latest_detail = %s, + latest_recommended_action = %s, + review_required = %s, + review_owner = %s, + review_status = %s, + source_payload_json = %s, + updated_at = %s + WHERE episode_id = %s + """, + ( + row.get("incident_key"), + row.get("incident_scope"), + row.get("incident_code"), + row.get("machine_id"), + row.get("context_key"), + row.get("production_order"), + row.get("odoo_workorder_id"), + row.get("operator_name"), + row.get("severity"), + next_max, + observed_at, + row.get("detail"), + row.get( + "recommended_action" + ), + next_required, + review_owner, + next_status, + Json(payload), + db_now, + existing["episode_id"], + ), + ) + updated += 1 + + cursor.execute( + """ + UPDATE + mv_loss_intelligence + .shopfloor_context_incident_episodes + SET + episode_status = 'CLOSED', + ended_at = last_seen_at, + closed_detected_at = %s, + close_reason = + 'INCIDENT_NO_LONGER_PRESENT', + current_severity = + max_severity, + updated_at = %s + WHERE tenant = %s + AND site = %s + AND episode_status = 'OPEN' + AND last_seen_at + < %s + - ( + %s + * interval '1 second' + ) + AND NOT ( + incident_fingerprint + = ANY(%s::text[]) + ) + RETURNING episode_id + """, + ( + db_now, + db_now, + config.tenant, + config.site, + db_now, + int(policy["close_grace_s"]), + current_fingerprints, + ), + ) + closed = len(cursor.fetchall()) + + cursor.execute( + """ + INSERT INTO + mv_loss_intelligence + .shopfloor_context_incident_collector_state ( + tenant, + site, + collector_name, + last_success_at, + last_current_count, + last_opened_count, + last_updated_count, + last_closed_count, + last_error, + updated_at + ) + VALUES ( + %s, %s, %s, %s, + %s, %s, %s, %s, + NULL, %s + ) + ON CONFLICT ( + tenant, + site, + collector_name + ) + DO UPDATE SET + last_success_at = + EXCLUDED.last_success_at, + last_current_count = + EXCLUDED.last_current_count, + last_opened_count = + EXCLUDED.last_opened_count, + last_updated_count = + EXCLUDED.last_updated_count, + last_closed_count = + EXCLUDED.last_closed_count, + last_error = NULL, + updated_at = + EXCLUDED.updated_at + """, + ( + config.tenant, + config.site, + COLLECTOR_NAME, + db_now, + len(deduplicated), + opened, + updated, + closed, + db_now, + ), + ) + + return { + "version": VERSION, + "status": "ok", + "db_now": db_now.isoformat(), + "current_incident_count": len(deduplicated), + "opened_count": opened, + "updated_count": updated, + "closed_count": closed, + "close_grace_s": int( + policy["close_grace_s"] + ), + "collector_interval_s": int( + policy["collector_interval_s"] + ), + } + + +def record_error( + config: Config, + error: Exception, +) -> None: + try: + with connect(config) as connection: + with connection.cursor() as cursor: + cursor.execute( + """ + INSERT INTO + mv_loss_intelligence + .shopfloor_context_incident_collector_state ( + tenant, + site, + collector_name, + last_error, + updated_at + ) + VALUES ( + %s, %s, %s, %s, now() + ) + ON CONFLICT ( + tenant, + site, + collector_name + ) + DO UPDATE SET + last_error = + EXCLUDED.last_error, + updated_at = now() + """, + ( + config.tenant, + config.site, + COLLECTOR_NAME, + str(error)[:4000], + ), + ) + except Exception as state_error: + print( + json.dumps( + { + "status": + "collector_error_state_write_failed", + "error": str(error), + "state_error": str(state_error), + }, + ensure_ascii=False, + ), + file=sys.stderr, + flush=True, + ) + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser( + description=( + "Persistencia SHADOW de incidencias " + "Shop Floor." + ) + ) + mode = parser.add_mutually_exclusive_group( + required=True + ) + mode.add_argument( + "--once", + action="store_true", + help="Ejecuta un ciclo y termina.", + ) + mode.add_argument( + "--watch", + action="store_true", + help="Ejecuta ciclos continuos.", + ) + parser.add_argument( + "--interval", + type=int, + default=None, + help=( + "Intervalo del modo watch. " + "Por defecto usa la política SQL." + ), + ) + parser.add_argument( + "--dry-run", + action="store_true", + help="Lee y muestra, pero no escribe.", + ) + return parser.parse_args() + + +def main() -> int: + args = parse_args() + config = load_config() + + while RUNNING: + try: + result = process_cycle( + config, + dry_run=args.dry_run, + ) + print( + json.dumps( + result, + ensure_ascii=False, + default=str, + ), + flush=True, + ) + + if args.once or args.dry_run: + return 0 + + interval = args.interval + if interval is None: + interval = int( + result.get( + "collector_interval_s", + 15, + ) + or 15 + ) + for _ in range(max(interval, 1)): + if not RUNNING: + break + time.sleep(1) + + except Exception as error: + record_error(config, error) + print( + json.dumps( + { + "version": VERSION, + "status": "error", + "error": str(error), + }, + ensure_ascii=False, + ), + file=sys.stderr, + flush=True, + ) + if args.once or args.dry_run: + return 1 + for _ in range( + max(args.interval or 15, 1) + ): + if not RUNNING: + break + time.sleep(1) + + return 0 + + +if __name__ == "__main__": + raise SystemExit(main())