557 lines
16 KiB
Python
Executable File
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",
|
|
"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()
|