This commit is contained in:
jhartworks
2026-06-25 14:35:57 +02:00
parent bf1214239f
commit a12ef41d1b
116 changed files with 978 additions and 59769 deletions
+365 -30
View File
@@ -11,6 +11,7 @@ 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;
@@ -143,7 +144,7 @@ function normalizeUnit(value) {
function normalizePointIndex(value, fallback) {
const numeric = Number(value);
return Number.isInteger(numeric) && numeric >= 1 && numeric <= 200 ? numeric : fallback;
return Number.isInteger(numeric) && numeric >= 1 && numeric <= LOCAL_POINT_LIMIT ? numeric : fallback;
}
function buildSourceTableName(body) {
@@ -222,6 +223,9 @@ function parseModbusCsv(payload) {
unit: normalizeUnit(get("unit")),
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 || "",
@@ -269,6 +273,9 @@ function parseKnxXml(payload) {
unit: normalizeUnit(""),
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(),
@@ -302,12 +309,17 @@ function createManualDatapoint(store, payload) {
alias,
unit: normalizeUnit(payload.unit),
kind: normalizePointKind(payload.kind),
enabled: payload.enabled !== false,
writeMode: normalizeWriteMode(payload.writeMode || source.writeMode),
condition: String(payload.condition || "").trim(),
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 || "",
currentValue: payload.currentValue ?? null,
lastReadAt: payload.lastReadAt || null,
createdAt: new Date().toISOString(),
};
@@ -318,6 +330,44 @@ function createManualDatapoint(store, payload) {
return datapoint;
}
function createScanDatapoint(source, payload, index) {
const pointIndex = normalizePointIndex(payload.pointIndex, index + 1);
const protocol = normalizeProtocol(payload.protocol || source.protocol);
const alias = String(payload.alias || payload.name || payload.displayName || `Punkt ${pointIndex}`).trim() || `Punkt ${pointIndex}`;
return {
id: payload.id || `dp-${protocol}-${Date.now()}-${Math.round(Math.random() * 100000)}-${pointIndex}`,
sourceId: source.id,
protocol,
pointIndex,
name: String(payload.name || alias).trim() || alias,
alias,
unit: normalizeUnit(payload.unit),
kind: normalizePointKind(payload.kind),
enabled: payload.enabled === true,
writeMode: normalizeWriteMode(payload.writeMode || source.writeMode),
condition: String(payload.condition || "").trim(),
dataType: String(payload.dataType || "").trim(),
address: String(payload.address || "").trim(),
nodeId: String(payload.nodeId || "").trim(),
objectType: payload.objectType,
objectInstance: payload.objectInstance,
propertyId: payload.propertyId || 85,
scale: Number(payload.scale || 1) || 1,
host: payload.host || source.host || "",
tableName: source.tableName || "",
currentValue: payload.currentValue ?? null,
lastReadAt: payload.currentValue === undefined ? null : new Date().toISOString(),
createdAt: new Date().toISOString(),
};
}
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) {
@@ -476,31 +526,83 @@ function buildOpcuaEndpoint(sourceOrBody) {
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: 10000,
requestedSessionTimeout: 15000,
});
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?.() || "",
}));
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" }],
sampleNodes,
tree,
nodes: state.nodes,
sampleNodes: state.nodes,
scannedAt: new Date().toISOString(),
};
} finally {
@@ -539,7 +641,10 @@ async function readOpcuaValue(source, datapoint) {
}
function parseBacnetObject(datapoint) {
const dataType = String(datapoint.dataType || "analogInput").trim();
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,
@@ -555,29 +660,104 @@ function parseBacnetObject(datapoint) {
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);
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 client = new bacnet({ apduTimeout: Number(body.timeout || 5000), port: Number(body.port || 47808) });
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(() => {
client.close();
resolve({ devices, scannedAt: new Date().toISOString() });
}, Number(body.timeout || 5000));
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) => {
devices.push({
host: device.address,
port: Number(body.port || 47808),
reachable: true,
deviceId: device.deviceId,
name: `BACnet Gerät ${device.deviceId}`,
});
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 {
@@ -670,6 +850,46 @@ async function readKnxValue(source, datapoint) {
});
}
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 || {}) };
@@ -678,13 +898,24 @@ async function pollGenericSource(source, datapoints, readValue) {
const value = await readValue(source, datapoint);
const key = String(datapoint.pointIndex);
const serialized = String(value);
if (source.writeMode !== "cov" || nextValues[key] !== serialized) {
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(),
@@ -694,12 +925,36 @@ async function pollGenericSource(source, datapoints, readValue) {
}
function buildValueColumns() {
const columns = [];
for (let index = 1; index <= 200; index += 1) {
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(`
@@ -711,6 +966,7 @@ async function ensureTrendTable(tableName) {
PRIMARY KEY (\`id\`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci
`);
await ensureTrendColumns(table);
return table;
}
@@ -774,13 +1030,24 @@ async function pollModbusSource(source, datapoints) {
const value = await readModbusTcp(source, datapoint);
const key = String(datapoint.pointIndex);
const serialized = String(value);
if (source.writeMode !== "cov" || nextValues[key] !== serialized) {
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(),
@@ -790,7 +1057,7 @@ async function pollModbusSource(source, datapoints) {
}
async function pollSource(source, datapoints) {
const sourcePoints = datapoints.filter((datapoint) => datapoint.sourceId === source.id);
const sourcePoints = datapoints.filter((datapoint) => datapoint.sourceId === source.id && datapoint.enabled !== false);
if (!sourcePoints.length) {
updateSourceState(source.id, {
lastPollAt: new Date().toISOString(),
@@ -970,7 +1237,7 @@ const server = createServer(async (req, res) => {
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);
upsertDatapoints(store, [datapoint]);
return store;
});
const created = next.datapoints.find((item) => item.sourceId === body.sourceId && item.pointIndex === normalizePointIndex(body.pointIndex, 1));
@@ -981,6 +1248,58 @@ const server = createServer(async (req, res) => {
return;
}
if (req.method === "PATCH" && url.pathname.startsWith("/datapoints/")) {
try {
const datapointId = 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: body.writeMode ? 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,
unit: body.unit ? normalizeUnit(body.unit) : 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;
});
replyJson(res, 200, { datapoint: next.datapoints.find((item) => item.id === datapointId), 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 = 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);
@@ -989,7 +1308,7 @@ const server = createServer(async (req, res) => {
if (body.sourceId) {
store.datapoints = store.datapoints.filter((item) => !(item.sourceId === body.sourceId && item.protocol === "modbus-tcp"));
}
store.datapoints.push(...imported);
upsertDatapoints(store, imported);
return store;
});
replyJson(res, 200, { imported: imported.length, datapoints: imported, updatedAt: next.updatedAt });
@@ -1007,7 +1326,7 @@ const server = createServer(async (req, res) => {
if (body.sourceId) {
store.datapoints = store.datapoints.filter((item) => !(item.sourceId === body.sourceId && item.protocol === "knx-ip"));
}
store.datapoints.push(...imported);
upsertDatapoints(store, imported);
return store;
});
replyJson(res, 200, { imported: imported.length, datapoints: imported, updatedAt: next.updatedAt });
@@ -1026,6 +1345,14 @@ const server = createServer(async (req, res) => {
}
const result = await scanOpcuaEndpoint(body);
touchStore((store) => {
if (body.sourceId) {
const source = store.sources.find((item) => item.id === body.sourceId);
if (source) {
const existingCount = store.datapoints.filter((item) => item.sourceId === source.id).length;
const candidates = (result.nodes || []).filter((node) => node.nodeId && node.hasValue).map((node, index) => createScanDatapoint(source, { ...node, protocol: "opc-ua", pointIndex: existingCount + index + 1, alias: node.path || node.displayName, name: node.displayName, address: node.nodeId, nodeId: node.nodeId }, index));
upsertDatapoints(store, candidates);
}
}
store.scans.opcua.unshift(result);
store.scans.opcua = store.scans.opcua.slice(0, 20);
return store;
@@ -1042,6 +1369,14 @@ const server = createServer(async (req, res) => {
const body = await readBody(req);
const payload = await scanBacnetNetwork(body);
touchStore((store) => {
if (body.sourceId) {
const source = store.sources.find((item) => item.id === body.sourceId);
if (source) {
const existingCount = store.datapoints.filter((item) => item.sourceId === source.id).length;
const candidates = (payload.points || []).map((point, index) => createScanDatapoint(source, { ...point, protocol: "bacnet-ip", pointIndex: existingCount + index + 1 }, index));
upsertDatapoints(store, candidates);
}
}
store.scans.bacnet.unshift(payload);
store.scans.bacnet = store.scans.bacnet.slice(0, 20);
return store;