Files
egorin/lora/provisioner.js
T

233 lines
9.1 KiB
JavaScript

/**
* ================================================================
* provisioner.js — сервис учёта и добавления устройств
* Егорин / ORBIT · Node.js 20+ · npm i mqtt pg express
* ================================================================
* Делает:
* 1. Слушает nodes/+/info → неизвестный узел = запись в devices
* со status='pending' («Ожидают добавления» в UI)
* 2. Слушает nodes/+/data → телеметрия в Postgres, online-флаг,
* события смены реле
* 3. Слушает nodes/+/status (LWT LAN-узлов) → online/offline
* 4. LoRa-узлы: offline, если молчат > OFFLINE_AFTER_S
* 5. REST API для интерфейса (см. README)
* 6. При подтверждении устройства публикует MQTT Discovery
* для Home Assistant — сущности появляются в HA сами
*
* Деплой: LXC/Docker на Proxmox, systemd unit или compose.
* ENV: MQTT_URL, PG_URL, PORT
* ================================================================
*/
const mqtt = require("mqtt");
const { Pool } = require("pg");
const express = require("express");
const MQTT_URL = process.env.MQTT_URL || "mqtt://192.168.1.200:1883";
const PG_URL = process.env.PG_URL || "postgres://energy:energy@127.0.0.1:5432/energy";
const PORT = process.env.PORT || 8090;
const OFFLINE_AFTER_S = 120; // 4 пропущенных периода LoRa-узла
const pg = new Pool({ connectionString: PG_URL });
const mq = mqtt.connect(MQTT_URL, { will: {
topic: "services/provisioner/status", payload: "offline", retain: true
}});
const lastRelays = new Map(); // node_id -> [..] для детекта фронтов
// ---------------- MQTT ----------------
mq.on("connect", () => {
mq.publish("services/provisioner/status", "online", { retain: true });
mq.subscribe(["nodes/+/info", "nodes/+/data", "nodes/+/status"]);
console.log("MQTT connected:", MQTT_URL);
});
mq.on("message", async (topic, payload) => {
const [, nodeId, kind] = topic.split("/");
try {
if (kind === "info") await onInfo(nodeId, JSON.parse(payload));
if (kind === "data") await onData(nodeId, JSON.parse(payload));
if (kind === "status") await onStatus(nodeId, payload.toString());
} catch (e) {
console.error(`${topic}: ${e.message}`);
}
});
// узел представился — регистрируем или обновляем паспорт
async function onInfo(id, m) {
await pg.query(`
INSERT INTO devices (node_id, transport, caps, fw, ip, mac, gw, last_seen, online)
VALUES ($1,$2,$3,$4,$5,$6,$7,now(),TRUE)
ON CONFLICT (node_id) DO UPDATE SET
transport=$2, caps=$3, fw=$4,
ip=COALESCE($5, devices.ip), mac=COALESCE($6, devices.mac),
gw=COALESCE($7, devices.gw), last_seen=now(), online=TRUE`,
[id, m.hw || "lora", m.caps || {}, m.fw || null,
m.ip || null, m.mac || null, m.gw || null]);
console.log(`info: ${id} (${m.hw})`);
}
// телеметрия
async function onData(id, m) {
// авторегистрация, если data пришла раньше info
await pg.query(`
INSERT INTO devices (node_id, transport) VALUES ($1,$2)
ON CONFLICT (node_id) DO NOTHING`, [id, m.rssi != null ? "lora" : "lan"]);
await pg.query(`
UPDATE devices SET last_seen=now(), online=TRUE,
last_rssi=$2, last_snr=$3, gw=COALESCE($4, gw)
WHERE node_id=$1`,
[id, m.rssi ?? null, m.snr ?? null, m.gw ?? null]);
await pg.query(`
INSERT INTO telemetry (node_id, v, a, w, kwh, pf, hz, relays, rssi, snr)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10)`,
[id, m.v ?? null, m.a ?? null, m.w ?? null, m.kwh ?? null,
m.pf ?? null, m.hz ?? null, JSON.stringify(m.r ?? []),
m.rssi ?? null, m.snr ?? null]);
// фронты реле → журнал
const prev = lastRelays.get(id);
if (Array.isArray(m.r)) {
if (prev) {
for (let i = 0; i < m.r.length; i++) {
if (prev[i] !== m.r[i]) {
await pg.query(
`INSERT INTO relay_events (node_id, relay_n, state) VALUES ($1,$2,$3)`,
[id, i, !!m.r[i]]);
}
}
}
lastRelays.set(id, m.r.slice());
}
}
// LWT LAN-узлов
async function onStatus(id, s) {
await pg.query(`UPDATE devices SET online=$2 WHERE node_id=$1`,
[id, s === "online"]);
}
// LoRa-узлы без LWT: оффлайн по таймауту
setInterval(async () => {
await pg.query(`
UPDATE devices SET online=FALSE
WHERE online AND transport='lora'
AND last_seen < now() - make_interval(secs => $1)`, [OFFLINE_AFTER_S]);
}, 30000);
// ---------------- HA MQTT Discovery ----------------
// вызывается при approve: сущности появляются в HA автоматически
async function publishHaDiscovery(d) {
const dev = {
identifiers: [`energy_${d.node_id}`],
name: d.name || d.node_id,
manufacturer: "МПО-Информ",
model: d.transport === "lan" ? "Energy Node LAN" : "Energy Node LoRa",
sw_version: d.fw || "",
};
const base = `nodes/${d.node_id}`;
const pub = (t, cfg) =>
mq.publish(t, JSON.stringify(cfg), { retain: true });
const sensors = [
["v", "Напряжение", "voltage", "V", null],
["a", "Ток", "current", "A", null],
["w", "Мощность", "power", "W", null],
["kwh", "Энергия", "energy", "kWh", "total_increasing"],
];
for (const [key, name, cls, unit, stateClass] of sensors) {
const cfg = {
name, unique_id: `${d.node_id}_${key}`,
state_topic: `${base}/data`,
value_template: `{{ value_json.${key} }}`,
device_class: cls, unit_of_measurement: unit, device: dev,
};
if (stateClass) cfg.state_class = stateClass;
pub(`homeassistant/sensor/${d.node_id}_${key}/config`, cfg);
}
const nRelays = d.caps?.relays || 0;
const names = d.relay_names || [];
for (let i = 0; i < nRelays; i++) {
pub(`homeassistant/switch/${d.node_id}_r${i}/config`, {
name: names[i] || `Реле ${i + 1}`,
unique_id: `${d.node_id}_r${i}`,
state_topic: `${base}/data`,
value_template: `{{ 'ON' if value_json.r[${i}] == 1 else 'OFF' }}`,
command_topic: `${base}/cmd`,
payload_on: JSON.stringify({ cmd: "relay", n: i, s: 1 }),
payload_off: JSON.stringify({ cmd: "relay", n: i, s: 0 }),
device: dev,
});
}
}
// ---------------- REST API для интерфейса ----------------
const app = express();
app.use(express.json());
// список устройств (+ последняя телеметрия); ?status=pending|active
app.get("/api/devices", async (req, res) => {
const { status, project } = req.query;
const cond = [], args = [];
if (status) { args.push(status); cond.push(`status=$${args.length}`); }
if (project) { args.push(project); cond.push(`project=$${args.length}`); }
const where = cond.length ? "WHERE " + cond.join(" AND ") : "";
const { rows } = await pg.query(
`SELECT * FROM v_devices_board ${where} ORDER BY project, name NULLS LAST, node_id`, args);
res.json(rows);
});
// подтвердить устройство (добавить на сервер)
app.post("/api/devices/:id/approve", async (req, res) => {
const { name, project, object, relay_names, approved_by } = req.body || {};
const { rows } = await pg.query(`
UPDATE devices SET status='active', name=$2, project=$3, object=$4,
relay_names=$5, approved_at=now(), approved_by=$6
WHERE node_id=$1 RETURNING *`,
[req.params.id, name || req.params.id, project || null, object || null,
JSON.stringify(relay_names || []), approved_by || "ui"]);
if (!rows.length) return res.status(404).json({ error: "not found" });
await publishHaDiscovery(rows[0]);
res.json(rows[0]);
});
// отклонить / отключить
app.post("/api/devices/:id/disable", async (req, res) => {
await pg.query(`UPDATE devices SET status='disabled' WHERE node_id=$1`,
[req.params.id]);
res.json({ ok: true });
});
// команда узлу (реле, poll) + журнал
app.post("/api/devices/:id/cmd", async (req, res) => {
const payload = req.body || {};
mq.publish(`nodes/${req.params.id}/cmd`, JSON.stringify(payload));
await pg.query(
`INSERT INTO commands (node_id, payload, source) VALUES ($1,$2,'ui')`,
[req.params.id, JSON.stringify(payload)]);
res.json({ ok: true });
});
// история для графиков: ?hours=24
app.get("/api/devices/:id/history", async (req, res) => {
const hours = Math.min(parseInt(req.query.hours || "24", 10), 720);
const { rows } = await pg.query(`
SELECT ts, v, a, w, kwh FROM telemetry
WHERE node_id=$1 AND ts > now() - make_interval(hours => $2)
ORDER BY ts`, [req.params.id, hours]);
res.json(rows);
});
// журнал реле
app.get("/api/devices/:id/relay-events", async (req, res) => {
const { rows } = await pg.query(`
SELECT ts, relay_n, state FROM relay_events
WHERE node_id=$1 ORDER BY ts DESC LIMIT 100`, [req.params.id]);
res.json(rows);
});
app.listen(PORT, () => console.log(`API on :${PORT}`));