Files
SE-LTD/collector/src/index.js
T
2026-06-23 14:48:36 +02:00

547 lines
18 KiB
JavaScript

import { createServer } from "node:http";
import { mkdirSync, readFileSync, writeFileSync, existsSync } from "node:fs";
import { dirname, join } from "node:path";
import { fileURLToPath } from "node:url";
import { Buffer } from "node:buffer";
import net from "node:net";
const __dirname = dirname(fileURLToPath(import.meta.url));
const dataDir = join(__dirname, "..", "data");
const storeFile = join(dataDir, "sources.json");
const port = Number(process.env.PORT || 18110);
mkdirSync(dataDir, { recursive: true });
function baseStore() {
return {
sources: [],
datapoints: [],
scans: { opcua: [], bacnet: [] },
updatedAt: null,
};
}
function readStore() {
if (!existsSync(storeFile)) {
return baseStore();
}
try {
const parsed = JSON.parse(readFileSync(storeFile, "utf8"));
return {
...baseStore(),
...parsed,
scans: { ...baseStore().scans, ...(parsed.scans || {}) },
sources: Array.isArray(parsed.sources) ? parsed.sources : [],
datapoints: Array.isArray(parsed.datapoints) ? parsed.datapoints : [],
};
} catch {
return baseStore();
}
}
function writeStore(payload) {
writeFileSync(storeFile, JSON.stringify(payload, null, 2), "utf8");
}
function touchStore(mutator) {
const store = readStore();
const next = mutator(structuredClone(store)) || store;
next.updatedAt = new Date().toISOString();
writeStore(next);
return next;
}
function replyJson(res, statusCode, body) {
res.writeHead(statusCode, {
"Content-Type": "application/json; charset=utf-8",
"Access-Control-Allow-Origin": "*",
"Access-Control-Allow-Methods": "GET,POST,PUT,DELETE,OPTIONS",
"Access-Control-Allow-Headers": "Content-Type, Authorization",
});
res.end(JSON.stringify(body));
}
function readBody(req) {
return new Promise((resolve, reject) => {
const chunks = [];
req.on("data", (chunk) => chunks.push(chunk));
req.on("end", () => {
try {
const text = Buffer.concat(chunks).toString("utf8");
resolve(text ? JSON.parse(text) : {});
} catch (error) {
reject(error);
}
});
req.on("error", reject);
});
}
function slugify(value) {
return String(value || "")
.toLowerCase()
.replace(/[^a-z0-9]+/g, "_")
.replace(/^_+|_+$/g, "") || "quelle";
}
function normalizeProtocol(value) {
const protocol = String(value || "").toLowerCase();
if (["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip"].includes(protocol)) {
return protocol;
}
return "unknown";
}
function normalizeWriteMode(value) {
return String(value || "").toLowerCase() === "cov" ? "cov" : "interval";
}
function normalizePollInterval(value) {
const numeric = Number(value);
return Number.isFinite(numeric) && numeric > 0 ? Math.round(numeric) : 60;
}
function normalizePointKind(value) {
return String(value || "").toLowerCase() === "digital" ? "digital" : "analog";
}
function normalizeUnit(value) {
const raw = String(value ?? "").trim();
if (!raw) {
return "°C";
}
return raw
.replace(/°/g, "°")
.replace(/^\?C$/i, "°C")
.replace(/^°\s*C$/i, "°C");
}
function normalizePointIndex(value, fallback) {
const numeric = Number(value);
return Number.isInteger(numeric) && numeric >= 1 && numeric <= 200 ? numeric : fallback;
}
function ensureSource(store, body) {
const now = new Date().toISOString();
const source = {
id: body.id || `src-${Date.now()}`,
protocol: normalizeProtocol(body.protocol),
name: body.name || "Neue Quelle",
host: body.host || "",
port: body.port || "",
deviceId: body.deviceId || "",
displayName: body.displayName || body.name || "Neue Quelle",
tableName: body.tableName || `trend_${slugify(body.name || body.protocol || "daten")}`,
writeMode: normalizeWriteMode(body.writeMode),
pollIntervalSeconds: normalizePollInterval(body.pollIntervalSeconds),
options: {
...(body.options || {}),
writeMode: normalizeWriteMode(body.writeMode),
pollIntervalSeconds: normalizePollInterval(body.pollIntervalSeconds),
},
lastPollAt: body.lastPollAt || null,
lastPollStatus: body.lastPollStatus || "idle",
lastPollError: body.lastPollError || "",
createdAt: body.createdAt || now,
updatedAt: now,
};
const existingIndex = store.sources.findIndex((item) => item.id === source.id);
if (existingIndex >= 0) {
store.sources[existingIndex] = { ...store.sources[existingIndex], ...source };
} else {
store.sources.push(source);
}
return source;
}
function splitDelimited(text, delimiter) {
return text
.split(/\r?\n/)
.map((line) => line.trim())
.filter(Boolean)
.map((line) => line.split(delimiter).map((part) => part.trim().replace(/^"|"$/g, "")));
}
function parseModbusCsv(payload) {
const delimiter = payload.delimiter || ";";
const rows = splitDelimited(payload.csv || "", delimiter);
if (rows.length < 2) {
return [];
}
const header = rows[0].map((value) => value.toLowerCase());
return rows.slice(1).map((row, index) => {
const get = (name, fallback = "") => {
const columnIndex = header.indexOf(name);
return columnIndex >= 0 ? row[columnIndex] || fallback : fallback;
};
const address = get("address", get("register", String(index + 1)));
const pointIndex = normalizePointIndex(get("pointindex", index + 1), index + 1);
return {
id: `dp-modbus-${Date.now()}-${index}`,
sourceId: payload.sourceId || "",
protocol: "modbus-tcp",
address,
pointIndex,
name: get("name", `Register ${index + 1}`),
alias: get("alias", get("name", `Register ${index + 1}`)),
unit: normalizeUnit(get("unit")),
dataType: get("datatype", "holding-register"),
kind: normalizePointKind(get("kind", get("type", "analog"))),
scale: Number(get("scale", 1)) || 1,
host: payload.host || "",
tableName: payload.tableName || "",
createdAt: new Date().toISOString(),
};
});
}
function decodeXmlEntities(value) {
return String(value || "")
.replaceAll("&amp;", "&")
.replaceAll("&quot;", '"')
.replaceAll("&lt;", "<")
.replaceAll("&gt;", ">")
.replaceAll("&apos;", "'");
}
function parseKnxXml(payload) {
const xml = String(payload.xml || "");
const datapoints = [];
const pattern = /<GroupAddress\b([^>]*)\/?>(?:<\/GroupAddress>)?/gi;
let match;
let index = 0;
while ((match = pattern.exec(xml))) {
const attrs = match[1];
const getAttr = (name) => {
const attrMatch = attrs.match(new RegExp(`${name}="([^"]*)"`, "i"));
return attrMatch ? decodeXmlEntities(attrMatch[1]) : "";
};
const address = getAttr("Address") || getAttr("address");
const name = getAttr("Name") || getAttr("name") || `KNX ${index + 1}`;
const dpt = getAttr("DatapointType") || getAttr("DPTs") || getAttr("DPT") || "";
if (address) {
datapoints.push({
id: `dp-knx-${Date.now()}-${index}`,
sourceId: payload.sourceId || "",
protocol: "knx-ip",
address,
pointIndex: index + 1,
name,
alias: name,
unit: normalizeUnit(""),
dataType: dpt || "group-address",
kind: dpt.startsWith("1.") ? "digital" : "analog",
host: payload.host || "",
tableName: payload.tableName || "",
createdAt: new Date().toISOString(),
});
index += 1;
}
}
return datapoints;
}
function createManualDatapoint(store, payload) {
const source = store.sources.find((item) => item.id === payload.sourceId);
if (!source) {
throw new Error("Quelle nicht gefunden.");
}
const existing = store.datapoints.filter((item) => item.sourceId === source.id);
const fallbackIndex = existing.length + 1;
const pointIndex = normalizePointIndex(payload.pointIndex, fallbackIndex);
const protocol = normalizeProtocol(payload.protocol || source.protocol);
const alias = String(payload.alias || payload.name || `Punkt ${pointIndex}`).trim() || `Punkt ${pointIndex}`;
const name = String(payload.name || payload.alias || alias).trim() || alias;
const datapoint = {
id: payload.id || `dp-${protocol}-${Date.now()}-${pointIndex}`,
sourceId: source.id,
protocol,
pointIndex,
name,
alias,
unit: normalizeUnit(payload.unit),
kind: normalizePointKind(payload.kind),
dataType: String(payload.dataType || (protocol === "knx-ip" ? "group-address" : "holding-register")).trim(),
address: String(payload.address || payload.register || "").trim(),
nodeId: String(payload.nodeId || "").trim(),
scale: Number(payload.scale || 1) || 1,
host: source.host || "",
tableName: source.tableName || "",
createdAt: new Date().toISOString(),
};
if (!datapoint.address && protocol !== "opc-ua") {
datapoint.address = String(pointIndex);
}
return datapoint;
}
function tryTcpConnect(host, portNumber, timeout = 1200) {
return new Promise((resolve) => {
if (!host || !portNumber) {
resolve(false);
return;
}
const socket = new net.Socket();
let settled = false;
const finish = (value) => {
if (!settled) {
settled = true;
socket.destroy();
resolve(value);
}
};
socket.setTimeout(timeout);
socket.once("connect", () => finish(true));
socket.once("timeout", () => finish(false));
socket.once("error", () => finish(false));
socket.connect(portNumber, host);
});
}
const server = createServer(async (req, res) => {
const url = new URL(req.url || "/", `http://${req.headers.host || "localhost"}`);
if (req.method === "OPTIONS") {
replyJson(res, 204, {});
return;
}
if (req.method === "GET" && url.pathname === "/health") {
replyJson(res, 200, { ok: true, service: "collector", timestamp: new Date().toISOString() });
return;
}
if (req.method === "GET" && url.pathname === "/capabilities") {
replyJson(res, 200, {
protocols: ["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip"],
imports: ["modbus-csv", "knx-xml"],
scans: ["opc-ua", "bacnet-ip"],
manualDatapoints: ["modbus-tcp", "knx-ip"],
note: "Importe, manuelle Datenpunkte und Diagnose sind aktiv. Live-Polling und echte Feldbus-Kommunikation sind noch im nächsten Ausbauschritt.",
});
return;
}
if (req.method === "GET" && url.pathname === "/sources") {
replyJson(res, 200, readStore());
return;
}
if (req.method === "GET" && url.pathname === "/datapoints") {
const store = readStore();
const sourceId = url.searchParams.get("sourceId");
const items = sourceId ? store.datapoints.filter((item) => item.sourceId === sourceId) : store.datapoints;
replyJson(res, 200, { datapoints: items, total: items.length, updatedAt: store.updatedAt });
return;
}
if (req.method === "GET" && url.pathname === "/debug") {
const store = readStore();
replyJson(res, 200, {
updatedAt: store.updatedAt,
sourceCount: store.sources.length,
datapointCount: store.datapoints.length,
sources: store.sources.map((source) => ({
...source,
datapointCount: store.datapoints.filter((item) => item.sourceId === source.id).length,
})),
latestScans: {
opcua: store.scans.opcua[0] || null,
bacnet: store.scans.bacnet[0] || null,
},
driverState: {
pollingImplemented: false,
note: "Der Collector speichert Quellen, Importe und Scan-Ergebnisse bereits. Echtes zyklisches Polling auf Bus-Ebene ist aktuell noch nicht implementiert.",
},
});
return;
}
if (req.method === "POST" && url.pathname === "/sources") {
try {
const body = await readBody(req);
const next = touchStore((store) => {
const source = ensureSource(store, body);
return { ...store, lastSource: source.id };
});
replyJson(res, 201, next.sources.at(-1));
} catch {
replyJson(res, 400, { error: "Ungültige JSON-Daten." });
}
return;
}
if (req.method === "PUT" && url.pathname.startsWith("/sources/")) {
try {
const sourceId = url.pathname.split("/").at(-1);
const body = await readBody(req);
const next = touchStore((store) => {
const existing = store.sources.find((item) => item.id === sourceId);
if (!existing) {
throw new Error("not-found");
}
ensureSource(store, { ...existing, ...body, id: sourceId });
return store;
});
replyJson(res, 200, next.sources.find((item) => item.id === sourceId));
} catch (error) {
replyJson(res, error.message === "not-found" ? 404 : 400, {
error: error.message === "not-found" ? "Quelle nicht gefunden." : "Ungültige JSON-Daten.",
});
}
return;
}
if (req.method === "DELETE" && url.pathname.startsWith("/sources/")) {
const sourceId = url.pathname.split("/").at(-1);
const next = touchStore((store) => ({
...store,
sources: store.sources.filter((item) => item.id !== sourceId),
datapoints: store.datapoints.filter((item) => item.sourceId !== sourceId),
}));
replyJson(res, 200, { ok: true, sources: next.sources.length, datapoints: next.datapoints.length });
return;
}
if (req.method === "POST" && url.pathname === "/datapoints/manual") {
try {
const body = await readBody(req);
const next = touchStore((store) => {
const datapoint = createManualDatapoint(store, body);
store.datapoints = store.datapoints.filter((item) => !(item.sourceId === datapoint.sourceId && item.pointIndex === datapoint.pointIndex));
store.datapoints.push(datapoint);
return store;
});
const created = next.datapoints.find((item) => item.sourceId === body.sourceId && item.pointIndex === normalizePointIndex(body.pointIndex, 1));
replyJson(res, 201, { ok: true, datapoint: created || null, updatedAt: next.updatedAt });
} catch (error) {
replyJson(res, 400, { error: error.message || "Datenpunkt konnte nicht angelegt werden." });
}
return;
}
if (req.method === "POST" && url.pathname === "/imports/modbus-csv") {
try {
const body = await readBody(req);
const imported = parseModbusCsv(body);
const next = touchStore((store) => {
if (body.sourceId) {
store.datapoints = store.datapoints.filter((item) => !(item.sourceId === body.sourceId && item.protocol === "modbus-tcp"));
}
store.datapoints.push(...imported);
return store;
});
replyJson(res, 200, { imported: imported.length, datapoints: imported, updatedAt: next.updatedAt });
} catch {
replyJson(res, 400, { error: "CSV-Import fehlgeschlagen." });
}
return;
}
if (req.method === "POST" && url.pathname === "/imports/knx-xml") {
try {
const body = await readBody(req);
const imported = parseKnxXml(body);
const next = touchStore((store) => {
if (body.sourceId) {
store.datapoints = store.datapoints.filter((item) => !(item.sourceId === body.sourceId && item.protocol === "knx-ip"));
}
store.datapoints.push(...imported);
return store;
});
replyJson(res, 200, { imported: imported.length, datapoints: imported, updatedAt: next.updatedAt });
} catch {
replyJson(res, 400, { error: "KNX-XML-Import fehlgeschlagen." });
}
return;
}
if (req.method === "POST" && url.pathname === "/scan/opcua") {
try {
const body = await readBody(req);
const host = body.host || "";
if (!host) {
replyJson(res, 400, { error: "Für den OPC-UA-Scan fehlt Host oder IP." });
return;
}
const portNumber = Number(body.port || 4840);
const reachable = await tryTcpConnect(host, portNumber);
const result = {
host,
port: portNumber,
reachable,
endpoints: reachable
? [
{ url: `opc.tcp://${host}:${portNumber}`, securityMode: "None", securityPolicy: "None" },
]
: [],
sampleNodes: reachable
? [
{ nodeId: "ns=0;i=2258", browseName: "Server", displayName: "Server" },
{ nodeId: "ns=0;i=2267", browseName: "ServerStatus", displayName: "ServerStatus" },
]
: [],
scannedAt: new Date().toISOString(),
};
touchStore((store) => {
store.scans.opcua.unshift(result);
store.scans.opcua = store.scans.opcua.slice(0, 20);
return store;
});
replyJson(res, 200, result);
} catch {
replyJson(res, 400, { error: "OPC-UA-Scan fehlgeschlagen." });
}
return;
}
if (req.method === "POST" && url.pathname === "/scan/bacnet") {
try {
const body = await readBody(req);
const hosts = Array.isArray(body.hosts) ? body.hosts : [body.host].filter(Boolean);
const results = [];
for (const host of hosts) {
const reachable = await tryTcpConnect(host, Number(body.port || 47808), 650);
results.push({
host,
port: Number(body.port || 47808),
reachable,
deviceId: body.deviceId || "",
name: reachable ? `BACnet Gerät ${host}` : "Nicht erreichbar",
});
}
const payload = { devices: results, scannedAt: new Date().toISOString() };
touchStore((store) => {
store.scans.bacnet.unshift(payload);
store.scans.bacnet = store.scans.bacnet.slice(0, 20);
return store;
});
replyJson(res, 200, payload);
} catch {
replyJson(res, 400, { error: "BACnet-Scan fehlgeschlagen." });
}
return;
}
replyJson(res, 404, { error: "Nicht gefunden." });
});
server.listen(port, () => {
console.log(`SE Local Trenddata collector listening on port ${port}`);
});