From 7c5ad7f5a397392e8502ce2faffa8f0345006628 Mon Sep 17 00:00:00 2001 From: Victor Fraile Garcia Date: Mon, 22 Jun 2026 12:54:08 +0200 Subject: [PATCH] feat(ucepsa): persist Odoo order_state for Grafana --- .../docker-compose.order-state-pg.yml | 28 ++ .../sql/081_odoo_order_state_pg.sql | 125 +++++++++ .../tools/order_state_pg_sink.py | 247 ++++++++++++++++++ 3 files changed, 400 insertions(+) create mode 100644 ucepsa/edge-oee-demo/deploy/ucepsa-local/docker-compose.order-state-pg.yml create mode 100644 ucepsa/edge-oee-demo/sql/081_odoo_order_state_pg.sql create mode 100755 ucepsa/edge-oee-demo/tools/order_state_pg_sink.py diff --git a/ucepsa/edge-oee-demo/deploy/ucepsa-local/docker-compose.order-state-pg.yml b/ucepsa/edge-oee-demo/deploy/ucepsa-local/docker-compose.order-state-pg.yml new file mode 100644 index 0000000..6c76f7b --- /dev/null +++ b/ucepsa/edge-oee-demo/deploy/ucepsa-local/docker-compose.order-state-pg.yml @@ -0,0 +1,28 @@ +services: + mv_ucepsa_odoo_order_state_pg_sink: + image: edge-oee-ucepsa-mv_ucepsa_edge_oee_mqtt_pg_sink + container_name: mv_ucepsa_odoo_order_state_pg_sink + restart: unless-stopped + env_file: + - .env + environment: + TENANT: ucepsa + SITE: ucepsa_onpremise + MQTT_HOST: mv_ucepsa_mosquitto + MQTT_PORT: 1883 + MQTT_TOPIC: vertical/ucepsa/ucepsa_onpremise/edge_oee/order_state/+ + PGHOST: mv_ucepsa_postgres_hot + PGPORT: 5432 + PGDATABASE: ${POSTGRES_DB} + PGUSER: ${POSTGRES_USER} + PGPASSWORD: ${POSTGRES_PASSWORD} + volumes: + - ./tools:/app/tools:ro + working_dir: /app + command: python tools/order_state_pg_sink.py + networks: + - mv_ucepsa_net + +networks: + mv_ucepsa_net: + external: true diff --git a/ucepsa/edge-oee-demo/sql/081_odoo_order_state_pg.sql b/ucepsa/edge-oee-demo/sql/081_odoo_order_state_pg.sql new file mode 100644 index 0000000..3708697 --- /dev/null +++ b/ucepsa/edge-oee-demo/sql/081_odoo_order_state_pg.sql @@ -0,0 +1,125 @@ +CREATE TABLE IF NOT EXISTS mv_hot.odoo_order_state_current ( + tenant text NOT NULL, + site text NOT NULL, + machine_id text NOT NULL, + topic text, + mqtt_ts timestamptz, + received_at timestamptz NOT NULL DEFAULT now(), + + source text, + workcenter_id bigint, + workcenter_name text, + workcenter_display_name text, + + order_present boolean, + order_coherent boolean, + odoo_state text, + reason text, + + workorder_id text, + production_order text, + product text, + + waiting_count integer, + ready_count integer, + progress_count integer, + + payload_hash text, + payload_json jsonb NOT NULL, + + PRIMARY KEY (tenant, site, machine_id) +); + +CREATE TABLE IF NOT EXISTS mv_hot.odoo_order_state_events ( + id bigserial PRIMARY KEY, + tenant text NOT NULL, + site text NOT NULL, + machine_id text NOT NULL, + topic text, + mqtt_ts timestamptz, + received_at timestamptz NOT NULL DEFAULT now(), + + odoo_state text, + order_present boolean, + order_coherent boolean, + production_order text, + product text, + reason text, + + waiting_count integer, + ready_count integer, + progress_count integer, + + payload_hash text NOT NULL, + payload_json jsonb NOT NULL +); + +CREATE UNIQUE INDEX IF NOT EXISTS ux_odoo_order_state_events_state_hash +ON mv_hot.odoo_order_state_events (tenant, site, machine_id, payload_hash); + +DROP VIEW IF EXISTS mv_reports_ucepsa_demo.odoo_order_state_machine_grafana; + +CREATE VIEW mv_reports_ucepsa_demo.odoo_order_state_machine_grafana AS +SELECT + COALESCE(m.machine_id, o.machine_id) AS machine_id, + + o.received_at AS odoo_state_received_at, + o.mqtt_ts AS odoo_state_ts, + EXTRACT(EPOCH FROM (now() - o.received_at))::double precision AS order_age_s, + + o.workcenter_id, + o.workcenter_name, + o.workcenter_display_name, + + o.production_order AS order_ref, + o.product AS product_name, + o.odoo_state, + o.order_present, + o.order_coherent, + o.reason, + + o.waiting_count, + o.ready_count, + o.progress_count, + + m.ingest_ts AS telemetry_ingest_ts, + m.age_s AS telemetry_age_s, + m.comm_ok, + m.sample_valid, + m.machine_running, + m.cycle_rate_ppm, + m.cycle_pulse_count, + m.digital_inputs, + + CASE + WHEN o.machine_id IS NULL THEN 'SIN_ESTADO_ODOO' + WHEN EXTRACT(EPOCH FROM (now() - o.received_at)) > 30 THEN 'ODOO_STATE_ANTIGUO' + WHEN o.odoo_state = 'multiple_loaded_orders' THEN 'MULTIPLES_ORDENES_EN_PROGRESS' + WHEN o.order_present = true + AND o.order_coherent = true + AND COALESCE(m.machine_running, false) = true + THEN 'ORDEN_CARGADA_MAQUINA_EN_MARCHA' + WHEN o.order_present = true + AND o.order_coherent = true + AND COALESCE(m.machine_running, false) = false + THEN 'ORDEN_CARGADA_MAQUINA_PARADA' + WHEN COALESCE(o.order_present, false) = false + AND COALESCE(m.machine_running, false) = true + THEN 'MAQUINA_EN_MARCHA_SIN_ORDEN_ODOO' + WHEN COALESCE(o.order_present, false) = false + AND COALESCE(m.machine_running, false) = false + THEN 'SIN_ORDEN_MAQUINA_PARADA' + ELSE 'REVISAR' + END AS edge_odoo_status, + + o.payload_json AS odoo_payload_json +FROM mv_reports_ucepsa_demo.machine_state_latest_multi m +FULL OUTER JOIN mv_hot.odoo_order_state_current o + ON o.tenant = 'ucepsa' + AND o.site = 'ucepsa_onpremise' + AND o.machine_id = m.machine_id +WHERE COALESCE(m.site, o.site) = 'ucepsa_onpremise'; + +GRANT SELECT ON mv_hot.odoo_order_state_current TO grafana_ucepsa_ro; +GRANT SELECT ON mv_hot.odoo_order_state_events TO grafana_ucepsa_ro; +GRANT SELECT ON mv_reports_ucepsa_demo.odoo_order_state_machine_grafana TO grafana_ucepsa_ro; diff --git a/ucepsa/edge-oee-demo/tools/order_state_pg_sink.py b/ucepsa/edge-oee-demo/tools/order_state_pg_sink.py new file mode 100755 index 0000000..72684a1 --- /dev/null +++ b/ucepsa/edge-oee-demo/tools/order_state_pg_sink.py @@ -0,0 +1,247 @@ +#!/usr/bin/env python3 +import hashlib +import json +import os +import signal +import sys +from datetime import datetime, timezone + +import paho.mqtt.client as mqtt +import psycopg2 +from psycopg2.extras import Json + + +RUNNING = True + + +def stop_handler(signum, frame): + global RUNNING + RUNNING = False + + +signal.signal(signal.SIGTERM, stop_handler) +signal.signal(signal.SIGINT, stop_handler) + + +def env_first(*names, default=None): + for name in names: + value = os.getenv(name) + if value not in (None, ""): + return value + return default + + +def to_bool(value): + if isinstance(value, bool): + return value + if value is None: + return None + return str(value).strip().lower() in ("true", "t", "1", "yes", "y", "si", "sí") + + +def to_int(value): + if value in (None, ""): + return None + try: + return int(value) + except Exception: + return None + + +def parse_ts(value): + if not value: + return None + try: + return datetime.fromisoformat(str(value).replace("Z", "+00:00")) + except Exception: + return None + + +def state_hash(payload): + diagnostics = payload.get("diagnostics") or {} + relevant = { + "machine_id": payload.get("machine_id"), + "workcenter_id": payload.get("workcenter_id"), + "order_present": payload.get("order_present"), + "order_coherent": payload.get("order_coherent"), + "odoo_state": payload.get("odoo_state"), + "reason": payload.get("reason"), + "workorder_id": payload.get("workorder_id"), + "production_order": payload.get("production_order"), + "product": payload.get("product"), + "waiting_count": diagnostics.get("waiting_count"), + "ready_count": diagnostics.get("ready_count"), + "progress_count": diagnostics.get("progress_count"), + } + raw = json.dumps(relevant, sort_keys=True, ensure_ascii=False) + return hashlib.sha256(raw.encode("utf-8")).hexdigest() + + +TENANT = env_first("TENANT", "TENANT_OVERRIDE", default="ucepsa") +SITE = env_first("SITE", "SITE_OVERRIDE", default="ucepsa_onpremise") + +MQTT_HOST = env_first("MQTT_HOST", "MQTT_HOST_OVERRIDE", default="mv_ucepsa_mosquitto") +MQTT_PORT = int(env_first("MQTT_PORT", default="1883")) +MQTT_TOPIC = env_first( + "MQTT_TOPIC", + "ORDER_STATE_TOPIC", + default=f"vertical/{TENANT}/{SITE}/edge_oee/order_state/+", +) + +MQTT_USER = env_first("MQTT_USER", "MQTT_USER_REVPI") +MQTT_PASSWORD = env_first("MQTT_PASSWORD", "MQTT_PASSWORD_REVPI") + +PGHOST = env_first("PGHOST", default="mv_ucepsa_postgres_hot") +PGPORT = int(env_first("PGPORT", default="5432")) +PGDATABASE = env_first("PGDATABASE", "POSTGRES_DB") +PGUSER = env_first("PGUSER", "POSTGRES_USER") +PGPASSWORD = env_first("PGPASSWORD", "POSTGRES_PASSWORD") + +if not PGDATABASE or not PGUSER or not PGPASSWORD: + print("ERROR: faltan credenciales PostgreSQL", file=sys.stderr) + sys.exit(1) + +conn = psycopg2.connect( + host=PGHOST, + port=PGPORT, + dbname=PGDATABASE, + user=PGUSER, + password=PGPASSWORD, +) +conn.autocommit = True + + +def upsert_order_state(topic, payload): + machine_id = payload.get("machine_id") or topic.rstrip("/").split("/")[-1] + diagnostics = payload.get("diagnostics") or {} + + mqtt_ts = parse_ts(payload.get("ts")) + h = state_hash(payload) + + row = { + "tenant": TENANT, + "site": SITE, + "machine_id": machine_id, + "topic": topic, + "mqtt_ts": mqtt_ts, + "source": payload.get("source"), + "workcenter_id": to_int(payload.get("workcenter_id")), + "workcenter_name": payload.get("workcenter_name"), + "workcenter_display_name": payload.get("workcenter_display_name"), + "order_present": to_bool(payload.get("order_present")), + "order_coherent": to_bool(payload.get("order_coherent")), + "odoo_state": payload.get("odoo_state"), + "reason": payload.get("reason"), + "workorder_id": str(payload.get("workorder_id")) if payload.get("workorder_id") is not None else None, + "production_order": payload.get("production_order"), + "product": payload.get("product"), + "waiting_count": to_int(diagnostics.get("waiting_count")), + "ready_count": to_int(diagnostics.get("ready_count")), + "progress_count": to_int(diagnostics.get("progress_count")), + "payload_hash": h, + "payload_json": payload, + } + + with conn.cursor() as cur: + cur.execute( + """ + INSERT INTO mv_hot.odoo_order_state_current ( + tenant, site, machine_id, topic, mqtt_ts, received_at, + source, workcenter_id, workcenter_name, workcenter_display_name, + order_present, order_coherent, odoo_state, reason, + workorder_id, production_order, product, + waiting_count, ready_count, progress_count, + payload_hash, payload_json + ) + VALUES ( + %(tenant)s, %(site)s, %(machine_id)s, %(topic)s, %(mqtt_ts)s, now(), + %(source)s, %(workcenter_id)s, %(workcenter_name)s, %(workcenter_display_name)s, + %(order_present)s, %(order_coherent)s, %(odoo_state)s, %(reason)s, + %(workorder_id)s, %(production_order)s, %(product)s, + %(waiting_count)s, %(ready_count)s, %(progress_count)s, + %(payload_hash)s, %(payload_json)s + ) + ON CONFLICT (tenant, site, machine_id) + DO UPDATE SET + topic = EXCLUDED.topic, + mqtt_ts = EXCLUDED.mqtt_ts, + received_at = now(), + source = EXCLUDED.source, + workcenter_id = EXCLUDED.workcenter_id, + workcenter_name = EXCLUDED.workcenter_name, + workcenter_display_name = EXCLUDED.workcenter_display_name, + order_present = EXCLUDED.order_present, + order_coherent = EXCLUDED.order_coherent, + odoo_state = EXCLUDED.odoo_state, + reason = EXCLUDED.reason, + workorder_id = EXCLUDED.workorder_id, + production_order = EXCLUDED.production_order, + product = EXCLUDED.product, + waiting_count = EXCLUDED.waiting_count, + ready_count = EXCLUDED.ready_count, + progress_count = EXCLUDED.progress_count, + payload_hash = EXCLUDED.payload_hash, + payload_json = EXCLUDED.payload_json + """, + {**row, "payload_json": Json(payload)}, + ) + + cur.execute( + """ + INSERT INTO mv_hot.odoo_order_state_events ( + tenant, site, machine_id, topic, mqtt_ts, received_at, + odoo_state, order_present, order_coherent, + production_order, product, reason, + waiting_count, ready_count, progress_count, + payload_hash, payload_json + ) + VALUES ( + %(tenant)s, %(site)s, %(machine_id)s, %(topic)s, %(mqtt_ts)s, now(), + %(odoo_state)s, %(order_present)s, %(order_coherent)s, + %(production_order)s, %(product)s, %(reason)s, + %(waiting_count)s, %(ready_count)s, %(progress_count)s, + %(payload_hash)s, %(payload_json)s + ) + ON CONFLICT (tenant, site, machine_id, payload_hash) + DO NOTHING + """, + {**row, "payload_json": Json(payload)}, + ) + + +def on_connect(client, userdata, flags, rc, properties=None): + if rc != 0: + print(f"MQTT connect failed rc={rc}", flush=True) + return + print(f"MQTT connected. Subscribing to {MQTT_TOPIC}", flush=True) + client.subscribe(MQTT_TOPIC, qos=1) + + +def on_message(client, userdata, msg): + try: + payload = json.loads(msg.payload.decode("utf-8")) + upsert_order_state(msg.topic, payload) + machine_id = payload.get("machine_id") or msg.topic.rstrip("/").split("/")[-1] + print( + f"ORDER_STATE {machine_id} " + f"state={payload.get('odoo_state')} " + f"present={payload.get('order_present')} " + f"coherent={payload.get('order_coherent')} " + f"mo={payload.get('production_order')}", + flush=True, + ) + except Exception as exc: + print(f"ERROR processing message topic={msg.topic}: {exc}", file=sys.stderr, flush=True) + + +client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2) + +if MQTT_USER: + client.username_pw_set(MQTT_USER, MQTT_PASSWORD) + +client.on_connect = on_connect +client.on_message = on_message + +print(f"Starting order_state_pg_sink MQTT={MQTT_HOST}:{MQTT_PORT} topic={MQTT_TOPIC}", flush=True) +client.connect(MQTT_HOST, MQTT_PORT, keepalive=30) +client.loop_forever()