#!/usr/bin/env python3 import hashlib import json import os import signal import sys from datetime import datetime 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_id": payload.get("production_id"), "production_order": payload.get("production_order"), "product_id": payload.get("product_id"), "product": payload.get("product"), "qty_production": payload.get("qty_production"), "date_start": payload.get("date_start"), "date_finished": payload.get("date_finished"), "workorder_name": payload.get("workorder_name"), "workorder_state": payload.get("workorder_state"), "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 = False 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")) payload_hash = 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 ), "workorder_name": payload.get( "workorder_name" ), "workorder_state": payload.get( "workorder_state" ), "production_id": to_int( payload.get("production_id") ), "production_order": payload.get( "production_order" ), "product_id": to_int( payload.get("product_id") ), "product": payload.get("product"), "qty_production": payload.get( "qty_production" ), "date_start": parse_ts( payload.get("date_start") ), "date_finished": parse_ts( payload.get("date_finished") ), "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": payload_hash, "payload_json": Json(payload), } event_inserted = False try: with conn: with conn.cursor() as cur: # Serializa los mensajes de una misma máquina incluso si, # por error, llegaran a existir dos instancias del sink. lock_key = ( f"{TENANT}|{SITE}|{machine_id}" ) cur.execute( """ SELECT pg_advisory_xact_lock( hashtextextended(%s, 0) ) """, (lock_key,), ) # La deduplicación es secuencial: # A → A no inserta un evento nuevo. # A → B inserta B. # A → B → A vuelve a insertar A. 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, workorder_id, workorder_name, workorder_state, production_id, production_order, product_id, product, qty_production, date_start, date_finished, reason, waiting_count, ready_count, progress_count, payload_hash, payload_json ) SELECT %(tenant)s, %(site)s, %(machine_id)s, %(topic)s, %(mqtt_ts)s, now(), %(odoo_state)s, %(order_present)s, %(order_coherent)s, %(workorder_id)s, %(workorder_name)s, %(workorder_state)s, %(production_id)s, %(production_order)s, %(product_id)s, %(product)s, %(qty_production)s, %(date_start)s, %(date_finished)s, %(reason)s, %(waiting_count)s, %(ready_count)s, %(progress_count)s, %(payload_hash)s, %(payload_json)s WHERE ( SELECT e.payload_hash FROM mv_hot.odoo_order_state_events e WHERE e.tenant = %(tenant)s AND e.site = %(site)s AND e.machine_id = %(machine_id)s ORDER BY e.id DESC LIMIT 1 ) IS DISTINCT FROM %(payload_hash)s RETURNING id """, row, ) event_inserted = ( cur.fetchone() is not None ) # CURRENT se refresca con cada mensaje, aunque no haya # transición, para conservar frescura y payload vivo. 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, workorder_name, workorder_state, production_id, production_order, product_id, product, qty_production, date_start, date_finished, 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, %(workorder_name)s, %(workorder_state)s, %(production_id)s, %(production_order)s, %(product_id)s, %(product)s, %(qty_production)s, %(date_start)s, %(date_finished)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, workorder_name = EXCLUDED.workorder_name, workorder_state = EXCLUDED.workorder_state, production_id = EXCLUDED.production_id, production_order = EXCLUDED.production_order, product_id = EXCLUDED.product_id, product = EXCLUDED.product, qty_production = EXCLUDED.qty_production, date_start = EXCLUDED.date_start, date_finished = EXCLUDED.date_finished, 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, ) return event_inserted except Exception: if not conn.closed: conn.rollback() raise 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") ) event_inserted = upsert_order_state( msg.topic, payload, ) machine_id = ( payload.get("machine_id") or msg.topic.rstrip("/").split("/")[-1] ) event_action = ( "transition_inserted" if event_inserted else "same_as_last" ) 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')} " f"event={event_action}", flush=True, ) except Exception as exc: print( "ERROR processing message " f"topic={msg.topic}: " f"{type(exc).__name__}: {exc!r}", 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( "Starting order_state_pg_sink " f"MQTT={MQTT_HOST}:{MQTT_PORT} " f"topic={MQTT_TOPIC}", flush=True, ) client.connect( MQTT_HOST, MQTT_PORT, keepalive=30, ) client.loop_forever()