40-clients/ucepsa/edge-oee-demo/tools/order_state_pg_sink.py
2026-07-19 09:40:32 +02:00

557 lines
16 KiB
Python
Executable File

#!/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",
"",
)
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()