diff --git a/api/src/index.js b/api/src/index.js index 1dcfc04..a59c752 100644 --- a/api/src/index.js +++ b/api/src/index.js @@ -1,4 +1,4 @@ -import bcrypt from "bcryptjs"; +import bcrypt from "bcryptjs"; import cors from "cors"; import express from "express"; import { authenticateToken, authenticateTokenOptional, createToken, requireRoles } from "./auth.js"; @@ -214,11 +214,22 @@ function sanitizeApiSettings(body = {}) { }; } -function sanitizeRunnerJob(body = {}, existing = {}) { - const time = String(body.time ?? existing.time ?? "06:00").trim(); - if (!/^\d{2}:\d{2}$/.test(time)) { - throw new Error("Uhrzeit bitte im Format HH:MM angeben."); +function normalizeRunnerTime(value) { + const raw = String(value ?? "").trim().replace(/\s*uhr$/i, "").replace(".", ":"); + const match = raw.match(/^(\d{1,2}):(\d{2})$/); + if (!match) { + throw new Error("Uhrzeit bitte im deutschen Format HH:MM Uhr angeben."); } + const hours = Number(match[1]); + const minutes = Number(match[2]); + if (hours > 23 || minutes > 59) { + throw new Error("Uhrzeit ist ungültig."); + } + return `${String(hours).padStart(2, "0")}:${String(minutes).padStart(2, "0")}`; +} + +function sanitizeRunnerJob(body = {}, existing = {}) { + const time = normalizeRunnerTime(body.time ?? existing.time ?? "06:00"); return { ...existing, id: existing.id || body.id || `job-${Date.now()}`, @@ -246,6 +257,26 @@ async function listRunnerJobs() { })); return jobs.filter(Boolean).sort((left, right) => String(left.time).localeCompare(String(right.time))); } +async function listRunnerQueue() { + const entries = await redis.lrange("runner:run-now", 0, 99); + const jobs = await listRunnerJobs(); + const jobsById = new Map(jobs.map((job) => [job.id, job])); + return entries.map((raw, index) => { + try { + const queued = JSON.parse(raw); + const job = jobsById.get(queued.jobId); + return { + position: index + 1, + jobId: queued.jobId || "", + jobName: job?.name || queued.jobId || "Unbekannter Job", + requestedAt: queued.requestedAt || "", + requestedBy: queued.requestedBy || "", + }; + } catch { + return { position: index + 1, jobId: "", jobName: "Ungültiger Auftrag", requestedAt: "", requestedBy: "" }; + } + }); +} async function resolveSelectionPoints(body, user) { if (body.selectionId) { const selection = await getSelection(redis, body.selectionId); @@ -803,8 +834,9 @@ app.post("/api/dashboard/data", authenticateTokenOptional, wrap(async (req, res) app.get("/api/integrations", authenticateToken, requirePermission("manage_integrations"), wrap(async (_req, res) => { const current = JSON.parse(await redis.get("integration:api") || "{}"); + const settings = sanitizeApiSettings(current); res.json({ - ...sanitizeApiSettings(current), + settings, mailConfigured: Boolean(current.mailApiKey), whatsappConfigured: Boolean(current.whatsappApiKey), voipConfigured: Boolean(current.voipApiKey), @@ -814,11 +846,16 @@ app.get("/api/integrations", authenticateToken, requirePermission("manage_integr app.put("/api/integrations", authenticateToken, requirePermission("manage_integrations"), wrap(async (req, res) => { const payload = sanitizeApiSettings(req.body || {}); await redis.set("integration:api", JSON.stringify(payload)); - res.json({ ...payload, mailConfigured: Boolean(payload.mailApiKey), whatsappConfigured: Boolean(payload.whatsappApiKey), voipConfigured: Boolean(payload.voipApiKey) }); + res.json({ + settings: payload, + mailConfigured: Boolean(payload.mailApiKey), + whatsappConfigured: Boolean(payload.whatsappApiKey), + voipConfigured: Boolean(payload.voipApiKey), + }); })); app.get("/api/runner/jobs", authenticateToken, requirePermission("manage_runner"), wrap(async (_req, res) => { - res.json({ jobs: await listRunnerJobs() }); + res.json({ jobs: await listRunnerJobs(), queue: await listRunnerQueue() }); })); app.post("/api/runner/jobs", authenticateToken, requirePermission("manage_runner"), wrap(async (req, res) => { @@ -839,6 +876,14 @@ app.patch("/api/runner/jobs/:jobId", authenticateToken, requirePermission("manag res.json(job); })); +app.post("/api/runner/jobs/:jobId/test", authenticateToken, requirePermission("manage_runner"), wrap(async (req, res) => { + const id = req.params.jobId; + const existing = JSON.parse(await redis.get(`runner:job:${id}`) || "null"); + if (!existing) return res.status(404).json({ error: "Runner-Job nicht gefunden." }); + await redis.rpush("runner:run-now", JSON.stringify({ jobId: id, requestedAt: new Date().toISOString(), requestedBy: req.user.sub })); + res.status(202).json({ queued: true, jobId: id }); +})); + app.delete("/api/runner/jobs/:jobId", authenticateToken, requirePermission("manage_runner"), wrap(async (req, res) => { const id = req.params.jobId; await redis.del(`runner:job:${id}`, `runner:last:${id}`, `runner:last-result:${id}`); diff --git a/collector/src/index.js b/collector/src/index.js index 3ae26f0..d9705a5 100644 --- a/collector/src/index.js +++ b/collector/src/index.js @@ -1,4 +1,4 @@ -import { createServer } from "node:http"; +import { createServer } from "node:http"; import { mkdirSync, readFileSync, writeFileSync, existsSync } from "node:fs"; import { dirname, join } from "node:path"; import { fileURLToPath } from "node:url"; @@ -33,7 +33,7 @@ function baseStore() { return { sources: [], datapoints: [], - scans: { opcua: [], bacnet: [] }, + scans: { opcua: [], bacnet: [], ipsymcon: [] }, updatedAt: null, }; } @@ -112,12 +112,26 @@ function sanitizeTrendTable(value) { function normalizeProtocol(value) { const protocol = String(value || "").toLowerCase(); - if (["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip"].includes(protocol)) { + if (["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip", "ipsymcon"].includes(protocol)) { return protocol; } return "unknown"; } +function redactSourceSecrets(source) { + const options = { ...(source?.options || {}) }; + const hasIpsymconApiKey = Boolean(options.ipsymconApiKey); + delete options.ipsymconApiKey; + return { ...source, options, hasIpsymconApiKey }; +} + +function redactStoreSecrets(store) { + return { + ...store, + sources: (store.sources || []).map(redactSourceSecrets), + }; +} + function normalizeWriteMode(value) { return String(value || "").toLowerCase() === "cov" ? "cov" : "interval"; } @@ -204,6 +218,9 @@ function ensureSource(store, body) { const bacnetBroadcastAddress = normalizeText(body.bacnetBroadcastAddress || body.options?.bacnetBroadcastAddress || existing?.options?.bacnetBroadcastAddress || "255.255.255.255") || "255.255.255.255"; const bacnetForeignDeviceTtl = normalizeBacnetTtl(body.bacnetForeignDeviceTtl || body.options?.bacnetForeignDeviceTtl || existing?.options?.bacnetForeignDeviceTtl); const bacnetScanSubnets = normalizeText(body.bacnetScanSubnets || body.options?.bacnetScanSubnets || existing?.options?.bacnetScanSubnets || ""); + const ipsymconCategoryIds = normalizeText(body.ipsymconCategoryIds ?? body.options?.ipsymconCategoryIds ?? existing?.options?.ipsymconCategoryIds ?? ""); + const requestedIpsymconApiKey = normalizeText(body.ipsymconApiKey ?? body.options?.ipsymconApiKey ?? ""); + const ipsymconApiKey = requestedIpsymconApiKey || normalizeText(existing?.options?.ipsymconApiKey); const tableName = sanitizeTrendTable(existing?.tableName || buildSourceTableName(body)); const source = { id: body.id || `src-${Date.now()}`, @@ -225,6 +242,8 @@ function ensureSource(store, body) { bacnetBroadcastAddress, bacnetForeignDeviceTtl, bacnetScanSubnets, + ipsymconCategoryIds, + ipsymconApiKey, }, lastValues: body.lastValues || {}, lastPollAt: body.lastPollAt || null, @@ -461,6 +480,9 @@ function buildScannedDatapoints(store, source, rawPoints, toPayload) { unit: existing.enabled === true ? normalizePointUnit(existing.unit || scanned.unit, existing.kind || scanned.kind) : normalizePointUnit(scanned.unit || existing.unit, scanned.kind || existing.kind), kind: existing.kind || scanned.kind, enabled: existing.enabled === true, + alarmEnabled: existing.alarmEnabled === true, + alarmCondition: normalizeText(existing.alarmCondition || scanned.alarmCondition), + alarmGroupId: normalizeText(existing.alarmGroupId || scanned.alarmGroupId), writeMode: normalizeWriteMode(existing.writeMode || scanned.writeMode), condition: normalizeText(existing.condition || scanned.condition), dataType: normalizeText(existing.dataType || scanned.dataType), @@ -476,7 +498,7 @@ function buildScannedDatapoints(store, source, rawPoints, toPayload) { 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 = store.datapoints.filter((item) => !(item.sourceId === datapoint.sourceId && item.id === datapoint.id)); store.datapoints.push(datapoint); }); } @@ -795,6 +817,175 @@ async function readOpcuaValue(source, datapoint) { } } +function parseIpsymconCategoryIds(value) { + const unique = new Set(); + for (const entry of String(value || "").split(/[;,\s]+/)) { + const id = Number(entry.trim()); + if (Number.isInteger(id) && id > 0) { + unique.add(id); + } + } + return Array.from(unique); +} + +function buildIpsymconEndpoint(source) { + const rawHost = normalizeText(source?.host); + if (!rawHost) { + throw new Error("IP-Symcon Host oder URL fehlt."); + } + + const endpoint = /^https?:\/\//i.test(rawHost) ? rawHost : `http://${rawHost}`; + let url; + try { + url = new URL(endpoint); + } catch { + throw new Error("IP-Symcon Host oder URL ist ungültig."); + } + + const port = Number(source?.port || 3777) || 3777; + if (!url.port) { + url.port = String(port); + } + const path = url.pathname.replace(/\/+$/, ""); + url.pathname = !path || path === "/" ? "/hook/api" : (/\/hook\/api$/i.test(path) ? path : `${path}/hook/api`); + url.search = ""; + url.hash = ""; + return url.toString(); +} + +async function requestIpsymconApi(source, payload) { + const apiKey = normalizeText(source?.options?.ipsymconApiKey || source?.ipsymconApiKey); + if (!apiKey) { + throw new Error("Für diese IP-Symcon-Quelle fehlt der API-Schlüssel."); + } + + const controller = new AbortController(); + const timer = setTimeout(() => controller.abort(), 30000); + try { + const response = await fetch(buildIpsymconEndpoint(source), { + method: "POST", + headers: { "Content-Type": "application/json; charset=utf-8" }, + body: JSON.stringify({ + api_key: apiKey, + include_hidden: true, + follow_links: true, + stale_after_seconds: 240, + ...payload, + }), + signal: controller.signal, + }); + const data = await response.json().catch(() => ({})); + if (!response.ok || data?.ok !== true) { + throw new Error(data?.error?.message || data?.error || `IP-Symcon API antwortet mit HTTP ${response.status}.`); + } + return data; + } catch (error) { + if (error?.name === "AbortError") { + throw new Error("IP-Symcon API hat nicht innerhalb von 30 Sekunden geantwortet."); + } + throw error; + } finally { + clearTimeout(timer); + } +} + +function mapIpsymconVariable(source, variable, categoryId, index) { + const variableId = Number(variable?.id); + const variableType = normalizeText(variable?.variable_type || variable?.type || ""); + const kind = variableType === "boolean" || typeof variable?.value === "boolean" ? "digital" : "analog"; + const name = normalizeText(variable?.full_name || variable?.name || `Variable ${variableId}`) || `Variable ${variableId}`; + return { + id: `${source.id}::ipsymcon::${variableId}`, + protocol: "ipsymcon", + pointIndex: index + 1, + alias: name, + name, + address: String(variableId), + nodeId: String(variableId), + dataType: variableType ? `ipsymcon-${variableType}` : "ipsymcon-variable", + kind, + unit: normalizePointUnit(variable?.unit, kind), + currentValue: variable?.value ?? null, + quality: normalizeText(variable?.quality || "unknown"), + path: normalizeText(variable?.actual_path || ""), + categoryId: Number(categoryId) || null, + }; +} + +async function scanIpsymconSource(source) { + const categoryIds = parseIpsymconCategoryIds(source?.options?.ipsymconCategoryIds || source?.ipsymconCategoryIds); + if (!categoryIds.length) { + throw new Error("Bitte mindestens eine IP-Symcon Kategorie-ID angeben."); + } + + const variablesById = new Map(); + const categories = []; + for (const categoryId of categoryIds) { + let offset = 0; + let pageCount = 0; + do { + const response = await requestIpsymconApi(source, { + function: "getcategory", + id: categoryId, + value: 0, + max_depth: 30, + limit: 1000, + offset, + }); + const result = response.result || {}; + const variables = Array.isArray(result.variables) ? result.variables : []; + categories.push({ + id: categoryId, + name: normalizeText(result.root?.name || `Kategorie ${categoryId}`), + path: normalizeText(result.root?.path || ""), + returned: variables.length, + total: Number(result.pagination?.total || variables.length), + }); + for (const variable of variables) { + const variableId = Number(variable?.id); + if (Number.isInteger(variableId) && variableId > 0 && !variablesById.has(variableId)) { + variablesById.set(variableId, { variable, categoryId }); + } + } + const pagination = result.pagination || {}; + if (!pagination.has_more || variables.length === 0 || pageCount >= 20) { + break; + } + offset = Number(pagination.next_offset); + if (!Number.isFinite(offset) || offset < 0) { + break; + } + pageCount += 1; + } while (variablesById.size < LOCAL_POINT_LIMIT); + } + + const points = Array.from(variablesById.values()) + .slice(0, LOCAL_POINT_LIMIT) + .map(({ variable, categoryId }, index) => mapIpsymconVariable(source, variable, categoryId, index)); + return { + host: buildIpsymconEndpoint(source), + categoryIds, + categories, + points, + total: variablesById.size, + truncated: variablesById.size > points.length, + scannedAt: new Date().toISOString(), + }; +} + +async function readIpsymconValue(source, datapoint) { + const variableId = Number(datapoint.nodeId || datapoint.address); + if (!Number.isInteger(variableId) || variableId <= 0) { + throw new Error(`IP-Symcon Variablen-ID fehlt für ${datapoint.alias || datapoint.name}.`); + } + const response = await requestIpsymconApi(source, { function: "getvalue", id: variableId }); + const raw = response.value ?? response.variable?.value; + if (typeof raw === "boolean") { + return raw ? 1 : 0; + } + const numeric = Number(raw); + return Number.isFinite(numeric) ? numeric * (Number(datapoint.scale || 1) || 1) : raw === null || raw === undefined ? null : String(raw); +} function parseBacnetObject(datapoint) { const rawAddress = String(datapoint.address || "").trim(); const shorthand = rawAddress.match(/^(ai|ao|av|bi|bo|bv)(\d+)$/i); @@ -1445,7 +1636,9 @@ async function readDatapointCurrent(source, datapoint) { if (source.protocol === "knx-ip") { return await readKnxValue(source, datapoint); } - throw new Error("Protokoll wird nicht unterstützt."); + if (source.protocol === "ipsymcon") { + return await readIpsymconValue(source, datapoint); + } throw new Error("Protokoll wird nicht unterstützt."); } async function pollGenericSource(source, datapoints, readValue) { @@ -1596,13 +1789,32 @@ async function ensureSourceDatapointRegistryTable() { \`protocol\` VARCHAR(32) NOT NULL, \`unit\` VARCHAR(64) NULL, \`kind\` VARCHAR(32) NULL, + \`enabled\` TINYINT(1) NOT NULL DEFAULT 0, + \`alarm_enabled\` TINYINT(1) NOT NULL DEFAULT 0, + \`last_seen_at\` TIMESTAMP 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_point\` (\`source_id\`, \`point_index\`), KEY \`idx_source_datapoint_table\` (\`table_name\`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci `); + + const migrations = [ + "ALTER TABLE `source_datapoint_registry` ADD COLUMN IF NOT EXISTS `enabled` TINYINT(1) NOT NULL DEFAULT 0 AFTER `kind`", + "ALTER TABLE `source_datapoint_registry` ADD COLUMN IF NOT EXISTS `alarm_enabled` TINYINT(1) NOT NULL DEFAULT 0 AFTER `enabled`", + "ALTER TABLE `source_datapoint_registry` ADD COLUMN IF NOT EXISTS `last_seen_at` TIMESTAMP NULL AFTER `alarm_enabled`", + "ALTER TABLE `source_datapoint_registry` DROP INDEX IF EXISTS `uniq_source_point`", + "ALTER TABLE `source_datapoint_registry` ADD INDEX IF NOT EXISTS `idx_source_point` (`source_id`, `point_index`)", + ]; + + for (const statement of migrations) { + try { + await db.query(statement); + } catch (error) { + console.warn("Registry-Migration übersprungen:", error.message); + } + } } async function syncDatapointRegistryEntries(datapoints) { @@ -1617,16 +1829,6 @@ async function syncDatapointRegistryEntries(datapoints) { } 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\` ( @@ -1638,8 +1840,23 @@ async function syncDatapointRegistryEntries(datapoints) { \`name\`, \`protocol\`, \`unit\`, - \`kind\` - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) + \`kind\`, + \`enabled\`, + \`alarm_enabled\`, + \`last_seen_at\` + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP) + ON DUPLICATE KEY UPDATE + \`source_id\` = VALUES(\`source_id\`), + \`table_name\` = VALUES(\`table_name\`), + \`point_index\` = VALUES(\`point_index\`), + \`alias\` = VALUES(\`alias\`), + \`name\` = VALUES(\`name\`), + \`protocol\` = VALUES(\`protocol\`), + \`unit\` = VALUES(\`unit\`), + \`kind\` = VALUES(\`kind\`), + \`enabled\` = VALUES(\`enabled\`), + \`alarm_enabled\` = VALUES(\`alarm_enabled\`), + \`last_seen_at\` = CURRENT_TIMESTAMP `, [ String(datapoint.id), @@ -1651,11 +1868,12 @@ async function syncDatapointRegistryEntries(datapoints) { String(datapoint.protocol || "unknown"), normalizeUnit(datapoint.unit || ""), String(datapoint.kind || "analog"), + datapoint.enabled === true ? 1 : 0, + datapoint.alarmEnabled === true ? 1 : 0, ] ); } } - async function deleteDatapointRegistryForSource(sourceId) { if (!sourceId) { return; @@ -1673,7 +1891,7 @@ async function rebuildRegistryTablesFromStore() { for (const source of store.sources || []) { await syncSourceRegistryEntry(source); } - await syncDatapointRegistryEntries((store.datapoints || []).filter((datapoint) => datapoint.enabled === true)); + await syncDatapointRegistryEntries(store.datapoints || []); } @@ -1911,6 +2129,9 @@ async function pollDueSources() { const store = readStore(); const now = Date.now(); for (const source of store.sources) { + if (source.protocol === "ipsymcon") { + continue; + } if (shouldPollSource(source, now)) { await pollSource(source, store.datapoints); } @@ -1924,6 +2145,28 @@ setInterval(() => { pollDueSources().catch((error) => console.error("Collector polling failed", error)); }, Math.max(500, pollTickMs)); +async function applyExternalWorkerStates(sources) { + try { + const [rows] = await db.query( + "SELECT source_id, worker, last_poll_at, last_poll_status, last_poll_error FROM source_worker_state WHERE worker = 'seltd-symcon-collector'", + ); + const states = new Map(rows.map((row) => [String(row.source_id), row])); + return (sources || []).map((source) => { + const state = source.protocol === "ipsymcon" ? states.get(String(source.id)) : null; + if (!state) return source; + return { + ...source, + lastPollAt: state.last_poll_at ? new Date(state.last_poll_at).toISOString() : source.lastPollAt, + lastPollStatus: state.last_poll_status || source.lastPollStatus, + lastPollError: state.last_poll_error || "", + worker: state.worker, + }; + }); + } catch { + return sources || []; + } +} + const server = createServer(async (req, res) => { const url = new URL(req.url || "/", `http://${req.headers.host || "localhost"}`); @@ -1939,18 +2182,20 @@ const server = createServer(async (req, res) => { if (req.method === "GET" && url.pathname === "/capabilities") { replyJson(res, 200, { - protocols: ["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip"], + protocols: ["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip", "ipsymcon"], 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.", + scans: ["opc-ua", "bacnet-ip", "ipsymcon"], + manualDatapoints: ["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip", "ipsymcon"], + polling: ["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip", "ipsymcon"], + note: "Modbus TCP, OPC UA, BACnet IP, KNX IP und IP-Symcon werden zyklisch gelesen und in Trendtabellen geschrieben.", }); return; } if (req.method === "GET" && url.pathname === "/sources") { - replyJson(res, 200, readStore()); + const store = readStore(); + store.sources = await applyExternalWorkerStates(store.sources); + replyJson(res, 200, redactStoreSecrets(store)); return; } @@ -1964,23 +2209,25 @@ const server = createServer(async (req, res) => { if (req.method === "GET" && url.pathname === "/debug") { const store = readStore(); + store.sources = await applyExternalWorkerStates(store.sources); replyJson(res, 200, { updatedAt: store.updatedAt, sourceCount: store.sources.length, datapointCount: store.datapoints.length, alarmDatapointCount: store.datapoints.filter((item) => item.alarmEnabled === true).length, sources: store.sources.map((source) => ({ - ...source, + ...redactSourceSecrets(source), datapointCount: store.datapoints.filter((item) => item.sourceId === source.id).length, alarmDatapointCount: store.datapoints.filter((item) => item.sourceId === source.id && item.alarmEnabled === true).length, })), latestScans: { opcua: store.scans.opcua[0] || null, bacnet: store.scans.bacnet[0] || null, + ipsymcon: store.scans.ipsymcon[0] || null, }, driverState: { pollingImplemented: true, - pollingProtocols: ["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip"], + pollingProtocols: ["modbus-tcp", "opc-ua", "bacnet-ip", "knx-ip", "ipsymcon"], note: "Alle konfigurierten Protokollquellen pollten über native Treiber. Fehler stehen direkt an der Quelle.", }, }); @@ -1996,7 +2243,7 @@ const server = createServer(async (req, res) => { }); const createdSource = next.sources.at(-1); await ensureSourceTrendTable(createdSource); - replyJson(res, 201, createdSource); + replyJson(res, 201, redactSourceSecrets(createdSource)); } catch { replyJson(res, 400, { error: "Ungültige JSON-Daten." }); } @@ -2017,7 +2264,7 @@ const server = createServer(async (req, res) => { }); const updatedSource = next.sources.find((item) => item.id === sourceId); await ensureSourceTrendTable(updatedSource); - replyJson(res, 200, updatedSource); + replyJson(res, 200, redactSourceSecrets(updatedSource)); } catch (error) { replyJson(res, error.message === "not-found" ? 404 : 400, { error: error.message === "not-found" ? "Quelle nicht gefunden." : "Ungültige JSON-Daten.", @@ -2195,6 +2442,42 @@ const server = createServer(async (req, res) => { return; } + if (req.method === "POST" && url.pathname === "/scan/ipsymcon") { + try { + const body = await readBody(req); + const store = readStore(); + const source = store.sources.find((item) => item.id === body.sourceId); + if (!source || source.protocol !== "ipsymcon") { + replyJson(res, 404, { error: "IP-Symcon-Quelle nicht gefunden." }); + return; + } + + const result = await scanIpsymconSource(source); + let scannedCandidates = []; + touchStore((current) => { + const currentSource = current.sources.find((item) => item.id === source.id); + if (currentSource) { + scannedCandidates = buildScannedDatapoints(current, currentSource, result.points || [], (point, index) => ({ + ...point, + id: point.id || `${currentSource.id}::ipsymcon::${point.address || index + 1}`, + protocol: "ipsymcon", + pointIndex: index + 1, + alias: point.alias || point.name || `Variable ${point.address || index + 1}`, + name: point.name || point.alias || `Variable ${point.address || index + 1}`, + })); + upsertDatapoints(current, scannedCandidates); + } + current.scans.ipsymcon.unshift(result); + current.scans.ipsymcon = current.scans.ipsymcon.slice(0, 20); + return current; + }); + await syncDatapointRegistryEntries(scannedCandidates); + replyJson(res, 200, result); + } catch (error) { + replyJson(res, 400, { error: error.message || "IP-Symcon-Scan fehlgeschlagen." }); + } + return; + } if (req.method === "POST" && url.pathname === "/scan/bacnet") { try { const body = await readBody(req); @@ -2240,21 +2523,3 @@ rebuildRegistryTablesFromStore().catch((error) => console.error("Registry sync f server.listen(port, () => { console.log(`SE Local Trenddata collector listening on port ${port}`); }); - - - - - - - - - - - - - - - - - - diff --git a/docker-compose.yml b/docker-compose.yml index 5f26070..f9bdf59 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,4 +1,4 @@ -services: +services: mariadb: image: mariadb:11.4 container_name: seltd-mariadb @@ -102,6 +102,28 @@ bacnet-collector: condition: service_started + symcon-collector: + build: + context: ./symcon-collector + container_name: seltd-symcon-collector + restart: unless-stopped + environment: + DB_HOST: mariadb + DB_PORT: 3306 + DB_NAME: wago + DB_USER: root + DB_PASSWORD: ${MARIADB_ROOT_PASSWORD:-SE3112} + SOURCES_FILE: /app/collector-data/sources.json + SYMCON_POLL_TICK_MS: ${SYMCON_POLL_TICK_MS:-1000} + SYMCON_MAX_CONCURRENCY: ${SYMCON_MAX_CONCURRENCY:-4} + volumes: + - collector_data:/app/collector-data:ro + depends_on: + mariadb: + condition: service_healthy + collector: + condition: service_started + bacnet-collector: build: context: ./bacnet-collector @@ -128,7 +150,8 @@ DB_USER: root DB_PASSWORD: ${MARIADB_ROOT_PASSWORD:-SE3112} EXPORT_DIR: /exports - RUNNER_TICK_MS: ${RUNNER_TICK_MS:-30000} + RUNNER_TICK_MS: ${RUNNER_TICK_MS:-1000} + RUNNER_TIMEZONE: ${RUNNER_TIMEZONE:-Europe/Berlin} volumes: - runner_exports:/exports depends_on: @@ -141,7 +164,7 @@ container_name: seltd-dockhand restart: unless-stopped ports: - - "3000:3000" + - "3008:3000" volumes: - /var/run/docker.sock:/var/run/docker.sock - dockhand_data:/app/data diff --git a/runner/src/index.js b/runner/src/index.js index ffb6a92..aa31a45 100644 --- a/runner/src/index.js +++ b/runner/src/index.js @@ -1,4 +1,4 @@ -import dns from "node:dns"; +import dns from "node:dns"; import { createReadStream, mkdirSync, statSync } from "node:fs"; import { writeFile } from "node:fs/promises"; import https from "node:https"; @@ -18,26 +18,51 @@ const db = mysql.createPool({ connectionLimit: 3, }); const exportDir = process.env.EXPORT_DIR || "/exports"; -const tickMs = Number(process.env.RUNNER_TICK_MS || 30000); +const tickMs = Math.max(500, Number(process.env.RUNNER_TICK_MS || 1000)); +const timeZone = process.env.RUNNER_TIMEZONE || "Europe/Berlin"; function safeParse(value, fallback = null) { try { return value ? JSON.parse(value) : fallback; } catch { return fallback; } } +function formatExportValue(value) { + if (value === null || value === undefined) return null; + if (typeof value === "boolean") return value ? "1" : "0"; + const text = String(value).trim(); + if (!text) return null; + const numeric = Number(text); + // German decimals prevent Excel from treating values such as 24.2 as dates. + if (Number.isFinite(numeric)) { + return new Intl.NumberFormat("de-DE", { useGrouping: false, maximumFractionDigits: 10 }).format(numeric); + } + return text; +} function csvCell(value, delimiter) { const text = String(value ?? ""); return /["\r\n;]/.test(text) || text.includes(delimiter) ? `"${text.replace(/"/g, '""')}"` : text; } +function localDateParts(date = new Date()) { + const parts = new Intl.DateTimeFormat("en-CA", { + timeZone, + year: "numeric", + month: "2-digit", + day: "2-digit", + hour: "2-digit", + minute: "2-digit", + hourCycle: "h23", + }).formatToParts(date); + return Object.fromEntries(parts.filter((part) => part.type !== "literal").map((part) => [part.type, part.value])); +} + function localDateKey(date = new Date()) { - const y = date.getFullYear(); - const m = String(date.getMonth() + 1).padStart(2, "0"); - const d = String(date.getDate()).padStart(2, "0"); - return `${y}-${m}-${d}`; + const parts = localDateParts(date); + return `${parts.year}-${parts.month}-${parts.day}`; } function localTime(date = new Date()) { - return `${String(date.getHours()).padStart(2, "0")}:${String(date.getMinutes()).padStart(2, "0")}`; + const parts = localDateParts(date); + return `${parts.hour}:${parts.minute}`; } function sanitizeTableName(value) { @@ -57,7 +82,11 @@ async function isNarrowTable(table) { async function getPointMeta(isp, pointIndex) { const raw = await redis.get(`meta:isp:${isp}`); const meta = safeParse(raw, {}); - const point = (meta.points || []).find((item) => Number(item.pointIndex) === Number(pointIndex)); + const points = meta.points || {}; + // Legacy metadata stores values by index, newer sources use an array. + const point = Array.isArray(points) + ? points.find((item) => Number(item.pointIndex) === Number(pointIndex)) + : points[String(pointIndex)] || points[Number(pointIndex)]; return { alias: point?.alias || `Wert${pointIndex}`, unit: point?.unit || "" }; } @@ -76,14 +105,22 @@ async function exportSelection(job, selection) { `SELECT datum, alias, value FROM \`${table}\` WHERE point_index = ? AND datum >= ? ORDER BY datum ASC`, [pointIndex, since], ); - for (const item of data) rows.push([item.datum?.toISOString?.() || item.datum, table, pointIndex, item.alias || meta.alias, item.value, meta.unit]); + for (const item of data) { + const value = formatExportValue(item.value); + if (value === null) continue; + rows.push([item.datum?.toISOString?.() || item.datum, table, pointIndex, item.alias || meta.alias, value, meta.unit]); + } } else { const column = `Wert${pointIndex}`; const [data] = await db.query( `SELECT datum, \`${column}\` AS value FROM \`${table}\` WHERE datum >= ? ORDER BY datum ASC`, [since], ); - for (const item of data) rows.push([item.datum?.toISOString?.() || item.datum, table, pointIndex, meta.alias, item.value, meta.unit]); + for (const item of data) { + const value = formatExportValue(item.value); + if (value === null) continue; + rows.push([item.datum?.toISOString?.() || item.datum, table, pointIndex, meta.alias, value, meta.unit]); + } } } @@ -91,7 +128,8 @@ async function exportSelection(job, selection) { const stamp = new Date().toISOString().replace(/[:.]/g, "-"); const fileName = `${selection.name || selection.id || "selection"}-${stamp}.csv`.replace(/[^a-zA-Z0-9_.-]+/g, "_"); const filePath = path.join(exportDir, fileName); - await writeFile(filePath, rows.map((row) => row.map((cell) => csvCell(cell, delimiter)).join(delimiter)).join("\r\n"), "utf8"); + // Excel uses the BOM to reliably recognize UTF-8 when opening CSV files directly. + await writeFile(filePath, `\uFEFF${rows.map((row) => row.map((cell) => csvCell(cell, delimiter)).join(delimiter)).join("\r\n")}`, "utf8"); return filePath; } @@ -144,16 +182,24 @@ async function postMultipartIpv4(url, fields, attachmentPath = "") { }); } +function parseMailRecipients(value) { + return [...new Set(String(value || "").split(/[;,]/).map((entry) => entry.trim()).filter(Boolean))]; +} + async function maybeMail(job, filePath) { if (!job.mailEnabled || !job.mailTo) return; const settings = safeParse(await redis.get("integration:api"), {}); if (!settings.mailApiKey) return; - await postMultipartIpv4(settings.mailUrl || "https://mailapi.se-inno.de", { - api_key: settings.mailApiKey, - address: job.mailTo, - subject: job.mailSubject || "SE-LTD CSV Export", - message: job.mailMessage || "Automatischer CSV Export aus SE Local Trenddata.", - }, filePath); + const recipients = parseMailRecipients(job.mailTo); + if (!recipients.length) return; + for (const address of recipients) { + await postMultipartIpv4(settings.mailUrl || "https://mailapi.se-inno.de", { + api_key: settings.mailApiKey, + address, + subject: job.mailSubject || "SE-LTD CSV Export", + message: job.mailMessage || "Automatischer CSV Export aus SE Local Trenddata.", + }, filePath); + } } async function runJob(job) { @@ -164,7 +210,24 @@ async function runJob(job) { await redis.set(`runner:last-result:${job.id}`, JSON.stringify({ ok: true, filePath, ranAt: new Date().toISOString() })); } +async function runQueuedJobs() { + for (let index = 0; index < 20; index += 1) { + const queued = safeParse(await redis.lpop("runner:run-now")); + if (!queued?.jobId) return; + const job = safeParse(await redis.get(`runner:job:${queued.jobId}`)); + if (!job) continue; + try { + await runJob(job); + console.log(`Runner-Test ausgeführt: ${job.name || job.id}`); + } catch (error) { + await redis.set(`runner:last-result:${job.id}`, JSON.stringify({ ok: false, error: error.message, ranAt: new Date().toISOString() })); + console.error(`Runner-Test fehlgeschlagen (${job.name || job.id}):`, error.message); + } + } +} + async function tick() { + await runQueuedJobs(); const ids = await redis.smembers("runner:job:ids"); const now = new Date(); const time = localTime(now); @@ -182,6 +245,6 @@ async function tick() { } } -console.log("SE Local Trenddata runner started"); +console.log(`SE Local Trenddata runner started (${timeZone}, ${tickMs} ms)`); setInterval(() => tick().catch((error) => console.error(error)), tickMs); tick().catch((error) => console.error(error)); diff --git a/symcon-collector/Dockerfile b/symcon-collector/Dockerfile new file mode 100644 index 0000000..da6fccf --- /dev/null +++ b/symcon-collector/Dockerfile @@ -0,0 +1,10 @@ +FROM node:20-alpine + +WORKDIR /app + +COPY package*.json ./ +RUN npm install + +COPY src ./src + +CMD ["npm", "start"] \ No newline at end of file diff --git a/symcon-collector/package.json b/symcon-collector/package.json new file mode 100644 index 0000000..0010238 --- /dev/null +++ b/symcon-collector/package.json @@ -0,0 +1,12 @@ +{ + "name": "seltd-symcon-collector", + "version": "1.0.0", + "private": true, + "type": "module", + "scripts": { + "start": "node src/index.js" + }, + "dependencies": { + "mysql2": "^3.12.0" + } +} \ No newline at end of file diff --git a/symcon-collector/src/index.js b/symcon-collector/src/index.js new file mode 100644 index 0000000..caec383 --- /dev/null +++ b/symcon-collector/src/index.js @@ -0,0 +1,265 @@ +import { existsSync, readFileSync } from "node:fs"; +import mysql from "mysql2/promise"; + +const sourcesFile = process.env.SOURCES_FILE || "/app/collector-data/sources.json"; +const pollTickMs = Math.max(500, Number(process.env.SYMCON_POLL_TICK_MS || 1000)); +const maxConcurrency = Math.max(1, Math.min(16, Number(process.env.SYMCON_MAX_CONCURRENCY || 4))); +const sourcePollState = new Map(); +const valueCache = new Map(); +let polling = false; + +const db = mysql.createPool({ + host: process.env.DB_HOST || "mariadb", + port: Number(process.env.DB_PORT || 3306), + database: process.env.DB_NAME || "wago", + user: process.env.DB_USER || "root", + password: process.env.DB_PASSWORD || "SE3112", + charset: "utf8mb4", + waitForConnections: true, + connectionLimit: Math.max(2, maxConcurrency + 1), +}); + +function normalizeText(value) { + return String(value ?? "").trim(); +} + +function normalizeInterval(value) { + const number = Number(value); + return Number.isFinite(number) && number > 0 ? Math.round(number) : 60; +} + +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; +} + +async function ensureWorkerStateTable() { + await db.query(` + CREATE TABLE IF NOT EXISTS \`source_worker_state\` ( + \`source_id\` VARCHAR(191) NOT NULL, + \`worker\` VARCHAR(64) NOT NULL, + \`last_poll_at\` TIMESTAMP NULL, + \`last_poll_status\` VARCHAR(255) NOT NULL DEFAULT '', + \`last_poll_error\` TEXT NULL, + \`updated_at\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, + PRIMARY KEY (\`source_id\`) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci + `); +} + +async function updateWorkerState(source, status, error = "") { + await db.query( + ` + INSERT INTO \`source_worker_state\` (\`source_id\`, \`worker\`, \`last_poll_at\`, \`last_poll_status\`, \`last_poll_error\`) + VALUES (?, 'seltd-symcon-collector', CURRENT_TIMESTAMP, ?, ?) + ON DUPLICATE KEY UPDATE + \`worker\` = VALUES(\`worker\`), + \`last_poll_at\` = CURRENT_TIMESTAMP, + \`last_poll_status\` = VALUES(\`last_poll_status\`), + \`last_poll_error\` = VALUES(\`last_poll_error\`) + `, + [String(source.id), String(status || ""), String(error || "")], + ); +} + +function readStore() { + if (!existsSync(sourcesFile)) { + return { sources: [], datapoints: [] }; + } + try { + const store = JSON.parse(readFileSync(sourcesFile, "utf8")); + return { + sources: Array.isArray(store.sources) ? store.sources : [], + datapoints: Array.isArray(store.datapoints) ? store.datapoints : [], + }; + } catch (error) { + console.warn("Quellenkonfiguration konnte nicht gelesen werden:", error.message); + return { sources: [], datapoints: [] }; + } +} + +function shouldPoll(source, now) { + const lastPoll = sourcePollState.get(source.id) || 0; + return now - lastPoll >= normalizeInterval(source.pollIntervalSeconds) * 1000; +} + +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 (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; +} + +function buildEndpoint(source) { + const rawHost = normalizeText(source.host); + if (!rawHost) throw new Error("IP-Symcon Host oder URL fehlt."); + const input = /^https?:\/\//i.test(rawHost) ? rawHost : `http://${rawHost}`; + const url = new URL(input); + if (!url.port) url.port = String(Number(source.port || 3777) || 3777); + const path = url.pathname.replace(/\/+$/, ""); + url.pathname = !path || path === "/" ? "/hook/api" : (/\/hook\/api$/i.test(path) ? path : `${path}/hook/api`); + url.search = ""; + url.hash = ""; + return url.toString(); +} + +async function requestApi(source, payload) { + const apiKey = normalizeText(source.options?.ipsymconApiKey || source.ipsymconApiKey); + if (!apiKey) throw new Error("IP-Symcon API-Schlüssel fehlt."); + const controller = new AbortController(); + const timeout = setTimeout(() => controller.abort(), 30000); + try { + const response = await fetch(buildEndpoint(source), { + method: "POST", + headers: { "Content-Type": "application/json; charset=utf-8" }, + body: JSON.stringify({ api_key: apiKey, ...payload }), + signal: controller.signal, + }); + const data = await response.json().catch(() => ({})); + if (!response.ok || data.ok !== true) { + throw new Error(data?.error?.message || data?.error || `IP-Symcon API HTTP ${response.status}`); + } + return data; + } finally { + clearTimeout(timeout); + } +} + +async function readValue(source, datapoint) { + const variableId = Number(datapoint.nodeId || datapoint.address); + if (!Number.isInteger(variableId) || variableId <= 0) { + throw new Error(`Ungültige Variablen-ID bei ${datapoint.alias || datapoint.name}.`); + } + const response = await requestApi(source, { function: "getvalue", id: variableId }); + const raw = response.value ?? response.variable?.value; + if (typeof raw === "boolean") return raw ? 1 : 0; + const numeric = Number(raw); + return Number.isFinite(numeric) ? numeric * (Number(datapoint.scale || 1) || 1) : raw === null || raw === undefined ? null : String(raw); +} + +async function mapLimit(items, limit, task) { + const result = []; + let nextIndex = 0; + const workers = Array.from({ length: Math.min(limit, items.length) }, async () => { + while (nextIndex < items.length) { + const item = items[nextIndex++]; + result.push(await task(item)); + } + }); + await Promise.all(workers); + return result; +} + +async function ensureNarrowTrendTable(tableName) { + const table = sanitizeTrendTable(tableName); + const [columns] = await db.query( + "SELECT COLUMN_NAME AS name FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ?", + [table], + ); + const names = new Set(columns.map((entry) => String(entry.name).toLowerCase())); + if (names.size && Array.from(names).some((name) => /^wert\d+$/.test(name))) { + throw new Error(`${table} ist eine Legacy-Tabelle und keine IP-Symcon-Trendtabelle.`); + } + if (!names.size) { + await db.query(` + CREATE TABLE \`${table}\` ( + \`id\` BIGINT NOT NULL AUTO_INCREMENT, + \`datum\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + \`userlevel\` INT 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_${table}_datum\` (\`datum\`), + KEY \`idx_${table}_point_datum\` (\`point_index\`, \`datum\`) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci + `); + } + return table; +} + +async function writeReadings(source, readings) { + if (!readings.length) return 0; + const table = await ensureNarrowTrendTable(source.tableName); + const placeholders = readings.map(() => "(?, ?, ?, ?, ?)").join(", "); + const values = readings.flatMap((reading) => [0, String(source.id), Number(reading.pointIndex), String(reading.alias), String(reading.value)]); + await db.query(`INSERT INTO \`${table}\` (\`userlevel\`, \`source_id\`, \`point_index\`, \`alias\`, \`value\`) VALUES ${placeholders}`, values); + return readings.length; +} + +async function pollSource(source, points) { + const errors = []; + const results = await mapLimit(points, maxConcurrency, async (datapoint) => { + try { + return { datapoint, value: await readValue(source, datapoint) }; + } catch (error) { + errors.push(`${datapoint.alias || datapoint.address}: ${error.message}`); + return null; + } + }); + + const readings = []; + for (const result of results) { + if (!result) continue; + const { datapoint, value } = result; + const key = `${source.id}:${datapoint.id}`; + const serialized = String(value); + const writeMode = datapoint.writeMode || source.writeMode || "interval"; + if (matchesLogCondition(datapoint.condition, value) && (writeMode !== "cov" || valueCache.get(key) !== serialized)) { + readings.push({ pointIndex: datapoint.pointIndex, alias: datapoint.alias || datapoint.name || `Wert${datapoint.pointIndex}`, value }); + } + valueCache.set(key, serialized); + } + + const written = await writeReadings(source, readings); + const suffix = errors.length ? `, ${errors.length} Fehler` : ""; + const status = errors.length ? `teilweise ok: ${written} Wert(e) geschrieben${suffix}` : `ok: ${written} Wert(e) geschrieben`; + await updateWorkerState(source, status, errors.slice(0, 5).join(" | ")); + console.log(`[${source.name || source.id}] ${written} Wert(e) geschrieben${suffix}`); + if (errors.length) console.warn(`[${source.name || source.id}]`, errors.slice(0, 5).join(" | ")); +} + +async function pollDueSources() { + if (polling) return; + polling = true; + try { + const store = readStore(); + const now = Date.now(); + const sources = store.sources.filter((source) => source.protocol === "ipsymcon" && shouldPoll(source, now)); + for (const source of sources) { + const points = store.datapoints.filter((point) => point.sourceId === source.id && point.protocol === "ipsymcon" && point.enabled === true); + sourcePollState.set(source.id, now); + if (!points.length) { + await updateWorkerState(source, "idle: keine Datenpunkte"); + continue; + } + try { + await pollSource(source, points); + } catch (error) { + await updateWorkerState(source, "error", error.message || "Polling fehlgeschlagen.").catch(() => undefined); + console.error(`[${source.name || source.id}] Polling fehlgeschlagen:`, error.message); + } + } + } finally { + polling = false; + } +} + +setInterval(() => pollDueSources().catch((error) => console.error("IP-Symcon Polling fehlgeschlagen:", error.message)), pollTickMs); +ensureWorkerStateTable() + .then(() => pollDueSources()) + .catch((error) => console.error("IP-Symcon Start-Polling fehlgeschlagen:", error.message)); +console.log(`SE LTD Symcon Collector gestartet (Intervall-Prüfung ${pollTickMs} ms, Parallelität ${maxConcurrency}).`); \ No newline at end of file diff --git a/web/public/bg-darkmode.png b/web/public/bg-darkmode.png new file mode 100644 index 0000000..64d4fea Binary files /dev/null and b/web/public/bg-darkmode.png differ diff --git a/web/public/bg-whitemode.png b/web/public/bg-whitemode.png new file mode 100644 index 0000000..c14d036 Binary files /dev/null and b/web/public/bg-whitemode.png differ diff --git a/web/src/App.jsx b/web/src/App.jsx index 2ad77df..ae11f0d 100644 --- a/web/src/App.jsx +++ b/web/src/App.jsx @@ -1,4 +1,4 @@ - + import { useEffect, useMemo, useState } from "react"; import { api } from "./api.js"; import Gauge from "./components/Gauge.jsx"; @@ -123,6 +123,7 @@ const PROTOCOL_DEFAULT_PORTS = { "opc-ua": "4840", "bacnet-ip": "47808", "knx-ip": "3671", + ipsymcon: "3777", }; const defaultSourceDraft = { id: "", @@ -139,6 +140,8 @@ const defaultSourceDraft = { bacnetBroadcastAddress: "auto", bacnetScanSubnets: "", bacnetForeignDeviceTtl: "120", + ipsymconCategoryIds: "", + ipsymconApiKey: "", }; @@ -190,6 +193,10 @@ function getOpcuaScanDatapointId(sourceId, node) { return sourceId && node?.nodeId ? `${sourceId}::${node.nodeId}` : ""; } +function getIpsymconScanDatapointId(sourceId, point) { + return sourceId && point?.address ? `${sourceId}::ipsymcon::${point.address}` : ""; +} + function getBacnetScanDatapointId(sourceId, sourceHost, point, index = 0) { if (!sourceId || !point) { return ""; @@ -732,10 +739,11 @@ export default function App() { const [collectorPointFilter, setCollectorPointFilter] = useState(""); const [alarmPointFilter, setAlarmPointFilter] = useState(""); const [manualDatapointDraft, setManualDatapointDraft] = useState(emptyManualDatapointDraft()); - const [scanResults, setScanResults] = useState({ opcua: null, bacnet: null }); + const [scanResults, setScanResults] = useState({ opcua: null, bacnet: null, ipsymcon: null }); const [debugSnapshot, setDebugSnapshot] = useState(null); const [integrationSettings, setIntegrationSettings] = useState({ mailApiKey: "", whatsappApiKey: "", phoneApiKey: "" }); const [runnerJobs, setRunnerJobs] = useState([]); + const [runnerQueue, setRunnerQueue] = useState([]); const [runnerDraft, setRunnerDraft] = useState({ name: "", selectionId: "", enabled: true, time: "06:00", rangeHours: "24", delimiter: ";", mailEnabled: false, mailTo: "", mailSubject: "Trenddaten", mailMessage: "Automatischer Export aus SE Local Trenddata." }); const trendIsps = useMemo(() => isps.filter((isp) => Number(isp.rowCount) > 0), [isps]); @@ -811,18 +819,36 @@ export default function App() { setIsps(nextIsps); return nextIsps; }; + const scannedLoggedPointIds = useMemo(() => new Set( + collectorDatapoints + .filter((item) => item.sourceId === selectedSourceId && item.enabled === true) + .map((item) => String(item.id)) + ), [collectorDatapoints, selectedSourceId]); const filteredOpcuaNodes = useMemo( - () => (scanResults.opcua?.nodes || []).filter((node) => matchesSearch(collectorPointFilter, [ + () => (scanResults.opcua?.nodes || []).filter((node) => !scannedLoggedPointIds.has(getOpcuaScanDatapointId(selectedSourceId, node)) && matchesSearch(collectorPointFilter, [ node.displayName, node.browseName, node.nodeId, node.path, node.currentValue, ])), - [collectorPointFilter, scanResults.opcua] + [collectorPointFilter, scanResults.opcua, scannedLoggedPointIds, selectedSourceId] + ); + const filteredIpsymconPoints = useMemo( + () => (scanResults.ipsymcon?.points || []).filter((point) => !scannedLoggedPointIds.has(getIpsymconScanDatapointId(selectedSourceId, point)) && matchesSearch(collectorPointFilter, [ + point.alias, + point.name, + point.address, + point.nodeId, + point.path, + point.categoryId, + point.currentValue, + point.unit, + ])), + [collectorPointFilter, scanResults.ipsymcon, scannedLoggedPointIds, selectedSourceId] ); const filteredBacnetPoints = useMemo( - () => (scanResults.bacnet?.points || []).filter((point) => matchesSearch(collectorPointFilter, [ + () => (scanResults.bacnet?.points || []).filter((point, index) => !scannedLoggedPointIds.has(getBacnetScanDatapointId(selectedSourceId, currentSource?.host, point, index)) && matchesSearch(collectorPointFilter, [ point.alias, point.name, point.id, @@ -831,7 +857,7 @@ export default function App() { point.objectInstance, point.currentValue, ])), - [collectorPointFilter, scanResults.bacnet] + [collectorPointFilter, currentSource?.host, scanResults.bacnet, scannedLoggedPointIds, selectedSourceId] ); const filteredBacnetDevices = useMemo( () => (scanResults.bacnet?.devices || []).filter((device) => matchesSearch(collectorPointFilter, [ @@ -842,7 +868,7 @@ export default function App() { device.vendorId, device.pointCount, ])), - [collectorPointFilter, scanResults.bacnet] + [collectorPointFilter, currentSource?.host, scanResults.bacnet, scannedLoggedPointIds, selectedSourceId] ); const filteredAvailableCollectorPoints = useMemo( () => availableCollectorPoints.filter((item) => matchesSearch(collectorPointFilter, [ @@ -1000,7 +1026,7 @@ export default function App() { }, [collectorSources, isCreatingSource, selectedSourceId]); useEffect(() => { - setScanResults({ opcua: null, bacnet: null }); + setScanResults({ opcua: null, bacnet: null, ipsymcon: null }); setCollectorPointFilter(""); }, [selectedSourceId]); @@ -1123,11 +1149,12 @@ export default function App() { useEffect(() => { if (!token || !canManageRunner) { setRunnerJobs([]); + setRunnerQueue([]); return; } api.getRunnerJobs(token) - .then((result) => setRunnerJobs(result.jobs || [])) + .then((result) => { setRunnerJobs(result.jobs || []); setRunnerQueue(result.queue || []); }) .catch((error) => setBootError(error.message)); }, [canManageRunner, token]); @@ -1187,6 +1214,8 @@ export default function App() { bacnetBroadcastAddress: currentSource.options?.bacnetBroadcastAddress || "auto", bacnetScanSubnets: currentSource.options?.bacnetScanSubnets || "", bacnetForeignDeviceTtl: String(currentSource.options?.bacnetForeignDeviceTtl || 120), + ipsymconCategoryIds: currentSource.options?.ipsymconCategoryIds || "", + ipsymconApiKey: "", }); setManualDatapointDraft((current) => ({ ...current, @@ -1722,6 +1751,8 @@ export default function App() { bacnetBroadcastAddress: sourceDraft.bacnetBroadcastAddress, bacnetScanSubnets: sourceDraft.bacnetScanSubnets, bacnetForeignDeviceTtl: sourceDraft.bacnetForeignDeviceTtl, + ipsymconCategoryIds: sourceDraft.ipsymconCategoryIds, + ipsymconApiKey: sourceDraft.ipsymconApiKey, }, }; @@ -1906,6 +1937,25 @@ export default function App() { } }; + const runIpsymconScan = async () => { + const activeSourceId = selectedSourceId || sourceDraft.id; + if (!activeSourceId) { + setBootError("Bitte zuerst eine IP-Symcon-Quelle auswählen oder anlegen."); + return; + } + + resetFlash(); + try { + const result = await api.scanIpsymcon({ sourceId: activeSourceId }); + setScanResults((current) => ({ ...current, ipsymcon: result })); + await refreshCollectorData(); + await refreshCollectorDebug().catch(() => undefined); + setSuccessMessage(`IP-Symcon-Scan abgeschlossen: ${result.points?.length || 0} Datenpunkt(e) aus ${result.categories?.length || 0} Kategorie(n).`); + } catch (error) { + setBootError(error.message); + } + }; + const selectBacnetDevice = (device) => { const host = device.host || device.address || ""; const name = normalizeDisplayText(device.name, host || "BACnet Quelle"); @@ -2085,6 +2135,7 @@ export default function App() { } const refreshed = await api.getRunnerJobs(token); setRunnerJobs(refreshed.jobs || []); + setRunnerQueue(refreshed.queue || []); resetRunnerDraft(); setSuccessMessage("Runner-Job gespeichert."); } catch (error) { @@ -2092,6 +2143,26 @@ export default function App() { } }; + const testRunnerJob = async (job) => { + if (!window.confirm(`Test für "${job.name || job.id}" jetzt ausführen? Der CSV-Export wird sofort erstellt und bei aktivierter Mailzustellung verschickt.`)) { + return; + } + resetFlash(); + try { + await api.testRunnerJob(token, job.id); + setSuccessMessage("Testauftrag eingereiht. Das Ergebnis wird in wenigen Sekunden am Job angezeigt."); + window.setTimeout(async () => { + try { + const refreshed = await api.getRunnerJobs(token); + setRunnerJobs(refreshed.jobs || []); + setRunnerQueue(refreshed.queue || []); + } catch {} + }, 1800); + } catch (error) { + setBootError(error.message); + } + }; + const deleteRunnerJob = async (job) => { resetFlash(); if (!window.confirm(`Runner-Job "${job.name || job.id}" wirklich löschen?`)) { @@ -2101,6 +2172,7 @@ export default function App() { await api.deleteRunnerJob(token, job.id); const refreshed = await api.getRunnerJobs(token); setRunnerJobs(refreshed.jobs || []); + setRunnerQueue(refreshed.queue || []); if (runnerDraft.id === job.id) { resetRunnerDraft(); } @@ -2309,7 +2381,7 @@ export default function App() {
Noch kein automatischer Export angelegt.
Noch kein automatischer Export angelegt.
Keine Aufträge warten auf die Ausführung.