fix(ucepsa): preserve repeated order-state transitions

This commit is contained in:
Victor Fraile Garcia 2026-07-19 09:40:32 +02:00
parent e7aa13c433
commit 43795f8038
4 changed files with 780 additions and 128 deletions

View File

@ -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.

View File

@ -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;

View File

@ -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;

View File

@ -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", "")
return str(value).strip().lower() in (
"true",
"t",
"1",
"yes",
"y",
"si",
"",
)
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()