This commit is contained in:
jhartworks
2026-06-25 10:20:38 +02:00
parent f0cb94ce4e
commit bf1214239f
5 changed files with 680 additions and 109 deletions
+560 -45
View File
@@ -1,14 +1,29 @@
import { createServer } from "node:http";
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";
import mysql from "mysql2/promise";
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);
const pollTickMs = Number(process.env.POLL_TICK_MS || 1000);
let transactionId = 1;
let polling = false;
const db = mysql.createPool({
host: process.env.DB_HOST || "localhost",
port: Number(process.env.DB_PORT || 3306),
user: process.env.DB_USER || "root",
password: process.env.DB_PASSWORD || "SE3112",
database: process.env.DB_NAME || "wago",
waitForConnections: true,
connectionLimit: 4,
charset: "utf8mb4",
});
mkdirSync(dataDir, { recursive: true });
@@ -85,6 +100,14 @@ function slugify(value) {
.replace(/^_+|_+$/g, "") || "quelle";
}
function sanitizeTrendTable(value) {
const table = String(value || "").toLowerCase();
if (!/^(isp[a-z0-9_]*|trend_[a-z0-9_]+)$/.test(table)) {
throw new Error("Ungültiger Trendtabellenname.");
}
return table;
}
function normalizeProtocol(value) {
const protocol = String(value || "").toLowerCase();
if (["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip"].includes(protocol)) {
@@ -123,8 +146,15 @@ function normalizePointIndex(value, fallback) {
return Number.isInteger(numeric) && numeric >= 1 && numeric <= 200 ? numeric : fallback;
}
function buildSourceTableName(body) {
return `trend_${slugify(body.displayName || body.name || body.protocol || "daten")}`;
}
function ensureSource(store, body) {
const now = new Date().toISOString();
const existingIndex = store.sources.findIndex((item) => item.id === body.id);
const existing = existingIndex >= 0 ? store.sources[existingIndex] : null;
const tableName = sanitizeTrendTable(existing?.tableName || buildSourceTableName(body));
const source = {
id: body.id || `src-${Date.now()}`,
protocol: normalizeProtocol(body.protocol),
@@ -133,7 +163,7 @@ function ensureSource(store, body) {
port: body.port || "",
deviceId: body.deviceId || "",
displayName: body.displayName || body.name || "Neue Quelle",
tableName: body.tableName || `trend_${slugify(body.name || body.protocol || "daten")}`,
tableName,
writeMode: normalizeWriteMode(body.writeMode),
pollIntervalSeconds: normalizePollInterval(body.pollIntervalSeconds),
options: {
@@ -141,6 +171,7 @@ function ensureSource(store, body) {
writeMode: normalizeWriteMode(body.writeMode),
pollIntervalSeconds: normalizePollInterval(body.pollIntervalSeconds),
},
lastValues: body.lastValues || {},
lastPollAt: body.lastPollAt || null,
lastPollStatus: body.lastPollStatus || "idle",
lastPollError: body.lastPollError || "",
@@ -148,7 +179,6 @@ function ensureSource(store, body) {
updatedAt: now,
};
const existingIndex = store.sources.findIndex((item) => item.id === source.id);
if (existingIndex >= 0) {
store.sources[existingIndex] = { ...store.sources[existingIndex], ...source };
} else {
@@ -313,6 +343,520 @@ function tryTcpConnect(host, portNumber, timeout = 1200) {
});
}
function inferModbusAddress(datapoint) {
const rawAddress = Number(String(datapoint.address || datapoint.register || "").replace(/[^0-9]/g, ""));
const dataType = String(datapoint.dataType || "").toLowerCase();
const pointIndex = normalizePointIndex(datapoint.pointIndex, 1);
const raw = Number.isFinite(rawAddress) && rawAddress > 0 ? rawAddress : pointIndex;
let functionCode = dataType.includes("coil") ? 1 : dataType.includes("discrete") ? 2 : dataType.includes("input") ? 4 : 3;
let address = Math.max(0, raw - 1);
if (raw >= 40001) {
functionCode = 3;
address = raw - 40001;
} else if (raw >= 30001) {
functionCode = 4;
address = raw - 30001;
} else if (raw >= 10001) {
functionCode = 2;
address = raw - 10001;
} else if (raw <= 9999 && functionCode > 2) {
address = raw - 1;
}
return { functionCode, address };
}
function readModbusTcp(source, datapoint) {
return new Promise((resolve, reject) => {
const host = String(source.host || "").trim();
const portNumber = Number(source.port || 502);
if (!host) {
reject(new Error("Modbus Host fehlt."));
return;
}
const { functionCode, address } = inferModbusAddress(datapoint);
const unitId = Number(source.deviceId || 1) & 0xff;
const currentTransaction = transactionId;
transactionId = transactionId >= 0xffff ? 1 : transactionId + 1;
const request = Buffer.alloc(12);
request.writeUInt16BE(currentTransaction, 0);
request.writeUInt16BE(0, 2);
request.writeUInt16BE(6, 4);
request.writeUInt8(unitId, 6);
request.writeUInt8(functionCode, 7);
request.writeUInt16BE(address, 8);
request.writeUInt16BE(1, 10);
const socket = new net.Socket();
const chunks = [];
let settled = false;
const finish = (error, value) => {
if (settled) {
return;
}
settled = true;
socket.destroy();
if (error) {
reject(error);
} else {
resolve(value);
}
};
socket.setTimeout(1800);
socket.once("timeout", () => finish(new Error("Modbus Timeout.")));
socket.once("error", (error) => finish(error));
socket.on("data", (chunk) => {
chunks.push(chunk);
const response = Buffer.concat(chunks);
if (response.length < 9) {
return;
}
const expectedLength = response.readUInt16BE(4) + 6;
if (response.length < expectedLength) {
return;
}
const responseFunction = response.readUInt8(7);
if (responseFunction & 0x80) {
finish(new Error(`Modbus Exception ${response.readUInt8(8)}`));
return;
}
if (responseFunction !== functionCode) {
finish(new Error("Modbus Antwort passt nicht zur Anfrage."));
return;
}
const scale = Number(datapoint.scale || 1) || 1;
if (functionCode === 1 || functionCode === 2) {
const digital = (response.readUInt8(9) & 0x01) ? 1 : 0;
finish(null, digital);
return;
}
if (response.length < 11) {
finish(new Error("Modbus Antwort zu kurz."));
return;
}
const dataType = String(datapoint.dataType || "").toLowerCase();
const raw = dataType.includes("int16") || dataType.includes("signed")
? response.readInt16BE(9)
: response.readUInt16BE(9);
finish(null, raw * scale);
});
socket.connect(portNumber, host, () => socket.write(request));
});
}
async function loadOptionalPackage(packageName) {
try {
const module = await import(packageName);
return module.default || module;
} catch (error) {
throw new Error(`${packageName} ist nicht installiert. Container bitte neu bauen, damit der Treiber verfügbar ist.`);
}
}
function buildOpcuaEndpoint(sourceOrBody) {
const host = String(sourceOrBody.host || "").trim();
if (!host) {
throw new Error("OPC-UA Host fehlt.");
}
if (host.startsWith("opc.tcp://")) {
return host;
}
return `opc.tcp://${host}:${Number(sourceOrBody.port || 4840)}`;
}
async function scanOpcuaEndpoint(body) {
const { OPCUAClient } = await loadOptionalPackage("node-opcua");
const endpointUrl = buildOpcuaEndpoint(body);
const client = OPCUAClient.create({
endpointMustExist: false,
connectionStrategy: { initialDelay: 250, maxRetry: 0 },
requestedSessionTimeout: 10000,
});
try {
await client.connect(endpointUrl);
const session = await client.createSession();
const browseResult = await session.browse("ObjectsFolder");
const sampleNodes = (browseResult.references || []).slice(0, 80).map((reference) => ({
nodeId: reference.nodeId?.toString?.() || String(reference.nodeId || ""),
browseName: reference.browseName?.toString?.() || String(reference.browseName || ""),
displayName: reference.displayName?.text || reference.displayName?.toString?.() || "",
}));
await session.close();
return {
host: body.host || "",
port: Number(body.port || 4840),
reachable: true,
endpoints: [{ url: endpointUrl, securityMode: "auto", securityPolicy: "auto" }],
sampleNodes,
scannedAt: new Date().toISOString(),
};
} finally {
await client.disconnect().catch(() => undefined);
}
}
async function readOpcuaValue(source, datapoint) {
const { OPCUAClient } = await loadOptionalPackage("node-opcua");
const endpointUrl = buildOpcuaEndpoint(source);
const nodeId = datapoint.nodeId || datapoint.address;
if (!nodeId) {
throw new Error(`OPC-UA NodeId fehlt für ${datapoint.alias || datapoint.name}.`);
}
const client = OPCUAClient.create({
endpointMustExist: false,
connectionStrategy: { initialDelay: 250, maxRetry: 0 },
requestedSessionTimeout: 10000,
});
try {
await client.connect(endpointUrl);
const session = await client.createSession();
const dataValue = await session.readVariableValue(nodeId);
await session.close();
const value = dataValue?.value?.value;
if (typeof value === "boolean") {
return value ? 1 : 0;
}
const numeric = Number(value);
return Number.isFinite(numeric) ? numeric * (Number(datapoint.scale || 1) || 1) : String(value ?? "");
} finally {
await client.disconnect().catch(() => undefined);
}
}
function parseBacnetObject(datapoint) {
const dataType = String(datapoint.dataType || "analogInput").trim();
const objectTypeMap = {
analoginput: 0,
analogoutput: 1,
analogvalue: 2,
binaryinput: 3,
binaryoutput: 4,
binaryvalue: 5,
multiinput: 13,
multioutput: 14,
multivalue: 19,
device: 8,
};
const objectType = Number.isFinite(Number(datapoint.objectType))
? Number(datapoint.objectType)
: objectTypeMap[dataType.toLowerCase().replace(/[^a-z]/g, "")] ?? 0;
const instance = Number(datapoint.objectInstance || datapoint.address || datapoint.pointIndex || 0);
return { type: objectType, instance };
}
async function scanBacnetNetwork(body) {
const bacnet = await loadOptionalPackage("node-bacnet");
const client = new bacnet({ apduTimeout: Number(body.timeout || 5000), port: Number(body.port || 47808) });
const devices = [];
return await new Promise((resolve, reject) => {
const timer = setTimeout(() => {
client.close();
resolve({ devices, scannedAt: new Date().toISOString() });
}, Number(body.timeout || 5000));
client.on("iAm", (device) => {
devices.push({
host: device.address,
port: Number(body.port || 47808),
reachable: true,
deviceId: device.deviceId,
name: `BACnet Gerät ${device.deviceId}`,
});
});
try {
if (body.deviceId) {
client.whoIs(Number(body.deviceId), Number(body.deviceId));
} else {
client.whoIs();
}
} catch (error) {
clearTimeout(timer);
client.close();
reject(error);
}
});
}
async function readBacnetValue(source, datapoint) {
const bacnet = await loadOptionalPackage("node-bacnet");
const client = new bacnet({ apduTimeout: 5000, port: Number(source.port || 47808) });
const propertyId = Number(datapoint.propertyId || 85);
const objectId = parseBacnetObject(datapoint);
const host = String(source.host || datapoint.host || "").trim();
if (!host) {
throw new Error("BACnet Host fehlt.");
}
return await new Promise((resolve, reject) => {
client.readProperty(host, objectId, propertyId, (error, value) => {
client.close();
if (error) {
reject(error);
return;
}
const raw = value?.values?.[0]?.value;
if (typeof raw === "boolean") {
resolve(raw ? 1 : 0);
return;
}
const numeric = Number(raw);
resolve(Number.isFinite(numeric) ? numeric * (Number(datapoint.scale || 1) || 1) : String(raw ?? ""));
});
});
}
async function readKnxValue(source, datapoint) {
const knx = await loadOptionalPackage("knx");
const groupAddress = String(datapoint.address || "").trim();
if (!groupAddress) {
throw new Error(`KNX Gruppenadresse fehlt für ${datapoint.alias || datapoint.name}.`);
}
return await new Promise((resolve, reject) => {
let done = false;
let connection;
const finish = (error, value) => {
if (done) {
return;
}
done = true;
clearTimeout(timer);
try { connection?.Disconnect?.(); } catch {}
if (error) {
reject(error);
} else {
const numeric = Number(value);
resolve(typeof value === "boolean" ? (value ? 1 : 0) : Number.isFinite(numeric) ? numeric : String(value ?? ""));
}
};
const timer = setTimeout(() => finish(new Error("KNX Timeout.")), 7000);
connection = new knx.Connection({
ipAddr: source.host,
ipPort: Number(source.port || 3671),
handlers: {
connected: () => {
try {
connection.read(groupAddress);
} catch (error) {
finish(error);
}
},
event: (_event, _src, dest, value) => {
if (!dest || dest === groupAddress) {
finish(null, value);
}
},
error: (connectionError) => finish(connectionError instanceof Error ? connectionError : new Error(String(connectionError))),
},
});
});
}
async function pollGenericSource(source, datapoints, readValue) {
const readings = [];
const nextValues = { ...(source.lastValues || {}) };
for (const datapoint of datapoints) {
const value = await readValue(source, datapoint);
const key = String(datapoint.pointIndex);
const serialized = String(value);
if (source.writeMode !== "cov" || nextValues[key] !== serialized) {
readings.push({ pointIndex: datapoint.pointIndex, value });
}
nextValues[key] = serialized;
}
const written = await writeTrendRow(source, readings);
updateSourceState(source.id, {
lastValues: nextValues,
lastPollAt: new Date().toISOString(),
lastPollStatus: written ? `ok: ${readings.length} Wert(e) geschrieben` : "ok: keine Änderung",
lastPollError: "",
});
}
function buildValueColumns() {
const columns = [];
for (let index = 1; index <= 200; index += 1) {
columns.push('`Wert' + index + '` TEXT NULL');
}
return columns.join(", ");
}
async function ensureTrendTable(tableName) {
const table = sanitizeTrendTable(tableName);
await db.query(`
CREATE TABLE IF NOT EXISTS \`${table}\` (
\`id\` INT(11) NOT NULL AUTO_INCREMENT,
\`datum\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
\`userlevel\` INT(1) NOT NULL,
${buildValueColumns()},
PRIMARY KEY (\`id\`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci
`);
return table;
}
async function writeTrendRow(source, readings) {
if (!readings.length) {
return false;
}
const table = await ensureTrendTable(source.tableName);
const columns = ["`userlevel`"];
const placeholders = ["?"];
const values = [0];
readings.forEach((reading) => {
const pointIndex = normalizePointIndex(reading.pointIndex, 0);
if (!pointIndex) {
return;
}
columns.push("`Wert" + pointIndex + "`");
placeholders.push("?");
values.push(String(reading.value));
});
if (values.length === 1) {
return false;
}
await db.query(`INSERT INTO \`${table}\` (${columns.join(", ")}) VALUES (${placeholders.join(", ")})`, values);
return true;
}
function shouldPollSource(source, now) {
if (!source.id || !source.protocol || !source.tableName) {
return false;
}
const intervalMs = normalizePollInterval(source.pollIntervalSeconds) * 1000;
const lastPoll = source.lastPollAt ? new Date(source.lastPollAt).getTime() : 0;
return !lastPoll || now - lastPoll >= intervalMs;
}
function updateSourceState(sourceId, changes) {
touchStore((store) => {
const index = store.sources.findIndex((source) => source.id === sourceId);
if (index >= 0) {
store.sources[index] = {
...store.sources[index],
...changes,
updatedAt: new Date().toISOString(),
};
}
return store;
});
}
async function pollModbusSource(source, datapoints) {
const readings = [];
const nextValues = { ...(source.lastValues || {}) };
for (const datapoint of datapoints) {
const value = await readModbusTcp(source, datapoint);
const key = String(datapoint.pointIndex);
const serialized = String(value);
if (source.writeMode !== "cov" || nextValues[key] !== serialized) {
readings.push({ pointIndex: datapoint.pointIndex, value });
}
nextValues[key] = serialized;
}
const written = await writeTrendRow(source, readings);
updateSourceState(source.id, {
lastValues: nextValues,
lastPollAt: new Date().toISOString(),
lastPollStatus: written ? `ok: ${readings.length} Wert(e) geschrieben` : "ok: keine Änderung",
lastPollError: "",
});
}
async function pollSource(source, datapoints) {
const sourcePoints = datapoints.filter((datapoint) => datapoint.sourceId === source.id);
if (!sourcePoints.length) {
updateSourceState(source.id, {
lastPollAt: new Date().toISOString(),
lastPollStatus: "idle: keine Datenpunkte",
lastPollError: "",
});
return;
}
try {
if (source.protocol === "modbus-tcp") {
await pollModbusSource(source, sourcePoints);
return;
}
if (source.protocol === "opc-ua") {
await pollGenericSource(source, sourcePoints, readOpcuaValue);
return;
}
if (source.protocol === "bacnet-ip") {
await pollGenericSource(source, sourcePoints, readBacnetValue);
return;
}
if (source.protocol === "knx-ip") {
await pollGenericSource(source, sourcePoints, readKnxValue);
return;
}
updateSourceState(source.id, {
lastPollAt: new Date().toISOString(),
lastPollStatus: "unsupported",
lastPollError: `${source.protocol} wird vom Collector nicht unterstützt.`,
});
} catch (error) {
updateSourceState(source.id, {
lastPollAt: new Date().toISOString(),
lastPollStatus: "error",
lastPollError: error.message || "Polling fehlgeschlagen.",
});
}
}
async function pollDueSources() {
if (polling) {
return;
}
polling = true;
try {
const store = readStore();
const now = Date.now();
for (const source of store.sources) {
if (shouldPollSource(source, now)) {
await pollSource(source, store.datapoints);
}
}
} finally {
polling = false;
}
}
setInterval(() => {
pollDueSources().catch((error) => console.error("Collector polling failed", error));
}, Math.max(500, pollTickMs));
const server = createServer(async (req, res) => {
const url = new URL(req.url || "/", `http://${req.headers.host || "localhost"}`);
@@ -331,8 +875,9 @@ const server = createServer(async (req, res) => {
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.",
manualDatapoints: ["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip"],
polling: ["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip"],
note: "Modbus TCP, OPC UA, BACnet IP und KNX IP werden zyklisch gelesen und in Trendtabellen geschrieben.",
});
return;
}
@@ -365,8 +910,9 @@ const server = createServer(async (req, res) => {
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.",
pollingImplemented: true,
pollingProtocols: ["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip"],
note: "Alle konfigurierten Protokollquellen pollten über native Treiber. Fehler stehen direkt an der Quelle.",
},
});
return;
@@ -474,38 +1020,19 @@ const server = createServer(async (req, res) => {
if (req.method === "POST" && url.pathname === "/scan/opcua") {
try {
const body = await readBody(req);
const host = body.host || "";
if (!host) {
if (!body.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(),
};
const result = await scanOpcuaEndpoint(body);
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." });
} catch (error) {
replyJson(res, 400, { error: error.message || "OPC-UA-Scan fehlgeschlagen." });
}
return;
}
@@ -513,27 +1040,15 @@ const server = createServer(async (req, res) => {
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() };
const payload = await scanBacnetNetwork(body);
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." });
} catch (error) {
replyJson(res, 400, { error: error.message || "BACnet-Scan fehlgeschlagen." });
}
return;
}