From 43795f803814b20b1803a5fdcaa5a22f7daa38c0 Mon Sep 17 00:00:00 2001 From: Victor Fraile Garcia Date: Sun, 19 Jul 2026 09:40:32 +0200 Subject: [PATCH] fix(ucepsa): preserve repeated order-state transitions --- ...der_state_transition_history_fix_v0.3.2.md | 70 +++ ..._order_state_event_transition_fix_v032.sql | 135 +++++ ...rder_state_event_transition_dedup_v032.sql | 167 ++++++ .../tools/order_state_pg_sink.py | 536 +++++++++++++----- 4 files changed, 780 insertions(+), 128 deletions(-) create mode 100644 ucepsa/edge-oee-demo/docs/runbooks/096_order_state_transition_history_fix_v0.3.2.md create mode 100644 ucepsa/edge-oee-demo/ops/validate_order_state_event_transition_fix_v032.sql create mode 100644 ucepsa/edge-oee-demo/sql/versions/103_fix_order_state_event_transition_dedup_v032.sql diff --git a/ucepsa/edge-oee-demo/docs/runbooks/096_order_state_transition_history_fix_v0.3.2.md b/ucepsa/edge-oee-demo/docs/runbooks/096_order_state_transition_history_fix_v0.3.2.md new file mode 100644 index 0000000..6851365 --- /dev/null +++ b/ucepsa/edge-oee-demo/docs/runbooks/096_order_state_transition_history_fix_v0.3.2.md @@ -0,0 +1,70 @@ +# UCEPSA — Hotfix historial de transiciones Odoo v0.3.2 + +## Causa confirmada + +La tabla `mv_hot.odoo_order_state_events` tenía un índice único global: + +```text +tenant + site + machine_id + payload_hash +``` + +Esto impedía representar una secuencia válida: + +```text +A → B → A +``` + +El segundo `A` se descartaba aunque fuera una transición nueva en otro momento. + +## Síntoma observado + +CORT-02: + +```text +WH/MO/00081-002 +→ no_loaded_order +``` + +El estado `current` se actualizó correctamente, pero el evento histórico no se +insertó porque el hash de `no_loaded_order` ya existía de un hueco anterior. + +## Corrección + +- Se elimina la unicidad global por hash. +- El sink compara con el último evento de la máquina. +- Solo evita repeticiones consecutivas: + - A → A: no inserta; + - A → B: inserta B; + - A → B → A: inserta de nuevo A. +- CURRENT sigue refrescándose con cada mensaje. +- Se usa un advisory lock transaccional por máquina. +- Se repara la transición conocida de CORT-02 a las 15:51:03.861315589+02. + +## Orden de despliegue obligatorio + +1. Detener `mv_ucepsa_odoo_order_state_pg_sink`. +2. Aplicar migración 103. +3. Instalar el script corregido. +4. Arrancar el sink. +5. Ejecutar validación. + +No se debe retirar el índice con el sink antiguo en funcionamiento, porque su +`ON CONFLICT` depende de ese índice. + +## Alcance + +No modifica: + +- Odoo; +- publicador Odoo; +- MQTT; +- RevPi; +- WISE; +- balizas; +- tabla raw de paros; +- cola oficial; +- Loss Ledger. + +La reparación conocida solo corrige el intervalo de WH/MO/00081-002. +Los huecos históricos anteriores que no puedan reconstruirse con evidencia deben +mantenerse con confianza reducida. diff --git a/ucepsa/edge-oee-demo/ops/validate_order_state_event_transition_fix_v032.sql b/ucepsa/edge-oee-demo/ops/validate_order_state_event_transition_fix_v032.sql new file mode 100644 index 0000000..0feeb69 --- /dev/null +++ b/ucepsa/edge-oee-demo/ops/validate_order_state_event_transition_fix_v032.sql @@ -0,0 +1,135 @@ +\pset pager off + +\echo '=== 1. Índice global eliminado ===' + +SELECT COUNT(*) AS global_unique_hash_index_count +FROM pg_indexes +WHERE schemaname = 'mv_hot' + AND tablename = 'odoo_order_state_events' + AND indexname = + 'ux_odoo_order_state_events_state_hash'; + +\echo '=== 2. Índices nuevos ===' + +SELECT + indexname, + indexdef +FROM pg_indexes +WHERE schemaname = 'mv_hot' + AND tablename = 'odoo_order_state_events' + AND indexname IN ( + 'idx_odoo_order_state_events_machine_last', + 'idx_odoo_order_state_events_machine_hash' + ) +ORDER BY indexname; + +\echo '=== 3. Transición reparada ===' + +SELECT + id, + machine_id, + odoo_state, + order_present, + order_coherent, + production_order, + mqtt_ts, + received_at, + payload_hash, + payload_json -> '_mesavault_repair' + AS repair_metadata +FROM mv_hot.odoo_order_state_events +WHERE tenant = 'ucepsa' + AND site = 'ucepsa_onpremise' + AND machine_id = 'CORT-02' + AND odoo_state = 'no_loaded_order' + AND COALESCE(mqtt_ts, received_at) + >= '2026-07-17 15:50:59+02' + AND COALESCE(mqtt_ts, received_at) + < '2026-07-17 15:52:00+02' +ORDER BY id; + +\echo '=== 4. WH/MO/00081-002 debe quedar cerrado ===' + +SELECT + production_order, + production_id, + workorder_id, + workorder_state, + effective_from, + effective_to +FROM mv_reports_ucepsa_prod.li_order_state_intervals_v1 +WHERE machine_id = 'CORT-02' + AND production_order = 'WH/MO/00081-002' +ORDER BY effective_from; + +\echo '=== 5. Estado posterior sin orden ===' + +SELECT + odoo_state, + order_present, + order_coherent, + production_order, + effective_from, + effective_to +FROM mv_reports_ucepsa_prod.li_order_state_intervals_v1 +WHERE machine_id = 'CORT-02' + AND effective_from >= '2026-07-17 15:50:59+02' +ORDER BY effective_from +LIMIT 10; + +\echo '=== 6. No deben quedar candidatos Odoo posteriores al cierre ===' + +SELECT COUNT(*) AS false_odoo_candidates_after_close +FROM mv_reports_ucepsa_prod.li_stop_candidates_context_v1 +WHERE machine_id = 'CORT-02' + AND order_ref = 'WH/MO/00081-002' + AND started_at >= + '2026-07-17 15:51:03.861315589+02'; + +\echo '=== 7. Muestra de paros reclasificados ===' + +SELECT + source_stop_id, + started_at, + ended_at, + context_status, + primary_context_source, + order_ref, + eligibility_status, + technically_eligible, + shadow_capacity_loss_eur, + official_ledger_eligible +FROM mv_reports_ucepsa_prod.li_stop_candidates_context_v1 +WHERE source_stop_id IN ( + 2636, + 2646, + 2649, + 2657, + 2665, + 2694 +) +ORDER BY source_stop_id, started_at; + +\echo '=== 8. Cola oficial sigue vacía ===' + +SELECT COUNT(*) AS official_operator_queue_rows +FROM mv_reports_ucepsa_prod.li_operator_stop_queue_v2; + +\echo '=== 9. No debe haber eventos consecutivos con el mismo hash desde el fix ===' + +WITH ordered AS ( + SELECT + id, + machine_id, + payload_hash, + LAG(payload_hash) OVER ( + PARTITION BY tenant, site, machine_id + ORDER BY id + ) AS previous_hash + FROM mv_hot.odoo_order_state_events + WHERE tenant = 'ucepsa' + AND site = 'ucepsa_onpremise' +) +SELECT COUNT(*) AS consecutive_duplicate_event_count +FROM ordered +WHERE payload_hash = previous_hash; diff --git a/ucepsa/edge-oee-demo/sql/versions/103_fix_order_state_event_transition_dedup_v032.sql b/ucepsa/edge-oee-demo/sql/versions/103_fix_order_state_event_transition_dedup_v032.sql new file mode 100644 index 0000000..165150d --- /dev/null +++ b/ucepsa/edge-oee-demo/sql/versions/103_fix_order_state_event_transition_dedup_v032.sql @@ -0,0 +1,167 @@ +BEGIN; + +-- La unicidad global por hash impide representar A → B → A. +-- Se sustituye por deduplicación secuencial en el sink. +DROP INDEX IF EXISTS + mv_hot.ux_odoo_order_state_events_state_hash; + +CREATE INDEX IF NOT EXISTS + idx_odoo_order_state_events_machine_last +ON mv_hot.odoo_order_state_events ( + tenant, + site, + machine_id, + id DESC +); + +CREATE INDEX IF NOT EXISTS + idx_odoo_order_state_events_machine_hash +ON mv_hot.odoo_order_state_events ( + tenant, + site, + machine_id, + payload_hash +); + +COMMENT ON INDEX +mv_hot.idx_odoo_order_state_events_machine_last IS +'Soporta la comparación con el último evento de una máquina para deduplicación secuencial.'; + +COMMENT ON INDEX +mv_hot.idx_odoo_order_state_events_machine_hash IS +'Índice no único para diagnóstico de estados que pueden reaparecer en momentos distintos.'; + +-- Reparación puntual demostrada: +-- WH/MO/00081-002 terminó en Odoo a las 15:50:59+02. +-- El sink recibió no_loaded_order por primera vez a las +-- 15:51:03.861315589+02, pero el índice global por hash +-- impidió insertar la nueva transición porque el mismo hash +-- ya había aparecido a las 13:48:59+02. +WITH source_event AS ( + SELECT e.* + FROM mv_hot.odoo_order_state_events e + WHERE e.tenant = 'ucepsa' + AND e.site = 'ucepsa_onpremise' + AND e.machine_id = 'CORT-02' + AND e.odoo_state = 'no_loaded_order' + AND e.payload_hash = + '6d10a29b35613a69e761dd2c1a1c7cb6324e913950593298a19a70b0e1901c38' + ORDER BY COALESCE(e.mqtt_ts, e.received_at) + LIMIT 1 +), +repair_candidate AS ( + SELECT + s.tenant, + s.site, + s.machine_id, + s.topic, + '2026-07-17 15:51:03.861315589+02'::timestamptz + AS mqtt_ts, + now() AS received_at, + s.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, + jsonb_set( + s.payload_json + || jsonb_build_object( + '_mesavault_repair', + jsonb_build_object( + 'reason', + 'Reinserción de transición A-B-A omitida por unicidad global del payload_hash', + 'source_event_id', + s.id, + 'repaired_at', + now() + ) + ), + '{ts}', + to_jsonb( + '2026-07-17T13:51:03.861315589Z'::text + ), + true + ) AS payload_json + FROM source_event s + WHERE NOT EXISTS ( + SELECT 1 + FROM mv_hot.odoo_order_state_events e + WHERE e.tenant = 'ucepsa' + AND e.site = 'ucepsa_onpremise' + AND e.machine_id = 'CORT-02' + AND e.odoo_state = 'no_loaded_order' + AND COALESCE(e.mqtt_ts, e.received_at) + >= '2026-07-17 15:50:59+02'::timestamptz + AND COALESCE(e.mqtt_ts, e.received_at) + < '2026-07-17 15:52:00+02'::timestamptz + ) +) +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, + 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 +FROM repair_candidate; + +COMMIT; diff --git a/ucepsa/edge-oee-demo/tools/order_state_pg_sink.py b/ucepsa/edge-oee-demo/tools/order_state_pg_sink.py index 715f26e..8f4d5ef 100755 --- a/ucepsa/edge-oee-demo/tools/order_state_pg_sink.py +++ b/ucepsa/edge-oee-demo/tools/order_state_pg_sink.py @@ -4,7 +4,7 @@ import json import os import signal import sys -from datetime import datetime, timezone +from datetime import datetime import paho.mqtt.client as mqtt import psycopg2 @@ -36,7 +36,15 @@ def to_bool(value): return value if value is None: return None - return str(value).strip().lower() in ("true", "t", "1", "yes", "y", "si", "sí") + return str(value).strip().lower() in ( + "true", + "t", + "1", + "yes", + "y", + "si", + "sí", + ) def to_int(value): @@ -52,7 +60,9 @@ def parse_ts(value): if not value: return None try: - return datetime.fromisoformat(str(value).replace("Z", "+00:00")) + return datetime.fromisoformat( + str(value).replace("Z", "+00:00") + ) except Exception: return None @@ -80,32 +90,72 @@ def state_hash(payload): "ready_count": diagnostics.get("ready_count"), "progress_count": diagnostics.get("progress_count"), } - raw = json.dumps(relevant, sort_keys=True, ensure_ascii=False) + 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") +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_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/+", + 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") +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") +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") +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) + print( + "ERROR: faltan credenciales PostgreSQL", + file=sys.stderr, + ) sys.exit(1) conn = psycopg2.connect( @@ -115,15 +165,18 @@ conn = psycopg2.connect( user=PGUSER, password=PGPASSWORD, ) -conn.autocommit = True +conn.autocommit = False def upsert_order_state(topic, payload): - machine_id = payload.get("machine_id") or topic.rstrip("/").split("/")[-1] + 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) + payload_hash = state_hash(payload) row = { "tenant": TENANT, @@ -132,145 +185,372 @@ def upsert_order_state(topic, payload): "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")), + "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")), + "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": h, - "payload_json": payload, + "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), } - 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, 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, "payload_json": Json(payload)}, - ) + event_inserted = False - 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 - ) - VALUES ( - %(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 - ) - ON CONFLICT (tenant, site, machine_id, payload_hash) - DO NOTHING - """, - {**row, "payload_json": Json(payload)}, - ) + 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): +def on_connect( + client, + userdata, + flags, + rc, + properties=None, +): if rc != 0: - print(f"MQTT connect failed rc={rc}", flush=True) + print( + f"MQTT connect failed rc={rc}", + flush=True, + ) return - print(f"MQTT connected. Subscribing to {MQTT_TOPIC}", flush=True) + 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] + 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"mo={payload.get('production_order')} " + f"event={event_action}", flush=True, ) except Exception as exc: - print(f"ERROR processing message topic={msg.topic}: {exc}", file=sys.stderr, flush=True) + 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) +client = mqtt.Client( + mqtt.CallbackAPIVersion.VERSION2 +) if MQTT_USER: - client.username_pw_set(MQTT_USER, MQTT_PASSWORD) + 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) +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()