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); const LOCAL_POINT_LIMIT = 2000; 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 }); 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,PATCH,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 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)) { 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 repairText(value) { let current = String(value ?? ""); for (let attempt = 0; attempt < 3; attempt += 1) { const currentScore = (current.match(/[\u00C3\u00C2\uFFFD]/g) || []).length; if (!currentScore) { break; } const repaired = Buffer.from(current, "latin1").toString("utf8"); const repairedScore = (repaired.match(/[\u00C3\u00C2\uFFFD]/g) || []).length; if (repairedScore < currentScore) { current = repaired; continue; } break; } return current.replace(/\u00A0/g, " ").trim(); } function normalizeText(value) { return repairText(String(value === undefined || value === null ? "" : value).trim()); } function normalizeUnit(value) { const raw = normalizeText(value); const collapsed = raw.replace(/\s+/g, ""); if (!collapsed) { return "°C"; } if (/^(?:\?C|\u00B0C)$/i.test(collapsed)) { return "°C"; } return raw; } function normalizePointUnit(value, kind = "analog") { const normalizedKind = normalizePointKind(kind); const raw = normalizeText(value); if (!raw) { return normalizedKind === "digital" ? "bool" : "°C"; } return normalizeUnit(raw); } function normalizePointIndex(value, fallback) { const numeric = Number(value); return Number.isInteger(numeric) && numeric >= 1 && numeric <= LOCAL_POINT_LIMIT ? 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), name: normalizeText(body.name) || "Neue Quelle", host: normalizeText(body.host), port: normalizeText(body.port), deviceId: normalizeText(body.deviceId), displayName: normalizeText(body.displayName || body.name) || "Neue Quelle", tableName, writeMode: normalizeWriteMode(body.writeMode), pollIntervalSeconds: normalizePollInterval(body.pollIntervalSeconds), options: { ...(body.options || {}), writeMode: normalizeWriteMode(body.writeMode), pollIntervalSeconds: normalizePollInterval(body.pollIntervalSeconds), }, lastValues: body.lastValues || {}, lastPollAt: body.lastPollAt || null, lastPollStatus: body.lastPollStatus || "idle", lastPollError: body.lastPollError || "", createdAt: body.createdAt || now, updatedAt: now, }; 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: normalizeText(get("name", `Register ${index + 1}`)) || `Register ${index + 1}`, alias: normalizeText(get("alias", get("name", `Register ${index + 1}`))) || normalizeText(get("name", `Register ${index + 1}`)) || `Register ${index + 1}`, unit: normalizePointUnit(get("unit"), get("kind", get("type", "analog"))), dataType: get("datatype", "holding-register"), kind: normalizePointKind(get("kind", get("type", "analog"))), enabled: payload.enabled === true, writeMode: normalizeWriteMode(get("writemode", payload.writeMode || "")), condition: get("condition", payload.condition || ""), scale: Number(get("scale", 1)) || 1, host: payload.host || "", tableName: payload.tableName || "", createdAt: new Date().toISOString(), }; }); } function decodeXmlEntities(value) { return String(value || "") .replaceAll("&", "&") .replaceAll(""", '"') .replaceAll("<", "<") .replaceAll(">", ">") .replaceAll("'", "'"); } function parseKnxXml(payload) { const xml = String(payload.xml || ""); const datapoints = []; const pattern = /]*)\/?>(?:<\/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 = normalizeText(getAttr("Name") || getAttr("name") || `KNX ${index + 1}`) || `KNX ${index + 1}`; const dpt = normalizeText(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: normalizePointUnit("", dpt.startsWith("1.") ? "digital" : "analog"), dataType: dpt || "group-address", kind: dpt.startsWith("1.") ? "digital" : "analog", enabled: payload.enabled === true, writeMode: normalizeWriteMode(payload.writeMode || ""), condition: payload.condition || "", 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 = normalizeText(payload.alias || payload.name || `Punkt ${pointIndex}`) || `Punkt ${pointIndex}`; const name = normalizeText(payload.name || payload.alias || alias) || alias; const datapoint = { id: payload.id || `dp-${protocol}-${Date.now()}-${pointIndex}`, sourceId: source.id, protocol, pointIndex, name, alias, unit: normalizePointUnit(payload.unit, payload.kind), kind: normalizePointKind(payload.kind), enabled: payload.enabled !== false, writeMode: normalizeWriteMode(payload.writeMode || source.writeMode), condition: normalizeText(payload.condition), dataType: normalizeText(payload.dataType || (protocol === "knx-ip" ? "group-address" : "holding-register")), address: normalizeText(payload.address || payload.register || ""), nodeId: normalizeText(payload.nodeId || ""), scale: Number(payload.scale || 1) || 1, host: source.host || "", tableName: source.tableName || "", currentValue: payload.currentValue ?? null, lastReadAt: payload.lastReadAt || null, createdAt: new Date().toISOString(), }; if (!datapoint.address && protocol !== "opc-ua") { datapoint.address = String(pointIndex); } return datapoint; } function createScanDatapoint(source, payload, index) { const pointIndex = normalizePointIndex(payload.pointIndex, index + 1); const protocol = normalizeProtocol(payload.protocol || source.protocol); const alias = normalizeText(payload.alias || payload.name || payload.displayName || `Punkt ${pointIndex}`) || `Punkt ${pointIndex}`; return { id: payload.id || `dp-${protocol}-${Date.now()}-${Math.round(Math.random() * 100000)}-${pointIndex}`, sourceId: source.id, protocol, pointIndex, name: normalizeText(payload.name || alias) || alias, alias, unit: normalizePointUnit(payload.unit, payload.kind), kind: normalizePointKind(payload.kind), enabled: payload.enabled === true, writeMode: normalizeWriteMode(payload.writeMode || source.writeMode), condition: normalizeText(payload.condition), dataType: normalizeText(payload.dataType), address: normalizeText(payload.address), nodeId: normalizeText(payload.nodeId), objectType: payload.objectType, objectInstance: payload.objectInstance, propertyId: payload.propertyId || 85, scale: Number(payload.scale || 1) || 1, host: normalizeText(payload.host || source.host), tableName: source.tableName || "", currentValue: payload.currentValue ?? null, lastReadAt: payload.currentValue === undefined ? null : new Date().toISOString(), createdAt: new Date().toISOString(), }; } function buildScannedDatapoints(store, source, rawPoints, toPayload) { const existingPoints = store.datapoints.filter((item) => item.sourceId === source.id); const existingById = new Map(existingPoints.map((item) => [String(item.id), item])); const seenIds = new Set(); let nextPointIndex = existingPoints.reduce((max, item) => Math.max(max, Number(item.pointIndex) || 0), 0) + 1; return rawPoints.reduce((items, rawPoint, index) => { const payload = toPayload(rawPoint, index); const stableId = String(payload.id || ""); if (stableId) { if (seenIds.has(stableId)) { return items; } seenIds.add(stableId); } const existing = stableId ? existingById.get(stableId) : null; const pointIndex = existing ? existing.pointIndex : normalizePointIndex(payload.pointIndex, nextPointIndex); if (!existing) { nextPointIndex = Math.max(nextPointIndex, pointIndex + 1); } const scanned = createScanDatapoint(source, { ...payload, pointIndex }, index); items.push(!existing ? scanned : { ...existing, ...scanned, pointIndex: existing.pointIndex, alias: existing.alias || scanned.alias, name: existing.name || scanned.name, unit: normalizePointUnit(existing.unit || scanned.unit, existing.kind || scanned.kind), kind: existing.kind || scanned.kind, enabled: existing.enabled === true, writeMode: normalizeWriteMode(existing.writeMode || scanned.writeMode), condition: normalizeText(existing.condition || scanned.condition), dataType: normalizeText(existing.dataType || scanned.dataType), address: normalizeText(existing.address || scanned.address), nodeId: normalizeText(existing.nodeId || scanned.nodeId), scale: Number(existing.scale || scanned.scale || 1) || 1, host: normalizeText(existing.host || scanned.host), createdAt: existing.createdAt || scanned.createdAt, }); return items; }, []); } function upsertDatapoints(store, datapoints) { datapoints.forEach((datapoint) => { store.datapoints = store.datapoints.filter((item) => item.id !== datapoint.id && !(item.sourceId === datapoint.sourceId && item.pointIndex === datapoint.pointIndex)); store.datapoints.push(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); }); } function inferModbusAddress(datapoint) { const addressText = String(datapoint.address || datapoint.register || "").trim(); const rawAddress = Number(addressText.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 (functionCode === 3 && raw >= 40001 && raw <= 49999) { address = raw - 40001; } else if (functionCode === 4 && raw >= 30001 && raw <= 39999) { address = raw - 30001; } else if (functionCode === 2 && raw >= 10001 && raw <= 19999) { address = raw - 10001; } else if (!addressText.match(/^[134]\d{4,}$/) && raw <= 9999 && functionCode > 2) { address = raw - 1; } return { functionCode, address }; } function swapModbusWords(buffer) { if (!Buffer.isBuffer(buffer) || buffer.length < 4) { return buffer; } return Buffer.from([buffer[2], buffer[3], buffer[0], buffer[1]]); } function resolveModbusReadDetails(datapoint) { const dataType = String(datapoint.dataType || "").toLowerCase(); const scale = Number(datapoint.scale || 1) || 1; const quantity = /(?:float32|real|int32|uint32|dint|udint)/.test(dataType) ? 2 : 1; const wantsSwap = /(?:swap|swapped|cdab|badc)/.test(dataType); return { quantity, decode(responseBuffer, functionCode) { if (functionCode === 1 || functionCode === 2) { return (responseBuffer.readUInt8(0) & 0x01) ? 1 : 0; } if (responseBuffer.length < quantity * 2) { throw new Error("Modbus Antwort enth\u00e4lt zu wenige Register."); } const baseBuffer = responseBuffer.subarray(0, quantity * 2); const buffer = wantsSwap && quantity === 2 ? swapModbusWords(baseBuffer) : baseBuffer; if (/(?:float32|real)/.test(dataType)) { return buffer.readFloatBE(0) * scale; } if (/(?:uint32|udint)/.test(dataType)) { return buffer.readUInt32BE(0) * scale; } if (/(?:int32|dint)/.test(dataType)) { return buffer.readInt32BE(0) * scale; } if (/(?:int16|signed)/.test(dataType)) { return buffer.readInt16BE(0) * scale; } return buffer.readUInt16BE(0) * scale; }, }; } 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 { quantity, decode } = resolveModbusReadDetails(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(quantity, 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.once("close", (hadError) => { if (!settled) { finish(new Error(hadError ? "Modbus Verbindung wurde fehlerhaft geschlossen." : "Modbus Verbindung wurde ohne Antwort geschlossen.")); } }); 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 byteCount = Number(response.readUInt8(8) || 0); const startOffset = 9; if (response.length < startOffset + byteCount) { finish(new Error("Modbus Antwort zu kurz.")); return; } try { finish(null, decode(response.subarray(startOffset, startOffset + byteCount), functionCode)); } catch (error) { finish(error instanceof Error ? error : new Error(String(error))); } }); 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 readOpcuaNodeValue(session, nodeId) { try { const dataValue = await session.readVariableValue(nodeId); const value = dataValue?.value?.value; return value === undefined ? null : value; } catch { return null; } } async function browseOpcuaChildren(session, nodeId, pathParts, state) { if (state.nodes.length >= state.maxNodes || pathParts.length > state.maxDepth) { return []; } const browseResult = await session.browse(nodeId); const children = []; for (const reference of browseResult.references || []) { if (state.nodes.length >= state.maxNodes) { break; } const refNodeId = reference.nodeId?.toString?.() || String(reference.nodeId || ""); const browseName = reference.browseName?.toString?.() || String(reference.browseName || ""); const displayName = reference.displayName?.text || reference.displayName?.toString?.() || browseName || refNodeId; const nodeClass = reference.nodeClass?.key || reference.nodeClass?.toString?.() || String(reference.nodeClass || ""); const isVariable = /Variable/i.test(nodeClass) || Number(reference.nodeClass?.value ?? reference.nodeClass) === 2; const nodePath = pathParts.concat(displayName).filter(Boolean); const currentValue = isVariable ? await readOpcuaNodeValue(session, refNodeId) : null; const item = { id: refNodeId, nodeId: refNodeId, browseName, displayName, path: nodePath.join(" / "), nodeClass, kind: typeof currentValue === "boolean" ? "digital" : "analog", dataType: "node-id", currentValue, hasValue: currentValue !== null, children: [], }; state.nodes.push(item); children.push(item); if (!isVariable && pathParts.length + 1 < state.maxDepth) { item.children = await browseOpcuaChildren(session, refNodeId, nodePath, state); } } return children; } 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: 15000, }); try { await client.connect(endpointUrl); const session = await client.createSession(); const state = { nodes: [], maxNodes: Math.min(Number(body.maxNodes || 500), LOCAL_POINT_LIMIT), maxDepth: Math.min(Number(body.maxDepth || 5), 8), }; const tree = await browseOpcuaChildren(session, "ObjectsFolder", [], state); await session.close(); return { host: body.host || "", port: Number(body.port || 4840), reachable: true, endpoints: [{ url: endpointUrl, securityMode: "auto", securityPolicy: "auto" }], tree, nodes: state.nodes, sampleNodes: state.nodes, 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 rawAddress = String(datapoint.address || "").trim(); const shorthand = rawAddress.match(/^(ai|ao|av|bi|bo|bv)(\d+)$/i); const shorthandMap = { ai: "analogInput", ao: "analogOutput", av: "analogValue", bi: "binaryInput", bo: "binaryOutput", bv: "binaryValue" }; const dataType = String(datapoint.dataType || (shorthand ? shorthandMap[shorthand[1].toLowerCase()] : "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 || (shorthand ? shorthand[2] : datapoint.address) || datapoint.pointIndex || 0); return { type: objectType, instance }; } function bacnetObjectLabel(objectId) { const typeMap = { 0: "analogInput", 1: "analogOutput", 2: "analogValue", 3: "binaryInput", 4: "binaryOutput", 5: "binaryValue", 8: "device", 13: "multiStateInput", 14: "multiStateOutput", 19: "multiStateValue", }; const prefix = typeMap[objectId.type] || "object" + objectId.type; return prefix + objectId.instance; } function readBacnetProperty(client, host, objectId, propertyId) { return new Promise((resolve) => { client.readProperty(host, objectId, propertyId, (error, value) => { if (error) { resolve(null); return; } resolve(value); }); }); } async function scanBacnetObjects(client, device, portNumber) { const host = device.address; const deviceObject = { type: 8, instance: Number(device.deviceId) }; const objectList = await readBacnetProperty(client, host, deviceObject, 76); const values = objectList?.values || []; const points = []; for (const entry of values.slice(0, LOCAL_POINT_LIMIT)) { const objectId = entry.value || entry; if (!objectId || objectId.type === 8 || !Number.isFinite(Number(objectId.instance))) { continue; } const presentValue = await readBacnetProperty(client, host, objectId, 85); const raw = presentValue?.values?.[0]?.value ?? null; const label = bacnetObjectLabel(objectId); points.push({ id: host + "-" + objectId.type + "-" + objectId.instance, host, port: portNumber, objectType: objectId.type, objectInstance: objectId.instance, propertyId: 85, address: String(objectId.instance), dataType: label.replace(/[0-9]+$/, ""), name: label, alias: label, kind: typeof raw === "boolean" || objectId.type >= 3 && objectId.type <= 5 ? "digital" : "analog", currentValue: raw, hasValue: raw !== null, }); } return points; } async function scanBacnetNetwork(body) { const bacnet = await loadOptionalPackage("node-bacnet"); const portNumber = Number(body.port || 47808); const timeout = Number(body.timeout || 6000); const client = new bacnet({ apduTimeout: timeout, port: portNumber }); const devices = []; return await new Promise((resolve, reject) => { const timer = setTimeout(async () => { try { for (const device of devices) { device.points = await scanBacnetObjects(client, device, portNumber); } client.close(); resolve({ devices, points: devices.flatMap((device) => device.points || []), scannedAt: new Date().toISOString() }); } catch (error) { client.close(); reject(error); } }, timeout); client.on("iAm", (device) => { if (!devices.some((item) => item.host === device.address && item.deviceId === device.deviceId)) { devices.push({ host: device.address, port: portNumber, reachable: true, deviceId: device.deviceId, name: "BACnet Gerät " + device.deviceId, points: [], }); } }); 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 || "")); }); }); } function normalizeKnxDatapointType(dataType) { const raw = String(dataType || "").trim(); if (!raw) { return ""; } const match = raw.match(/(?:dpt\s*)?(\d+)(?:[.-](\d+))?/i); if (!match) { return raw.toLowerCase(); } return match[2] ? `${match[1]}.${match[2]}` : match[1]; } function decodeKnxPayload(value, dataType) { if (value === undefined || value === null) { return null; } if (typeof value === "boolean") { return value ? 1 : 0; } if (typeof value === "number") { return value; } if (typeof value === "object" && !Buffer.isBuffer(value)) { if (Array.isArray(value.data)) { return decodeKnxPayload(Buffer.from(value.data), dataType); } if (Array.isArray(value.buffer)) { return decodeKnxPayload(Buffer.from(value.buffer), dataType); } if ("value" in value && value.value !== value) { return decodeKnxPayload(value.value, dataType); } } const normalizedType = normalizeKnxDatapointType(dataType); const mainType = normalizedType.split(/[.-]/)[0]; const buffer = Buffer.isBuffer(value) ? value : value instanceof Uint8Array ? Buffer.from(value) : Array.isArray(value) ? Buffer.from(value) : Buffer.from(String(value), "binary"); if (mainType === "1") { return (buffer[buffer.length - 1] & 0x01) ? 1 : 0; } if (mainType === "5") { return buffer[buffer.length - 1] || 0; } if (mainType === "9" && buffer.length >= 2) { const hi = buffer[0]; const lo = buffer[1]; const sign = (hi & 0x80) ? -1 : 1; const exponent = (hi & 0x78) >> 3; let mantissa = ((hi & 0x07) << 8) | lo; if (sign === -1) { mantissa = -(~(mantissa - 1) & 0x07ff); } return 0.01 * mantissa * Math.pow(2, exponent); } if (mainType === "13" && buffer.length >= 4) { return buffer.readInt32BE(0); } if (mainType === "14" && buffer.length >= 4) { return buffer.readFloatBE(0); } const numeric = Number(String(value).replace(",", ".")); if (Number.isFinite(numeric)) { return numeric; } return repairText(String(value)); } 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\u00fcr ${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 decoded = decodeKnxPayload(value, datapoint.dataType); if (typeof decoded === "number" && Number.isFinite(decoded)) { const scale = Number(datapoint.scale || 1) || 1; resolve(decoded * scale); } else { resolve(decoded === null ? "" : decoded); } } }; 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))), }, }); }); } function matchesLogCondition(condition, value) { const text = String(condition || "").trim(); if (!text) { return true; } const numeric = Number(value); const match = text.match(/^(?:value\s*)?(>=|<=|>|<|==|!=)\s*(-?\d+(?:[.,]\d+)?)$/i); if (!match || !Number.isFinite(numeric)) { return true; } const target = Number(match[2].replace(",", ".")); if (!Number.isFinite(target)) { return true; } if (match[1] === ">") return numeric > target; if (match[1] === ">=") return numeric >= target; if (match[1] === "<") return numeric < target; if (match[1] === "<=") return numeric <= target; if (match[1] === "==") return numeric === target; if (match[1] === "!=") return numeric !== target; return true; } async function readDatapointCurrent(source, datapoint) { if (source.protocol === "modbus-tcp") { return await readModbusTcp(source, datapoint); } if (source.protocol === "opc-ua") { return await readOpcuaValue(source, datapoint); } if (source.protocol === "bacnet-ip") { return await readBacnetValue(source, datapoint); } if (source.protocol === "knx-ip") { return await readKnxValue(source, datapoint); } throw new Error("Protokoll wird nicht unterstützt."); } 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); const effectiveWriteMode = datapoint.writeMode || source.writeMode; const conditionMatches = matchesLogCondition(datapoint.condition, value); if (conditionMatches && (effectiveWriteMode !== "cov" || nextValues[key] !== serialized)) { readings.push({ pointIndex: datapoint.pointIndex, value }); } nextValues[key] = serialized; datapoint.currentValue = value; datapoint.lastReadAt = new Date().toISOString(); } const written = await writeTrendRow(source, readings); touchStore((store) => { store.datapoints = store.datapoints.map((item) => { const updated = datapoints.find((datapoint) => datapoint.id === item.id); return updated ? { ...item, currentValue: updated.currentValue, lastReadAt: updated.lastReadAt } : item; }); return store; }); 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 <= LOCAL_POINT_LIMIT; index += 1) { columns.push('`Wert' + index + '` TEXT NULL'); } return columns.join(", "); } async function ensureTrendColumns(table) { const [columns] = await db.query( ` SELECT COLUMN_NAME AS columnName FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND COLUMN_NAME LIKE 'Wert%' `, [table] ); const existing = new Set(columns.map((column) => column.columnName)); const missing = []; for (let index = 1; index <= LOCAL_POINT_LIMIT; index += 1) { const columnName = `Wert${index}`; if (!existing.has(columnName)) { missing.push(`ADD COLUMN \`${columnName}\` TEXT NULL`); } } for (let index = 0; index < missing.length; index += 50) { await db.query(`ALTER TABLE \`${table}\` ${missing.slice(index, index + 50).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 `); await ensureTrendColumns(table); return table; } async function ensureSourceRegistryTable() { await db.query(` CREATE TABLE IF NOT EXISTS \`source_registry\` ( \`source_id\` VARCHAR(191) NOT NULL, \`table_name\` VARCHAR(191) NOT NULL, \`source_name\` VARCHAR(255) NOT NULL, \`display_name\` VARCHAR(255) NOT NULL, \`protocol\` VARCHAR(32) NOT NULL, \`created_at\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, \`updated_at\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (\`source_id\`), KEY \`idx_source_registry_table_name\` (\`table_name\`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci `); } async function syncSourceRegistryEntry(source) { if (!source?.id || !source?.tableName) { return; } await ensureSourceRegistryTable(); await db.query( ` INSERT INTO \`source_registry\` ( \`source_id\`, \`table_name\`, \`source_name\`, \`display_name\`, \`protocol\` ) VALUES (?, ?, ?, ?, ?) ON DUPLICATE KEY UPDATE \`table_name\` = VALUES(\`table_name\`), \`source_name\` = VALUES(\`source_name\`), \`display_name\` = VALUES(\`display_name\`), \`protocol\` = VALUES(\`protocol\`) `, [ String(source.id), sanitizeTrendTable(source.tableName), String(source.name || source.displayName || source.id), String(source.displayName || source.name || source.id), String(source.protocol || "unknown"), ] ); } async function deleteSourceRegistryEntry(sourceId) { if (!sourceId) { return; } await ensureSourceRegistryTable(); await db.query("DELETE FROM `source_registry` WHERE `source_id` = ?", [String(sourceId)]); } async function ensureSourceDatapointRegistryTable() { await db.query(` CREATE TABLE IF NOT EXISTS \`source_datapoint_registry\` ( \`datapoint_id\` VARCHAR(255) NOT NULL, \`source_id\` VARCHAR(191) NOT NULL, \`table_name\` VARCHAR(191) NOT NULL, \`point_index\` INT NOT NULL, \`alias\` VARCHAR(255) NOT NULL, \`name\` VARCHAR(255) NOT NULL, \`protocol\` VARCHAR(32) NOT NULL, \`unit\` VARCHAR(64) NULL, \`kind\` VARCHAR(32) NULL, \`created_at\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, \`updated_at\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (\`datapoint_id\`), UNIQUE KEY \`uniq_source_point\` (\`source_id\`, \`point_index\`), KEY \`idx_source_datapoint_table\` (\`table_name\`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci `); } async function syncDatapointRegistryEntries(datapoints) { if (!Array.isArray(datapoints) || !datapoints.length) { return; } await ensureSourceDatapointRegistryTable(); for (const datapoint of datapoints) { if (!datapoint?.id || !datapoint?.sourceId || !datapoint?.tableName) { continue; } const pointIndex = normalizePointIndex(datapoint.pointIndex, 1); await db.query("DELETE FROM `source_datapoint_registry` WHERE `datapoint_id` = ? OR (`source_id` = ? AND `point_index` = ?)", [ String(datapoint.id), String(datapoint.sourceId), pointIndex, ]); if (datapoint.enabled !== true) { continue; } await db.query( ` INSERT INTO \`source_datapoint_registry\` ( \`datapoint_id\`, \`source_id\`, \`table_name\`, \`point_index\`, \`alias\`, \`name\`, \`protocol\`, \`unit\`, \`kind\` ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) `, [ String(datapoint.id), String(datapoint.sourceId), sanitizeTrendTable(datapoint.tableName), pointIndex, normalizeText(datapoint.alias || datapoint.name || `Wert${datapoint.pointIndex}`), normalizeText(datapoint.name || datapoint.alias || `Wert${datapoint.pointIndex}`), String(datapoint.protocol || "unknown"), normalizeUnit(datapoint.unit || ""), String(datapoint.kind || "analog"), ] ); } } async function deleteDatapointRegistryForSource(sourceId) { if (!sourceId) { return; } await ensureSourceDatapointRegistryTable(); await db.query("DELETE FROM `source_datapoint_registry` WHERE `source_id` = ?", [String(sourceId)]); } async function rebuildRegistryTablesFromStore() { const store = readStore(); await ensureSourceRegistryTable(); await ensureSourceDatapointRegistryTable(); await db.query("DELETE FROM `source_datapoint_registry`"); for (const source of store.sources || []) { await syncSourceRegistryEntry(source); } await syncDatapointRegistryEntries((store.datapoints || []).filter((datapoint) => datapoint.enabled === true)); } async function getTrendTableInfo(tableName) { const table = sanitizeTrendTable(tableName); const [columns] = await db.query( ` SELECT COLUMN_NAME AS columnName FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? ORDER BY ORDINAL_POSITION `, [table] ); const names = columns.map((column) => column.columnName); const lower = new Set(names.map((name) => String(name).toLowerCase())); const wideColumnCount = names.filter((name) => /^Wert\d+$/i.test(name)).length; if (wideColumnCount > 0) { return { table, mode: "wide", exists: true }; } if (lower.has("source_id") && lower.has("point_index") && lower.has("value") && lower.has("datum")) { return { table, mode: "narrow", exists: true }; } return { table, mode: null, exists: names.length > 0 }; } async function ensureSourceTrendTable(source) { const info = await getTrendTableInfo(source.tableName); if (info.mode === "wide") { return info; } if (info.mode === "narrow") { await syncSourceRegistryEntry(source); return info; } if (info.exists) { throw new Error(`Trendtabelle ${info.table} hat ein unbekanntes Format.`); } await db.query(` CREATE TABLE IF NOT EXISTS \`${info.table}\` ( \`id\` BIGINT NOT NULL AUTO_INCREMENT, \`datum\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, \`userlevel\` INT(1) NOT NULL DEFAULT 0, \`source_id\` VARCHAR(191) NOT NULL, \`point_index\` INT NOT NULL, \`alias\` VARCHAR(255) NOT NULL, \`value\` TEXT NULL, PRIMARY KEY (\`id\`), KEY \`idx_${info.table}_datum\` (\`datum\`), KEY \`idx_${info.table}_point_datum\` (\`point_index\`, \`datum\`), KEY \`idx_${info.table}_source_point\` (\`source_id\`, \`point_index\`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci `); await syncSourceRegistryEntry(source); return { ...info, mode: "narrow", exists: true }; } async function writeTrendRow(source, readings) { if (!readings.length) { return false; } const tableInfo = await ensureSourceTrendTable(source); const table = tableInfo.table; if (tableInfo.mode === "narrow") { const statement = ` INSERT INTO \`${table}\` ( \`userlevel\`, \`source_id\`, \`point_index\`, \`alias\`, \`value\` ) VALUES (?, ?, ?, ?, ?) `; for (const reading of readings) { const pointIndex = normalizePointIndex(reading.pointIndex, 0); if (!pointIndex) { continue; } const datapoint = source.__activeDatapoints?.find((item) => item.pointIndex === pointIndex); await db.query(statement, [ 0, String(source.id), pointIndex, String(datapoint?.alias || `Wert${pointIndex}`), String(reading.value), ]); } return true; } 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); const effectiveWriteMode = datapoint.writeMode || source.writeMode; const conditionMatches = matchesLogCondition(datapoint.condition, value); if (conditionMatches && (effectiveWriteMode !== "cov" || nextValues[key] !== serialized)) { readings.push({ pointIndex: datapoint.pointIndex, value }); } nextValues[key] = serialized; datapoint.currentValue = value; datapoint.lastReadAt = new Date().toISOString(); } const written = await writeTrendRow(source, readings); touchStore((store) => { store.datapoints = store.datapoints.map((item) => { const updated = datapoints.find((datapoint) => datapoint.id === item.id); return updated ? { ...item, currentValue: updated.currentValue, lastReadAt: updated.lastReadAt } : item; }); return store; }); 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 && datapoint.enabled === true); if (!sourcePoints.length) { updateSourceState(source.id, { lastPollAt: new Date().toISOString(), lastPollStatus: "idle: keine Datenpunkte", lastPollError: "", }); return; } try { source.__activeDatapoints = sourcePoints; 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\u00fctzt.`, }); } catch (error) { updateSourceState(source.id, { lastPollAt: new Date().toISOString(), lastPollStatus: "error", lastPollError: error.message || "Polling fehlgeschlagen.", }); } finally { delete source.__activeDatapoints; } } 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"}`); 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", "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; } 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: 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; } 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 }; }); const createdSource = next.sources.at(-1); await ensureSourceTrendTable(createdSource); replyJson(res, 201, createdSource); } 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; }); const updatedSource = next.sources.find((item) => item.id === sourceId); await ensureSourceTrendTable(updatedSource); replyJson(res, 200, updatedSource); } 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), })); await deleteSourceRegistryEntry(sourceId); await deleteDatapointRegistryForSource(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)); upsertDatapoints(store, [datapoint]); return store; }); const created = next.datapoints.find((item) => item.sourceId === body.sourceId && item.pointIndex === normalizePointIndex(body.pointIndex, 1)); await syncDatapointRegistryEntries(created ? [created] : []); 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 === "PATCH" && url.pathname.startsWith("/datapoints/")) { try { const datapointId = decodeURIComponent(url.pathname.split("/").at(-1) || ""); const body = await readBody(req); const next = touchStore((store) => { const index = store.datapoints.findIndex((item) => item.id === datapointId); if (index < 0) { throw new Error("not-found"); } store.datapoints[index] = { ...store.datapoints[index], enabled: typeof body.enabled === "boolean" ? body.enabled : store.datapoints[index].enabled, writeMode: typeof body.writeMode === "string" ? normalizeWriteMode(body.writeMode) : store.datapoints[index].writeMode, condition: typeof body.condition === "string" ? body.condition.trim() : store.datapoints[index].condition, pointIndex: body.pointIndex ? normalizePointIndex(body.pointIndex, store.datapoints[index].pointIndex) : store.datapoints[index].pointIndex, alias: typeof body.alias === "string" && body.alias.trim() ? body.alias.trim() : store.datapoints[index].alias, name: typeof body.name === "string" && body.name.trim() ? body.name.trim() : store.datapoints[index].name, address: typeof body.address === "string" ? normalizeText(body.address) : store.datapoints[index].address, nodeId: typeof body.nodeId === "string" ? normalizeText(body.nodeId) : store.datapoints[index].nodeId, dataType: typeof body.dataType === "string" && body.dataType.trim() ? normalizeText(body.dataType) : store.datapoints[index].dataType, unit: body.unit ? normalizePointUnit(body.unit, body.kind || store.datapoints[index].kind) : store.datapoints[index].unit, kind: body.kind ? normalizePointKind(body.kind) : store.datapoints[index].kind, scale: Number.isFinite(Number(body.scale)) ? Number(body.scale) : store.datapoints[index].scale, }; return store; }); const updatedDatapoint = next.datapoints.find((item) => item.id === datapointId); await syncDatapointRegistryEntries(updatedDatapoint ? [updatedDatapoint] : []); replyJson(res, 200, { datapoint: updatedDatapoint, updatedAt: next.updatedAt }); } catch (error) { replyJson(res, error.message === "not-found" ? 404 : 400, { error: error.message === "not-found" ? "Datenpunkt nicht gefunden." : error.message }); } return; } if (req.method === "POST" && url.pathname.match(/^\/datapoints\/[^/]+\/read$/)) { try { const datapointId = decodeURIComponent(url.pathname.split("/").at(-2) || ""); const store = readStore(); const datapoint = store.datapoints.find((item) => item.id === datapointId); const source = datapoint ? store.sources.find((item) => item.id === datapoint.sourceId) : null; if (!datapoint || !source) { replyJson(res, 404, { error: "Datenpunkt oder Quelle nicht gefunden." }); return; } const value = await readDatapointCurrent(source, datapoint); const next = touchStore((current) => { current.datapoints = current.datapoints.map((item) => item.id === datapointId ? { ...item, currentValue: value, lastReadAt: new Date().toISOString() } : item); return current; }); replyJson(res, 200, { value, datapoint: next.datapoints.find((item) => item.id === datapointId), readAt: new Date().toISOString() }); } catch (error) { replyJson(res, 400, { error: error.message || "Live-Wert konnte nicht gelesen 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")); } upsertDatapoints(store, imported); return store; }); await syncDatapointRegistryEntries(imported); 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")); } upsertDatapoints(store, imported); return store; }); await syncDatapointRegistryEntries(imported); 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); if (!body.host) { replyJson(res, 400, { error: "Für den OPC-UA-Scan fehlt Host oder IP." }); return; } const result = await scanOpcuaEndpoint(body); let scannedCandidates = []; touchStore((store) => { if (body.sourceId) { const source = store.sources.find((item) => item.id === body.sourceId); if (source) { const rawNodes = (result.nodes || []).filter((node) => node.nodeId); const candidates = buildScannedDatapoints(store, source, rawNodes, (node, index) => ({ ...node, id: `${source.id}::${node.nodeId}`, protocol: "opc-ua", pointIndex: index + 1, alias: node.displayName || node.browseName || node.nodeId, name: node.displayName || node.browseName || node.nodeId, address: node.nodeId, nodeId: node.nodeId, })); scannedCandidates = candidates; upsertDatapoints(store, candidates); } } store.scans.opcua.unshift(result); store.scans.opcua = store.scans.opcua.slice(0, 20); return store; }); await syncDatapointRegistryEntries(scannedCandidates); replyJson(res, 200, result); } catch (error) { replyJson(res, 400, { error: error.message || "OPC-UA-Scan fehlgeschlagen." }); } return; } if (req.method === "POST" && url.pathname === "/scan/bacnet") { try { const body = await readBody(req); const payload = await scanBacnetNetwork(body); let scannedCandidates = []; touchStore((store) => { if (body.sourceId) { const source = store.sources.find((item) => item.id === body.sourceId); if (source) { const rawPoints = payload.points || []; const candidates = buildScannedDatapoints(store, source, rawPoints, (point, index) => ({ ...point, id: point.id || `${source.id}::${point.host || source.host}::${point.dataType || point.objectType || "object"}::${point.objectInstance || index + 1}`, protocol: "bacnet-ip", pointIndex: index + 1, alias: point.alias || point.name || `${point.dataType || "obj"}${point.objectInstance || index + 1}`, name: point.name || point.alias || `${point.dataType || "obj"}${point.objectInstance || index + 1}`, })); scannedCandidates = candidates; upsertDatapoints(store, candidates); } } store.scans.bacnet.unshift(payload); store.scans.bacnet = store.scans.bacnet.slice(0, 20); return store; }); await syncDatapointRegistryEntries(scannedCandidates); replyJson(res, 200, payload); } catch (error) { replyJson(res, 400, { error: error.message || "BACnet-Scan fehlgeschlagen." }); } return; } replyJson(res, 404, { error: "Nicht gefunden." }); }); rebuildRegistryTablesFromStore().catch((error) => console.error("Registry sync failed", error)); server.listen(port, () => { console.log(`SE Local Trenddata collector listening on port ${port}`); });