411 lines
11 KiB
Python

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