#!/usr/bin/env python3 import json import os import signal import sys from datetime import datetime, timezone 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 getenv(name, default=None, required=False): value = os.getenv(name, default) if required and not value: print(f"[ERROR] Missing required env var: {name}", file=sys.stderr) sys.exit(2) return value MQTT_HOST = getenv("MQTT_HOST", required=True) MQTT_PORT = int(getenv("MQTT_PORT", "1883")) MQTT_TOPIC = getenv("MQTT_TOPIC", required=True) MQTT_USER = getenv("MQTT_USER") MQTT_PASSWORD = getenv("MQTT_PASSWORD") PGHOST = getenv("PGHOST", required=True) PGPORT = int(getenv("PGPORT", "5432")) PGDATABASE = getenv("PGDATABASE", required=True) PGUSER = getenv("PGUSER", required=True) PGPASSWORD = getenv("PGPASSWORD", required=True) DEFAULT_TENANT = getenv("DEFAULT_TENANT", "ucepsa") DEFAULT_SITE = getenv("DEFAULT_SITE", "demo_edge_oee") DEFAULT_ASSET = getenv("DEFAULT_ASSET", "revpi_oee_node_01") def num(payload, key): value = payload.get(key) if value is None or value == "": return None try: return float(value) except Exception: return None def bool_or_false(payload, key): value = payload.get(key) if isinstance(value, bool): return value if value is None: return False if isinstance(value, str): return value.lower() in ("true", "1", "yes", "ok") return bool(value) def parse_ts(payload): raw = payload.get("ts") if not raw: return datetime.now(timezone.utc) try: return datetime.fromisoformat(raw.replace("Z", "+00:00")) except Exception: return datetime.now(timezone.utc) def connect_pg(): return psycopg2.connect( host=PGHOST, port=PGPORT, dbname=PGDATABASE, user=PGUSER, password=PGPASSWORD, connect_timeout=5, ) PG_CONN = connect_pg() PG_CONN.autocommit = True def ensure_pg(): global PG_CONN try: with PG_CONN.cursor() as cur: cur.execute("SELECT 1;") except Exception: try: PG_CONN.close() except Exception: pass PG_CONN = connect_pg() PG_CONN.autocommit = True def insert_payload(topic, payload): ensure_pg() tenant = payload.get("tenant") or DEFAULT_TENANT site = payload.get("site") or DEFAULT_SITE asset = payload.get("asset") or DEFAULT_ASSET machine_id = payload.get("machine_id") or None order_id = payload.get("order_id") or None comm_ok = bool_or_false(payload, "comm_ok") error_text = payload.get("error_text") sql = """ INSERT INTO mv_hot.edge_oee_eastron_readings ( ts, tenant, site, asset, source, meter_model, meter_id, machine_id, order_id, voltage_l1_v, voltage_l2_v, voltage_l3_v, current_l1_a, current_l2_a, current_l3_a, power_l1_kw, power_l2_kw, power_l3_kw, power_total_kw, frequency_hz, import_kwh, export_kwh, comm_ok, error_text, payload_json ) VALUES ( %(ts)s, %(tenant)s, %(site)s, %(asset)s, %(source)s, %(meter_model)s, %(meter_id)s, %(machine_id)s, %(order_id)s, %(voltage_l1_v)s, %(voltage_l2_v)s, %(voltage_l3_v)s, %(current_l1_a)s, %(current_l2_a)s, %(current_l3_a)s, %(power_l1_kw)s, %(power_l2_kw)s, %(power_l3_kw)s, %(power_total_kw)s, %(frequency_hz)s, %(import_kwh)s, %(export_kwh)s, %(comm_ok)s, %(error_text)s, %(payload_json)s ); """ power_total_kw = num(payload, "power_total_kw") params = { "ts": parse_ts(payload), "tenant": tenant, "site": site, "asset": asset, "source": payload.get("source") or "mqtt_edge_oee_sink", "meter_model": payload.get("meter_model"), "meter_id": payload.get("meter_id"), "machine_id": machine_id, "order_id": order_id, "voltage_l1_v": num(payload, "voltage_l1_v"), "voltage_l2_v": num(payload, "voltage_l2_v"), "voltage_l3_v": num(payload, "voltage_l3_v"), "current_l1_a": num(payload, "current_l1_a"), "current_l2_a": num(payload, "current_l2_a"), "current_l3_a": num(payload, "current_l3_a"), "power_l1_kw": num(payload, "power_l1_kw") if payload.get("power_l1_kw") is not None else power_total_kw, "power_l2_kw": num(payload, "power_l2_kw"), "power_l3_kw": num(payload, "power_l3_kw"), "power_total_kw": power_total_kw, "frequency_hz": num(payload, "frequency_hz"), "import_kwh": num(payload, "import_kwh"), "export_kwh": num(payload, "export_kwh"), "comm_ok": comm_ok, "error_text": error_text, "payload_json": Json({ **payload, "_mqtt_topic": topic, "_ingested_by": "edge_oee_mqtt_pg_sink", }), } with PG_CONN.cursor() as cur: cur.execute(sql, params) def insert_event(topic, payload): ensure_pg() tenant = payload.get("tenant") or DEFAULT_TENANT site = payload.get("site") or DEFAULT_SITE vertical = payload.get("vertical") or "edge_oee" asset = payload.get("asset") or DEFAULT_ASSET raw_ts = payload.get("event_ts") or payload.get("ts") try: event_ts = datetime.fromisoformat(raw_ts.replace("Z", "+00:00")) if raw_ts else datetime.now(timezone.utc) except Exception: event_ts = datetime.now(timezone.utc) def to_bool(value): if isinstance(value, bool): return value if value is None: return None if isinstance(value, str): return value.lower() in ("true", "1", "yes", "ok") return bool(value) def to_int(value): if value is None or value == "": return None try: return int(value) except Exception: return None def to_float(value): if value is None or value == "": return None try: return float(value) except Exception: return None sql = """ INSERT INTO mv_hot.edge_oee_node_events ( event_ts, topic, tenant, site, vertical, asset, machine_id, order_id, event_type, event_name, previous_auto_signal, new_auto_signal, machine_running, reason_pending, cycle_pulse_count, cycle_rate_ppm, last_cycle_age_s, payload_json ) VALUES ( %(event_ts)s, %(topic)s, %(tenant)s, %(site)s, %(vertical)s, %(asset)s, %(machine_id)s, %(order_id)s, %(event_type)s, %(event_name)s, %(previous_auto_signal)s, %(new_auto_signal)s, %(machine_running)s, %(reason_pending)s, %(cycle_pulse_count)s, %(cycle_rate_ppm)s, %(last_cycle_age_s)s, %(payload_json)s ); """ params = { "event_ts": event_ts, "topic": topic, "tenant": tenant, "site": site, "vertical": vertical, "asset": asset, "machine_id": payload.get("machine_id") or None, "order_id": payload.get("order_id") or None, "event_type": payload.get("event_type") or "unknown", "event_name": payload.get("event_name") or "unknown", "previous_auto_signal": to_bool(payload.get("previous_auto_signal")), "new_auto_signal": to_bool(payload.get("new_auto_signal")), "machine_running": to_bool(payload.get("machine_running")), "reason_pending": to_bool(payload.get("reason_pending")), "cycle_pulse_count": to_int(payload.get("cycle_pulse_count")), "cycle_rate_ppm": to_float(payload.get("cycle_rate_ppm")), "last_cycle_age_s": to_float(payload.get("last_cycle_age_s")), "payload_json": Json({ **payload, "_mqtt_topic": topic, "_ingested_by": "edge_oee_mqtt_pg_sink", }), } with PG_CONN.cursor() as cur: cur.execute(sql, params) def process_machine_stop_event(topic, payload): if payload.get("event_type") != "machine_state_changed": return event_name = payload.get("event_name") if event_name not in ("machine_stopped", "machine_started"): return ensure_pg() tenant = payload.get("tenant") or DEFAULT_TENANT site = payload.get("site") or DEFAULT_SITE vertical = payload.get("vertical") or "edge_oee" asset = payload.get("asset") or DEFAULT_ASSET machine_id = payload.get("machine_id") or None order_id = payload.get("order_id") or None raw_ts = payload.get("event_ts") or payload.get("ts") try: event_ts = datetime.fromisoformat(raw_ts.replace("Z", "+00:00")) if raw_ts else datetime.now(timezone.utc) except Exception: event_ts = datetime.now(timezone.utc) if event_name == "machine_stopped": sql = """ INSERT INTO mv_hot.edge_oee_machine_stops ( tenant, site, vertical, asset, machine_id, order_id, started_at, status, classification_status, start_event_name, start_payload_json ) SELECT %(tenant)s, %(site)s, %(vertical)s, %(asset)s, %(machine_id)s, %(order_id)s, %(started_at)s, 'OPEN', 'PENDING', %(event_name)s, %(payload_json)s WHERE NOT EXISTS ( SELECT 1 FROM mv_hot.edge_oee_machine_stops WHERE asset = %(asset)s AND COALESCE(machine_id, '') = COALESCE(%(machine_id)s, '') AND status = 'OPEN' ); """ params = { "tenant": tenant, "site": site, "vertical": vertical, "asset": asset, "machine_id": machine_id, "order_id": order_id, "started_at": event_ts, "event_name": event_name, "payload_json": Json({ **payload, "_mqtt_topic": topic, "_stop_logic": "open_stop", }), } with PG_CONN.cursor() as cur: cur.execute(sql, params) return if event_name == "machine_started": sql = """ WITH open_stop AS ( SELECT id, started_at, value_hour_eur FROM mv_hot.edge_oee_machine_stops WHERE asset = %(asset)s AND COALESCE(machine_id, '') = COALESCE(%(machine_id)s, '') AND status = 'OPEN' ORDER BY started_at DESC LIMIT 1 ) UPDATE mv_hot.edge_oee_machine_stops s SET ended_at = %(ended_at)s, status = 'CLOSED', duration_s = EXTRACT(EPOCH FROM (%(ended_at)s - s.started_at))::numeric(14,2), cost_eur = ROUND(((EXTRACT(EPOCH FROM (%(ended_at)s - s.started_at)) / 3600.0) * s.value_hour_eur)::numeric, 2), end_event_name = %(event_name)s, end_payload_json = %(payload_json)s, updated_at = now() FROM open_stop WHERE s.id = open_stop.id; """ params = { "asset": asset, "machine_id": machine_id, "ended_at": event_ts, "event_name": event_name, "payload_json": Json({ **payload, "_mqtt_topic": topic, "_stop_logic": "close_stop", }), } with PG_CONN.cursor() as cur: cur.execute(sql, params) return def on_connect(client, userdata, flags, reason_code, properties=None): print(f"[MQTT] connected reason_code={reason_code}") client.subscribe(MQTT_TOPIC) print(f"[MQTT] subscribed topic={MQTT_TOPIC}") def on_message(client, userdata, msg): try: body = msg.payload.decode("utf-8") payload = json.loads(body) if msg.topic.endswith("/telemetry"): insert_payload(msg.topic, payload) asset = payload.get("asset") voltage = payload.get("voltage_l1_v") power = payload.get("power_total_kw") comm_ok = payload.get("comm_ok") pulses = payload.get("cycle_pulse_count") rate = payload.get("cycle_rate_ppm") print( f"[OK][telemetry] topic={msg.topic} asset={asset} " f"voltage={voltage} power_kw={power} pulses={pulses} rate_ppm={rate} comm_ok={comm_ok}", flush=True, ) elif msg.topic.endswith("/event"): insert_event(msg.topic, payload) process_machine_stop_event(msg.topic, payload) print( f"[OK][event] topic={msg.topic} " f"event_name={payload.get('event_name')} " f"asset={payload.get('asset')} " f"previous={payload.get('previous_auto_signal')} " f"new={payload.get('new_auto_signal')} " f"pulses={payload.get('cycle_pulse_count')}", flush=True, ) else: print(f"[WARN] Ignored topic={msg.topic}", flush=True) except Exception as exc: print(f"[ERROR] processing topic={msg.topic}: {type(exc).__name__}: {exc}", file=sys.stderr, flush=True) def main(): print("[START] edge_oee_mqtt_pg_sink") print(f"[CONFIG] MQTT {MQTT_HOST}:{MQTT_PORT} topic={MQTT_TOPIC}") print(f"[CONFIG] PostgreSQL {PGHOST}:{PGPORT}/{PGDATABASE} user={PGUSER}") client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2) if MQTT_USER: client.username_pw_set(MQTT_USER, MQTT_PASSWORD) print(f"[CONFIG] MQTT auth user={MQTT_USER}", flush=True) else: print("[WARN] MQTT auth disabled: MQTT_USER not set", flush=True) client.on_connect = on_connect client.on_message = on_message client.connect(MQTT_HOST, MQTT_PORT, keepalive=60) client.loop_start() try: while RUNNING: signal.pause() except AttributeError: import time while RUNNING: time.sleep(1) finally: client.loop_stop() client.disconnect() try: PG_CONN.close() except Exception: pass print("[STOP] edge_oee_mqtt_pg_sink") return 0 if __name__ == "__main__": sys.exit(main())