update
This commit is contained in:
@@ -0,0 +1,7 @@
|
||||
FROM node:20-alpine
|
||||
WORKDIR /app
|
||||
COPY package.json ./
|
||||
COPY src ./src
|
||||
RUN mkdir -p /app/data
|
||||
EXPOSE 18110
|
||||
CMD ["npm", "run", "start"]
|
||||
@@ -0,0 +1,9 @@
|
||||
{
|
||||
"name": "seltd-collector",
|
||||
"version": "1.0.0",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"scripts": {
|
||||
"start": "node src/index.js"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,436 @@
|
||||
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";
|
||||
|
||||
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);
|
||||
|
||||
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,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 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 ensureSource(store, body) {
|
||||
const now = new Date().toISOString();
|
||||
const source = {
|
||||
id: body.id || `src-${Date.now()}`,
|
||||
protocol: normalizeProtocol(body.protocol),
|
||||
name: body.name || "Neue Quelle",
|
||||
host: body.host || "",
|
||||
port: body.port || "",
|
||||
deviceId: body.deviceId || "",
|
||||
displayName: body.displayName || body.name || "Neue Quelle",
|
||||
tableName: body.tableName || `trend_${slugify(body.name || body.protocol || "daten")}`,
|
||||
writeMode: normalizeWriteMode(body.writeMode),
|
||||
pollIntervalSeconds: normalizePollInterval(body.pollIntervalSeconds),
|
||||
options: {
|
||||
...(body.options || {}),
|
||||
writeMode: normalizeWriteMode(body.writeMode),
|
||||
pollIntervalSeconds: normalizePollInterval(body.pollIntervalSeconds),
|
||||
},
|
||||
createdAt: body.createdAt || now,
|
||||
updatedAt: now,
|
||||
};
|
||||
|
||||
const existingIndex = store.sources.findIndex((item) => item.id === source.id);
|
||||
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;
|
||||
};
|
||||
|
||||
return {
|
||||
id: `dp-modbus-${Date.now()}-${index}`,
|
||||
sourceId: payload.sourceId || "",
|
||||
protocol: "modbus-tcp",
|
||||
address: get("address", get("register", String(index + 1))),
|
||||
pointIndex: Number(get("pointindex", index + 1)) || index + 1,
|
||||
name: get("name", `Register ${index + 1}`),
|
||||
alias: get("alias", get("name", `Register ${index + 1}`)),
|
||||
unit: get("unit"),
|
||||
dataType: get("datatype", "holding-register"),
|
||||
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\b([^>]*)\/?>(?:<\/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 = getAttr("Name") || getAttr("name") || `KNX ${index + 1}`;
|
||||
const dpt = 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: "",
|
||||
dataType: dpt || "group-address",
|
||||
host: payload.host || "",
|
||||
tableName: payload.tableName || "",
|
||||
createdAt: new Date().toISOString(),
|
||||
});
|
||||
index += 1;
|
||||
}
|
||||
}
|
||||
|
||||
return datapoints;
|
||||
}
|
||||
|
||||
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);
|
||||
});
|
||||
}
|
||||
|
||||
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"],
|
||||
note: "Importe und erste Scan-Stubs sind aktiv. Poll-Intervall und Schreibmodus werden bereits pro Quelle gespeichert; die echte Live-Kommunikation folgt als nächste Ausbaustufe.",
|
||||
});
|
||||
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 === "POST" && url.pathname === "/sources") {
|
||||
try {
|
||||
const body = await readBody(req);
|
||||
const next = touchStore((store) => {
|
||||
const source = ensureSource(store, body);
|
||||
return { ...store, lastSource: source.id };
|
||||
});
|
||||
replyJson(res, 201, next.sources.at(-1));
|
||||
} 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;
|
||||
});
|
||||
replyJson(res, 200, next.sources.find((item) => item.id === sourceId));
|
||||
} 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),
|
||||
}));
|
||||
replyJson(res, 200, { ok: true, sources: next.sources.length, datapoints: next.datapoints.length });
|
||||
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"));
|
||||
}
|
||||
store.datapoints.push(...imported);
|
||||
return store;
|
||||
});
|
||||
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"));
|
||||
}
|
||||
store.datapoints.push(...imported);
|
||||
return store;
|
||||
});
|
||||
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);
|
||||
const host = body.host || "";
|
||||
const portNumber = Number(body.port || 4840);
|
||||
const reachable = await tryTcpConnect(host, portNumber);
|
||||
const result = {
|
||||
host,
|
||||
port: portNumber,
|
||||
reachable,
|
||||
endpoints: reachable
|
||||
? [
|
||||
{ url: `opc.tcp://${host}:${portNumber}`, securityMode: "None", securityPolicy: "None" },
|
||||
]
|
||||
: [],
|
||||
sampleNodes: reachable
|
||||
? [
|
||||
{ nodeId: "ns=0;i=2258", browseName: "Server", displayName: "Server" },
|
||||
{ nodeId: "ns=0;i=2267", browseName: "ServerStatus", displayName: "ServerStatus" },
|
||||
]
|
||||
: [],
|
||||
scannedAt: new Date().toISOString(),
|
||||
};
|
||||
touchStore((store) => {
|
||||
store.scans.opcua.unshift(result);
|
||||
store.scans.opcua = store.scans.opcua.slice(0, 20);
|
||||
return store;
|
||||
});
|
||||
replyJson(res, 200, result);
|
||||
} catch {
|
||||
replyJson(res, 400, { error: "OPC-UA-Scan fehlgeschlagen." });
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (req.method === "POST" && url.pathname === "/scan/bacnet") {
|
||||
try {
|
||||
const body = await readBody(req);
|
||||
const hosts = Array.isArray(body.hosts) ? body.hosts : [body.host].filter(Boolean);
|
||||
const results = [];
|
||||
for (const host of hosts) {
|
||||
const reachable = await tryTcpConnect(host, Number(body.port || 47808), 650);
|
||||
results.push({
|
||||
host,
|
||||
port: Number(body.port || 47808),
|
||||
reachable,
|
||||
deviceId: body.deviceId || "",
|
||||
name: reachable ? `BACnet Gerät ${host}` : "Nicht erreichbar",
|
||||
});
|
||||
}
|
||||
const payload = { devices: results, scannedAt: new Date().toISOString() };
|
||||
touchStore((store) => {
|
||||
store.scans.bacnet.unshift(payload);
|
||||
store.scans.bacnet = store.scans.bacnet.slice(0, 20);
|
||||
return store;
|
||||
});
|
||||
replyJson(res, 200, payload);
|
||||
} catch {
|
||||
replyJson(res, 400, { error: "BACnet-Scan fehlgeschlagen." });
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
replyJson(res, 404, { error: "Nicht gefunden." });
|
||||
});
|
||||
|
||||
server.listen(port, () => {
|
||||
console.log(`SE Local Trenddata collector listening on port ${port}`);
|
||||
});
|
||||
Reference in New Issue
Block a user