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
{metadata?.displayName || "ISP"}
+{normalizeDisplayText(metadata?.displayName, "ISP")}
Trendtabellen, SQL-Importe und Protokollquellen an einer Stelle verwalten.
+Einheiten zentral verwalten und danach bei Aliasen oder Quellen auswählen.
+Quellen oben auswählen, per Modal bearbeiten und unten gefundene Punkte ins Logging übernehmen.
+{collectorCapabilities?.protocols?.join(" · ") || "Modbus TCP, OPC UA, BACnet IP und KNX IP"}
-{collectorCapabilities?.protocols?.join(" · ") || "Modbus TCP · OPC UA · BACnet IP · KNX IP"}
Noch keine Protokollquelle angelegt. Über Neue Quelle startest du sauber mit leerem Formular.
Host, Port, Schreibmodus und Poll-Intervall werden hier zentral verwaltet.
+Links landen Scan-Ergebnisse und importierte Punkte. Von hier schiebst du Werte ins Logging.
+{collectorPointFilter.trim() ? "Keine Treffer für die aktuelle Suche." : "Nach Scan oder Import erscheinen die gefundenen Punkte genau hier links."}
Nur diese rechte Liste wird vom Collector gepollt und in die Trendtabelle geschrieben.
+Hier erscheinen nur die Punkte, die wirklich in die DB geschrieben werden.
{scanResults.opcua?.nodes?.length || scanResults.bacnet?.points?.length || 0} gefundene Datenpunkte oder Objekte.
-{collectorSourcePoints.length} importierte, gescannte oder manuelle Punkte. Aktivierte Punkte werden geloggt.
-