beta. OPC fixed + tested

This commit is contained in:
jhartworks
2026-06-29 15:24:26 +02:00
parent a12ef41d1b
commit 2ce3be068e
9 changed files with 1614 additions and 399 deletions
+43 -19
View File
@@ -5,7 +5,7 @@ import { authenticateToken, createToken, requireRoles } from "./auth.js";
import { config } from "./config.js"; import { config } from "./config.js";
import { pool } from "./db.js"; import { pool } from "./db.js";
import { import {
createWideTrendTable, createTrendTable,
ensureIspMetadata, ensureIspMetadata,
LOCAL_POINT_LIMIT, LOCAL_POINT_LIMIT,
fetchLatestValues, fetchLatestValues,
@@ -13,6 +13,7 @@ import {
getAvailableIsps, getAvailableIsps,
normalizeEngineeringUnit, normalizeEngineeringUnit,
normalizeRange, normalizeRange,
resolveTrendTablePointCount,
sanitizeIspName, sanitizeIspName,
sanitizePointIndex, sanitizePointIndex,
updateIspMetadata, updateIspMetadata,
@@ -25,10 +26,12 @@ import {
getSelection, getSelection,
listRoles, listRoles,
listSelections, listSelections,
listUnits,
saveDashboard, saveDashboard,
savePreferences, savePreferences,
saveRole, saveRole,
saveSelection, saveSelection,
saveUnit,
} from "./lib/store.js"; } from "./lib/store.js";
import { redis } from "./redis.js"; import { redis } from "./redis.js";
@@ -378,25 +381,14 @@ app.get("/api/auth/me", authenticateToken, wrap(async (req, res) => {
})); }));
async function resolveIspPointCount(isp) { async function resolveIspPointCount(isp) {
const safeIsp = sanitizeIspName(isp); return await resolveTrendTablePointCount(pool, config.dbName, 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);
} }
app.get("/api/isps", authenticateToken, wrap(async (_req, res) => { app.get("/api/isps", authenticateToken, wrap(async (_req, res) => {
const isps = await getAvailableIsps(pool, config.dbName); const isps = await getAvailableIsps(pool, config.dbName);
const response = await Promise.all( const response = await Promise.all(
isps.map(async (isp) => { isps.map(async (isp) => {
const metadata = await ensureIspMetadata(redis, isp.tableName, isp.pointCount); const metadata = await ensureIspMetadata(pool, redis, isp.tableName, isp.pointCount);
return { return {
...isp, ...isp,
displayName: metadata.displayName, 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) => { app.get("/api/isps/:isp/aliases", authenticateToken, wrap(async (req, res) => {
const pointCount = await resolveIspPointCount(req.params.isp); 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); res.json(metadata);
})); }));
@@ -421,7 +413,7 @@ app.put(
wrap(async (req, res) => { wrap(async (req, res) => {
const { displayName, description, points } = req.body || {}; const { displayName, description, points } = req.body || {};
const pointCount = await resolveIspPointCount(req.params.isp); 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") { if (typeof displayName === "string") {
metadata.displayName = displayName.trim() || metadata.displayName; 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) => { app.post("/api/trend-tables", authenticateToken, requirePermission("manage_sources"), wrap(async (req, res) => {
const rawName = String(req.body.tableName || "").trim().toLowerCase(); const rawName = String(req.body.tableName || "").trim().toLowerCase();
const displayName = String(req.body.displayName || rawName).trim(); const displayName = String(req.body.displayName || rawName).trim();
const safeName = await createWideTrendTable(pool, rawName.startsWith("trend_") ? rawName : `trend_${rawName}`); const safeName = await createTrendTable(pool, rawName.startsWith("trend_") ? rawName : `trend_${rawName}`);
const metadata = await updateIspMetadata(redis, safeName, (current) => ({ const metadata = await updateIspMetadata(pool, redis, safeName, (current) => ({
...current, ...current,
displayName: displayName || current.displayName, displayName: displayName || current.displayName,
}), LOCAL_POINT_LIMIT); }), LOCAL_POINT_LIMIT);
@@ -482,7 +503,7 @@ app.post(
const normalized = []; const normalized = [];
for (const isp of targets) { 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 }); normalized.push({ isp: isp.tableName, metadata });
} }
@@ -740,3 +761,6 @@ app.listen(config.port, () => {
+301 -58
View File
@@ -116,21 +116,20 @@ function createRandomPointColor(seed) {
export function normalizeEngineeringUnit(value) { export function normalizeEngineeringUnit(value) {
const raw = String(value ?? "").trim(); const raw = String(value ?? "").trim();
if (!raw) { if (!raw) {
return "°C"; return "\u00B0C";
} }
const normalized = raw if (/^\?C$/i.test(raw) || /^\u00B0\s*C$/i.test(raw) || (/[^\x00-\x7F]/.test(raw) && /C/i.test(raw))) {
.replace(/°/g, "°") return "\u00B0C";
.replace(/^\?C$/i, "°C") }
.replace(/^°\s*C$/i, "°C");
return normalized || "°C"; return raw;
} }
function createDefaultPoint(index) { function createDefaultPoint(index) {
return { return {
alias: `Wert${index}`, alias: `Wert${index}`,
unit: "°C", unit: "\u00B0C",
factor: 1, factor: 1,
kind: "analog", kind: "analog",
min: 0, 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 safeTable = sanitizeIspName(tableName);
const sql = ` const sql = `
CREATE TABLE IF NOT EXISTS \`${safeTable}\` ( CREATE TABLE IF NOT EXISTS \`${safeTable}\` (
@@ -181,18 +180,145 @@ export async function createWideTrendTable(pool, tableName) {
return safeTable; 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) { export async function getAvailableIsps(pool, dbName) {
const [tables] = await pool.query( const [tables] = await pool.query(
` `
SELECT SELECT t.TABLE_NAME AS tableName
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
FROM information_schema.TABLES t FROM information_schema.TABLES t
WHERE t.TABLE_SCHEMA = ? WHERE t.TABLE_SCHEMA = ?
AND (t.TABLE_NAME LIKE 'isp%' OR t.TABLE_NAME LIKE 'trend_%') AND (t.TABLE_NAME LIKE 'isp%' OR t.TABLE_NAME LIKE 'trend_%')
@@ -201,36 +327,17 @@ export async function getAvailableIsps(pool, dbName) {
[dbName] [dbName]
); );
const stats = new Map(); const described = await Promise.all(tables.map((table) => describeTrendTable(pool, dbName, table.tableName)));
return described.filter(Boolean);
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,
};
});
} }
export function createDefaultMetadata(isp, pointCount = LEGACY_POINT_LIMIT) { export function createDefaultMetadata(isp, pointCount = LEGACY_POINT_LIMIT) {
const points = {}; const points = {};
const pointOrder = [];
for (let index = 1; index <= pointCount; index += 1) { 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 { return {
@@ -238,12 +345,18 @@ export function createDefaultMetadata(isp, pointCount = LEGACY_POINT_LIMIT) {
displayName: isp.toUpperCase(), displayName: isp.toUpperCase(),
description: "", description: "",
points, 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 safeIsp = sanitizeIspName(isp);
const key = `${META_PREFIX}${safeIsp}`; 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); const existing = await redis.get(key);
if (existing) { if (existing) {
@@ -255,19 +368,43 @@ export async function ensureIspMetadata(redis, isp, pointCount = LEGACY_POINT_LI
changed = true; changed = true;
} }
for (let index = 1; index <= pointCount; index += 1) { if (!Array.isArray(parsed.pointOrder)) {
const pointKey = String(index); 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 defaultPoint = createDefaultPoint(index);
const sourceSeed = seedMap.get(pointKey);
const initialAlias = sourceSeed?.alias || sourceSeed?.name || defaultPoint.alias;
if (!parsed.points[pointKey]) { 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; changed = true;
continue; return;
} }
const point = parsed.points[pointKey]; const point = parsed.points[pointKey];
if (!point.alias) { if (!point.alias || point.alias === defaultPoint.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; changed = true;
} }
const normalizedUnit = normalizeEngineeringUnit(point.unit); const normalizedUnit = normalizeEngineeringUnit(point.unit);
@@ -275,8 +412,15 @@ export async function ensureIspMetadata(redis, isp, pointCount = LEGACY_POINT_LI
point.unit = normalizedUnit; point.unit = normalizedUnit;
changed = true; 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) { if (!point.kind) {
point.kind = defaultPoint.kind; point.kind = sourceSeed?.kind || defaultPoint.kind;
changed = true; changed = true;
} }
if (!Number.isFinite(Number(point.factor))) { if (!Number.isFinite(Number(point.factor))) {
@@ -300,7 +444,7 @@ export async function ensureIspMetadata(redis, isp, pointCount = LEGACY_POINT_LI
point.color = defaultPoint.color; point.color = defaultPoint.color;
changed = true; changed = true;
} }
} });
if (changed) { if (changed) {
await redis.set(key, JSON.stringify(parsed)); 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); 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)); await redis.set(key, JSON.stringify(metadata));
return metadata; return metadata;
} }
export async function updateIspMetadata(redis, isp, updater, pointCount = LEGACY_POINT_LIMIT) { export async function updateIspMetadata(pool, redis, isp, updater, pointCount = LEGACY_POINT_LIMIT) {
const metadata = await ensureIspMetadata(redis, isp, pointCount); const metadata = await ensureIspMetadata(pool, redis, isp, pointCount);
const updated = updater(JSON.parse(JSON.stringify(metadata))); const updated = updater(JSON.parse(JSON.stringify(metadata)));
await redis.set(`${META_PREFIX}${sanitizeIspName(isp)}`, JSON.stringify(updated)); await redis.set(`${META_PREFIX}${sanitizeIspName(isp)}`, JSON.stringify(updated));
return updated; return updated;
} }
export async function getPointDescriptors(redis, points) { export async function getPointDescriptors(pool, redis, points) {
const grouped = new Map(); const grouped = new Map();
points.forEach((point) => { points.forEach((point) => {
@@ -338,10 +495,11 @@ export async function getPointDescriptors(redis, points) {
const descriptors = []; const descriptors = [];
for (const [isp, pointIndexes] of grouped.entries()) { 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) => { pointIndexes.forEach((pointIndex) => {
const pointMeta = metadata.points[String(pointIndex)]; const pointMeta = metadata.points[String(pointIndex)] || createDefaultPoint(pointIndex);
descriptors.push({ descriptors.push({
isp, isp,
pointIndex, pointIndex,
@@ -362,7 +520,7 @@ export async function getPointDescriptors(redis, points) {
} }
export async function fetchTrendSeries(pool, redis, points, range) { 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) => { const grouped = descriptors.reduce((accumulator, descriptor) => {
if (!accumulator.has(descriptor.isp)) { if (!accumulator.has(descriptor.isp)) {
accumulator.set(descriptor.isp, []); accumulator.set(descriptor.isp, []);
@@ -375,6 +533,43 @@ export async function fetchTrendSeries(pool, redis, points, range) {
const mergedRows = new Map(); const mergedRows = new Map();
for (const [isp, pointDescriptors] of grouped.entries()) { 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 const columns = pointDescriptors
.map((descriptor) => `\`Wert${descriptor.pointIndex}\``) .map((descriptor) => `\`Wert${descriptor.pointIndex}\``)
.join(", "); .join(", ");
@@ -410,7 +605,7 @@ export async function fetchTrendSeries(pool, redis, points, range) {
} }
export async function fetchLatestValues(pool, redis, points) { 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) => { const grouped = descriptors.reduce((accumulator, descriptor) => {
if (!accumulator.has(descriptor.isp)) { if (!accumulator.has(descriptor.isp)) {
accumulator.set(descriptor.isp, []); accumulator.set(descriptor.isp, []);
@@ -423,6 +618,52 @@ export async function fetchLatestValues(pool, redis, points) {
const values = []; const values = [];
for (const [isp, pointDescriptors] of grouped.entries()) { 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 const columns = pointDescriptors
.map((descriptor) => `\`Wert${descriptor.pointIndex}\``) .map((descriptor) => `\`Wert${descriptor.pointIndex}\``)
.join(", "); .join(", ");
@@ -459,3 +700,5 @@ export function normalizeRange(input = {}) {
} }
+96
View File
@@ -1,6 +1,7 @@
const SELECTION_SET_KEY = "selection:ids"; const SELECTION_SET_KEY = "selection:ids";
const SELECTION_SEQ_KEY = "selection:seq"; const SELECTION_SEQ_KEY = "selection:seq";
const ROLE_SET_KEY = "role:names"; const ROLE_SET_KEY = "role:names";
const UNIT_SET_KEY = "unit:symbols";
export const DEFAULT_ROLE_DEFINITIONS = [ 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) { function safeParse(jsonValue) {
return jsonValue ? JSON.parse(jsonValue) : null; 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) { function normalizeRoleName(value) {
const normalized = String(value || "") const normalized = String(value || "")
.trim() .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) { export async function listRoles(redis) {
await ensureDefaultRoles(redis); await ensureDefaultRoles(redis);
const names = await redis.smembers(ROLE_SET_KEY); const names = await redis.smembers(ROLE_SET_KEY);
@@ -118,6 +173,47 @@ export async function saveRole(redis, role, previousName = "") {
return merged; 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) { export async function listSelections(redis, user) {
const ids = await redis.smembers(SELECTION_SET_KEY); const ids = await redis.smembers(SELECTION_SET_KEY);
const items = await Promise.all(ids.map((id) => redis.get(`selection:${id}`))); const items = await Promise.all(ids.map((id) => redis.get(`selection:${id}`)));
+389 -43
View File
@@ -1,4 +1,4 @@
import { createServer } from "node:http"; import { createServer } from "node:http";
import { mkdirSync, readFileSync, writeFileSync, existsSync } from "node:fs"; import { mkdirSync, readFileSync, writeFileSync, existsSync } from "node:fs";
import { dirname, join } from "node:path"; import { dirname, join } from "node:path";
import { fileURLToPath } from "node:url"; import { fileURLToPath } from "node:url";
@@ -72,7 +72,7 @@ function replyJson(res, statusCode, body) {
res.writeHead(statusCode, { res.writeHead(statusCode, {
"Content-Type": "application/json; charset=utf-8", "Content-Type": "application/json; charset=utf-8",
"Access-Control-Allow-Origin": "*", "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", "Access-Control-Allow-Headers": "Content-Type, Authorization",
}); });
res.end(JSON.stringify(body)); res.end(JSON.stringify(body));
@@ -130,16 +130,40 @@ function normalizePointKind(value) {
return String(value || "").toLowerCase() === "digital" ? "digital" : "analog"; 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) { function normalizeUnit(value) {
const raw = String(value ?? "").trim(); const raw = normalizeText(value);
if (!raw) { const collapsed = raw.replace(/\s+/g, "");
return "°C"; if (!collapsed) {
return "\u00B0C";
} }
return raw if (/^\?C$/i.test(collapsed) || /^\u00B0C$/i.test(collapsed)) {
.replace(/°/g, "°") return "\u00B0C";
.replace(/^\?C$/i, "°C") }
.replace(/^°\s*C$/i, "°C");
return raw;
} }
function normalizePointIndex(value, fallback) { function normalizePointIndex(value, fallback) {
@@ -159,11 +183,11 @@ function ensureSource(store, body) {
const source = { const source = {
id: body.id || `src-${Date.now()}`, id: body.id || `src-${Date.now()}`,
protocol: normalizeProtocol(body.protocol), protocol: normalizeProtocol(body.protocol),
name: body.name || "Neue Quelle", name: normalizeText(body.name) || "Neue Quelle",
host: body.host || "", host: normalizeText(body.host),
port: body.port || "", port: normalizeText(body.port),
deviceId: body.deviceId || "", deviceId: normalizeText(body.deviceId),
displayName: body.displayName || body.name || "Neue Quelle", displayName: normalizeText(body.displayName || body.name) || "Neue Quelle",
tableName, tableName,
writeMode: normalizeWriteMode(body.writeMode), writeMode: normalizeWriteMode(body.writeMode),
pollIntervalSeconds: normalizePollInterval(body.pollIntervalSeconds), pollIntervalSeconds: normalizePollInterval(body.pollIntervalSeconds),
@@ -297,8 +321,8 @@ function createManualDatapoint(store, payload) {
const fallbackIndex = existing.length + 1; const fallbackIndex = existing.length + 1;
const pointIndex = normalizePointIndex(payload.pointIndex, fallbackIndex); const pointIndex = normalizePointIndex(payload.pointIndex, fallbackIndex);
const protocol = normalizeProtocol(payload.protocol || source.protocol); const protocol = normalizeProtocol(payload.protocol || source.protocol);
const alias = String(payload.alias || payload.name || `Punkt ${pointIndex}`).trim() || `Punkt ${pointIndex}`; const alias = normalizeText(payload.alias || payload.name || `Punkt ${pointIndex}`) || `Punkt ${pointIndex}`;
const name = String(payload.name || payload.alias || alias).trim() || alias; const name = normalizeText(payload.name || payload.alias || alias) || alias;
const datapoint = { const datapoint = {
id: payload.id || `dp-${protocol}-${Date.now()}-${pointIndex}`, id: payload.id || `dp-${protocol}-${Date.now()}-${pointIndex}`,
@@ -311,10 +335,10 @@ function createManualDatapoint(store, payload) {
kind: normalizePointKind(payload.kind), kind: normalizePointKind(payload.kind),
enabled: payload.enabled !== false, enabled: payload.enabled !== false,
writeMode: normalizeWriteMode(payload.writeMode || source.writeMode), writeMode: normalizeWriteMode(payload.writeMode || source.writeMode),
condition: String(payload.condition || "").trim(), condition: normalizeText(payload.condition),
dataType: String(payload.dataType || (protocol === "knx-ip" ? "group-address" : "holding-register")).trim(), dataType: normalizeText(payload.dataType || (protocol === "knx-ip" ? "group-address" : "holding-register")),
address: String(payload.address || payload.register || "").trim(), address: normalizeText(payload.address || payload.register || ""),
nodeId: String(payload.nodeId || "").trim(), nodeId: normalizeText(payload.nodeId || ""),
scale: Number(payload.scale || 1) || 1, scale: Number(payload.scale || 1) || 1,
host: source.host || "", host: source.host || "",
tableName: source.tableName || "", tableName: source.tableName || "",
@@ -333,27 +357,27 @@ function createManualDatapoint(store, payload) {
function createScanDatapoint(source, payload, index) { function createScanDatapoint(source, payload, index) {
const pointIndex = normalizePointIndex(payload.pointIndex, index + 1); const pointIndex = normalizePointIndex(payload.pointIndex, index + 1);
const protocol = normalizeProtocol(payload.protocol || source.protocol); 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 { return {
id: payload.id || `dp-${protocol}-${Date.now()}-${Math.round(Math.random() * 100000)}-${pointIndex}`, id: payload.id || `dp-${protocol}-${Date.now()}-${Math.round(Math.random() * 100000)}-${pointIndex}`,
sourceId: source.id, sourceId: source.id,
protocol, protocol,
pointIndex, pointIndex,
name: String(payload.name || alias).trim() || alias, name: normalizeText(payload.name || alias) || alias,
alias, alias,
unit: normalizeUnit(payload.unit), unit: normalizeUnit(payload.unit),
kind: normalizePointKind(payload.kind), kind: normalizePointKind(payload.kind),
enabled: payload.enabled === true, enabled: payload.enabled === true,
writeMode: normalizeWriteMode(payload.writeMode || source.writeMode), writeMode: normalizeWriteMode(payload.writeMode || source.writeMode),
condition: String(payload.condition || "").trim(), condition: normalizeText(payload.condition),
dataType: String(payload.dataType || "").trim(), dataType: normalizeText(payload.dataType),
address: String(payload.address || "").trim(), address: normalizeText(payload.address),
nodeId: String(payload.nodeId || "").trim(), nodeId: normalizeText(payload.nodeId),
objectType: payload.objectType, objectType: payload.objectType,
objectInstance: payload.objectInstance, objectInstance: payload.objectInstance,
propertyId: payload.propertyId || 85, propertyId: payload.propertyId || 85,
scale: Number(payload.scale || 1) || 1, scale: Number(payload.scale || 1) || 1,
host: payload.host || source.host || "", host: normalizeText(payload.host || source.host),
tableName: source.tableName || "", tableName: source.tableName || "",
currentValue: payload.currentValue ?? null, currentValue: payload.currentValue ?? null,
lastReadAt: payload.currentValue === undefined ? null : new Date().toISOString(), 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) { function upsertDatapoints(store, datapoints) {
datapoints.forEach((datapoint) => { datapoints.forEach((datapoint) => {
store.datapoints = store.datapoints.filter((item) => item.id !== datapoint.id && !(item.sourceId === datapoint.sourceId && item.pointIndex === datapoint.pointIndex)); store.datapoints = store.datapoints.filter((item) => item.id !== datapoint.id && !(item.sourceId === datapoint.sourceId && item.pointIndex === datapoint.pointIndex));
@@ -634,7 +702,7 @@ async function readOpcuaValue(source, datapoint) {
return value ? 1 : 0; return value ? 1 : 0;
} }
const numeric = Number(value); 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 { } finally {
await client.disconnect().catch(() => undefined); await client.disconnect().catch(() => undefined);
} }
@@ -797,7 +865,7 @@ async function readBacnetValue(source, datapoint) {
return; return;
} }
const numeric = Number(raw); 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); reject(error);
} else { } else {
const numeric = Number(value); 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; 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) { async function writeTrendRow(source, readings) {
if (!readings.length) { if (!readings.length) {
return false; 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 columns = ["`userlevel`"];
const placeholders = ["?"]; const placeholders = ["?"];
const values = [0]; const values = [0];
@@ -1057,7 +1364,7 @@ async function pollModbusSource(source, datapoints) {
} }
async function pollSource(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) { if (!sourcePoints.length) {
updateSourceState(source.id, { updateSourceState(source.id, {
lastPollAt: new Date().toISOString(), lastPollAt: new Date().toISOString(),
@@ -1068,6 +1375,7 @@ async function pollSource(source, datapoints) {
} }
try { try {
source.__activeDatapoints = sourcePoints;
if (source.protocol === "modbus-tcp") { if (source.protocol === "modbus-tcp") {
await pollModbusSource(source, sourcePoints); await pollModbusSource(source, sourcePoints);
return; return;
@@ -1091,7 +1399,7 @@ async function pollSource(source, datapoints) {
updateSourceState(source.id, { updateSourceState(source.id, {
lastPollAt: new Date().toISOString(), lastPollAt: new Date().toISOString(),
lastPollStatus: "unsupported", lastPollStatus: "unsupported",
lastPollError: `${source.protocol} wird vom Collector nicht unterstützt.`, lastPollError: `${source.protocol} wird vom Collector nicht unterst?tzt.`,
}); });
} catch (error) { } catch (error) {
updateSourceState(source.id, { updateSourceState(source.id, {
@@ -1099,6 +1407,8 @@ async function pollSource(source, datapoints) {
lastPollStatus: "error", lastPollStatus: "error",
lastPollError: error.message || "Polling fehlgeschlagen.", lastPollError: error.message || "Polling fehlgeschlagen.",
}); });
} finally {
delete source.__activeDatapoints;
} }
} }
@@ -1192,7 +1502,9 @@ const server = createServer(async (req, res) => {
const source = ensureSource(store, body); const source = ensureSource(store, body);
return { ...store, lastSource: source.id }; 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 { } catch {
replyJson(res, 400, { error: "Ungültige JSON-Daten." }); replyJson(res, 400, { error: "Ungültige JSON-Daten." });
} }
@@ -1211,7 +1523,9 @@ const server = createServer(async (req, res) => {
ensureSource(store, { ...existing, ...body, id: sourceId }); ensureSource(store, { ...existing, ...body, id: sourceId });
return store; 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) { } catch (error) {
replyJson(res, error.message === "not-found" ? 404 : 400, { replyJson(res, error.message === "not-found" ? 404 : 400, {
error: error.message === "not-found" ? "Quelle nicht gefunden." : "Ungültige JSON-Daten.", 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), sources: store.sources.filter((item) => item.id !== sourceId),
datapoints: store.datapoints.filter((item) => item.sourceId !== 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 }); replyJson(res, 200, { ok: true, sources: next.sources.length, datapoints: next.datapoints.length });
return; return;
} }
@@ -1241,6 +1557,7 @@ const server = createServer(async (req, res) => {
return store; return store;
}); });
const created = next.datapoints.find((item) => item.sourceId === body.sourceId && item.pointIndex === normalizePointIndex(body.pointIndex, 1)); 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 }); replyJson(res, 201, { ok: true, datapoint: created || null, updatedAt: next.updatedAt });
} catch (error) { } catch (error) {
replyJson(res, 400, { error: error.message || "Datenpunkt konnte nicht angelegt werden." }); 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/")) { if (req.method === "PATCH" && url.pathname.startsWith("/datapoints/")) {
try { try {
const datapointId = url.pathname.split("/").at(-1); const datapointId = decodeURIComponent(url.pathname.split("/").at(-1) || "");
const body = await readBody(req); const body = await readBody(req);
const next = touchStore((store) => { const next = touchStore((store) => {
const index = store.datapoints.findIndex((item) => item.id === datapointId); const index = store.datapoints.findIndex((item) => item.id === datapointId);
@@ -1270,7 +1587,9 @@ const server = createServer(async (req, res) => {
}; };
return store; 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) { } catch (error) {
replyJson(res, error.message === "not-found" ? 404 : 400, { error: error.message === "not-found" ? "Datenpunkt nicht gefunden." : error.message }); 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$/)) { if (req.method === "POST" && url.pathname.match(/^\/datapoints\/[^/]+\/read$/)) {
try { try {
const datapointId = url.pathname.split("/").at(-2); const datapointId = decodeURIComponent(url.pathname.split("/").at(-2) || "");
const store = readStore(); const store = readStore();
const datapoint = store.datapoints.find((item) => item.id === datapointId); const datapoint = store.datapoints.find((item) => item.id === datapointId);
const source = datapoint ? store.sources.find((item) => item.id === datapoint.sourceId) : null; const source = datapoint ? store.sources.find((item) => item.id === datapoint.sourceId) : null;
@@ -1299,7 +1618,6 @@ const server = createServer(async (req, res) => {
return; return;
} }
if (req.method === "POST" && url.pathname === "/imports/modbus-csv") { if (req.method === "POST" && url.pathname === "/imports/modbus-csv") {
try { try {
const body = await readBody(req); const body = await readBody(req);
@@ -1311,6 +1629,7 @@ const server = createServer(async (req, res) => {
upsertDatapoints(store, imported); upsertDatapoints(store, imported);
return store; return store;
}); });
await syncDatapointRegistryEntries(imported);
replyJson(res, 200, { imported: imported.length, datapoints: imported, updatedAt: next.updatedAt }); replyJson(res, 200, { imported: imported.length, datapoints: imported, updatedAt: next.updatedAt });
} catch { } catch {
replyJson(res, 400, { error: "CSV-Import fehlgeschlagen." }); replyJson(res, 400, { error: "CSV-Import fehlgeschlagen." });
@@ -1329,6 +1648,7 @@ const server = createServer(async (req, res) => {
upsertDatapoints(store, imported); upsertDatapoints(store, imported);
return store; return store;
}); });
await syncDatapointRegistryEntries(imported);
replyJson(res, 200, { imported: imported.length, datapoints: imported, updatedAt: next.updatedAt }); replyJson(res, 200, { imported: imported.length, datapoints: imported, updatedAt: next.updatedAt });
} catch { } catch {
replyJson(res, 400, { error: "KNX-XML-Import fehlgeschlagen." }); replyJson(res, 400, { error: "KNX-XML-Import fehlgeschlagen." });
@@ -1344,12 +1664,23 @@ const server = createServer(async (req, res) => {
return; return;
} }
const result = await scanOpcuaEndpoint(body); const result = await scanOpcuaEndpoint(body);
let scannedCandidates = [];
touchStore((store) => { touchStore((store) => {
if (body.sourceId) { if (body.sourceId) {
const source = store.sources.find((item) => item.id === body.sourceId); const source = store.sources.find((item) => item.id === body.sourceId);
if (source) { if (source) {
const existingCount = store.datapoints.filter((item) => item.sourceId === source.id).length; const rawNodes = (result.nodes || []).filter((node) => node.nodeId);
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 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); upsertDatapoints(store, candidates);
} }
} }
@@ -1357,6 +1688,7 @@ const server = createServer(async (req, res) => {
store.scans.opcua = store.scans.opcua.slice(0, 20); store.scans.opcua = store.scans.opcua.slice(0, 20);
return store; return store;
}); });
await syncDatapointRegistryEntries(scannedCandidates);
replyJson(res, 200, result); replyJson(res, 200, result);
} catch (error) { } catch (error) {
replyJson(res, 400, { error: error.message || "OPC-UA-Scan fehlgeschlagen." }); replyJson(res, 400, { error: error.message || "OPC-UA-Scan fehlgeschlagen." });
@@ -1368,12 +1700,21 @@ const server = createServer(async (req, res) => {
try { try {
const body = await readBody(req); const body = await readBody(req);
const payload = await scanBacnetNetwork(body); const payload = await scanBacnetNetwork(body);
let scannedCandidates = [];
touchStore((store) => { touchStore((store) => {
if (body.sourceId) { if (body.sourceId) {
const source = store.sources.find((item) => item.id === body.sourceId); const source = store.sources.find((item) => item.id === body.sourceId);
if (source) { if (source) {
const existingCount = store.datapoints.filter((item) => item.sourceId === source.id).length; const rawPoints = payload.points || [];
const candidates = (payload.points || []).map((point, index) => createScanDatapoint(source, { ...point, protocol: "bacnet-ip", pointIndex: existingCount + index + 1 }, index)); 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); upsertDatapoints(store, candidates);
} }
} }
@@ -1381,6 +1722,7 @@ const server = createServer(async (req, res) => {
store.scans.bacnet = store.scans.bacnet.slice(0, 20); store.scans.bacnet = store.scans.bacnet.slice(0, 20);
return store; return store;
}); });
await syncDatapointRegistryEntries(scannedCandidates);
replyJson(res, 200, payload); replyJson(res, 200, payload);
} catch (error) { } catch (error) {
replyJson(res, 400, { error: error.message || "BACnet-Scan fehlgeschlagen." }); 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." }); replyJson(res, 404, { error: "Nicht gefunden." });
}); });
rebuildRegistryTablesFromStore().catch((error) => console.error("Registry sync failed", error));
server.listen(port, () => { server.listen(port, () => {
console.log(`SE Local Trenddata collector listening on port ${port}`); console.log(`SE Local Trenddata collector listening on port ${port}`);
}); });
+22 -2
View File
@@ -5,9 +5,10 @@ services:
restart: unless-stopped restart: unless-stopped
environment: environment:
MARIADB_ROOT_PASSWORD: ${MARIADB_ROOT_PASSWORD:-SE3112} MARIADB_ROOT_PASSWORD: ${MARIADB_ROOT_PASSWORD:-SE3112}
MARIADB_ROOT_HOST: mariadb
MARIADB_DATABASE: wago MARIADB_DATABASE: wago
MARIADB_USER: app MARIADB_USER: wago
MARIADB_PASSWORD: app MARIADB_PASSWORD: ${MARIADB_PASSWORD:-SE3112}
ports: ports:
- "3306:3306" - "3306:3306"
volumes: volumes:
@@ -105,9 +106,28 @@ services:
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock - /var/run/docker.sock:/var/run/docker.sock
- dockhand_data:/app/data - 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: volumes:
mariadb_data: mariadb_data:
redis_data: redis_data:
collector_data: collector_data:
dockhand_data: dockhand_data:
+7 -1
View File
@@ -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`; USE `wago`;
DELIMITER // DELIMITER //
@@ -55,3 +55,9 @@ CREATE TABLE IF NOT EXISTS `app_users` (
UNIQUE KEY `uniq_app_users_username` (`username`), UNIQUE KEY `uniq_app_users_username` (`username`),
UNIQUE KEY `uniq_app_users_email` (`email`) UNIQUE KEY `uniq_app_users_email` (`email`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; ) 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;
+516 -170
View File
File diff suppressed because it is too large Load Diff
+3
View File
@@ -56,6 +56,9 @@ export const api = {
getIsps: (token) => apiRequest("/isps", { token }), getIsps: (token) => apiRequest("/isps", { token }),
getAliases: (token, isp) => apiRequest(`/isps/${isp}/aliases`, { token }), getAliases: (token, isp) => apiRequest(`/isps/${isp}/aliases`, { token }),
saveAliases: (token, isp, body) => apiRequest(`/isps/${isp}/aliases`, { method: "PUT", token, body }), 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 }), createTrendTable: (token, body) => apiRequest("/trend-tables", { method: "POST", token, body }),
normalizeAliasMetadata: (token, body) => apiRequest("/admin/normalize-alias-metadata", { 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 }), importSql: (token, body) => apiRequest("/admin/import-sql", { method: "POST", token, body }),
+131
View File
@@ -919,3 +919,134 @@ span {
font-weight: 700; font-weight: 700;
white-space: nowrap; 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%;
}
}