From 2ce3be068ecdcf4dd1aa2ca0ef9c11628623df7d Mon Sep 17 00:00:00 2001 From: jhartworks Date: Mon, 29 Jun 2026 15:24:26 +0200 Subject: [PATCH] beta. OPC fixed + tested --- api/src/index.js | 62 ++- api/src/lib/isp.js | 359 ++++++++++--- api/src/lib/store.js | 96 ++++ collector/src/index.js | 434 ++++++++++++++-- docker-compose.yml | 24 +- docker/mariadb/init/01-init.sql | 8 +- web/src/App.jsx | 896 ++++++++++++++++++++++---------- web/src/api.js | 3 + web/src/styles.css | 131 +++++ 9 files changed, 1614 insertions(+), 399 deletions(-) diff --git a/api/src/index.js b/api/src/index.js index dd720d0..101fdb3 100644 --- a/api/src/index.js +++ b/api/src/index.js @@ -5,7 +5,7 @@ import { authenticateToken, createToken, requireRoles } from "./auth.js"; import { config } from "./config.js"; import { pool } from "./db.js"; import { - createWideTrendTable, + createTrendTable, ensureIspMetadata, LOCAL_POINT_LIMIT, fetchLatestValues, @@ -13,6 +13,7 @@ import { getAvailableIsps, normalizeEngineeringUnit, normalizeRange, + resolveTrendTablePointCount, sanitizeIspName, sanitizePointIndex, updateIspMetadata, @@ -25,10 +26,12 @@ import { getSelection, listRoles, listSelections, + listUnits, saveDashboard, savePreferences, saveRole, saveSelection, + saveUnit, } from "./lib/store.js"; import { redis } from "./redis.js"; @@ -378,25 +381,14 @@ app.get("/api/auth/me", authenticateToken, wrap(async (req, res) => { })); async function resolveIspPointCount(isp) { - const safeIsp = sanitizeIspName(isp); - const [rows] = await pool.query( - ` - SELECT COUNT(*) AS pointCount - FROM information_schema.COLUMNS - WHERE TABLE_SCHEMA = ? - AND TABLE_NAME = ? - AND COLUMN_NAME LIKE 'Wert%' - `, - [config.dbName, safeIsp] - ); - return Number(rows[0]?.pointCount || 0) || (safeIsp.startsWith("trend_") ? LOCAL_POINT_LIMIT : 200); + return await resolveTrendTablePointCount(pool, config.dbName, isp); } app.get("/api/isps", authenticateToken, wrap(async (_req, res) => { const isps = await getAvailableIsps(pool, config.dbName); const response = await Promise.all( isps.map(async (isp) => { - const metadata = await ensureIspMetadata(redis, isp.tableName, isp.pointCount); + const metadata = await ensureIspMetadata(pool, redis, isp.tableName, isp.pointCount); return { ...isp, displayName: metadata.displayName, @@ -410,7 +402,7 @@ app.get("/api/isps", authenticateToken, wrap(async (_req, res) => { app.get("/api/isps/:isp/aliases", authenticateToken, wrap(async (req, res) => { const pointCount = await resolveIspPointCount(req.params.isp); - const metadata = await ensureIspMetadata(redis, req.params.isp, pointCount); + const metadata = await ensureIspMetadata(pool, redis, req.params.isp, pointCount); res.json(metadata); })); @@ -421,7 +413,7 @@ app.put( wrap(async (req, res) => { const { displayName, description, points } = req.body || {}; const pointCount = await resolveIspPointCount(req.params.isp); - const updated = await updateIspMetadata(redis, req.params.isp, (metadata) => { + const updated = await updateIspMetadata(pool, redis, req.params.isp, (metadata) => { if (typeof displayName === "string") { metadata.displayName = displayName.trim() || metadata.displayName; } @@ -454,11 +446,40 @@ app.put( }) ); +app.get("/api/units", authenticateToken, wrap(async (_req, res) => { + const units = await listUnits(redis); + res.json({ units }); +})); + +app.post("/api/units", authenticateToken, requirePermission("edit_aliases"), wrap(async (req, res) => { + const symbol = String(req.body.symbol || "").trim(); + const label = String(req.body.label || symbol).trim(); + + if (!symbol) { + return res.status(400).json({ error: "Einheit fehlt." }); + } + + const unit = await saveUnit(redis, { symbol, label }); + res.status(201).json({ unit }); +})); + +app.patch("/api/units/:symbol", authenticateToken, requirePermission("edit_aliases"), wrap(async (req, res) => { + const existingSymbol = decodeURIComponent(req.params.symbol || "").trim(); + if (!existingSymbol) { + return res.status(400).json({ error: "Einheit fehlt." }); + } + + const unit = await saveUnit(redis, { + symbol: String(req.body.symbol || existingSymbol).trim(), + label: String(req.body.label || req.body.symbol || existingSymbol).trim() || existingSymbol, + }, existingSymbol); + res.json({ unit }); +})); app.post("/api/trend-tables", authenticateToken, requirePermission("manage_sources"), wrap(async (req, res) => { const rawName = String(req.body.tableName || "").trim().toLowerCase(); const displayName = String(req.body.displayName || rawName).trim(); - const safeName = await createWideTrendTable(pool, rawName.startsWith("trend_") ? rawName : `trend_${rawName}`); - const metadata = await updateIspMetadata(redis, safeName, (current) => ({ + const safeName = await createTrendTable(pool, rawName.startsWith("trend_") ? rawName : `trend_${rawName}`); + const metadata = await updateIspMetadata(pool, redis, safeName, (current) => ({ ...current, displayName: displayName || current.displayName, }), LOCAL_POINT_LIMIT); @@ -482,7 +503,7 @@ app.post( const normalized = []; for (const isp of targets) { - const metadata = await ensureIspMetadata(redis, isp.tableName, isp.pointCount || (isp.tableName.startsWith("trend_") ? LOCAL_POINT_LIMIT : 200)); + const metadata = await ensureIspMetadata(pool, redis, isp.tableName, isp.pointCount || (isp.tableName.startsWith("trend_") ? LOCAL_POINT_LIMIT : 200)); normalized.push({ isp: isp.tableName, metadata }); } @@ -740,3 +761,6 @@ app.listen(config.port, () => { + + + diff --git a/api/src/lib/isp.js b/api/src/lib/isp.js index b58856d..3044291 100644 --- a/api/src/lib/isp.js +++ b/api/src/lib/isp.js @@ -116,21 +116,20 @@ function createRandomPointColor(seed) { export function normalizeEngineeringUnit(value) { const raw = String(value ?? "").trim(); if (!raw) { - return "°C"; + return "\u00B0C"; } - const normalized = raw - .replace(/°/g, "°") - .replace(/^\?C$/i, "°C") - .replace(/^°\s*C$/i, "°C"); + if (/^\?C$/i.test(raw) || /^\u00B0\s*C$/i.test(raw) || (/[^\x00-\x7F]/.test(raw) && /C/i.test(raw))) { + return "\u00B0C"; + } - return normalized || "°C"; + return raw; } function createDefaultPoint(index) { return { alias: `Wert${index}`, - unit: "°C", + unit: "\u00B0C", factor: 1, kind: "analog", min: 0, @@ -164,7 +163,7 @@ async function ensureWideTrendColumns(pool, tableName) { } } -export async function createWideTrendTable(pool, tableName) { +export async function createLegacyWideTrendTable(pool, tableName) { const safeTable = sanitizeIspName(tableName); const sql = ` CREATE TABLE IF NOT EXISTS \`${safeTable}\` ( @@ -181,18 +180,145 @@ export async function createWideTrendTable(pool, tableName) { return safeTable; } +export async function createTrendTable(pool, tableName) { + const safeTable = sanitizeIspName(tableName); + await pool.query(` + CREATE TABLE IF NOT EXISTS \`${safeTable}\` ( + \`id\` BIGINT NOT NULL AUTO_INCREMENT, + \`datum\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + \`userlevel\` INT(1) 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_${safeTable}_datum\` (\`datum\`), + KEY \`idx_${safeTable}_point_datum\` (\`point_index\`, \`datum\`), + KEY \`idx_${safeTable}_source_point\` (\`source_id\`, \`point_index\`) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci; + `); + return safeTable; +} + +async function loadSourcePointSeeds(pool, tableName) { + const safeTable = sanitizeIspName(tableName); + + try { + const [rows] = await pool.query( + ` + SELECT + point_index AS pointIndex, + alias, + name, + unit, + kind + FROM \`source_datapoint_registry\` + WHERE table_name = ? + ORDER BY point_index ASC, updated_at ASC + `, + [safeTable] + ); + + return rows.map((row) => ({ + pointIndex: sanitizePointIndex(row.pointIndex), + alias: String(row.alias || "").trim(), + name: String(row.name || "").trim(), + unit: String(row.unit || "").trim(), + kind: row.kind === "digital" ? "digital" : "analog", + })); + } catch (error) { + if (String(error?.message || "").includes("doesn't exist")) { + return []; + } + throw error; + } +} + +async function hasSourceRegistryEntry(pool, tableName) { + const safeTable = sanitizeIspName(tableName); + + try { + const [rows] = await pool.query( + "SELECT 1 FROM `source_registry` WHERE table_name = ? LIMIT 1", + [safeTable] + ); + return rows.length > 0; + } catch (error) { + if (String(error?.message || "").includes("doesn't exist")) { + return false; + } + throw error; + } +} + +async function describeTrendTable(pool, dbName, tableName) { + const safeTable = sanitizeIspName(tableName); + const usesExplicitSchema = Boolean(dbName); + const [columns] = await pool.query( + ` + SELECT COLUMN_NAME AS columnName + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = ${usesExplicitSchema ? "?" : "DATABASE()"} + AND TABLE_NAME = ? + ORDER BY ORDINAL_POSITION + `, + usesExplicitSchema ? [dbName, safeTable] : [safeTable] + ); + + const columnNames = columns.map((column) => column.columnName); + const lowerNames = new Set(columnNames.map((name) => String(name).toLowerCase())); + const widePointCount = columnNames.filter((name) => /^Wert\d+$/i.test(name)).length; + + if (widePointCount > 0) { + const [rows] = await pool.query(`SELECT COUNT(*) AS rowCount, MAX(datum) AS latestTimestamp FROM \`${safeTable}\``); + return { + tableName: safeTable, + storage: "wide", + pointCount: widePointCount, + rowCount: Number(rows[0]?.rowCount || 0), + latestTimestamp: rows[0]?.latestTimestamp || null, + }; + } + + if (lowerNames.has("datum") && lowerNames.has("point_index") && lowerNames.has("value")) { + const sourceSeeds = await loadSourcePointSeeds(pool, safeTable); + const sourceBacked = await hasSourceRegistryEntry(pool, safeTable); + const [rows] = await pool.query( + ` + SELECT + COUNT(*) AS rowCount, + MAX(datum) AS latestTimestamp, + COALESCE(MAX(point_index), 0) AS pointCount + FROM \`${safeTable}\` + ` + ); + const measuredPointCount = Number(rows[0]?.pointCount || 0); + const pointCount = sourceSeeds.length ? sourceSeeds.length : measuredPointCount; + return { + tableName: safeTable, + storage: "narrow", + pointCount: pointCount || (sourceBacked ? 0 : (safeTable.startsWith("trend_") ? LOCAL_POINT_LIMIT : LEGACY_POINT_LIMIT)), + rowCount: Number(rows[0]?.rowCount || 0), + latestTimestamp: rows[0]?.latestTimestamp || null, + }; + } + + return null; +} + +export async function resolveTrendTablePointCount(pool, dbName, tableName) { + const safeTable = sanitizeIspName(tableName); + const info = await describeTrendTable(pool, dbName, safeTable); + if (!info) { + return safeTable.startsWith("trend_") ? LOCAL_POINT_LIMIT : LEGACY_POINT_LIMIT; + } + return Number(info.pointCount || 0) || (safeTable.startsWith("trend_") ? LOCAL_POINT_LIMIT : LEGACY_POINT_LIMIT); +} + export async function getAvailableIsps(pool, dbName) { const [tables] = await pool.query( ` - SELECT - t.TABLE_NAME AS tableName, - ( - SELECT COUNT(*) - FROM information_schema.COLUMNS c - WHERE c.TABLE_SCHEMA = t.TABLE_SCHEMA - AND c.TABLE_NAME = t.TABLE_NAME - AND c.COLUMN_NAME LIKE 'Wert%' - ) AS pointCount + SELECT t.TABLE_NAME AS tableName FROM information_schema.TABLES t WHERE t.TABLE_SCHEMA = ? AND (t.TABLE_NAME LIKE 'isp%' OR t.TABLE_NAME LIKE 'trend_%') @@ -201,36 +327,17 @@ export async function getAvailableIsps(pool, dbName) { [dbName] ); - const stats = new Map(); - - await Promise.all( - tables.map(async (table) => { - const safeTable = sanitizeIspName(table.tableName); - const [rows] = await pool.query( - `SELECT COUNT(*) AS rowCount, MAX(datum) AS latestTimestamp FROM \`${safeTable}\`` - ); - stats.set(safeTable, rows[0] || { rowCount: 0, latestTimestamp: null }); - }) - ); - - return tables.map((table) => { - const safeTable = sanitizeIspName(table.tableName); - const info = stats.get(safeTable) || {}; - - return { - tableName: safeTable, - pointCount: Number(table.pointCount || 0), - rowCount: Number(info.rowCount || 0), - latestTimestamp: info.latestTimestamp || null, - }; - }); + const described = await Promise.all(tables.map((table) => describeTrendTable(pool, dbName, table.tableName))); + return described.filter(Boolean); } - export function createDefaultMetadata(isp, pointCount = LEGACY_POINT_LIMIT) { const points = {}; + const pointOrder = []; for (let index = 1; index <= pointCount; index += 1) { - points[String(index)] = createDefaultPoint(index); + const key = String(index); + points[key] = createDefaultPoint(index); + pointOrder.push(key); } return { @@ -238,12 +345,18 @@ export function createDefaultMetadata(isp, pointCount = LEGACY_POINT_LIMIT) { displayName: isp.toUpperCase(), description: "", points, + pointOrder, }; } -export async function ensureIspMetadata(redis, isp, pointCount = LEGACY_POINT_LIMIT) { +export async function ensureIspMetadata(pool, redis, isp, pointCount = LEGACY_POINT_LIMIT) { const safeIsp = sanitizeIspName(isp); const key = `${META_PREFIX}${safeIsp}`; + const sourceSeeds = await loadSourcePointSeeds(pool, safeIsp); + const seedMap = new Map(sourceSeeds.map((seed) => [String(seed.pointIndex), seed])); + const pointOrder = sourceSeeds.length + ? sourceSeeds.map((seed) => String(seed.pointIndex)) + : Array.from({ length: pointCount }, (_value, index) => String(index + 1)); const existing = await redis.get(key); if (existing) { @@ -255,19 +368,43 @@ export async function ensureIspMetadata(redis, isp, pointCount = LEGACY_POINT_LI changed = true; } - for (let index = 1; index <= pointCount; index += 1) { - const pointKey = String(index); + if (!Array.isArray(parsed.pointOrder)) { + parsed.pointOrder = []; + changed = true; + } + + if (JSON.stringify(parsed.pointOrder) !== JSON.stringify(pointOrder)) { + parsed.pointOrder = pointOrder; + changed = true; + } + + pointOrder.forEach((pointKey) => { + const index = Number(pointKey); const defaultPoint = createDefaultPoint(index); + const sourceSeed = seedMap.get(pointKey); + const initialAlias = sourceSeed?.alias || sourceSeed?.name || defaultPoint.alias; if (!parsed.points[pointKey]) { - parsed.points[pointKey] = defaultPoint; + parsed.points[pointKey] = { + ...defaultPoint, + alias: initialAlias, + unit: normalizeEngineeringUnit(sourceSeed?.unit || defaultPoint.unit), + kind: sourceSeed?.kind || defaultPoint.kind, + }; changed = true; - continue; + return; } const point = parsed.points[pointKey]; - if (!point.alias) { - point.alias = defaultPoint.alias; + if (!point.alias || point.alias === defaultPoint.alias) { + const nextAlias = initialAlias || defaultPoint.alias; + if (point.alias !== nextAlias) { + point.alias = nextAlias; + changed = true; + } + } + if (sourceSeed?.kind && !point.kind) { + point.kind = sourceSeed.kind; changed = true; } const normalizedUnit = normalizeEngineeringUnit(point.unit); @@ -275,8 +412,15 @@ export async function ensureIspMetadata(redis, isp, pointCount = LEGACY_POINT_LI point.unit = normalizedUnit; changed = true; } + if ((!point.unit || point.unit === normalizeEngineeringUnit(defaultPoint.unit)) && sourceSeed?.unit) { + const seededUnit = normalizeEngineeringUnit(sourceSeed.unit); + if (point.unit !== seededUnit) { + point.unit = seededUnit; + changed = true; + } + } if (!point.kind) { - point.kind = defaultPoint.kind; + point.kind = sourceSeed?.kind || defaultPoint.kind; changed = true; } if (!Number.isFinite(Number(point.factor))) { @@ -300,7 +444,7 @@ export async function ensureIspMetadata(redis, isp, pointCount = LEGACY_POINT_LI point.color = defaultPoint.color; changed = true; } - } + }); if (changed) { await redis.set(key, JSON.stringify(parsed)); @@ -310,18 +454,31 @@ export async function ensureIspMetadata(redis, isp, pointCount = LEGACY_POINT_LI } const metadata = createDefaultMetadata(safeIsp, pointCount); + metadata.pointOrder = pointOrder; + pointOrder.forEach((pointKey) => { + const sourceSeed = seedMap.get(pointKey); + if (!sourceSeed) { + return; + } + metadata.points[pointKey] = { + ...metadata.points[pointKey], + alias: sourceSeed.alias || sourceSeed.name || metadata.points[pointKey].alias, + unit: normalizeEngineeringUnit(sourceSeed.unit || metadata.points[pointKey].unit), + kind: sourceSeed.kind || metadata.points[pointKey].kind, + }; + }); await redis.set(key, JSON.stringify(metadata)); return metadata; } -export async function updateIspMetadata(redis, isp, updater, pointCount = LEGACY_POINT_LIMIT) { - const metadata = await ensureIspMetadata(redis, isp, pointCount); +export async function updateIspMetadata(pool, redis, isp, updater, pointCount = LEGACY_POINT_LIMIT) { + const metadata = await ensureIspMetadata(pool, redis, isp, pointCount); const updated = updater(JSON.parse(JSON.stringify(metadata))); await redis.set(`${META_PREFIX}${sanitizeIspName(isp)}`, JSON.stringify(updated)); return updated; } -export async function getPointDescriptors(redis, points) { +export async function getPointDescriptors(pool, redis, points) { const grouped = new Map(); points.forEach((point) => { @@ -338,10 +495,11 @@ export async function getPointDescriptors(redis, points) { const descriptors = []; for (const [isp, pointIndexes] of grouped.entries()) { - const metadata = await ensureIspMetadata(redis, isp, LOCAL_POINT_LIMIT); + const pointCount = await resolveTrendTablePointCount(pool, null, isp); + const metadata = await ensureIspMetadata(pool, redis, isp, pointCount); pointIndexes.forEach((pointIndex) => { - const pointMeta = metadata.points[String(pointIndex)]; + const pointMeta = metadata.points[String(pointIndex)] || createDefaultPoint(pointIndex); descriptors.push({ isp, pointIndex, @@ -362,7 +520,7 @@ export async function getPointDescriptors(redis, points) { } export async function fetchTrendSeries(pool, redis, points, range) { - const descriptors = await getPointDescriptors(redis, points); + const descriptors = await getPointDescriptors(pool, redis, points); const grouped = descriptors.reduce((accumulator, descriptor) => { if (!accumulator.has(descriptor.isp)) { accumulator.set(descriptor.isp, []); @@ -375,6 +533,43 @@ export async function fetchTrendSeries(pool, redis, points, range) { const mergedRows = new Map(); for (const [isp, pointDescriptors] of grouped.entries()) { + const tableInfo = await describeTrendTable(pool, null, isp); + if (!tableInfo) { + continue; + } + + if (tableInfo.storage === "narrow") { + const placeholders = pointDescriptors.map(() => "?").join(", "); + const limit = Math.max(pointDescriptors.length, (range.maxRows || 4000) * Math.max(1, pointDescriptors.length)); + const [rows] = await pool.query( + ` + SELECT datum, point_index, value + FROM \`${isp}\` + WHERE point_index IN (${placeholders}) + AND datum BETWEEN ? AND ? + ORDER BY datum DESC, id DESC + LIMIT ? + `, + [...pointDescriptors.map((descriptor) => descriptor.pointIndex), range.from, range.to, limit] + ); + + rows.reverse().forEach((row) => { + const timestamp = new Date(row.datum).toISOString(); + if (!mergedRows.has(timestamp)) { + mergedRows.set(timestamp, { timestamp }); + } + const target = mergedRows.get(timestamp); + const descriptor = pointDescriptors.find((item) => item.pointIndex === Number(row.point_index)); + if (!descriptor) { + return; + } + const parsed = parseTrendValue(row.value, descriptor.factor); + target[descriptor.key] = parsed.numeric; + target[`${descriptor.key}:raw`] = parsed.raw; + }); + continue; + } + const columns = pointDescriptors .map((descriptor) => `\`Wert${descriptor.pointIndex}\``) .join(", "); @@ -410,7 +605,7 @@ export async function fetchTrendSeries(pool, redis, points, range) { } export async function fetchLatestValues(pool, redis, points) { - const descriptors = await getPointDescriptors(redis, points); + const descriptors = await getPointDescriptors(pool, redis, points); const grouped = descriptors.reduce((accumulator, descriptor) => { if (!accumulator.has(descriptor.isp)) { accumulator.set(descriptor.isp, []); @@ -423,6 +618,52 @@ export async function fetchLatestValues(pool, redis, points) { const values = []; for (const [isp, pointDescriptors] of grouped.entries()) { + const tableInfo = await describeTrendTable(pool, null, isp); + if (!tableInfo) { + pointDescriptors.forEach((descriptor) => { + values.push({ + ...descriptor, + timestamp: null, + raw: null, + value: null, + }); + }); + continue; + } + + if (tableInfo.storage === "narrow") { + const placeholders = pointDescriptors.map(() => "?").join(", "); + const [rows] = await pool.query( + ` + SELECT point_index, datum, value + FROM ( + SELECT + point_index, + datum, + value, + ROW_NUMBER() OVER (PARTITION BY point_index ORDER BY datum DESC, id DESC) AS point_rank + FROM \`${isp}\` + WHERE point_index IN (${placeholders}) + ) ranked + WHERE point_rank = 1 + `, + pointDescriptors.map((descriptor) => descriptor.pointIndex) + ); + const latestByPoint = new Map(rows.map((row) => [Number(row.point_index), row])); + + pointDescriptors.forEach((descriptor) => { + const row = latestByPoint.get(descriptor.pointIndex); + const parsed = parseTrendValue(row?.value, descriptor.factor); + values.push({ + ...descriptor, + timestamp: row?.datum ? new Date(row.datum).toISOString() : null, + raw: parsed.raw, + value: parsed.numeric, + }); + }); + continue; + } + const columns = pointDescriptors .map((descriptor) => `\`Wert${descriptor.pointIndex}\``) .join(", "); @@ -459,3 +700,5 @@ export function normalizeRange(input = {}) { } + + diff --git a/api/src/lib/store.js b/api/src/lib/store.js index 90858ed..b262be9 100644 --- a/api/src/lib/store.js +++ b/api/src/lib/store.js @@ -1,6 +1,7 @@ const SELECTION_SET_KEY = "selection:ids"; const SELECTION_SEQ_KEY = "selection:seq"; const ROLE_SET_KEY = "role:names"; +const UNIT_SET_KEY = "unit:symbols"; export const DEFAULT_ROLE_DEFINITIONS = [ { @@ -23,10 +24,48 @@ export const DEFAULT_ROLE_DEFINITIONS = [ }, ]; +export const DEFAULT_UNIT_DEFINITIONS = [ + { symbol: "\u00B0C", label: "Grad Celsius", isSystem: true }, + { symbol: "%", label: "Prozent", isSystem: true }, + { symbol: "K", label: "Kelvin", isSystem: true }, + { symbol: "V", label: "Volt", isSystem: true }, + { symbol: "A", label: "Ampere", isSystem: true }, + { symbol: "Pa", label: "Pascal", isSystem: true }, +]; + function safeParse(jsonValue) { return jsonValue ? JSON.parse(jsonValue) : null; } +function normalizeUnitSymbol(value, fallback = "") { + const raw = String(value ?? "").trim(); + const collapsed = raw.replace(/\s+/g, ""); + + if (!collapsed) { + return fallback; + } + + if (/^\?C$/i.test(collapsed) || /^Â?\u00B0C$/i.test(collapsed)) { + return "\u00B0C"; + } + + return raw; +} + +function unitStorageKey(symbol) { + return `unit:${encodeURIComponent(symbol)}`; +} + +function normalizeUnitRecord(unit) { + const symbol = normalizeUnitSymbol(unit?.symbol, "\u00B0C"); + return { + symbol, + label: String(unit?.label || symbol).trim() || symbol, + isSystem: Boolean(unit?.isSystem), + updatedAt: unit?.updatedAt || new Date().toISOString(), + }; +} + function normalizeRoleName(value) { const normalized = String(value || "") .trim() @@ -72,6 +111,22 @@ async function ensureDefaultRoles(redis) { } } +async function ensureDefaultUnits(redis) { + for (const definition of DEFAULT_UNIT_DEFINITIONS) { + const key = unitStorageKey(definition.symbol); + const existing = safeParse(await redis.get(key)); + const payload = normalizeUnitRecord({ + ...definition, + ...(existing || {}), + symbol: definition.symbol, + label: existing?.label || definition.label, + isSystem: true, + }); + await redis.set(key, JSON.stringify(payload)); + await redis.sadd(UNIT_SET_KEY, definition.symbol); + } +} + export async function listRoles(redis) { await ensureDefaultRoles(redis); const names = await redis.smembers(ROLE_SET_KEY); @@ -118,6 +173,47 @@ export async function saveRole(redis, role, previousName = "") { return merged; } +export async function listUnits(redis) { + await ensureDefaultUnits(redis); + const symbols = await redis.smembers(UNIT_SET_KEY); + const items = await Promise.all(symbols.map((symbol) => redis.get(unitStorageKey(symbol)))); + + return items + .map((item) => safeParse(item)) + .filter(Boolean) + .map((unit) => normalizeUnitRecord(unit)) + .sort((left, right) => left.symbol.localeCompare(right.symbol, "de-DE")); +} + +export async function saveUnit(redis, unit, previousSymbol = "") { + await ensureDefaultUnits(redis); + + const existingSymbol = normalizeUnitSymbol(previousSymbol || unit?.symbol, ""); + const payload = normalizeUnitRecord(unit); + const current = existingSymbol ? safeParse(await redis.get(unitStorageKey(existingSymbol))) : null; + + if (current?.isSystem && existingSymbol !== payload.symbol) { + throw new Error("Systemeinheiten können nicht umbenannt werden."); + } + + const merged = normalizeUnitRecord({ + ...(current || {}), + ...payload, + symbol: current?.isSystem ? existingSymbol : payload.symbol, + isSystem: Boolean(current?.isSystem || payload.isSystem), + updatedAt: new Date().toISOString(), + }); + + if (existingSymbol && existingSymbol !== merged.symbol) { + await redis.del(unitStorageKey(existingSymbol)); + await redis.srem(UNIT_SET_KEY, existingSymbol); + } + + await redis.set(unitStorageKey(merged.symbol), JSON.stringify(merged)); + await redis.sadd(UNIT_SET_KEY, merged.symbol); + return merged; +} + export async function listSelections(redis, user) { const ids = await redis.smembers(SELECTION_SET_KEY); const items = await Promise.all(ids.map((id) => redis.get(`selection:${id}`))); diff --git a/collector/src/index.js b/collector/src/index.js index f4c4535..0ea48b6 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"; @@ -72,7 +72,7 @@ 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-Methods": "GET,POST,PUT,PATCH,DELETE,OPTIONS", "Access-Control-Allow-Headers": "Content-Type, Authorization", }); res.end(JSON.stringify(body)); @@ -130,16 +130,40 @@ function normalizePointKind(value) { return String(value || "").toLowerCase() === "digital" ? "digital" : "analog"; } +function repairText(value) { + let current = String(value ?? ""); + for (let attempt = 0; attempt < 3; attempt += 1) { + const currentScore = (current.match(/[\u00C3\u00C2\uFFFD]/g) || []).length; + if (!currentScore) { + break; + } + const repaired = Buffer.from(current, "latin1").toString("utf8"); + const repairedScore = (repaired.match(/[\u00C3\u00C2\uFFFD]/g) || []).length; + if (repairedScore < currentScore) { + current = repaired; + continue; + } + break; + } + return current.replace(/\u00A0/g, " ").trim(); +} + +function normalizeText(value) { + return repairText(String(value === undefined || value === null ? "" : value).trim()); +} + function normalizeUnit(value) { - const raw = String(value ?? "").trim(); - if (!raw) { - return "°C"; + const raw = normalizeText(value); + const collapsed = raw.replace(/\s+/g, ""); + if (!collapsed) { + return "\u00B0C"; } - return raw - .replace(/°/g, "°") - .replace(/^\?C$/i, "°C") - .replace(/^°\s*C$/i, "°C"); + if (/^\?C$/i.test(collapsed) || /^\u00B0C$/i.test(collapsed)) { + return "\u00B0C"; + } + + return raw; } function normalizePointIndex(value, fallback) { @@ -159,11 +183,11 @@ function ensureSource(store, body) { 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", + name: normalizeText(body.name) || "Neue Quelle", + host: normalizeText(body.host), + port: normalizeText(body.port), + deviceId: normalizeText(body.deviceId), + displayName: normalizeText(body.displayName || body.name) || "Neue Quelle", tableName, writeMode: normalizeWriteMode(body.writeMode), pollIntervalSeconds: normalizePollInterval(body.pollIntervalSeconds), @@ -297,8 +321,8 @@ function createManualDatapoint(store, payload) { const fallbackIndex = existing.length + 1; const pointIndex = normalizePointIndex(payload.pointIndex, fallbackIndex); const protocol = normalizeProtocol(payload.protocol || source.protocol); - const alias = String(payload.alias || payload.name || `Punkt ${pointIndex}`).trim() || `Punkt ${pointIndex}`; - const name = String(payload.name || payload.alias || alias).trim() || alias; + const alias = normalizeText(payload.alias || payload.name || `Punkt ${pointIndex}`) || `Punkt ${pointIndex}`; + const name = normalizeText(payload.name || payload.alias || alias) || alias; const datapoint = { id: payload.id || `dp-${protocol}-${Date.now()}-${pointIndex}`, @@ -311,10 +335,10 @@ function createManualDatapoint(store, payload) { kind: normalizePointKind(payload.kind), enabled: payload.enabled !== false, writeMode: normalizeWriteMode(payload.writeMode || source.writeMode), - condition: String(payload.condition || "").trim(), - dataType: String(payload.dataType || (protocol === "knx-ip" ? "group-address" : "holding-register")).trim(), - address: String(payload.address || payload.register || "").trim(), - nodeId: String(payload.nodeId || "").trim(), + condition: normalizeText(payload.condition), + dataType: normalizeText(payload.dataType || (protocol === "knx-ip" ? "group-address" : "holding-register")), + address: normalizeText(payload.address || payload.register || ""), + nodeId: normalizeText(payload.nodeId || ""), scale: Number(payload.scale || 1) || 1, host: source.host || "", tableName: source.tableName || "", @@ -333,27 +357,27 @@ function createManualDatapoint(store, payload) { function createScanDatapoint(source, payload, index) { const pointIndex = normalizePointIndex(payload.pointIndex, index + 1); const protocol = normalizeProtocol(payload.protocol || source.protocol); - const alias = String(payload.alias || payload.name || payload.displayName || `Punkt ${pointIndex}`).trim() || `Punkt ${pointIndex}`; + const alias = normalizeText(payload.alias || payload.name || payload.displayName || `Punkt ${pointIndex}`) || `Punkt ${pointIndex}`; return { id: payload.id || `dp-${protocol}-${Date.now()}-${Math.round(Math.random() * 100000)}-${pointIndex}`, sourceId: source.id, protocol, pointIndex, - name: String(payload.name || alias).trim() || alias, + name: normalizeText(payload.name || alias) || alias, alias, unit: normalizeUnit(payload.unit), kind: normalizePointKind(payload.kind), enabled: payload.enabled === true, writeMode: normalizeWriteMode(payload.writeMode || source.writeMode), - condition: String(payload.condition || "").trim(), - dataType: String(payload.dataType || "").trim(), - address: String(payload.address || "").trim(), - nodeId: String(payload.nodeId || "").trim(), + condition: normalizeText(payload.condition), + dataType: normalizeText(payload.dataType), + address: normalizeText(payload.address), + nodeId: normalizeText(payload.nodeId), objectType: payload.objectType, objectInstance: payload.objectInstance, propertyId: payload.propertyId || 85, scale: Number(payload.scale || 1) || 1, - host: payload.host || source.host || "", + host: normalizeText(payload.host || source.host), tableName: source.tableName || "", currentValue: payload.currentValue ?? null, lastReadAt: payload.currentValue === undefined ? null : new Date().toISOString(), @@ -361,6 +385,50 @@ function createScanDatapoint(source, payload, index) { }; } +function buildScannedDatapoints(store, source, rawPoints, toPayload) { + const existingPoints = store.datapoints.filter((item) => item.sourceId === source.id); + const existingById = new Map(existingPoints.map((item) => [String(item.id), item])); + const seenIds = new Set(); + let nextPointIndex = existingPoints.reduce((max, item) => Math.max(max, Number(item.pointIndex) || 0), 0) + 1; + + return rawPoints.reduce((items, rawPoint, index) => { + const payload = toPayload(rawPoint, index); + const stableId = String(payload.id || ""); + if (stableId) { + if (seenIds.has(stableId)) { + return items; + } + seenIds.add(stableId); + } + + const existing = stableId ? existingById.get(stableId) : null; + const pointIndex = existing ? existing.pointIndex : normalizePointIndex(payload.pointIndex, nextPointIndex); + if (!existing) { + nextPointIndex = Math.max(nextPointIndex, pointIndex + 1); + } + const scanned = createScanDatapoint(source, { ...payload, pointIndex }, index); + items.push(!existing ? scanned : { + ...existing, + ...scanned, + pointIndex: existing.pointIndex, + alias: existing.alias || scanned.alias, + name: existing.name || scanned.name, + unit: normalizeUnit(existing.unit || scanned.unit), + kind: existing.kind || scanned.kind, + enabled: existing.enabled === true, + writeMode: normalizeWriteMode(existing.writeMode || scanned.writeMode), + condition: normalizeText(existing.condition || scanned.condition), + dataType: normalizeText(existing.dataType || scanned.dataType), + address: normalizeText(existing.address || scanned.address), + nodeId: normalizeText(existing.nodeId || scanned.nodeId), + scale: Number(existing.scale || scanned.scale || 1) || 1, + host: normalizeText(existing.host || scanned.host), + createdAt: existing.createdAt || scanned.createdAt, + }); + return items; + }, []); +} + 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)); @@ -634,7 +702,7 @@ async function readOpcuaValue(source, datapoint) { return value ? 1 : 0; } const numeric = Number(value); - return Number.isFinite(numeric) ? numeric * (Number(datapoint.scale || 1) || 1) : String(value ?? ""); + return Number.isFinite(numeric) ? numeric * (Number(datapoint.scale || 1) || 1) : String(value || ""); } finally { await client.disconnect().catch(() => undefined); } @@ -797,7 +865,7 @@ async function readBacnetValue(source, datapoint) { return; } const numeric = Number(raw); - resolve(Number.isFinite(numeric) ? numeric * (Number(datapoint.scale || 1) || 1) : String(raw ?? "")); + resolve(Number.isFinite(numeric) ? numeric * (Number(datapoint.scale || 1) || 1) : String(raw || "")); }); }); } @@ -823,7 +891,7 @@ async function readKnxValue(source, datapoint) { reject(error); } else { const numeric = Number(value); - resolve(typeof value === "boolean" ? (value ? 1 : 0) : Number.isFinite(numeric) ? numeric : String(value ?? "")); + resolve(typeof value === "boolean" ? (value ? 1 : 0) : Number.isFinite(numeric) ? numeric : String(value || "")); } }; @@ -970,12 +1038,251 @@ async function ensureTrendTable(tableName) { return table; } +async function ensureSourceRegistryTable() { + await db.query(` + CREATE TABLE IF NOT EXISTS \`source_registry\` ( + \`source_id\` VARCHAR(191) NOT NULL, + \`table_name\` VARCHAR(191) NOT NULL, + \`source_name\` VARCHAR(255) NOT NULL, + \`display_name\` VARCHAR(255) NOT NULL, + \`protocol\` VARCHAR(32) NOT NULL, + \`created_at\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + \`updated_at\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, + PRIMARY KEY (\`source_id\`), + KEY \`idx_source_registry_table_name\` (\`table_name\`) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci + `); +} + +async function syncSourceRegistryEntry(source) { + if (!source?.id || !source?.tableName) { + return; + } + + await ensureSourceRegistryTable(); + await db.query( + ` + INSERT INTO \`source_registry\` ( + \`source_id\`, + \`table_name\`, + \`source_name\`, + \`display_name\`, + \`protocol\` + ) VALUES (?, ?, ?, ?, ?) + ON DUPLICATE KEY UPDATE + \`table_name\` = VALUES(\`table_name\`), + \`source_name\` = VALUES(\`source_name\`), + \`display_name\` = VALUES(\`display_name\`), + \`protocol\` = VALUES(\`protocol\`) + `, + [ + String(source.id), + sanitizeTrendTable(source.tableName), + String(source.name || source.displayName || source.id), + String(source.displayName || source.name || source.id), + String(source.protocol || "unknown"), + ] + ); +} + +async function deleteSourceRegistryEntry(sourceId) { + if (!sourceId) { + return; + } + + await ensureSourceRegistryTable(); + await db.query("DELETE FROM `source_registry` WHERE `source_id` = ?", [String(sourceId)]); +} + +async function ensureSourceDatapointRegistryTable() { + await db.query(` + CREATE TABLE IF NOT EXISTS \`source_datapoint_registry\` ( + \`datapoint_id\` VARCHAR(255) NOT NULL, + \`source_id\` VARCHAR(191) NOT NULL, + \`table_name\` VARCHAR(191) NOT NULL, + \`point_index\` INT NOT NULL, + \`alias\` VARCHAR(255) NOT NULL, + \`name\` VARCHAR(255) NOT NULL, + \`protocol\` VARCHAR(32) NOT NULL, + \`unit\` VARCHAR(64) NULL, + \`kind\` VARCHAR(32) 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_datapoint_table\` (\`table_name\`) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci + `); +} + +async function syncDatapointRegistryEntries(datapoints) { + if (!Array.isArray(datapoints) || !datapoints.length) { + return; + } + + await ensureSourceDatapointRegistryTable(); + for (const datapoint of datapoints) { + if (!datapoint?.id || !datapoint?.sourceId || !datapoint?.tableName) { + continue; + } + + 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\` ( + \`datapoint_id\`, + \`source_id\`, + \`table_name\`, + \`point_index\`, + \`alias\`, + \`name\`, + \`protocol\`, + \`unit\`, + \`kind\` + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) + `, + [ + String(datapoint.id), + String(datapoint.sourceId), + sanitizeTrendTable(datapoint.tableName), + pointIndex, + normalizeText(datapoint.alias || datapoint.name || `Wert${datapoint.pointIndex}`), + normalizeText(datapoint.name || datapoint.alias || `Wert${datapoint.pointIndex}`), + String(datapoint.protocol || "unknown"), + normalizeUnit(datapoint.unit || ""), + String(datapoint.kind || "analog"), + ] + ); + } +} + +async function deleteDatapointRegistryForSource(sourceId) { + if (!sourceId) { + return; + } + + await ensureSourceDatapointRegistryTable(); + await db.query("DELETE FROM `source_datapoint_registry` WHERE `source_id` = ?", [String(sourceId)]); +} + +async function rebuildRegistryTablesFromStore() { + const store = readStore(); + await ensureSourceRegistryTable(); + await ensureSourceDatapointRegistryTable(); + await db.query("DELETE FROM `source_datapoint_registry`"); + for (const source of store.sources || []) { + await syncSourceRegistryEntry(source); + } + await syncDatapointRegistryEntries((store.datapoints || []).filter((datapoint) => datapoint.enabled === true)); +} + + +async function getTrendTableInfo(tableName) { + const table = sanitizeTrendTable(tableName); + const [columns] = await db.query( + ` + SELECT COLUMN_NAME AS columnName + FROM information_schema.COLUMNS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = ? + ORDER BY ORDINAL_POSITION + `, + [table] + ); + const names = columns.map((column) => column.columnName); + const lower = new Set(names.map((name) => String(name).toLowerCase())); + const wideColumnCount = names.filter((name) => /^Wert\d+$/i.test(name)).length; + + if (wideColumnCount > 0) { + return { table, mode: "wide", exists: true }; + } + + if (lower.has("source_id") && lower.has("point_index") && lower.has("value") && lower.has("datum")) { + return { table, mode: "narrow", exists: true }; + } + + return { table, mode: null, exists: names.length > 0 }; +} + +async function ensureSourceTrendTable(source) { + const info = await getTrendTableInfo(source.tableName); + if (info.mode === "wide") { + return info; + } + if (info.mode === "narrow") { + await syncSourceRegistryEntry(source); + return info; + } + if (info.exists) { + throw new Error(`Trendtabelle ${info.table} hat ein unbekanntes Format.`); + } + + await db.query(` + CREATE TABLE IF NOT EXISTS \`${info.table}\` ( + \`id\` BIGINT NOT NULL AUTO_INCREMENT, + \`datum\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + \`userlevel\` INT(1) 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_${info.table}_datum\` (\`datum\`), + KEY \`idx_${info.table}_point_datum\` (\`point_index\`, \`datum\`), + KEY \`idx_${info.table}_source_point\` (\`source_id\`, \`point_index\`) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci + `); + await syncSourceRegistryEntry(source); + return { ...info, mode: "narrow", exists: true }; +} + async function writeTrendRow(source, readings) { if (!readings.length) { return false; } - const table = await ensureTrendTable(source.tableName); + const tableInfo = await ensureSourceTrendTable(source); + const table = tableInfo.table; + + if (tableInfo.mode === "narrow") { + const statement = ` + INSERT INTO \`${table}\` ( + \`userlevel\`, + \`source_id\`, + \`point_index\`, + \`alias\`, + \`value\` + ) VALUES (?, ?, ?, ?, ?) + `; + + for (const reading of readings) { + const pointIndex = normalizePointIndex(reading.pointIndex, 0); + if (!pointIndex) { + continue; + } + + const datapoint = source.__activeDatapoints?.find((item) => item.pointIndex === pointIndex); + await db.query(statement, [ + 0, + String(source.id), + pointIndex, + String(datapoint?.alias || `Wert${pointIndex}`), + String(reading.value), + ]); + } + return true; + } + const columns = ["`userlevel`"]; const placeholders = ["?"]; const values = [0]; @@ -1057,7 +1364,7 @@ async function pollModbusSource(source, datapoints) { } async function pollSource(source, datapoints) { - const sourcePoints = datapoints.filter((datapoint) => datapoint.sourceId === source.id && datapoint.enabled !== false); + const sourcePoints = datapoints.filter((datapoint) => datapoint.sourceId === source.id && datapoint.enabled === true); if (!sourcePoints.length) { updateSourceState(source.id, { lastPollAt: new Date().toISOString(), @@ -1068,6 +1375,7 @@ async function pollSource(source, datapoints) { } try { + source.__activeDatapoints = sourcePoints; if (source.protocol === "modbus-tcp") { await pollModbusSource(source, sourcePoints); return; @@ -1091,7 +1399,7 @@ async function pollSource(source, datapoints) { updateSourceState(source.id, { lastPollAt: new Date().toISOString(), lastPollStatus: "unsupported", - lastPollError: `${source.protocol} wird vom Collector nicht unterstützt.`, + lastPollError: `${source.protocol} wird vom Collector nicht unterst?tzt.`, }); } catch (error) { updateSourceState(source.id, { @@ -1099,6 +1407,8 @@ async function pollSource(source, datapoints) { lastPollStatus: "error", lastPollError: error.message || "Polling fehlgeschlagen.", }); + } finally { + delete source.__activeDatapoints; } } @@ -1185,14 +1495,16 @@ const server = createServer(async (req, res) => { return; } - if (req.method === "POST" && url.pathname === "/sources") { + 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)); + const createdSource = next.sources.at(-1); + await ensureSourceTrendTable(createdSource); + replyJson(res, 201, createdSource); } catch { replyJson(res, 400, { error: "Ungültige JSON-Daten." }); } @@ -1211,7 +1523,9 @@ const server = createServer(async (req, res) => { ensureSource(store, { ...existing, ...body, id: sourceId }); return store; }); - replyJson(res, 200, next.sources.find((item) => item.id === sourceId)); + const updatedSource = next.sources.find((item) => item.id === sourceId); + await ensureSourceTrendTable(updatedSource); + replyJson(res, 200, updatedSource); } catch (error) { replyJson(res, error.message === "not-found" ? 404 : 400, { error: error.message === "not-found" ? "Quelle nicht gefunden." : "Ungültige JSON-Daten.", @@ -1227,6 +1541,8 @@ const server = createServer(async (req, res) => { sources: store.sources.filter((item) => item.id !== sourceId), datapoints: store.datapoints.filter((item) => item.sourceId !== sourceId), })); + await deleteSourceRegistryEntry(sourceId); + await deleteDatapointRegistryForSource(sourceId); replyJson(res, 200, { ok: true, sources: next.sources.length, datapoints: next.datapoints.length }); return; } @@ -1241,6 +1557,7 @@ const server = createServer(async (req, res) => { return store; }); const created = next.datapoints.find((item) => item.sourceId === body.sourceId && item.pointIndex === normalizePointIndex(body.pointIndex, 1)); + await syncDatapointRegistryEntries(created ? [created] : []); replyJson(res, 201, { ok: true, datapoint: created || null, updatedAt: next.updatedAt }); } catch (error) { replyJson(res, 400, { error: error.message || "Datenpunkt konnte nicht angelegt werden." }); @@ -1250,7 +1567,7 @@ const server = createServer(async (req, res) => { if (req.method === "PATCH" && url.pathname.startsWith("/datapoints/")) { try { - const datapointId = url.pathname.split("/").at(-1); + const datapointId = decodeURIComponent(url.pathname.split("/").at(-1) || ""); const body = await readBody(req); const next = touchStore((store) => { const index = store.datapoints.findIndex((item) => item.id === datapointId); @@ -1270,7 +1587,9 @@ const server = createServer(async (req, res) => { }; return store; }); - replyJson(res, 200, { datapoint: next.datapoints.find((item) => item.id === datapointId), updatedAt: next.updatedAt }); + const updatedDatapoint = next.datapoints.find((item) => item.id === datapointId); + await syncDatapointRegistryEntries(updatedDatapoint ? [updatedDatapoint] : []); + replyJson(res, 200, { datapoint: updatedDatapoint, updatedAt: next.updatedAt }); } catch (error) { replyJson(res, error.message === "not-found" ? 404 : 400, { error: error.message === "not-found" ? "Datenpunkt nicht gefunden." : error.message }); } @@ -1279,7 +1598,7 @@ const server = createServer(async (req, res) => { if (req.method === "POST" && url.pathname.match(/^\/datapoints\/[^/]+\/read$/)) { try { - const datapointId = url.pathname.split("/").at(-2); + const datapointId = decodeURIComponent(url.pathname.split("/").at(-2) || ""); const store = readStore(); const datapoint = store.datapoints.find((item) => item.id === datapointId); const source = datapoint ? store.sources.find((item) => item.id === datapoint.sourceId) : null; @@ -1299,7 +1618,6 @@ const server = createServer(async (req, res) => { return; } - if (req.method === "POST" && url.pathname === "/imports/modbus-csv") { try { const body = await readBody(req); @@ -1311,6 +1629,7 @@ const server = createServer(async (req, res) => { upsertDatapoints(store, imported); return store; }); + await syncDatapointRegistryEntries(imported); replyJson(res, 200, { imported: imported.length, datapoints: imported, updatedAt: next.updatedAt }); } catch { replyJson(res, 400, { error: "CSV-Import fehlgeschlagen." }); @@ -1329,6 +1648,7 @@ const server = createServer(async (req, res) => { upsertDatapoints(store, imported); return store; }); + await syncDatapointRegistryEntries(imported); replyJson(res, 200, { imported: imported.length, datapoints: imported, updatedAt: next.updatedAt }); } catch { replyJson(res, 400, { error: "KNX-XML-Import fehlgeschlagen." }); @@ -1344,12 +1664,23 @@ const server = createServer(async (req, res) => { return; } const result = await scanOpcuaEndpoint(body); + let scannedCandidates = []; touchStore((store) => { if (body.sourceId) { const source = store.sources.find((item) => item.id === body.sourceId); if (source) { - const existingCount = store.datapoints.filter((item) => item.sourceId === source.id).length; - const candidates = (result.nodes || []).filter((node) => node.nodeId && node.hasValue).map((node, index) => createScanDatapoint(source, { ...node, protocol: "opc-ua", pointIndex: existingCount + index + 1, alias: node.path || node.displayName, name: node.displayName, address: node.nodeId, nodeId: node.nodeId }, index)); + const rawNodes = (result.nodes || []).filter((node) => node.nodeId); + const candidates = buildScannedDatapoints(store, source, rawNodes, (node, index) => ({ + ...node, + id: `${source.id}::${node.nodeId}`, + protocol: "opc-ua", + pointIndex: index + 1, + alias: node.displayName || node.browseName || node.nodeId, + name: node.displayName || node.browseName || node.nodeId, + address: node.nodeId, + nodeId: node.nodeId, + })); + scannedCandidates = candidates; upsertDatapoints(store, candidates); } } @@ -1357,6 +1688,7 @@ const server = createServer(async (req, res) => { store.scans.opcua = store.scans.opcua.slice(0, 20); return store; }); + await syncDatapointRegistryEntries(scannedCandidates); replyJson(res, 200, result); } catch (error) { replyJson(res, 400, { error: error.message || "OPC-UA-Scan fehlgeschlagen." }); @@ -1368,12 +1700,21 @@ const server = createServer(async (req, res) => { try { const body = await readBody(req); const payload = await scanBacnetNetwork(body); + let scannedCandidates = []; touchStore((store) => { if (body.sourceId) { const source = store.sources.find((item) => item.id === body.sourceId); if (source) { - const existingCount = store.datapoints.filter((item) => item.sourceId === source.id).length; - const candidates = (payload.points || []).map((point, index) => createScanDatapoint(source, { ...point, protocol: "bacnet-ip", pointIndex: existingCount + index + 1 }, index)); + const rawPoints = payload.points || []; + const candidates = buildScannedDatapoints(store, source, rawPoints, (point, index) => ({ + ...point, + id: point.id || `${source.id}::${point.host || source.host}::${point.dataType || point.objectType || "object"}::${point.objectInstance || index + 1}`, + protocol: "bacnet-ip", + pointIndex: index + 1, + alias: point.alias || point.name || `${point.dataType || "obj"}${point.objectInstance || index + 1}`, + name: point.name || point.alias || `${point.dataType || "obj"}${point.objectInstance || index + 1}`, + })); + scannedCandidates = candidates; upsertDatapoints(store, candidates); } } @@ -1381,6 +1722,7 @@ const server = createServer(async (req, res) => { store.scans.bacnet = store.scans.bacnet.slice(0, 20); return store; }); + await syncDatapointRegistryEntries(scannedCandidates); replyJson(res, 200, payload); } catch (error) { replyJson(res, 400, { error: error.message || "BACnet-Scan fehlgeschlagen." }); @@ -1391,6 +1733,10 @@ const server = createServer(async (req, res) => { replyJson(res, 404, { error: "Nicht gefunden." }); }); +rebuildRegistryTablesFromStore().catch((error) => console.error("Registry sync failed", error)); + server.listen(port, () => { console.log(`SE Local Trenddata collector listening on port ${port}`); }); + + diff --git a/docker-compose.yml b/docker-compose.yml index 0e8dec9..3dd6e1d 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -5,9 +5,10 @@ services: restart: unless-stopped environment: MARIADB_ROOT_PASSWORD: ${MARIADB_ROOT_PASSWORD:-SE3112} + MARIADB_ROOT_HOST: mariadb MARIADB_DATABASE: wago - MARIADB_USER: app - MARIADB_PASSWORD: app + MARIADB_USER: wago + MARIADB_PASSWORD: ${MARIADB_PASSWORD:-SE3112} ports: - "3306:3306" volumes: @@ -105,9 +106,28 @@ services: volumes: - /var/run/docker.sock:/var/run/docker.sock - dockhand_data:/app/data + cadvisor: + image: ghcr.io/google/cadvisor:v0.60.2 + container_name: cadvisor + restart: unless-stopped + privileged: true + + ports: + - "9090:8080" + + volumes: + - /:/rootfs:ro + - /var/run:/var/run:ro + - /sys:/sys:ro + - /var/lib/docker/:/var/lib/docker:ro + - /dev/disk/:/dev/disk:ro + + devices: + - /dev/kmsg:/dev/kmsg volumes: mariadb_data: redis_data: collector_data: dockhand_data: + diff --git a/docker/mariadb/init/01-init.sql b/docker/mariadb/init/01-init.sql index 12b93b6..0f48cf0 100644 --- a/docker/mariadb/init/01-init.sql +++ b/docker/mariadb/init/01-init.sql @@ -1,4 +1,4 @@ -CREATE DATABASE IF NOT EXISTS `wago` CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci; +CREATE DATABASE IF NOT EXISTS `wago` CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci; USE `wago`; DELIMITER // @@ -55,3 +55,9 @@ CREATE TABLE IF NOT EXISTS `app_users` ( UNIQUE KEY `uniq_app_users_username` (`username`), UNIQUE KEY `uniq_app_users_email` (`email`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; + + +CREATE USER IF NOT EXISTS 'root'@'%' IDENTIFIED BY 'SE3112'; +GRANT ALL PRIVILEGES ON *.* TO 'root'@'%' WITH GRANT OPTION; +FLUSH PRIVILEGES; + diff --git a/web/src/App.jsx b/web/src/App.jsx index ed37921..94e9abb 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"; @@ -13,6 +13,7 @@ const views = [ const adminSections = [ { id: "aliases", label: "Aliase" }, + { id: "units", label: "Einheiten" }, { id: "sources", label: "Quellen" }, { id: "users", label: "Benutzer" }, { id: "roles", label: "Rollen" }, @@ -65,7 +66,15 @@ const pointColors = [ "#0891b2", ]; -const UNIT_PRESETS = ["°C", "%", "K", "V", "A", "Pa"]; +const DEFAULT_UNIT_OPTIONS = [ + { symbol: "\u00B0C", label: "Grad Celsius", isSystem: true }, + { symbol: "%", label: "Prozent", isSystem: true }, + { symbol: "K", label: "Kelvin", isSystem: true }, + { symbol: "V", label: "Volt", isSystem: true }, + { symbol: "A", label: "Ampere", isSystem: true }, + { symbol: "Pa", label: "Pascal", isSystem: true }, +]; +const UNIT_PRESETS = DEFAULT_UNIT_OPTIONS.map((unit) => unit.symbol); const CUSTOM_UNIT_VALUE = "__custom__"; const PROTOCOL_DEFAULT_PORTS = { "modbus-tcp": "502", @@ -120,6 +129,18 @@ function getDefaultPointDataType(protocol) { return "holding-register"; } +function getOpcuaScanDatapointId(sourceId, node) { + return sourceId && node?.nodeId ? `${sourceId}::${node.nodeId}` : ""; +} + +function getBacnetScanDatapointId(sourceId, sourceHost, point, index = 0) { + if (!sourceId || !point) { + return ""; + } + + return point.id || `${sourceId}::${point.host || sourceHost || ""}::${point.dataType || point.objectType || "object"}::${point.objectInstance || index + 1}`; +} + function getManualAddressLabel(protocol) { if (protocol === "knx-ip") { return "Gruppenadresse"; @@ -137,6 +158,43 @@ function formatValue(value) { return Number.isFinite(Number(value)) ? Number(value).toFixed(2) : "--"; } +function repairUiText(value) { + let current = String(value === undefined || value === null ? "" : value); + for (let attempt = 0; attempt < 3; attempt += 1) { + const currentScore = (current.match(/[\u00C3\u00C2\uFFFD]/g) || []).length; + if (!currentScore) { + break; + } + try { + const bytes = Uint8Array.from(current, (char) => char.charCodeAt(0) & 0xff); + const repaired = new TextDecoder("utf-8").decode(bytes); + const repairedScore = (repaired.match(/[\u00C3\u00C2\uFFFD]/g) || []).length; + if (repairedScore < currentScore) { + current = repaired; + continue; + } + } catch { + break; + } + break; + } + return current.replace(/\u00A0/g, " ").trim(); +} + +function normalizeDisplayText(value, fallback = "") { + const nextValue = repairUiText(value); + return nextValue || fallback; +} + +function normalizeDisplayUnit(value) { + const nextValue = normalizeDisplayText(value, "\u00B0C").replace(/°/g, "°"); + const collapsed = nextValue.replace(/\s+/g, ""); + if (!collapsed || /^\?C$/i.test(collapsed) || /^Â?°C$/i.test(collapsed) || (/[ÃÂ�]/.test(nextValue) && /(?:°|\?|C)/i.test(nextValue))) { + return "°C"; + } + return nextValue; +} + function formatCollectorValue(value) { if (value === null || value === undefined || value === "") { return "kein Wert"; @@ -144,7 +202,7 @@ function formatCollectorValue(value) { if (typeof value === "boolean") { return value ? "true" : "false"; } - return Number.isFinite(Number(value)) ? Number(value).toLocaleString("de-DE") : String(value); + return Number.isFinite(Number(value)) ? Number(value).toLocaleString("de-DE") : normalizeDisplayText(value, "--"); } function formatTimestampLabel(timestamp) { @@ -246,16 +304,30 @@ async function readAliasCsv(file) { const row = splitCsvLine(line); const item = {}; header.forEach((key, index) => { - item[key] = row[index] ?? ""; + item[key] = row[index] === undefined || row[index] === null ? "" : row[index]; }); return item; }); } +function getOrderedPointEntries(metadata) { + if (!metadata?.points) { + return []; + } + + const order = Array.isArray(metadata.pointOrder) && metadata.pointOrder.length + ? metadata.pointOrder + : Object.keys(metadata.points).sort((left, right) => Number(left) - Number(right)); + + return order + .filter((pointIndex) => metadata.points[String(pointIndex)]) + .map((pointIndex) => ({ pointIndex: Number(pointIndex), ...metadata.points[String(pointIndex)] })); +} + function buildAliasCsv(aliasEditor) { const header = ["pointIndex", "alias", "unit", "factor", "kind", "min", "max", "color"]; - const rows = Object.entries(aliasEditor.points).map(([pointIndex, point]) => [ - pointIndex, + const rows = getOrderedPointEntries(aliasEditor).map((point) => [ + point.pointIndex, point.alias, point.unit, point.factor, @@ -267,10 +339,19 @@ function buildAliasCsv(aliasEditor) { return [header] .concat(rows) - .map((line) => line.map((cell) => `"${String(cell ?? "").replaceAll('"', '""')}"`).join(";")) + .map((line) => line.map((cell) => `"${String(cell || "").replaceAll('"', '""')}"`).join(";")) .join("\n"); } +function matchesSearch(filterText, values) { + const search = String(filterText || "").trim().toLowerCase(); + if (!search) { + return true; + } + + return values.some((value) => normalizeDisplayText(value).toLowerCase().includes(search)); +} + function hashString(value) { return String(value) .split("") @@ -312,9 +393,16 @@ function SmallFootnote() { return
by J.H. | SE Inno GmbH
; } -function UnitField({ value, onChange }) { - const normalizedValue = value || "°C"; - const selectedValue = UNIT_PRESETS.includes(normalizedValue) ? normalizedValue : CUSTOM_UNIT_VALUE; +function UnitField({ units, value, onChange }) { + const options = units?.length ? units : DEFAULT_UNIT_OPTIONS; + const normalizedOptions = options.map((unit) => ({ + ...unit, + symbol: normalizeDisplayUnit(unit.symbol), + label: normalizeDisplayText(unit.label || unit.symbol, normalizeDisplayUnit(unit.symbol)), + })); + const normalizedValue = normalizeDisplayUnit(value || "\u00B0C"); + const hasPresetMatch = normalizedOptions.some((unit) => unit.symbol === normalizedValue); + const selectedValue = hasPresetMatch ? normalizedValue : CUSTOM_UNIT_VALUE; return (
@@ -323,13 +411,13 @@ function UnitField({ value, onChange }) { onChange={(event) => { const nextValue = event.target.value; if (nextValue === CUSTOM_UNIT_VALUE) { - onChange(UNIT_PRESETS.includes(normalizedValue) ? "" : normalizedValue); + onChange(hasPresetMatch ? "" : normalizedValue); return; } onChange(nextValue); }} > - {UNIT_PRESETS.map((unit) => )} + {normalizedOptions.map((unit) => )} {selectedValue === CUSTOM_UNIT_VALUE ? ( @@ -451,13 +539,12 @@ function PointLibrary({ metadata, filterText, onFilterTextChange, onAddPoint }) } const search = filterText.trim().toLowerCase(); - return Object.entries(metadata.points) - .map(([pointIndex, point]) => ({ pointIndex: Number(pointIndex), ...point })) + return getOrderedPointEntries(metadata) .filter((point) => { if (!search) { return true; } - const haystack = `${point.alias} ${point.unit} wert ${point.pointIndex}`.toLowerCase(); + const haystack = `${normalizeDisplayText(point.alias)} ${normalizeDisplayUnit(point.unit)} wert ${point.pointIndex}`.toLowerCase(); return haystack.includes(search); }); }, [filterText, metadata]); @@ -467,7 +554,7 @@ function PointLibrary({ metadata, filterText, onFilterTextChange, onAddPoint })

Datenpunkte

-

{metadata?.displayName || "ISP"}

+

{normalizeDisplayText(metadata?.displayName, "ISP")}

( @@ -511,6 +598,7 @@ export default function App() { const [user, setUser] = useState(null); const [loginError, setLoginError] = useState(""); const [bootError, setBootError] = useState(""); + const [collectorError, setCollectorError] = useState(""); const [successMessage, setSuccessMessage] = useState(""); const [busy, setBusy] = useState(false); const [sessionReady, setSessionReady] = useState(false); @@ -546,10 +634,15 @@ export default function App() { const [collectorCapabilities, setCollectorCapabilities] = useState(null); const [collectorSources, setCollectorSources] = useState([]); const [collectorDatapoints, setCollectorDatapoints] = useState([]); + const [units, setUnits] = useState(DEFAULT_UNIT_OPTIONS); + const [unitDraft, setUnitDraft] = useState({ symbol: "", label: "" }); + const [editingUnitSymbol, setEditingUnitSymbol] = useState(""); const [selectedSourceId, setSelectedSourceId] = useState(""); const [isCreatingSource, setIsCreatingSource] = useState(false); const [sourceDraft, setSourceDraft] = useState(defaultSourceDraft); + const [sourceModalOpen, setSourceModalOpen] = useState(false); const emptyManualDatapointDraft = (protocol = defaultSourceDraft.protocol) => ({ alias: "", address: "", pointIndex: "", unit: "°C", dataType: getDefaultPointDataType(protocol), kind: "analog", scale: "1", writeMode: "", condition: "" }); + const [collectorPointFilter, setCollectorPointFilter] = useState(""); const [manualDatapointDraft, setManualDatapointDraft] = useState(emptyManualDatapointDraft()); const [scanResults, setScanResults] = useState({ opcua: null, bacnet: null }); const [debugSnapshot, setDebugSnapshot] = useState(null); @@ -573,15 +666,60 @@ export default function App() { { name: "technician", label: "Techniker" }, { name: "admin", label: "Admin" }, ]; + const availableUnits = useMemo(() => (units.length ? units : DEFAULT_UNIT_OPTIONS).map((unit) => ({ + ...unit, + symbol: normalizeDisplayUnit(unit.symbol), + label: normalizeDisplayText(unit.label || unit.symbol, normalizeDisplayUnit(unit.symbol)), + })), [units]); const currentSource = useMemo(() => collectorSources.find((source) => source.id === selectedSourceId) || null, [collectorSources, selectedSourceId]); const collectorSourcePoints = useMemo(() => ( selectedSourceId ? collectorDatapoints.filter((item) => item.sourceId === selectedSourceId) : collectorDatapoints ), [collectorDatapoints, selectedSourceId]); + const collectorSourcePointById = useMemo(() => new Map(collectorSourcePoints.map((item) => [String(item.id), item])), [collectorSourcePoints]); + const availableCollectorPoints = useMemo(() => collectorSourcePoints.filter((item) => item.enabled !== true), [collectorSourcePoints]); + const loggedCollectorPoints = useMemo(() => collectorSourcePoints.filter((item) => item.enabled === true), [collectorSourcePoints]); + const loggedCollectorPointIds = useMemo(() => new Set(loggedCollectorPoints.map((item) => String(item.id))), [loggedCollectorPoints]); + const filteredOpcuaNodes = useMemo( + () => (scanResults.opcua?.nodes || []).filter((node) => matchesSearch(collectorPointFilter, [ + node.displayName, + node.browseName, + node.nodeId, + node.path, + node.currentValue, + ])), + [collectorPointFilter, scanResults.opcua] + ); + const filteredBacnetPoints = useMemo( + () => (scanResults.bacnet?.points || []).filter((point) => matchesSearch(collectorPointFilter, [ + point.alias, + point.name, + point.id, + point.host, + point.dataType, + point.objectInstance, + point.currentValue, + ])), + [collectorPointFilter, scanResults.bacnet] + ); + const filteredAvailableCollectorPoints = useMemo( + () => availableCollectorPoints.filter((item) => matchesSearch(collectorPointFilter, [ + item.alias, + item.name, + item.address, + item.nodeId, + item.dataType, + item.protocol, + item.pointIndex, + item.currentValue, + ])), + [availableCollectorPoints, collectorPointFilter] + ); const resetFlash = () => { setBootError(""); + setCollectorError(""); setSuccessMessage(""); }; @@ -632,6 +770,7 @@ export default function App() { setIsps(ispResult.isps); setSelections(selectionResult.selections); setDashboard(dashboardResult); + setActiveView("dashboard"); setPreferencesReady(true); } catch (error) { if (!ignore) { @@ -691,6 +830,11 @@ export default function App() { } }, [collectorSources, isCreatingSource, selectedSourceId]); + useEffect(() => { + setScanResults({ opcua: null, bacnet: null }); + setCollectorPointFilter(""); + }, [selectedSourceId]); + useEffect(() => { if (!token || !preferencesReady) { return; @@ -785,6 +929,17 @@ export default function App() { .catch((error) => setBootError(error.message)); }, [canManageRoles, token]); + useEffect(() => { + if (!token || !(canEditAliases || canManageSources)) { + setUnits(DEFAULT_UNIT_OPTIONS); + return; + } + + api.getUnits(token) + .then((result) => setUnits(result.units || DEFAULT_UNIT_OPTIONS)) + .catch((error) => setBootError(error.message)); + }, [canEditAliases, canManageSources, token]); + useEffect(() => { if (!token || !(canEditAliases || canManageSources || canViewDebug)) { setCollectorCapabilities(null); @@ -794,10 +949,10 @@ export default function App() { } api.getCollectorCapabilities() - .then((result) => setCollectorCapabilities(result)) - .catch((error) => setBootError(error.message)); + .then((result) => { setCollectorCapabilities(result); setCollectorError(""); }) + .catch((error) => { setCollectorCapabilities(null); setCollectorError(error.message); }); - refreshCollectorData().catch((error) => setBootError(error.message)); + refreshCollectorData().catch((error) => setCollectorError(error.message)); }, [canEditAliases, canManageSources, canViewDebug, token]); useEffect(() => { @@ -879,7 +1034,7 @@ export default function App() { return; } - refreshCollectorDebug().catch((error) => setBootError(error.message)); + refreshCollectorDebug().catch((error) => setCollectorError(error.message)); }, [adminSection, canViewDebug, token]); const handleLogin = async (credentials) => { @@ -889,6 +1044,8 @@ export default function App() { try { const result = await api.login(credentials); localStorage.setItem("seltd-token", result.token); + setActiveView("dashboard"); + setAdminSection("aliases"); setToken(result.token); setUser(result.user); } catch (error) { @@ -910,6 +1067,8 @@ export default function App() { setLatestValues([]); setWorkingPoints([]); setSelectionDraft(defaultSelectionDraft); + setActiveView("dashboard"); + setAdminSection("aliases"); localStorage.removeItem("seltd-token"); }; @@ -948,6 +1107,10 @@ export default function App() { }; const deleteSelection = async (selectionId) => { + if (!window.confirm("Auswahl wirklich löschen?")) { + return; + } + resetFlash(); try { await api.deleteSelection(token, selectionId); @@ -1032,6 +1195,10 @@ export default function App() { }; const removeWidget = (widgetId) => { + if (!window.confirm("Widget wirklich löschen?")) { + return; + } + setDashboard((current) => ({ ...current, widgets: current.widgets.filter((widget) => widget.id !== widgetId), @@ -1056,11 +1223,13 @@ export default function App() { ]); setCollectorSources(sourceResult.sources || []); setCollectorDatapoints(datapointResult.datapoints || []); + setCollectorError(""); }; const refreshCollectorDebug = async () => { const result = await api.getCollectorDebug(); setDebugSnapshot(result); + setCollectorError(""); }; const updateAliasPoint = (pointIndex, changes) => { @@ -1086,9 +1255,15 @@ export default function App() { const payload = { displayName: aliasEditor.displayName, description: aliasEditor.description, - points: Object.entries(aliasEditor.points).map(([pointIndex, point]) => ({ - pointIndex: Number(pointIndex), - ...point, + points: getOrderedPointEntries(aliasEditor).map((point) => ({ + pointIndex: Number(point.pointIndex), + alias: point.alias, + unit: point.unit, + factor: point.factor, + kind: point.kind, + min: point.min, + max: point.max, + color: point.color, })), }; const result = await api.saveAliases(token, adminIsp, payload); @@ -1162,6 +1337,41 @@ export default function App() { } }; + const resetUnitDraft = () => { + setEditingUnitSymbol(""); + setUnitDraft({ symbol: "", label: "" }); + }; + + const startUnitEdit = (unit) => { + setEditingUnitSymbol(unit.symbol); + setUnitDraft({ symbol: unit.symbol, label: unit.label }); + }; + + const saveUnitDefinition = async () => { + const symbol = normalizeDisplayUnit(unitDraft.symbol || unitDraft.label || "").trim(); + const label = normalizeDisplayText(unitDraft.label || symbol, symbol).trim(); + + if (!symbol) { + setBootError("Bitte ein Einheitssymbol angeben."); + return; + } + + resetFlash(); + try { + if (editingUnitSymbol) { + await api.updateUnit(token, editingUnitSymbol, { symbol, label }); + } else { + await api.createUnit(token, { symbol, label }); + } + const result = await api.getUnits(token); + setUnits(result.units || DEFAULT_UNIT_OPTIONS); + resetUnitDraft(); + setSuccessMessage("Einheit gespeichert."); + } catch (error) { + setBootError(error.message); + } + }; + const importSqlFile = async (event) => { const file = event.target.files?.[0]; if (!file) { @@ -1209,6 +1419,17 @@ export default function App() { ...defaultSourceDraft, port: getDefaultProtocolPort(defaultSourceDraft.protocol), }); + setSourceModalOpen(true); + }; + + const openSourceEditor = (sourceId) => { + setIsCreatingSource(false); + setSelectedSourceId(sourceId); + setSourceModalOpen(true); + }; + + const closeSourceEditor = () => { + setSourceModalOpen(false); }; const updateSourceProtocol = (nextProtocol) => { @@ -1260,6 +1481,7 @@ export default function App() { setIsps(refreshed.isps); setIsCreatingSource(false); setSelectedSourceId(saved.id); + setSourceModalOpen(false); setSuccessMessage("Protokollquelle gespeichert."); } catch (error) { setBootError(error.message); @@ -1271,6 +1493,10 @@ export default function App() { return; } + if (!window.confirm("Quelle wirklich löschen?")) { + return; + } + resetFlash(); try { await api.deleteCollectorSource(selectedSourceId); @@ -1297,7 +1523,7 @@ export default function App() { const saveManualCollectorDatapoint = async () => { const activeSourceId = selectedSourceId || sourceDraft.id; if (!activeSourceId) { - setBootError("Bitte zuerst eine Quelle speichern oder auswählen."); + setBootError("Bitte zuerst eine Quelle auswählen oder anlegen."); return; } @@ -1310,10 +1536,13 @@ export default function App() { name: manualDatapointDraft.alias, address: manualDatapointDraft.address, pointIndex: manualDatapointDraft.pointIndex, - unit: manualDatapointDraft.unit || "°C", + unit: manualDatapointDraft.unit || "\u00B0C", dataType: manualDatapointDraft.dataType, nodeId: sourceDraft.protocol === "opc-ua" ? manualDatapointDraft.address : "", kind: manualDatapointDraft.kind, + scale: manualDatapointDraft.scale, + writeMode: manualDatapointDraft.writeMode, + condition: manualDatapointDraft.condition, }); await refreshCollectorData(); await refreshCollectorDebug().catch(() => undefined); @@ -1332,7 +1561,7 @@ export default function App() { const activeSourceId = selectedSourceId || sourceDraft.id; if (!activeSourceId) { - setBootError("Bitte zuerst eine Quelle speichern oder auswählen."); + setBootError("Bitte zuerst eine Quelle auswählen oder anlegen."); event.target.value = ""; return; } @@ -1364,9 +1593,10 @@ export default function App() { const runOpcuaScan = async () => { const activeSourceId = selectedSourceId || sourceDraft.id; if (!activeSourceId) { - setBootError("Bitte zuerst eine Quelle speichern oder auswählen."); + setBootError("Bitte zuerst eine Quelle auswählen oder anlegen."); return; } + resetFlash(); try { const result = await api.scanOpcua({ @@ -1386,9 +1616,10 @@ export default function App() { const runBacnetScan = async () => { const activeSourceId = selectedSourceId || sourceDraft.id; if (!activeSourceId) { - setBootError("Bitte zuerst eine Quelle speichern oder auswählen."); + setBootError("Bitte zuerst eine Quelle auswählen oder anlegen."); return; } + resetFlash(); try { const hosts = sourceDraft.host @@ -1499,7 +1730,7 @@ export default function App() { } if (passwordDraft.nextPassword !== passwordDraft.repeatPassword) { - setBootError("Die neuen Passwörter stimmen nicht überein."); + setBootError("Die neuen Passw?rter stimmen nicht ?berein."); return; } @@ -1509,7 +1740,7 @@ export default function App() { nextPassword: passwordDraft.nextPassword, }); setPasswordDraft({ currentPassword: "", nextPassword: "", repeatPassword: "" }); - setSuccessMessage("Passwort geändert."); + setSuccessMessage("Passwort ge?ndert."); } catch (error) { setBootError(error.message); } @@ -1541,7 +1772,7 @@ export default function App() {
) : null} - setSidebarPinned((current) => !current)}> + setSidebarPinned((current) => !current)}> @@ -1598,6 +1829,7 @@ export default function App() { {bootError ?
{bootError}
: null} + {!bootError && collectorError && activeView === "admin" && (adminSection === "sources" || adminSection === "debug") ?
{collectorError}
: null} {successMessage ?
{successMessage}
: null} {activeView === "dashboard" ? ( @@ -1657,7 +1889,7 @@ export default function App() { ) : null} - +
@@ -1716,20 +1948,20 @@ export default function App() { ) : (
- {value.displayName} - {value.alias} + {normalizeDisplayText(value.displayName)} + {normalizeDisplayText(value.alias)} {formatTimestampLabel(value.timestamp)}
-
{formatValue(value.value)}{value.unit}
+
{formatValue(value.value)}{normalizeDisplayUnit(value.unit)}
))}
@@ -1754,7 +1986,7 @@ export default function App() { @@ -1781,8 +2013,8 @@ export default function App() { onClick={() => removePoint(point)} style={{ borderLeftColor: color, background: `${color}12` }} > - {metadataCache[point.isp]?.displayName || point.isp} - {meta?.alias || `Wert ${point.pointIndex}`} + {normalizeDisplayText(metadataCache[point.isp]?.displayName, point.isp)} + {normalizeDisplayText(meta?.alias, `Wert ${point.pointIndex}`)} ); })} @@ -1875,11 +2107,11 @@ export default function App() { {latestValues.map((value) => (
- {value.displayName} - {value.alias} + {normalizeDisplayText(value.displayName)} + {normalizeDisplayText(value.alias)} {formatTimestampLabel(value.timestamp)}
-
{value.value !== null ? formatValue(value.value) : value.raw || "--"}{value.unit}
+
{value.value !== null ? formatValue(value.value) : normalizeDisplayText(value.raw, "--")}{normalizeDisplayUnit(value.unit)}
))} @@ -1894,6 +2126,9 @@ export default function App() { if (section.id === "aliases") { return canEditAliases; } + if (section.id === "units") { + return canEditAliases; + } if (section.id === "sources") { return canManageSources || canEditAliases; } @@ -1939,7 +2174,7 @@ export default function App() { @@ -1965,17 +2200,17 @@ export default function App() {
Farbe
- {Object.entries(aliasEditor.points).map(([pointIndex, point]) => ( -
-
Wert {pointIndex}
- updateAliasPoint(pointIndex, { alias: event.target.value })} /> - updateAliasPoint(pointIndex, { unit: nextValue })} /> - updateAliasPoint(pointIndex, { factor: event.target.value })} /> - updateAliasPoint(String(point.pointIndex), { alias: event.target.value })} /> + updateAliasPoint(String(point.pointIndex), { unit: nextValue })} /> + updateAliasPoint(String(point.pointIndex), { factor: event.target.value })} /> + - updateAliasPoint(pointIndex, { color: nextValue })} /> + updateAliasPoint(String(point.pointIndex), { color: nextValue })} />
))}
@@ -1985,23 +2220,67 @@ export default function App() { ) : null} - {adminSection === "sources" && (canManageSources || canEditAliases) ? ( + {adminSection === "units" && canEditAliases ? (
-

Quellenverwaltung

-

Trendtabellen, SQL-Importe und Protokollquellen an einer Stelle verwalten.

+

Einheiten

+

Einheiten zentral verwalten und danach bei Aliasen oder Quellen auswählen.

+
+
+
-
-
+ ) : null} -
+ {adminSection === "sources" && (canManageSources || canEditAliases) ? ( +
+
+
+

Quellenverwaltung

+

Quellen oben auswählen, per Modal bearbeiten und unten gefundene Punkte ins Logging übernehmen.

+
+
+ + + {currentSource ? : null} +
+
+ +
+
+ +
-

Protokollquellen

-

{collectorCapabilities?.protocols?.join(" · ") || "Modbus TCP, OPC UA, BACnet IP und KNX IP"}

-
-
- - {selectedSourceId ? : null} +

{collectorCapabilities?.protocols?.join(" · ") || "Modbus TCP · OPC UA · BACnet IP · KNX IP"}

- - -
- - - - - - - - +
+ {collectorSources.map((source) => { + const isActive = source.id === selectedSourceId; + const pointCount = collectorDatapoints.filter((item) => item.sourceId === source.id).length; + return ( + + {source.protocol === "opc-ua" ? : null} + {source.protocol === "bacnet-ip" ? : null} + +
+ + ); + })} + {!collectorSources.length ?

Noch keine Protokollquelle angelegt. Über Neue Quelle startest du sauber mit leerem Formular.

: null}
- Das Poll-Intervall legt fest, wie oft die Quelle abgefragt wird. Mit COV werden Werte nur bei Änderung in die Datenbank geschrieben. - - - -
- {collectorSources.map((source) => ( -
- {collectorDatapoints.filter((item) => item.sourceId === source.id).length} DP - - ))} -
-
-
- Importe -
- - - -
- Zuletzt ausgewählte Quelle: {currentSource?.name || "noch keine"} -
- - {canManageSources ? ( -
- Einzelner Datenpunkt
- - - + + -
- -
- ) : null} -
- OPC UA -
- + Der Tabellenname wird automatisch passend erzeugt. So taucht die Quelle später sicher in der Trendansicht auf. + +
+ + {selectedSourceId ? : null} +
+
+
+ ) : null} + +
+
+
+
+

Gefunden / importiert

+

Links landen Scan-Ergebnisse und importierte Punkte. Von hier schiebst du Werte ins Logging.

+
+
+ + setCollectorPointFilter(event.target.value)} + /> + +
+
+ Importe +
+ + + +
+ Aktive Quelle: {normalizeDisplayText(currentSource?.name, "noch keine ausgewählt")} +
+ + {sourceDraft.protocol === "opc-ua" ? ( +
+ OPC UA +
+ +
+ Verwendet Host und Port der oben ausgewählten Quelle. + {scanResults.opcua ? {scanResults.opcua.reachable ? `${scanResults.opcua.nodes?.length || 0} Node(s), ${scanResults.opcua.endpoints?.length || 0} Endpunkt(e)` : "Host nicht erreichbar"} : null} +
+ ) : null} + + {sourceDraft.protocol === "bacnet-ip" ? ( +
+ BACnet IP +
+ +
+ Für mehrere Hosts im Host-Feld Komma, Leerzeichen oder Semikolon verwenden. + {scanResults.bacnet ? {`${scanResults.bacnet.devices?.length || 0} Gerät(e)`} : null} +
+ ) : null} +
+ +
+ {filteredOpcuaNodes.slice(0, 160).map((node, index) => { + const pointId = getOpcuaScanDatapointId(selectedSourceId, node); + const linkedPoint = collectorSourcePointById.get(pointId); + const isLogged = linkedPoint?.enabled === true || loggedCollectorPointIds.has(pointId); + return ( +
+
+ {normalizeDisplayText(node.displayName || node.browseName || node.nodeId)} + {normalizeDisplayText(node.path || node.nodeId)} + {node.currentValue === null || node.currentValue === undefined ? "kein Live-Wert" : formatCollectorValue(node.currentValue)} +
+ {linkedPoint?.pointIndex || index + 1} +
+ {linkedPoint ? : null} + +
+
+ ); + })} + {filteredBacnetPoints.slice(0, 160).map((point, index) => { + const pointId = getBacnetScanDatapointId(selectedSourceId, currentSource?.host, point, index); + const linkedPoint = collectorSourcePointById.get(pointId); + const isLogged = linkedPoint?.enabled === true || loggedCollectorPointIds.has(pointId); + return ( +
+
+ {normalizeDisplayText(point.alias || point.name)} + {`${normalizeDisplayText(point.host)} · ${normalizeDisplayText(point.dataType || "obj")}${point.objectInstance || ""}`} + {point.currentValue === null || point.currentValue === undefined ? "kein Live-Wert" : formatCollectorValue(point.currentValue)} +
+ {linkedPoint?.pointIndex || index + 1} +
+ {linkedPoint ? : null} + +
+
+ ); + })} + {(!scanResults.opcua && !scanResults.bacnet) && filteredAvailableCollectorPoints.length ? filteredAvailableCollectorPoints.slice(0, 2000).map((item) => ( +
+
+ {normalizeDisplayText(item.alias || item.name)} + {normalizeDisplayText(item.address || item.nodeId || item.dataType || item.protocol)} + {item.currentValue === null || item.currentValue === undefined ? "kein Live-Wert" : formatCollectorValue(item.currentValue)} +
+ {item.pointIndex} +
+ + +
+
+ )) : null} + {!filteredAvailableCollectorPoints.length && !filteredOpcuaNodes.length && !filteredBacnetPoints.length ?

{collectorPointFilter.trim() ? "Keine Treffer für die aktuelle Suche." : "Nach Scan oder Import erscheinen die gefundenen Punkte genau hier links."}

: null}
- Verwendet Host und Port aus der Quelle oberhalb des Speicher-Buttons. - {scanResults.opcua ? {scanResults.opcua.reachable ? `${scanResults.opcua.nodes?.length || 0} Node(s), ${scanResults.opcua.endpoints?.length || 0} Endpunkt(e)` : "Host nicht erreichbar"} : null}
-
- BACnet IP -
- +
+
+
+

Wird geloggt

+

Nur diese rechte Liste wird vom Collector gepollt und in die Trendtabelle geschrieben.

+
+
+ + {canManageSources ? ( +
+ Einzelnen Datenpunkt anlegen +
+ + + + + + + + + +
+ +
+ ) : null} + +
+ {loggedCollectorPoints.length ? loggedCollectorPoints.slice(0, 2000).map((item) => ( +
+
+ {normalizeDisplayText(item.alias || item.name)} + {`Wert ${item.pointIndex} · ${normalizeDisplayText(item.address || item.nodeId || item.dataType || item.protocol)}`} + {`${item.writeMode === "cov" ? "COV" : "Intervall"}${item.condition ? ` · ${normalizeDisplayText(item.condition)}` : ""}`} +
+ {formatCollectorValue(item.currentValue)}{item.unit ? ` ${normalizeDisplayUnit(item.unit)}` : ""} +
+ + +
+
+ )) :

Hier erscheinen nur die Punkte, die wirklich in die DB geschrieben werden.

}
- Für mehrere Hosts einfach im Feld oben Komma, Leerzeichen oder Semikolon verwenden. - {scanResults.bacnet ? {scanResults.bacnet.devices?.length || 0} Gerät(e), {scanResults.bacnet.points?.length || 0} Objekt(e) : null}
- - {(scanResults.opcua?.nodes?.length || scanResults.bacnet?.points?.length) ? ( -
-
-
-

Scan-Ergebnis

-

{scanResults.opcua?.nodes?.length || scanResults.bacnet?.points?.length || 0} gefundene Datenpunkte oder Objekte.

-
-
-
- {(scanResults.opcua?.nodes || []).filter((node) => node.hasValue).slice(0, 120).map((node) => ( -
-
- {node.displayName || node.browseName} - {node.path || node.nodeId} -
- {formatCollectorValue(node.currentValue)} -
- ))} - {(scanResults.bacnet?.points || []).slice(0, 120).map((point) => ( -
-
- {point.alias || point.name} - {point.host} · {point.dataType}{point.objectInstance} -
- {formatCollectorValue(point.currentValue)} -
- ))} -
-
- ) : null} - - {collectorSourcePoints.length ? ( -
-
-
-

Datenpunkte

-

{collectorSourcePoints.length} importierte, gescannte oder manuelle Punkte. Aktivierte Punkte werden geloggt.

-
-
-
- {collectorSourcePoints.slice(0, 2000).map((item) => ( -
-
- {item.alias || item.name} - Wert {item.pointIndex} · {item.address || item.nodeId || item.dataType || item.protocol} · {item.writeMode === "cov" ? "COV" : "Intervall"} - {item.condition ? {item.condition} : null} -
- {formatCollectorValue(item.currentValue)}{item.unit ? " " + item.unit : ""} -
- - -
-
- ))} -
-
- ) : null}
) : null} @@ -2390,21 +2730,24 @@ export default function App() {

Letzte Scans

OPC UA{debugSnapshot?.latestScans?.opcua?.reachable ? `${debugSnapshot.latestScans.opcua.endpoints.length} Endpunkte` : "kein Treffer"}
-
BACnet{debugSnapshot?.latestScans?.bacnet?.devices?.length || 0} Geräte
+
BACnet{debugSnapshot?.latestScans?.bacnet?.devices?.length || 0} Ger?te
- {(debugSnapshot?.sources || []).map((source) => ( -
-
- {source.name} - {source.protocol} · {source.host || "ohne Host"} + {(debugSnapshot?.sources || []).map((source) => { + const pointCount = collectorDatapoints.filter((item) => item.sourceId === source.id).length; + return ( +
+
+ {normalizeDisplayText(source.name, "Quelle")} + {`${source.protocol} · ${normalizeDisplayText(source.tableName)}`} +
+ {`${String(source.pollIntervalSeconds || 60)}s ? ${pointCount} DP`}
- {source.datapointCount} DP -
- ))} + ); + })}
) : null} @@ -2471,3 +2814,6 @@ export default function App() {
); } + + + diff --git a/web/src/api.js b/web/src/api.js index d337a32..6964296 100644 --- a/web/src/api.js +++ b/web/src/api.js @@ -56,6 +56,9 @@ export const api = { getIsps: (token) => apiRequest("/isps", { token }), getAliases: (token, isp) => apiRequest(`/isps/${isp}/aliases`, { token }), saveAliases: (token, isp, body) => apiRequest(`/isps/${isp}/aliases`, { method: "PUT", token, body }), + getUnits: (token) => apiRequest("/units", { token }), + createUnit: (token, body) => apiRequest("/units", { method: "POST", token, body }), + updateUnit: (token, symbol, body) => apiRequest(`/units/${encodeURIComponent(symbol)}`, { method: "PATCH", token, body }), createTrendTable: (token, body) => apiRequest("/trend-tables", { method: "POST", token, body }), normalizeAliasMetadata: (token, body) => apiRequest("/admin/normalize-alias-metadata", { method: "POST", token, body }), importSql: (token, body) => apiRequest("/admin/import-sql", { method: "POST", token, body }), diff --git a/web/src/styles.css b/web/src/styles.css index faaa7e4..e1db813 100644 --- a/web/src/styles.css +++ b/web/src/styles.css @@ -919,3 +919,134 @@ span { font-weight: 700; white-space: nowrap; } + +.source-split { + display: grid; + grid-template-columns: minmax(0, 1fr) minmax(0, 1fr); + gap: 0.65rem; + align-items: start; +} + +.source-split .scan-result-list, +.source-split .datapoint-table { + min-height: 320px; +} + +.source-admin-panel { + position: relative; +} + +.inline-create-table { + align-items: end; +} + +.inline-table-button { + display: flex; + align-items: end; +} + +.source-grid { + align-items: stretch; +} + +.source-tile { + display: grid; + grid-template-columns: minmax(0, 1fr) auto; + gap: 0.55rem; + align-items: start; +} + +.source-tile-copy { + display: grid; + gap: 0.18rem; + min-width: 0; +} + +.source-tile-copy small { + color: var(--muted); + font-size: 0.7rem; +} + +.source-tile-actions { + display: grid; + gap: 0.3rem; + justify-items: end; +} + +.source-toolbar-card { + display: flex; + align-items: center; + justify-content: space-between; + gap: 0.6rem; + border: 1px solid var(--border); + border-radius: 12px; + padding: 0.55rem 0.65rem; + background: color-mix(in srgb, var(--surface-strong) 90%, transparent); +} + +.source-toolbar-card > div:first-child { + display: grid; + gap: 0.14rem; +} + +.source-toolbar-card span { + color: var(--muted); + font-size: 0.74rem; +} + +.source-tools-grid { + grid-template-columns: repeat(3, minmax(0, 1fr)); +} + +.manual-point-card { + margin-bottom: 0.35rem; +} + +.modal-backdrop { + position: fixed; + inset: 0; + z-index: 30; + display: grid; + place-items: center; + padding: 1rem; + background: rgba(15, 23, 42, 0.38); + backdrop-filter: blur(8px); +} + +.modal-card { + width: min(760px, calc(100vw - 2rem)); + max-height: calc(100vh - 2rem); + overflow: auto; + border: 1px solid var(--border); + border-radius: 16px; + padding: 0.85rem; + background: var(--surface); + box-shadow: 0 28px 80px rgba(15, 23, 42, 0.22); +} + +.source-modal { + display: grid; + gap: 0.65rem; +} + +@media (max-width: 980px) { + .source-split, + .source-tools-grid, + .source-tile, + .source-toolbar-card { + grid-template-columns: 1fr; + } + + .source-toolbar-card { + align-items: start; + } + + .source-tile-actions { + width: 100%; + justify-items: stretch; + } + + .source-tile-actions .compact-button { + width: 100%; + } +}