beta. Modbus fixed and tested

This commit is contained in:
jhartworks
2026-06-30 09:42:55 +02:00
parent 2ce3be068e
commit e34626f15b
8 changed files with 1355 additions and 246 deletions
+162 -33
View File
@@ -242,8 +242,8 @@ function parseModbusCsv(payload) {
protocol: "modbus-tcp",
address,
pointIndex,
name: get("name", `Register ${index + 1}`),
alias: get("alias", get("name", `Register ${index + 1}`)),
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: normalizeUnit(get("unit")),
dataType: get("datatype", "holding-register"),
kind: normalizePointKind(get("kind", get("type", "analog"))),
@@ -282,8 +282,8 @@ function parseKnxXml(payload) {
};
const address = getAttr("Address") || getAttr("address");
const name = getAttr("Name") || getAttr("name") || `KNX ${index + 1}`;
const dpt = getAttr("DatapointType") || getAttr("DPTs") || getAttr("DPT") || "";
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({
@@ -462,29 +462,71 @@ function tryTcpConnect(host, portNumber, timeout = 1200) {
}
function inferModbusAddress(datapoint) {
const rawAddress = Number(String(datapoint.address || datapoint.register || "").replace(/[^0-9]/g, ""));
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 (raw >= 40001) {
functionCode = 3;
if (functionCode === 3 && raw >= 40001 && raw <= 49999) {
address = raw - 40001;
} else if (raw >= 30001) {
functionCode = 4;
} else if (functionCode === 4 && raw >= 30001 && raw <= 39999) {
address = raw - 30001;
} else if (raw >= 10001) {
functionCode = 2;
} else if (functionCode === 2 && raw >= 10001 && raw <= 19999) {
address = raw - 10001;
} else if (raw <= 9999 && functionCode > 2) {
} 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();
@@ -495,6 +537,7 @@ function readModbusTcp(source, datapoint) {
}
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;
@@ -506,7 +549,7 @@ function readModbusTcp(source, datapoint) {
request.writeUInt8(unitId, 6);
request.writeUInt8(functionCode, 7);
request.writeUInt16BE(address, 8);
request.writeUInt16BE(1, 10);
request.writeUInt16BE(quantity, 10);
const socket = new net.Socket();
const chunks = [];
@@ -528,6 +571,11 @@ function readModbusTcp(source, datapoint) {
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);
@@ -551,29 +599,23 @@ function readModbusTcp(source, datapoint) {
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) {
const byteCount = Number(response.readUInt8(8) || 0);
const startOffset = 9;
if (response.length < startOffset + byteCount) {
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);
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);
@@ -870,11 +912,87 @@ async function readBacnetValue(source, datapoint) {
});
}
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ür ${datapoint.alias || datapoint.name}.`);
throw new Error(`KNX Gruppenadresse fehlt f\u00fcr ${datapoint.alias || datapoint.name}.`);
}
return await new Promise((resolve, reject) => {
@@ -890,8 +1008,13 @@ async function readKnxValue(source, datapoint) {
if (error) {
reject(error);
} else {
const numeric = Number(value);
resolve(typeof value === "boolean" ? (value ? 1 : 0) : Number.isFinite(numeric) ? numeric : String(value || ""));
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);
}
}
};
@@ -917,7 +1040,6 @@ async function readKnxValue(source, datapoint) {
});
});
}
function matchesLogCondition(condition, value) {
const text = String(condition || "").trim();
if (!text) {
@@ -1399,7 +1521,7 @@ async function pollSource(source, datapoints) {
updateSourceState(source.id, {
lastPollAt: new Date().toISOString(),
lastPollStatus: "unsupported",
lastPollError: `${source.protocol} wird vom Collector nicht unterst?tzt.`,
lastPollError: `${source.protocol} wird vom Collector nicht unterst\u00fctzt.`,
});
} catch (error) {
updateSourceState(source.id, {
@@ -1577,10 +1699,14 @@ const server = createServer(async (req, res) => {
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,
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 ? 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,
@@ -1740,3 +1866,6 @@ server.listen(port, () => {
});