feat(bacnet-collector): add initial implementation of BACnet collector with scanning and reading capabilities

- Created package.json for the bacnet-collector project with necessary dependencies.
- Implemented main server logic in src/index.js to handle BACnet device scanning and reading.
- Added raw discovery functionality in src/raw-discover.js to facilitate device detection over the network.
- Introduced utility functions for handling IP addresses, broadcasting, and BACnet packet construction.
- Implemented JSON response handling for health checks and capabilities endpoints.
This commit is contained in:
jhartworks
2026-07-09 11:13:26 +02:00
parent 7a1fe401f2
commit f247d2827d
10 changed files with 1833 additions and 77233 deletions
+12
View File
@@ -0,0 +1,12 @@
FROM node:20-alpine
WORKDIR /app
COPY package*.json ./
RUN npm install
COPY src ./src
EXPOSE 18111
CMD ["npm", "start"]
+12
View File
@@ -0,0 +1,12 @@
{
"name": "seltd-bacnet-collector",
"version": "1.0.0",
"private": true,
"type": "module",
"scripts": {
"start": "node src/index.js"
},
"dependencies": {
"bacstack": "^0.0.1-beta.14"
}
}
+866
View File
@@ -0,0 +1,866 @@
import { createServer } from "node:http";
import { Buffer } from "node:buffer";
import { spawn } from "node:child_process";
import { fileURLToPath } from "node:url";
import { networkInterfaces } from "node:os";
import dgram from "node:dgram";
import Bacstack from "bacstack";
const port = Number(process.env.PORT || 18111);
const defaultTimeout = Number(process.env.BACNET_SCAN_TIMEOUT || 7000);
const pointLimit = Number(process.env.BACNET_POINT_LIMIT || 2000);
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,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 sleep(ms) {
return new Promise((resolve) => setTimeout(resolve, ms));
}
function normalizeText(value) {
return String(value ?? "").trim();
}
function normalizeBacnetMode(value) {
const mode = normalizeText(value || "broadcast").toLowerCase();
return ["broadcast", "bbmd", "foreign-device"].includes(mode) ? mode : "broadcast";
}
function normalizeBacnetTtl(value) {
const numeric = Number(value);
return Number.isFinite(numeric) && numeric > 0 ? Math.round(numeric) : 120;
}
function getRouting(body) {
return {
mode: normalizeBacnetMode(body?.bacnetMode || body?.options?.bacnetMode),
host: normalizeText(body?.host || body?.options?.bbmdHost || ""),
port: Number(body?.port || body?.options?.bbmdPort || 47808) || 47808,
broadcastAddress: normalizeText(body?.bacnetBroadcastAddress || body?.options?.bacnetBroadcastAddress || "255.255.255.255") || "255.255.255.255",
foreignDeviceTtl: normalizeBacnetTtl(body?.bacnetForeignDeviceTtl || body?.options?.bacnetForeignDeviceTtl),
};
}
function isSubnetMaskLike(address) {
return /^255\.255\.255\.(0|128|192|224|240|248|252|254)$/.test(String(address || ""));
}
function isAutoBroadcast(value) {
return !normalizeText(value) || ["auto", "automatisch"].includes(normalizeText(value).toLowerCase()) || isSubnetMaskLike(value);
}
function parseIpv4(value) {
const address = normalizeText(value).replace(/^::ffff:/, "");
const parts = address.split(".").map((part) => Number(part));
if (parts.length !== 4 || parts.some((part) => !Number.isInteger(part) || part < 0 || part > 255)) {
return null;
}
return parts;
}
function isUsableLanIpv4(value) {
const parts = parseIpv4(value);
if (!parts) return false;
if (parts[0] === 127 || parts[0] === 0 || parts[0] >= 224) return false;
if (parts[0] === 172 && parts[1] >= 17 && parts[1] <= 31) return false; // Docker bridge defaults.
return true;
}
function ipv4ToNumber(parts) {
return (((parts[0] << 24) >>> 0) + (parts[1] << 16) + (parts[2] << 8) + parts[3]) >>> 0;
}
function numberToIpv4(value) {
return [value >>> 24, (value >>> 16) & 255, (value >>> 8) & 255, value & 255].join(".");
}
function broadcastFromAddress(address, netmask = "255.255.255.0") {
const ipParts = parseIpv4(address);
const maskParts = parseIpv4(netmask) || [255, 255, 255, 0];
if (!ipParts || !isUsableLanIpv4(address)) return "";
const ip = ipv4ToNumber(ipParts);
const mask = ipv4ToNumber(maskParts);
return numberToIpv4((ip | (~mask >>> 0)) >>> 0);
}
function getInterfaceBroadcasts() {
return Object.values(networkInterfaces())
.flat()
.filter((item) => item && item.family === "IPv4" && !item.internal && isUsableLanIpv4(item.address))
.map((item) => broadcastFromAddress(item.address, item.netmask))
.filter(Boolean);
}
function getBroadcastTargets(routing, body = {}) {
const configured = normalizeText(routing.broadcastAddress || "auto");
const targets = [];
if (!isAutoBroadcast(configured)) {
targets.push(configured);
}
for (const subnet of parseHostList(body.scanSubnets || body.bacnetScanSubnets || body.options?.bacnetScanSubnets)) {
const broadcast = cidrBroadcast(subnet);
if (broadcast) targets.push(broadcast);
}
for (const candidate of [body.clientHost, body.browserHost, body.locationHost]) {
const broadcast = broadcastFromAddress(candidate);
if (broadcast) targets.push(broadcast);
}
targets.push(...getInterfaceBroadcasts());
targets.push(...parseHostList(process.env.BACNET_BROADCAST_FALLBACKS));
if (!targets.length) {
targets.push("255.255.255.255");
}
return Array.from(new Set(targets));
}
function expandCidrHosts(cidr) {
const match = String(cidr || "").trim().match(/^(\d+\.\d+\.\d+\.\d+)\/(\d{1,2})$/);
if (!match) return [];
const ipParts = parseIpv4(match[1]);
const prefix = Number(match[2]);
if (!ipParts || !Number.isInteger(prefix) || prefix < 16 || prefix > 32) return [];
if (prefix === 32) return [numberToIpv4(ipv4ToNumber(ipParts))];
const ip = ipv4ToNumber(ipParts);
const mask = prefix === 0 ? 0 : (0xffffffff << (32 - prefix)) >>> 0;
const network = (ip & mask) >>> 0;
const broadcast = (network | (~mask >>> 0)) >>> 0;
const limit = Math.min(broadcast - network - 1, 4094);
return Array.from({ length: Math.max(0, limit) }, (_, index) => numberToIpv4(network + index + 1));
}
function cidrBroadcast(cidr) {
const match = String(cidr || "").trim().match(/^(\d+\.\d+\.\d+\.\d+)\/(\d{1,2})$/);
if (!match) return "";
const ipParts = parseIpv4(match[1]);
const prefix = Number(match[2]);
if (!ipParts || !Number.isInteger(prefix) || prefix < 16 || prefix > 30) return "";
const ip = ipv4ToNumber(ipParts);
const mask = (0xffffffff << (32 - prefix)) >>> 0;
return numberToIpv4((ip | (~mask >>> 0)) >>> 0);
}
function expandDirectedBroadcast(address) {
const parts = String(address || "").split(".").map((part) => Number(part));
if (parts.length !== 4 || parts.some((part) => !Number.isInteger(part) || part < 0 || part > 255)) {
return [];
}
if (parts[3] !== 255) {
return [];
}
const prefix = parts.slice(0, 3).join(".");
return Array.from({ length: 254 }, (_, index) => prefix + "." + (index + 1));
}
function parseHostList(value) {
if (Array.isArray(value)) {
return value.flatMap((item) => parseHostList(item));
}
return String(value || "")
.split(/[;,\s]+/)
.map((item) => item.trim())
.filter(Boolean);
}
function buildUnicastScanHosts(body, routing) {
const hosts = [
...parseHostList(body.scanHosts),
...parseHostList(body.scanSubnets || body.bacnetScanSubnets || body.options?.bacnetScanSubnets).flatMap((subnet) => subnet.includes("/") ? expandCidrHosts(subnet) : (parseIpv4(subnet) ? [subnet] : [])),
...parseHostList(process.env.BACNET_SCAN_HOSTS),
];
if (routing.host && routing.mode !== "broadcast") {
hosts.push(routing.host);
}
for (const broadcastAddress of getBroadcastTargets(routing, body)) {
if (broadcastAddress !== "255.255.255.255") {
hosts.push(...expandDirectedBroadcast(broadcastAddress));
}
}
return Array.from(new Set(hosts)).slice(0, 512);
}
function sendUnicastWhoIs(client, host, deviceId) {
const packet = captureWhoIsPacket(client, deviceId);
packet[1] = 0x0a;
getTransport(client).send(packet, packet.length, String(host || "").replace(/:\d+$/, ""));
}
function createClient(body, timeout) {
const routing = getRouting(body);
const client = new Bacstack({
apduTimeout: timeout,
port: routing.port,
broadcastAddress: routing.broadcastAddress,
});
return { client, routing };
}
function closeClient(client) {
try {
client.close?.();
} catch {}
try {
client._transport?.close?.();
} catch {}
}
function getTransport(client) {
const transport = client?._transport;
if (!transport || typeof transport.send !== "function" || !transport._server) {
throw new Error("BACnet-Transport konnte nicht initialisiert werden.");
}
return transport;
}
function waitForSocket(client, timeout = 1200) {
const transport = getTransport(client);
try {
const address = transport._server.address?.();
if (address && typeof address.port === "number") {
return Promise.resolve();
}
} catch {}
return new Promise((resolve, reject) => {
let settled = false;
const finish = (error) => {
if (settled) return;
settled = true;
transport._server.off("listening", onListening);
transport._server.off("error", onError);
clearTimeout(timer);
if (error) {
reject(error);
} else {
resolve();
}
};
const onListening = () => finish();
const onError = (error) => finish(error);
const timer = setTimeout(() => finish(), timeout);
transport._server.once("listening", onListening);
transport._server.once("error", onError);
});
}
function captureWhoIsPacket(client, deviceId) {
const transport = getTransport(client);
const originalSend = transport.send.bind(transport);
let packet = null;
transport.send = (buffer, offset) => {
packet = Buffer.from(buffer.slice(0, offset));
};
try {
const numericDeviceId = Number(deviceId);
if (Number.isFinite(numericDeviceId)) {
client.whoIs({ lowLimit: numericDeviceId, highLimit: numericDeviceId });
} else {
client.whoIs();
}
} finally {
transport.send = originalSend;
}
if (!packet || packet.length < 4 || packet[0] !== 0x81) {
throw new Error("BACnet Who-Is konnte nicht aufgebaut werden.");
}
return packet;
}
function bvlcResultText(resultCode) {
const known = {
0x0000: "Erfolgreich",
0x0030: "Register-Foreign-Device NAK",
0x0060: "Distribute-Broadcast-To-Network NAK",
};
return known[resultCode] || `BVLC-Fehler ${resultCode}`;
}
async function sendControlPacket(client, receiver, packet, matcher, timeout = 3000) {
const transport = getTransport(client);
await waitForSocket(client);
return await new Promise((resolve, reject) => {
let settled = false;
const socket = transport._server;
const cleanup = () => {
socket.off("message", onMessage);
clearTimeout(timer);
};
const finish = (error, result) => {
if (settled) return;
settled = true;
cleanup();
if (error) {
reject(error instanceof Error ? error : new Error(String(error)));
} else {
resolve(result);
}
};
const onMessage = (message, rinfo) => {
try {
if (matcher(message, rinfo)) {
finish(null, { message, rinfo });
}
} catch (error) {
finish(error);
}
};
const timer = setTimeout(() => finish(new Error("BACnet-Steuertelegramm hat nicht rechtzeitig geantwortet.")), timeout);
socket.on("message", onMessage);
try {
transport.send(packet, packet.length, receiver);
} catch (error) {
finish(error);
}
});
}
async function registerForeignDevice(client, routing) {
if (routing.mode !== "foreign-device") {
return;
}
if (!routing.host) {
throw new Error("Für BACnet Foreign Device bitte BBMD Host/IP angeben.");
}
const packet = Buffer.alloc(6);
packet[0] = 0x81;
packet[1] = 0x05;
packet.writeUInt16BE(6, 2);
packet.writeUInt16BE(routing.foreignDeviceTtl, 4);
const { message } = await sendControlPacket(
client,
routing.host,
packet,
(candidate, rinfo) => candidate?.length >= 6 && candidate[0] === 0x81 && candidate[1] === 0x00 && Number(rinfo?.port) === routing.port,
3000,
);
const resultCode = message.readUInt16BE(4);
if (resultCode !== 0x0000) {
throw new Error(`Foreign-Device-Register wurde abgelehnt: ${bvlcResultText(resultCode)}.`);
}
}
function sendWhoIs(client, routing, deviceId) {
if (routing.mode === "bbmd" || routing.mode === "foreign-device") {
if (!routing.host) {
throw new Error("Für BACnet über BBMD bitte Host/IP angeben.");
}
const packet = captureWhoIsPacket(client, deviceId);
packet[1] = 0x09;
getTransport(client).send(packet, packet.length, routing.host);
return;
}
const numericDeviceId = Number(deviceId);
if (Number.isFinite(numericDeviceId)) {
client.whoIs({ lowLimit: numericDeviceId, highLimit: numericDeviceId });
} else {
client.whoIs();
}
}
function readProperty(client, host, objectId, propertyId) {
return new Promise((resolve) => {
client.readProperty(host, objectId, propertyId, (error, value) => {
resolve(error ? null : value);
});
});
}
function objectLabel(objectId) {
const typeMap = {
0: "analogInput",
1: "analogOutput",
2: "analogValue",
3: "binaryInput",
4: "binaryOutput",
5: "binaryValue",
8: "device",
13: "multiStateInput",
14: "multiStateOutput",
19: "multiStateValue",
};
return `${typeMap[objectId.type] || `object${objectId.type}`}${objectId.instance}`;
}
function extractValue(propertyValue) {
return propertyValue?.values?.[0]?.value ?? propertyValue?.values?.[0] ?? null;
}
async function scanObjects(client, device, portNumber) {
const host = device.address;
const deviceObject = { type: 8, instance: Number(device.deviceId) };
const objectList = await readProperty(client, host, deviceObject, 76);
const values = objectList?.values || [];
const points = [];
for (const entry of values.slice(0, pointLimit)) {
const objectId = entry?.value || entry;
if (!objectId || objectId.type === 8 || !Number.isFinite(Number(objectId.instance))) {
continue;
}
const objectName = await readProperty(client, host, objectId, 77);
const description = await readProperty(client, host, objectId, 28);
const presentValue = await readProperty(client, host, objectId, 85);
const raw = extractValue(presentValue);
const name = normalizeText(extractValue(objectName)) || objectLabel(objectId);
const alias = normalizeText(extractValue(description)) || name;
points.push({
id: `${host}-${objectId.type}-${objectId.instance}`,
host,
port: portNumber,
deviceId: device.deviceId,
objectType: objectId.type,
objectInstance: objectId.instance,
propertyId: 85,
address: String(objectId.instance),
dataType: objectLabel(objectId).replace(/[0-9]+$/, ""),
name,
description: alias,
alias,
kind: typeof raw === "boolean" || (objectId.type >= 3 && objectId.type <= 5) ? "digital" : "analog",
currentValue: raw,
hasValue: raw !== null,
});
}
return points;
}
function encodeBacnetContextUnsigned(tagNumber, value) {
const numeric = Number(value);
if (!Number.isFinite(numeric) || numeric < 0) {
return Buffer.alloc(0);
}
if (numeric <= 0xff) {
return Buffer.from([(tagNumber << 4) | 1, numeric]);
}
if (numeric <= 0xffff) {
const buffer = Buffer.alloc(3);
buffer[0] = (tagNumber << 4) | 2;
buffer.writeUInt16BE(numeric, 1);
return buffer;
}
const buffer = Buffer.alloc(5);
buffer[0] = (tagNumber << 4) | 4;
buffer.writeUInt32BE(numeric, 1);
return buffer;
}
function buildRawWhoIsPacket(deviceId, bvlcFunction) {
const numericDeviceId = Number(deviceId);
const payload = [Buffer.from([0x01, 0x00, 0x10, 0x08])];
if (Number.isFinite(numericDeviceId)) {
payload.push(encodeBacnetContextUnsigned(0, numericDeviceId));
payload.push(encodeBacnetContextUnsigned(1, numericDeviceId));
}
const body = Buffer.concat(payload);
const packet = Buffer.alloc(4 + body.length);
packet[0] = 0x81;
packet[1] = bvlcFunction;
packet.writeUInt16BE(packet.length, 2);
body.copy(packet, 4);
return packet;
}
function parseRawIAm(message) {
if (!message || message.length < 12 || message[0] !== 0x81) {
return null;
}
const start = message.indexOf(0xc4);
if (start < 0 || start + 4 >= message.length) {
return null;
}
const objectId = message.readUInt32BE(start + 1);
const objectType = objectId >>> 22;
const deviceId = objectId & 0x3fffff;
if (objectType !== 8) {
return null;
}
return { deviceId };
}
function runRawDiscoveryHelper(body, routing) {
const timeout = Number(body.timeout || defaultTimeout) + 3000;
const helperPath = fileURLToPath(new URL("./raw-discover.js", import.meta.url));
return new Promise((resolve, reject) => {
const child = spawn(process.execPath, [helperPath], {
env: {
...process.env,
RAW_BACNET_BODY: JSON.stringify({ body, routing }),
},
stdio: ["ignore", "pipe", "pipe"],
});
let stdout = "";
let stderr = "";
const timer = setTimeout(() => {
child.kill("SIGKILL");
reject(new Error("BACnet Raw-Discovery-Helper Timeout"));
}, Math.max(2500, timeout));
child.stdout.on("data", (chunk) => {
stdout += chunk.toString("utf8");
});
child.stderr.on("data", (chunk) => {
stderr += chunk.toString("utf8");
});
child.on("error", (error) => {
clearTimeout(timer);
reject(error);
});
child.on("close", (code) => {
clearTimeout(timer);
if (code !== 0) {
reject(new Error(stderr || `BACnet Raw-Discovery-Helper beendet mit Code ${code}.`));
return;
}
try {
resolve(JSON.parse(stdout || "{}"));
} catch (error) {
reject(new Error(`BACnet Raw-Discovery-Helper lieferte kein gültiges JSON: ${error.message}`));
}
});
});
}
async function rawBacnetDiscover(body, routing) {
if (routing.mode !== "broadcast") {
return [];
}
const hosts = buildUnicastScanHosts(body, routing);
const helperResult = await runRawDiscoveryHelper(body, routing).catch((error) => ({
helperError: error?.message || String(error),
devices: [],
rawTrace: [],
rawDebug: { hosts, sends: [], packets: [], helperError: error?.message || String(error) },
}));
const helperDevices = Array.isArray(helperResult.devices) ? helperResult.devices : [];
if (helperDevices.length) {
helperDevices.rawTrace = helperResult.rawTrace || [];
helperDevices.rawDebug = { ...(helperResult.rawDebug || {}), helper: "process" };
return helperDevices;
}
const devices = [];
const rawTrace = helperResult.rawTrace || [];
const rawDebug = { hosts, sends: [], packets: [], helper: "fallback", helperError: helperResult.helperError || "" };
const socket = dgram.createSocket({ type: "udp4", reuseAddr: true });
const addDevice = (deviceId, host, rinfo) => {
rawTrace.push({ deviceId: Number(deviceId), host, remote: rinfo ? rinfo.address + ":" + rinfo.port : "" });
if (!Number.isFinite(Number(deviceId)) || !host) {
return;
}
if (!devices.some((item) => item.deviceId === Number(deviceId) && item.address === host)) {
devices.push({
address: host,
host,
port: routing.port,
reachable: true,
deviceId: Number(deviceId),
name: "BACnet Gerät " + Number(deviceId),
vendorId: null,
maxApdu: null,
segmentation: null,
pointCount: 0,
points: [],
rawRemote: rinfo ? rinfo.address + ":" + rinfo.port : "",
});
}
};
await new Promise((resolve, reject) => {
socket.once("error", reject);
socket.bind(routing.port, "0.0.0.0", () => {
socket.setBroadcast(true);
resolve();
});
});
socket.on("message", (message, rinfo) => {
rawDebug.packets.push({ remote: rinfo.address + ":" + rinfo.port, length: message.length, hex: message.toString("hex") });
const parsed = parseRawIAm(message);
if (parsed) {
addDevice(parsed.deviceId, rinfo.address, rinfo);
}
});
try {
const broadcastPacket = buildRawWhoIsPacket(body.deviceId, 0x0b);
for (const broadcastAddress of getBroadcastTargets(routing, body)) {
socket.send(broadcastPacket, routing.port, broadcastAddress, (error) => rawDebug.sends.push({ host: broadcastAddress, error: error?.message || "" }));
await sleep(350);
}
const perHostDelay = Number(body.unicastDelayMs || process.env.BACNET_UNICAST_DELAY_MS || 20);
const packet = buildRawWhoIsPacket(body.deviceId, 0x0a);
for (const host of hosts) {
const before = devices.length;
socket.removeAllListeners("message");
socket.on("message", (message, rinfo) => {
rawDebug.packets.push({ remote: rinfo.address + ":" + rinfo.port, length: message.length, hex: message.toString("hex") });
const parsed = parseRawIAm(message);
if (parsed) {
addDevice(parsed.deviceId, host, rinfo);
}
});
socket.send(packet, routing.port, host, (error) => rawDebug.sends.push({ host, error: error?.message || "" }));
await sleep(Math.max(15, perHostDelay));
if (body.deviceId && devices.length > before && devices.some((device) => String(device.deviceId) === String(body.deviceId))) {
break;
}
}
await sleep(150);
} finally {
socket.close();
}
devices.rawTrace = rawTrace;
devices.rawDebug = rawDebug;
return devices;
}
function isLikelyDockerGatewayAddress(address) {
return /^172\.(1[7-9]|2\d|3[01])\.0\.1$/.test(String(address || ""));
}
function preferRealBacnetAddresses(devices) {
const list = Array.isArray(devices) ? devices : [];
const idsWithRealAddress = new Set(
list
.filter((device) => !isLikelyDockerGatewayAddress(device.address || device.host))
.map((device) => String(device.deviceId)),
);
return list.filter((device) => {
const address = device.address || device.host;
return !(idsWithRealAddress.has(String(device.deviceId)) && isLikelyDockerGatewayAddress(address));
});
}
async function scanNetwork(body) {
const timeout = Number(body.timeout || defaultTimeout);
const includeObjects = body.includeObjects === true || String(body.includeObjects || "").toLowerCase() === "true";
const objectReadTimeout = Number(body.objectReadTimeout || process.env.BACNET_OBJECT_READ_TIMEOUT || 1500);
const routing = getRouting(body);
const rawDevices = await rawBacnetDiscover(body, routing);
const { client } = createClient(body, includeObjects ? Math.min(timeout, objectReadTimeout) : timeout);
const devices = [...rawDevices];
const rawTrace = rawDevices.rawTrace || [];
const rawDebug = rawDevices.rawDebug || null;
let currentUnicastTarget = "";
const addDevice = (device, forcedAddress = "") => {
const address = normalizeText(forcedAddress || device?.address || "");
const deviceId = Number(device?.deviceId);
if (!address || !Number.isFinite(deviceId)) {
return;
}
const existing = devices.find((item) => item.deviceId === deviceId);
if (existing && forcedAddress) {
existing.address = address;
existing.host = address;
return;
}
if (!devices.some((item) => item.deviceId === deviceId && item.address === address)) {
devices.push({
address,
host: address,
port: routing.port,
reachable: true,
deviceId,
name: "BACnet Gerät " + deviceId,
vendorId: device?.vendorId ?? null,
maxApdu: device?.maxApdu ?? null,
segmentation: device?.segmentation ?? null,
pointCount: 0,
points: [],
});
}
};
client.on("iAm", (device) => addDevice(device, currentUnicastTarget));
try {
if ((routing.mode === "bbmd" || routing.mode === "foreign-device") && !routing.host) {
throw new Error("Für BACnet über BBMD bitte Host/IP angeben.");
}
await waitForSocket(client);
await registerForeignDevice(client, routing);
sendWhoIs(client, routing, body.deviceId);
await sleep(Math.min(timeout, 1800));
const hosts = routing.mode === "broadcast" ? buildUnicastScanHosts(body, routing) : [];
const needsUnicastMapping = routing.mode === "broadcast" && hosts.length && (!devices.length || devices.some((device) => !hosts.includes(device.address)));
if (needsUnicastMapping) {
const perHostDelay = Number(body.unicastDelayMs || process.env.BACNET_UNICAST_DELAY_MS || 20);
for (const host of hosts) {
currentUnicastTarget = host;
sendUnicastWhoIs(client, host, body.deviceId);
await sleep(Math.max(15, perHostDelay));
if (body.deviceId && devices.some((device) => String(device.deviceId) === String(body.deviceId))) {
break;
}
}
currentUnicastTarget = "";
await sleep(250);
}
let visibleDevices = preferRealBacnetAddresses(devices);
if (includeObjects) {
devices.splice(0, devices.length, ...visibleDevices);
for (const device of devices) {
const deviceObject = { type: 8, instance: Number(device.deviceId) };
const deviceName = await readProperty(client, device.address, deviceObject, 77);
const resolvedName = normalizeText(extractValue(deviceName));
if (resolvedName) {
device.name = resolvedName;
}
device.points = await scanObjects(client, device, routing.port);
device.pointCount = device.points.length;
}
}
visibleDevices = preferRealBacnetAddresses(devices);
return {
driver: "bacstack",
routing,
devices: visibleDevices,
rawTrace: body.debug === true ? rawTrace : undefined,
rawDebug: body.debug === true ? rawDebug : undefined,
points: visibleDevices.flatMap((device) => device.points || []),
scannedAt: new Date().toISOString(),
};
} finally {
closeClient(client);
}
}
function parseObjectId(datapoint) {
const typeMap = {
analogInput: 0,
analogOutput: 1,
analogValue: 2,
binaryInput: 3,
binaryOutput: 4,
binaryValue: 5,
multiStateInput: 13,
multiStateOutput: 14,
multiStateValue: 19,
};
if (Number.isFinite(Number(datapoint.objectType))) {
return { type: Number(datapoint.objectType), instance: Number(datapoint.objectInstance || datapoint.address || 0) };
}
const normalizedType = String(datapoint.dataType || "").replace(/[s_-]+/g, "");
const matchedType = Object.entries(typeMap).find(([name]) => name.toLowerCase() === normalizedType.toLowerCase())?.[1];
const shorthand = String(datapoint.address || datapoint.nodeId || "").match(/^(ai|ao|av|bi|bo|bv|msi|mso|msv)(\\d+)$/i);
const shorthandTypes = { ai: 0, ao: 1, av: 2, bi: 3, bo: 4, bv: 5, msi: 13, mso: 14, msv: 19 };
return {
type: matchedType ?? (shorthand ? shorthandTypes[shorthand[1].toLowerCase()] : 0),
instance: Number(datapoint.objectInstance || (shorthand ? shorthand[2] : datapoint.address) || datapoint.pointIndex || 0),
};
}
async function readCurrentValue(source, datapoint) {
const { client, routing } = createClient(source, Number(source.timeout || 5000));
try {
await waitForSocket(client);
await registerForeignDevice(client, routing);
const host = normalizeText(datapoint.host || source.host || "");
if (!host) {
throw new Error("BACnet Host fehlt.");
}
const propertyId = Number(datapoint.propertyId || 85);
const objectId = parseObjectId(datapoint);
const propertyValue = await new Promise((resolve, reject) => {
client.readProperty(host, objectId, propertyId, (error, value) => {
if (error) {
reject(error);
} else {
resolve(value);
}
});
});
const raw = extractValue(propertyValue);
if (typeof raw === "boolean") {
return raw ? 1 : 0;
}
const numeric = Number(raw);
return Number.isFinite(numeric) ? numeric * (Number(datapoint.scale || 1) || 1) : String(raw ?? "");
} finally {
closeClient(client);
}
}
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: "bacnet-collector", driver: "bacstack", timestamp: new Date().toISOString() });
return;
}
if (req.method === "GET" && url.pathname === "/capabilities") {
replyJson(res, 200, {
protocols: ["bacnet-ip"],
scans: ["bacnet-ip"],
driver: "bacstack",
modes: ["broadcast", "bbmd", "foreign-device"],
});
return;
}
if (req.method === "POST" && url.pathname === "/scan/bacnet") {
try {
const body = await readBody(req);
const payload = await scanNetwork(body);
replyJson(res, 200, payload);
} catch (error) {
replyJson(res, 400, { error: error.message || "BACnet-Scan fehlgeschlagen." });
}
return;
}
if (req.method === "POST" && url.pathname === "/read/bacnet") {
try {
const body = await readBody(req);
const value = await readCurrentValue(body.source || {}, body.datapoint || {});
replyJson(res, 200, { value, readAt: new Date().toISOString() });
} catch (error) {
replyJson(res, 400, { error: error.message || "BACnet-Livewert konnte nicht gelesen werden." });
}
return;
}
replyJson(res, 404, { error: "Nicht gefunden." });
});
server.listen(port, () => {
console.log(`SE Local Trenddata BACnet collector listening on port ${port}`);
});
+291
View File
@@ -0,0 +1,291 @@
import dgram from "node:dgram";
import { networkInterfaces } from "node:os";
function sleep(ms) {
return new Promise((resolve) => setTimeout(resolve, ms));
}
function normalizeText(value) {
return String(value ?? "").trim();
}
function parseHostList(value) {
if (Array.isArray(value)) {
return value.flatMap((item) => parseHostList(item));
}
return String(value || "")
.split(/[,\s;]+/)
.map((item) => item.trim())
.filter(Boolean);
}
function isSubnetMaskLike(address) {
return /^255\.255\.255\.(0|128|192|224|240|248|252|254)$/.test(String(address || ""));
}
function isAutoBroadcast(value) {
return !normalizeText(value) || ["auto", "automatisch"].includes(normalizeText(value).toLowerCase()) || isSubnetMaskLike(value);
}
function parseIpv4(value) {
const address = normalizeText(value).replace(/^::ffff:/, "");
const parts = address.split(".").map((part) => Number(part));
if (parts.length !== 4 || parts.some((part) => !Number.isInteger(part) || part < 0 || part > 255)) {
return null;
}
return parts;
}
function isUsableLanIpv4(value) {
const parts = parseIpv4(value);
if (!parts) return false;
if (parts[0] === 127 || parts[0] === 0 || parts[0] >= 224) return false;
if (parts[0] === 172 && parts[1] >= 17 && parts[1] <= 31) return false; // Docker bridge defaults.
return true;
}
function ipv4ToNumber(parts) {
return (((parts[0] << 24) >>> 0) + (parts[1] << 16) + (parts[2] << 8) + parts[3]) >>> 0;
}
function numberToIpv4(value) {
return [value >>> 24, (value >>> 16) & 255, (value >>> 8) & 255, value & 255].join(".");
}
function broadcastFromAddress(address, netmask = "255.255.255.0") {
const ipParts = parseIpv4(address);
const maskParts = parseIpv4(netmask) || [255, 255, 255, 0];
if (!ipParts || !isUsableLanIpv4(address)) return "";
const ip = ipv4ToNumber(ipParts);
const mask = ipv4ToNumber(maskParts);
return numberToIpv4((ip | (~mask >>> 0)) >>> 0);
}
function getInterfaceBroadcasts() {
return Object.values(networkInterfaces())
.flat()
.filter((item) => item && item.family === "IPv4" && !item.internal && isUsableLanIpv4(item.address))
.map((item) => broadcastFromAddress(item.address, item.netmask))
.filter(Boolean);
}
function getBroadcastTargets(routing, body = {}) {
const configured = normalizeText(routing.broadcastAddress || "auto");
const targets = [];
if (!isAutoBroadcast(configured)) {
targets.push(configured);
}
for (const subnet of parseHostList(body.scanSubnets || body.bacnetScanSubnets || body.options?.bacnetScanSubnets)) {
const broadcast = cidrBroadcast(subnet);
if (broadcast) targets.push(broadcast);
}
for (const candidate of [body.clientHost, body.browserHost, body.locationHost]) {
const broadcast = broadcastFromAddress(candidate);
if (broadcast) targets.push(broadcast);
}
targets.push(...getInterfaceBroadcasts());
targets.push(...parseHostList(process.env.BACNET_BROADCAST_FALLBACKS));
if (!targets.length) {
targets.push("255.255.255.255");
}
return Array.from(new Set(targets));
}
function expandCidrHosts(cidr) {
const match = String(cidr || "").trim().match(/^(\d+\.\d+\.\d+\.\d+)\/(\d{1,2})$/);
if (!match) return [];
const ipParts = parseIpv4(match[1]);
const prefix = Number(match[2]);
if (!ipParts || !Number.isInteger(prefix) || prefix < 16 || prefix > 32) return [];
if (prefix === 32) return [numberToIpv4(ipv4ToNumber(ipParts))];
const ip = ipv4ToNumber(ipParts);
const mask = prefix === 0 ? 0 : (0xffffffff << (32 - prefix)) >>> 0;
const network = (ip & mask) >>> 0;
const broadcast = (network | (~mask >>> 0)) >>> 0;
const limit = Math.min(broadcast - network - 1, 4094);
return Array.from({ length: Math.max(0, limit) }, (_, index) => numberToIpv4(network + index + 1));
}
function cidrBroadcast(cidr) {
const match = String(cidr || "").trim().match(/^(\d+\.\d+\.\d+\.\d+)\/(\d{1,2})$/);
if (!match) return "";
const ipParts = parseIpv4(match[1]);
const prefix = Number(match[2]);
if (!ipParts || !Number.isInteger(prefix) || prefix < 16 || prefix > 30) return "";
const ip = ipv4ToNumber(ipParts);
const mask = (0xffffffff << (32 - prefix)) >>> 0;
return numberToIpv4((ip | (~mask >>> 0)) >>> 0);
}
function expandDirectedBroadcast(address) {
const parts = String(address || "").split(".").map((part) => Number(part));
if (parts.length !== 4 || parts.some((part) => !Number.isInteger(part) || part < 0 || part > 255)) {
return [];
}
if (parts[3] !== 255) {
return [];
}
const prefix = parts.slice(0, 3).join(".");
return Array.from({ length: 254 }, (_, index) => prefix + "." + (index + 1));
}
function buildUnicastScanHosts(body, routing) {
const hosts = [
...parseHostList(body.scanHosts),
...parseHostList(body.scanSubnets || body.bacnetScanSubnets || body.options?.bacnetScanSubnets).flatMap((subnet) => subnet.includes("/") ? expandCidrHosts(subnet) : (parseIpv4(subnet) ? [subnet] : [])),
...parseHostList(process.env.BACNET_SCAN_HOSTS),
];
if (routing.host && routing.mode !== "broadcast") {
hosts.push(routing.host);
}
for (const broadcastAddress of getBroadcastTargets(routing, body)) {
if (broadcastAddress !== "255.255.255.255") {
hosts.push(...expandDirectedBroadcast(broadcastAddress));
}
}
return Array.from(new Set(hosts)).slice(0, 512);
}
function encodeBacnetContextUnsigned(tagNumber, value) {
const numeric = Number(value);
if (!Number.isFinite(numeric) || numeric < 0) {
return Buffer.alloc(0);
}
if (numeric <= 0xff) {
return Buffer.from([(tagNumber << 4) | 0x09, numeric]);
}
if (numeric <= 0xffff) {
const buffer = Buffer.alloc(3);
buffer[0] = (tagNumber << 4) | 0x0a;
buffer.writeUInt16BE(numeric, 1);
return buffer;
}
const buffer = Buffer.alloc(5);
buffer[0] = (tagNumber << 4) | 0x0c;
buffer.writeUInt32BE(numeric, 1);
return buffer;
}
function buildRawWhoIsPacket(deviceId, bvlcFunction) {
const numericDeviceId = Number(deviceId);
const payload = [Buffer.from([0x01, 0x00, 0x10, 0x08])];
if (Number.isFinite(numericDeviceId)) {
payload.push(encodeBacnetContextUnsigned(0, numericDeviceId));
payload.push(encodeBacnetContextUnsigned(1, numericDeviceId));
}
const body = Buffer.concat(payload);
const packet = Buffer.alloc(4 + body.length);
packet[0] = 0x81;
packet[1] = bvlcFunction;
packet.writeUInt16BE(packet.length, 2);
body.copy(packet, 4);
return packet;
}
function parseRawIAm(message) {
if (!message || message.length < 12 || message[0] !== 0x81) {
return null;
}
const start = message.indexOf(0xc4);
if (start < 0 || start + 4 >= message.length) {
return null;
}
const objectId = message.readUInt32BE(start + 1);
const objectType = objectId >>> 22;
const deviceId = objectId & 0x3fffff;
if (objectType !== 8) {
return null;
}
return { deviceId };
}
async function discover() {
const input = JSON.parse(process.env.RAW_BACNET_BODY || "{}");
const body = input.body || {};
const routing = input.routing || {};
const port = Number(routing.port || 47808) || 47808;
const hosts = buildUnicastScanHosts(body, routing);
const devices = [];
const rawTrace = [];
const rawDebug = { hosts, sends: [], packets: [] };
let currentUnicastHost = "";
if (routing.mode !== "broadcast") {
return { devices, rawTrace, rawDebug };
}
const socket = dgram.createSocket({ type: "udp4", reuseAddr: true });
const addDevice = (deviceId, host, rinfo) => {
rawTrace.push({ deviceId: Number(deviceId), host, remote: rinfo ? rinfo.address + ":" + rinfo.port : "" });
if (!Number.isFinite(Number(deviceId)) || !host) {
return;
}
if (!devices.some((item) => item.deviceId === Number(deviceId) && item.address === host)) {
devices.push({
address: host,
host,
port,
reachable: true,
deviceId: Number(deviceId),
name: "BACnet Gerät " + Number(deviceId),
vendorId: null,
maxApdu: null,
segmentation: null,
pointCount: 0,
points: [],
rawRemote: rinfo ? rinfo.address + ":" + rinfo.port : "",
});
}
};
await new Promise((resolve, reject) => {
socket.once("error", reject);
socket.bind(port, "0.0.0.0", () => {
socket.setBroadcast(true);
resolve();
});
});
socket.on("message", (message, rinfo) => {
rawDebug.packets.push({ remote: rinfo.address + ":" + rinfo.port, length: message.length, hex: message.toString("hex") });
const parsed = parseRawIAm(message);
if (parsed) {
addDevice(parsed.deviceId, currentUnicastHost || rinfo.address, rinfo);
}
});
try {
const broadcastPacket = buildRawWhoIsPacket(body.deviceId, 0x0b);
for (const broadcastAddress of getBroadcastTargets(routing, body)) {
socket.send(broadcastPacket, port, broadcastAddress, (error) => rawDebug.sends.push({ host: broadcastAddress, error: error?.message || "" }));
await sleep(350);
}
const perHostDelay = Number(body.unicastDelayMs || process.env.BACNET_UNICAST_DELAY_MS || 20);
const packet = buildRawWhoIsPacket(body.deviceId, 0x0a);
for (const host of hosts) {
const before = devices.length;
currentUnicastHost = host;
socket.send(packet, port, host, (error) => rawDebug.sends.push({ host, error: error?.message || "" }));
await sleep(Math.max(50, perHostDelay));
if (body.deviceId && devices.length > before && devices.some((device) => String(device.deviceId) === String(body.deviceId))) {
break;
}
}
currentUnicastHost = "";
await sleep(250);
} finally {
socket.close();
}
return { devices, rawTrace, rawDebug };
}
discover()
.then((result) => {
process.stdout.write(JSON.stringify(result));
})
.catch((error) => {
process.stderr.write(error?.stack || error?.message || String(error));
process.exitCode = 1;
});