initial commit
This commit is contained in:
@@ -0,0 +1,44 @@
|
||||
import jwt from "jsonwebtoken";
|
||||
import { config } from "./config.js";
|
||||
|
||||
export function createToken(user) {
|
||||
return jwt.sign(
|
||||
{
|
||||
sub: user.id,
|
||||
email: user.email,
|
||||
role: user.role,
|
||||
name: user.name,
|
||||
},
|
||||
config.jwtSecret,
|
||||
{
|
||||
expiresIn: "12h",
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
export function authenticateToken(req, res, next) {
|
||||
const header = req.headers.authorization || "";
|
||||
const token = header.startsWith("Bearer ") ? header.slice(7) : "";
|
||||
|
||||
if (!token) {
|
||||
return res.status(401).json({ error: "Authentication required." });
|
||||
}
|
||||
|
||||
try {
|
||||
req.user = jwt.verify(token, config.jwtSecret);
|
||||
next();
|
||||
} catch {
|
||||
return res.status(401).json({ error: "Token is invalid or expired." });
|
||||
}
|
||||
}
|
||||
|
||||
export function requireRoles(...roles) {
|
||||
return (req, res, next) => {
|
||||
if (!req.user || !roles.includes(req.user.role)) {
|
||||
return res.status(403).json({ error: "Insufficient permissions." });
|
||||
}
|
||||
|
||||
next();
|
||||
};
|
||||
}
|
||||
|
||||
@@ -0,0 +1,19 @@
|
||||
import dotenv from "dotenv";
|
||||
|
||||
dotenv.config();
|
||||
|
||||
export const config = {
|
||||
port: Number(process.env.PORT || 8080),
|
||||
dbHost: process.env.DB_HOST || "localhost",
|
||||
dbPort: Number(process.env.DB_PORT || 3306),
|
||||
dbName: process.env.DB_NAME || "wago",
|
||||
dbUser: process.env.DB_USER || "root",
|
||||
dbPassword: process.env.DB_PASSWORD || "root",
|
||||
redisUrl: process.env.REDIS_URL || "redis://127.0.0.1:6379",
|
||||
jwtSecret: process.env.JWT_SECRET || "change-me",
|
||||
corsOrigin: process.env.CORS_ORIGIN || "http://localhost:5173",
|
||||
adminUsername: process.env.ADMIN_USERNAME || "admin",
|
||||
adminEmail: process.env.ADMIN_EMAIL || "admin@example.com",
|
||||
adminPassword: process.env.ADMIN_PASSWORD || "ChangeMe123!",
|
||||
adminName: process.env.ADMIN_NAME || "System Admin",
|
||||
};
|
||||
@@ -0,0 +1,14 @@
|
||||
import mysql from "mysql2/promise";
|
||||
import { config } from "./config.js";
|
||||
|
||||
export const pool = mysql.createPool({
|
||||
host: config.dbHost,
|
||||
port: config.dbPort,
|
||||
user: config.dbUser,
|
||||
password: config.dbPassword,
|
||||
database: config.dbName,
|
||||
waitForConnections: true,
|
||||
connectionLimit: 10,
|
||||
namedPlaceholders: true,
|
||||
multipleStatements: true,
|
||||
});
|
||||
@@ -0,0 +1,514 @@
|
||||
import bcrypt from "bcryptjs";
|
||||
import cors from "cors";
|
||||
import express from "express";
|
||||
import { authenticateToken, createToken, requireRoles } from "./auth.js";
|
||||
import { config } from "./config.js";
|
||||
import { pool } from "./db.js";
|
||||
import {
|
||||
ensureIspMetadata,
|
||||
fetchLatestValues,
|
||||
fetchTrendSeries,
|
||||
getAvailableIsps,
|
||||
normalizeRange,
|
||||
sanitizeIspName,
|
||||
sanitizePointIndex,
|
||||
updateIspMetadata,
|
||||
} from "./lib/isp.js";
|
||||
import {
|
||||
deleteSelection,
|
||||
getDashboard,
|
||||
getPreferences,
|
||||
getSelection,
|
||||
listSelections,
|
||||
saveDashboard,
|
||||
savePreferences,
|
||||
saveSelection,
|
||||
} from "./lib/store.js";
|
||||
import { redis } from "./redis.js";
|
||||
|
||||
const app = express();
|
||||
const wrap = (handler) => (req, res, next) => Promise.resolve(handler(req, res, next)).catch(next);
|
||||
|
||||
app.use(
|
||||
cors({
|
||||
origin: config.corsOrigin,
|
||||
})
|
||||
);
|
||||
app.use(express.json({ limit: "50mb" }));
|
||||
|
||||
function serializeUser(row) {
|
||||
return {
|
||||
id: row.id,
|
||||
username: row.username,
|
||||
name: row.name,
|
||||
email: row.email,
|
||||
role: row.role,
|
||||
active: Boolean(row.active),
|
||||
createdAt: row.created_at,
|
||||
updatedAt: row.updated_at,
|
||||
};
|
||||
}
|
||||
|
||||
function assertPoints(points) {
|
||||
if (!Array.isArray(points) || points.length === 0) {
|
||||
throw new Error("At least one data point is required.");
|
||||
}
|
||||
|
||||
return points.map((point) => ({
|
||||
isp: sanitizeIspName(point.isp),
|
||||
pointIndex: sanitizePointIndex(point.pointIndex),
|
||||
}));
|
||||
}
|
||||
|
||||
function canUseSelection(selection, user) {
|
||||
return Boolean(selection) && (selection.shared || selection.ownerId === user.sub || user.role === "admin");
|
||||
}
|
||||
|
||||
function slugifyUsername(value) {
|
||||
const base = String(value || "")
|
||||
.toLowerCase()
|
||||
.replace(/[^a-z0-9]+/g, "-")
|
||||
.replace(/^-+|-+$/g, "")
|
||||
.slice(0, 40);
|
||||
|
||||
return base || "user";
|
||||
}
|
||||
|
||||
function nextAvailableUsername(seed, taken) {
|
||||
const base = slugifyUsername(seed);
|
||||
let candidate = base;
|
||||
let counter = 1;
|
||||
|
||||
while (taken.has(candidate)) {
|
||||
candidate = `${base}-${counter}`;
|
||||
counter += 1;
|
||||
}
|
||||
|
||||
taken.add(candidate);
|
||||
return candidate;
|
||||
}
|
||||
|
||||
function normalizeSqlDump(sql) {
|
||||
return String(sql || "")
|
||||
.replace(/^\uFEFF/, "")
|
||||
.split(/\r?\n/)
|
||||
.filter((line) => !line.trim().toUpperCase().startsWith("DELIMITER "))
|
||||
.join("\n")
|
||||
.trim();
|
||||
}
|
||||
|
||||
async function resolveSelectionPoints(body, user) {
|
||||
if (body.selectionId) {
|
||||
const selection = await getSelection(redis, body.selectionId);
|
||||
|
||||
if (!canUseSelection(selection, user)) {
|
||||
throw new Error("Selection not found or not accessible.");
|
||||
}
|
||||
|
||||
return assertPoints(selection.points || []);
|
||||
}
|
||||
|
||||
return assertPoints(body.points || []);
|
||||
}
|
||||
|
||||
async function ensureUserSchema() {
|
||||
await pool.query(
|
||||
`
|
||||
CREATE TABLE IF NOT EXISTS app_users (
|
||||
id INT NOT NULL AUTO_INCREMENT,
|
||||
username VARCHAR(80) NOT NULL,
|
||||
name VARCHAR(120) NOT NULL,
|
||||
email VARCHAR(160) NULL,
|
||||
password_hash VARCHAR(255) NOT NULL,
|
||||
role ENUM('viewer', 'technician', 'admin') NOT NULL DEFAULT 'viewer',
|
||||
active TINYINT(1) NOT NULL DEFAULT 1,
|
||||
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
|
||||
PRIMARY KEY (id),
|
||||
UNIQUE KEY uniq_app_users_username (username),
|
||||
UNIQUE KEY uniq_app_users_email (email)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
|
||||
`
|
||||
);
|
||||
|
||||
try {
|
||||
await pool.query("ALTER TABLE app_users ADD COLUMN username VARCHAR(80) NULL AFTER id");
|
||||
} catch {}
|
||||
|
||||
try {
|
||||
await pool.query("ALTER TABLE app_users MODIFY email VARCHAR(160) NULL");
|
||||
} catch {}
|
||||
|
||||
const [users] = await pool.query("SELECT id, username, name, email FROM app_users ORDER BY id");
|
||||
const taken = new Set(users.map((user) => user.username).filter(Boolean));
|
||||
|
||||
for (const user of users) {
|
||||
if (user.username) {
|
||||
continue;
|
||||
}
|
||||
|
||||
const seed = user.email?.split("@")[0] || user.name || `user-${user.id}`;
|
||||
const username = nextAvailableUsername(seed, taken);
|
||||
await pool.query("UPDATE app_users SET username = ? WHERE id = ?", [username, user.id]);
|
||||
}
|
||||
|
||||
try {
|
||||
await pool.query("ALTER TABLE app_users MODIFY username VARCHAR(80) NOT NULL");
|
||||
} catch {}
|
||||
|
||||
try {
|
||||
await pool.query("ALTER TABLE app_users ADD UNIQUE KEY uniq_app_users_username (username)");
|
||||
} catch {}
|
||||
|
||||
try {
|
||||
await pool.query("ALTER TABLE app_users ADD UNIQUE KEY uniq_app_users_email (email)");
|
||||
} catch {}
|
||||
}
|
||||
|
||||
async function ensureAdminUser() {
|
||||
const [rows] = await pool.query(
|
||||
"SELECT * FROM app_users WHERE username = ? OR email = ? LIMIT 1",
|
||||
[config.adminUsername.toLowerCase(), config.adminEmail.toLowerCase()]
|
||||
);
|
||||
|
||||
if (rows[0]) {
|
||||
const user = rows[0];
|
||||
await pool.query(
|
||||
`
|
||||
UPDATE app_users
|
||||
SET username = ?, name = ?, email = COALESCE(email, ?), role = 'admin', active = 1
|
||||
WHERE id = ?
|
||||
`,
|
||||
[config.adminUsername.toLowerCase(), config.adminName, config.adminEmail.toLowerCase(), user.id]
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
const passwordHash = await bcrypt.hash(config.adminPassword, 10);
|
||||
|
||||
await pool.query(
|
||||
`
|
||||
INSERT INTO app_users (username, name, email, password_hash, role, active)
|
||||
VALUES (?, ?, ?, ?, 'admin', 1)
|
||||
`,
|
||||
[config.adminUsername.toLowerCase(), config.adminName, config.adminEmail.toLowerCase(), passwordHash]
|
||||
);
|
||||
}
|
||||
|
||||
app.get("/api/health", wrap(async (_req, res) => {
|
||||
const [dbRows] = await pool.query("SELECT 1 AS ok");
|
||||
const redisOk = await redis.ping();
|
||||
res.json({
|
||||
ok: dbRows[0]?.ok === 1 && redisOk === "PONG",
|
||||
timestamp: new Date().toISOString(),
|
||||
});
|
||||
}));
|
||||
|
||||
app.post("/api/auth/login", wrap(async (req, res) => {
|
||||
const username = String(req.body.username || "").trim().toLowerCase();
|
||||
const password = String(req.body.password || "");
|
||||
|
||||
const [rows] = await pool.query(
|
||||
"SELECT * FROM app_users WHERE username = ? AND active = 1 LIMIT 1",
|
||||
[username]
|
||||
);
|
||||
|
||||
const user = rows[0];
|
||||
|
||||
if (!user) {
|
||||
return res.status(401).json({ error: "Invalid credentials." });
|
||||
}
|
||||
|
||||
const isValid = await bcrypt.compare(password, user.password_hash);
|
||||
|
||||
if (!isValid) {
|
||||
return res.status(401).json({ error: "Invalid credentials." });
|
||||
}
|
||||
|
||||
const serializedUser = serializeUser(user);
|
||||
|
||||
return res.json({
|
||||
token: createToken(serializedUser),
|
||||
user: serializedUser,
|
||||
});
|
||||
}));
|
||||
|
||||
app.get("/api/auth/me", authenticateToken, wrap(async (req, res) => {
|
||||
const [rows] = await pool.query("SELECT * FROM app_users WHERE id = ? LIMIT 1", [req.user.sub]);
|
||||
|
||||
if (!rows[0]) {
|
||||
return res.status(404).json({ error: "User not found." });
|
||||
}
|
||||
|
||||
return res.json({ user: serializeUser(rows[0]) });
|
||||
}));
|
||||
|
||||
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);
|
||||
return {
|
||||
...isp,
|
||||
displayName: metadata.displayName,
|
||||
description: metadata.description,
|
||||
};
|
||||
})
|
||||
);
|
||||
|
||||
res.json({ isps: response });
|
||||
}));
|
||||
|
||||
app.get("/api/isps/:isp/aliases", authenticateToken, wrap(async (req, res) => {
|
||||
const metadata = await ensureIspMetadata(redis, req.params.isp);
|
||||
res.json(metadata);
|
||||
}));
|
||||
|
||||
app.put(
|
||||
"/api/isps/:isp/aliases",
|
||||
authenticateToken,
|
||||
requireRoles("technician", "admin"),
|
||||
wrap(async (req, res) => {
|
||||
const { displayName, description, points } = req.body || {};
|
||||
const updated = await updateIspMetadata(redis, req.params.isp, (metadata) => {
|
||||
if (typeof displayName === "string") {
|
||||
metadata.displayName = displayName.trim() || metadata.displayName;
|
||||
}
|
||||
|
||||
if (typeof description === "string") {
|
||||
metadata.description = description.trim();
|
||||
}
|
||||
|
||||
if (Array.isArray(points)) {
|
||||
points.forEach((point) => {
|
||||
const pointIndex = sanitizePointIndex(point.pointIndex);
|
||||
const key = String(pointIndex);
|
||||
metadata.points[key] = {
|
||||
...metadata.points[key],
|
||||
alias: String(point.alias || metadata.points[key].alias).trim() || metadata.points[key].alias,
|
||||
unit: String(point.unit || metadata.points[key].unit).trim(),
|
||||
factor: Number.isFinite(Number(point.factor)) ? Number(point.factor) : metadata.points[key].factor,
|
||||
kind: point.kind === "digital" ? "digital" : "analog",
|
||||
min: Number.isFinite(Number(point.min)) ? Number(point.min) : metadata.points[key].min,
|
||||
max: Number.isFinite(Number(point.max)) ? Number(point.max) : metadata.points[key].max,
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
return metadata;
|
||||
});
|
||||
|
||||
res.json(updated);
|
||||
})
|
||||
);
|
||||
|
||||
app.post(
|
||||
"/api/admin/import-sql",
|
||||
authenticateToken,
|
||||
requireRoles("admin"),
|
||||
wrap(async (req, res) => {
|
||||
const sql = normalizeSqlDump(req.body.sql);
|
||||
|
||||
if (!sql) {
|
||||
throw new Error("SQL content is empty.");
|
||||
}
|
||||
|
||||
await pool.query(sql);
|
||||
const isps = await getAvailableIsps(pool, config.dbName);
|
||||
res.json({ ok: true, importedTables: isps.length, filename: req.body.filename || null });
|
||||
})
|
||||
);
|
||||
|
||||
app.get("/api/users", authenticateToken, requireRoles("admin"), wrap(async (_req, res) => {
|
||||
const [rows] = await pool.query(
|
||||
"SELECT id, username, name, email, role, active, created_at, updated_at FROM app_users ORDER BY username"
|
||||
);
|
||||
res.json({ users: rows.map(serializeUser) });
|
||||
}));
|
||||
|
||||
app.post("/api/users", authenticateToken, requireRoles("admin"), wrap(async (req, res) => {
|
||||
const username = String(req.body.username || "").trim().toLowerCase();
|
||||
const name = String(req.body.name || "").trim();
|
||||
const email = String(req.body.email || "").trim().toLowerCase() || null;
|
||||
const role = ["viewer", "technician", "admin"].includes(req.body.role) ? req.body.role : "viewer";
|
||||
const password = String(req.body.password || "");
|
||||
|
||||
if (!username || !name || password.length < 8) {
|
||||
return res.status(400).json({ error: "Username, name and password are required." });
|
||||
}
|
||||
|
||||
const passwordHash = await bcrypt.hash(password, 10);
|
||||
|
||||
await pool.query(
|
||||
`
|
||||
INSERT INTO app_users (username, name, email, password_hash, role, active)
|
||||
VALUES (?, ?, ?, ?, ?, 1)
|
||||
`,
|
||||
[username, name, email, passwordHash, role]
|
||||
);
|
||||
|
||||
const [rows] = await pool.query("SELECT * FROM app_users WHERE username = ? LIMIT 1", [username]);
|
||||
res.status(201).json({ user: serializeUser(rows[0]) });
|
||||
}));
|
||||
|
||||
app.patch("/api/users/:id", authenticateToken, requireRoles("admin"), wrap(async (req, res) => {
|
||||
const userId = Number(req.params.id);
|
||||
const payload = req.body || {};
|
||||
const [rows] = await pool.query("SELECT * FROM app_users WHERE id = ? LIMIT 1", [userId]);
|
||||
|
||||
if (!rows[0]) {
|
||||
return res.status(404).json({ error: "User not found." });
|
||||
}
|
||||
|
||||
const current = rows[0];
|
||||
const username = typeof payload.username === "string"
|
||||
? payload.username.trim().toLowerCase() || current.username
|
||||
: current.username;
|
||||
const name = typeof payload.name === "string" ? payload.name.trim() || current.name : current.name;
|
||||
const email = typeof payload.email === "string"
|
||||
? payload.email.trim().toLowerCase() || null
|
||||
: current.email;
|
||||
const role = ["viewer", "technician", "admin"].includes(payload.role) ? payload.role : current.role;
|
||||
const active = typeof payload.active === "boolean" ? Number(payload.active) : current.active;
|
||||
const passwordHash = payload.password
|
||||
? await bcrypt.hash(String(payload.password), 10)
|
||||
: current.password_hash;
|
||||
|
||||
await pool.query(
|
||||
`
|
||||
UPDATE app_users
|
||||
SET username = ?, name = ?, email = ?, role = ?, active = ?, password_hash = ?
|
||||
WHERE id = ?
|
||||
`,
|
||||
[username, name, email, role, active, passwordHash, userId]
|
||||
);
|
||||
|
||||
const [updatedRows] = await pool.query("SELECT * FROM app_users WHERE id = ? LIMIT 1", [userId]);
|
||||
res.json({ user: serializeUser(updatedRows[0]) });
|
||||
}));
|
||||
|
||||
app.get("/api/selections", authenticateToken, wrap(async (req, res) => {
|
||||
const selections = await listSelections(redis, req.user);
|
||||
res.json({ selections });
|
||||
}));
|
||||
|
||||
app.post("/api/selections", authenticateToken, wrap(async (req, res) => {
|
||||
const points = assertPoints(req.body.points || []);
|
||||
const selection = await saveSelection(redis, {
|
||||
ownerId: req.user.sub,
|
||||
ownerName: req.user.name,
|
||||
name: String(req.body.name || "").trim() || "Neue Auswahl",
|
||||
description: String(req.body.description || "").trim(),
|
||||
shared: Boolean(req.body.shared),
|
||||
points,
|
||||
});
|
||||
|
||||
res.status(201).json({ selection });
|
||||
}));
|
||||
|
||||
app.put("/api/selections/:selectionId", authenticateToken, wrap(async (req, res) => {
|
||||
const existing = await getSelection(redis, req.params.selectionId);
|
||||
|
||||
if (!existing) {
|
||||
return res.status(404).json({ error: "Selection not found." });
|
||||
}
|
||||
|
||||
if (existing.ownerId !== req.user.sub && req.user.role !== "admin") {
|
||||
return res.status(403).json({ error: "You can only edit your own selections." });
|
||||
}
|
||||
|
||||
const points = assertPoints(req.body.points || existing.points || []);
|
||||
const selection = await saveSelection(
|
||||
redis,
|
||||
{
|
||||
...existing,
|
||||
name: String(req.body.name || existing.name).trim(),
|
||||
description: String(req.body.description || existing.description || "").trim(),
|
||||
shared: typeof req.body.shared === "boolean" ? req.body.shared : existing.shared,
|
||||
points,
|
||||
},
|
||||
existing.id
|
||||
);
|
||||
|
||||
res.json({ selection });
|
||||
}));
|
||||
|
||||
app.delete("/api/selections/:selectionId", authenticateToken, wrap(async (req, res) => {
|
||||
const existing = await getSelection(redis, req.params.selectionId);
|
||||
|
||||
if (!existing) {
|
||||
return res.status(404).json({ error: "Selection not found." });
|
||||
}
|
||||
|
||||
if (existing.ownerId !== req.user.sub && req.user.role !== "admin") {
|
||||
return res.status(403).json({ error: "You can only delete your own selections." });
|
||||
}
|
||||
|
||||
await deleteSelection(redis, existing.id);
|
||||
res.status(204).send();
|
||||
}));
|
||||
|
||||
app.post("/api/trends/query", authenticateToken, wrap(async (req, res) => {
|
||||
const points = await resolveSelectionPoints(req.body, req.user);
|
||||
const trendData = await fetchTrendSeries(pool, redis, points, normalizeRange(req.body.range));
|
||||
res.json(trendData);
|
||||
}));
|
||||
|
||||
app.post("/api/points/latest", authenticateToken, wrap(async (req, res) => {
|
||||
const points = await resolveSelectionPoints(req.body, req.user);
|
||||
const values = await fetchLatestValues(pool, redis, points);
|
||||
res.json({ values });
|
||||
}));
|
||||
|
||||
app.get("/api/dashboard", authenticateToken, wrap(async (req, res) => {
|
||||
const dashboard = await getDashboard(redis, req.user.sub);
|
||||
res.json(dashboard);
|
||||
}));
|
||||
|
||||
app.put("/api/dashboard", authenticateToken, wrap(async (req, res) => {
|
||||
const dashboard = await saveDashboard(redis, req.user.sub, req.body || {});
|
||||
res.json(dashboard);
|
||||
}));
|
||||
|
||||
app.post("/api/dashboard/data", authenticateToken, wrap(async (req, res) => {
|
||||
const widgets = Array.isArray(req.body.widgets) ? req.body.widgets : [];
|
||||
const results = [];
|
||||
|
||||
for (const widget of widgets) {
|
||||
const points = widget.selectionId
|
||||
? await resolveSelectionPoints({ selectionId: widget.selectionId }, req.user)
|
||||
: assertPoints(widget.points || []);
|
||||
|
||||
if (widget.type === "chart") {
|
||||
const trendData = await fetchTrendSeries(pool, redis, points, normalizeRange(widget.range));
|
||||
results.push({ id: widget.id, type: widget.type, ...trendData });
|
||||
} else {
|
||||
const values = await fetchLatestValues(pool, redis, points);
|
||||
results.push({ id: widget.id, type: widget.type, values });
|
||||
}
|
||||
}
|
||||
|
||||
res.json({ widgets: results });
|
||||
}));
|
||||
|
||||
app.get("/api/preferences", authenticateToken, wrap(async (req, res) => {
|
||||
const preferences = await getPreferences(redis, req.user.sub);
|
||||
res.json(preferences);
|
||||
}));
|
||||
|
||||
app.put("/api/preferences", authenticateToken, wrap(async (req, res) => {
|
||||
const preferences = await savePreferences(redis, req.user.sub, req.body || {});
|
||||
res.json(preferences);
|
||||
}));
|
||||
|
||||
app.use((error, _req, res, _next) => {
|
||||
console.error(error);
|
||||
res.status(400).json({ error: error.message || "Unexpected error." });
|
||||
});
|
||||
|
||||
await ensureUserSchema();
|
||||
await ensureAdminUser();
|
||||
|
||||
app.listen(config.port, () => {
|
||||
console.log(`SE Local Trenddata API listening on port ${config.port}`);
|
||||
});
|
||||
@@ -0,0 +1,279 @@
|
||||
const META_PREFIX = "meta:isp:";
|
||||
|
||||
export function buildPointKey(isp, pointIndex) {
|
||||
return `${isp}:${pointIndex}`;
|
||||
}
|
||||
|
||||
export function sanitizeIspName(value) {
|
||||
const normalized = String(value || "").toLowerCase();
|
||||
|
||||
if (!/^isp\d+$/.test(normalized)) {
|
||||
throw new Error("Invalid ISP table name.");
|
||||
}
|
||||
|
||||
return normalized;
|
||||
}
|
||||
|
||||
export function sanitizePointIndex(value) {
|
||||
const pointIndex = Number(value);
|
||||
|
||||
if (!Number.isInteger(pointIndex) || pointIndex < 1 || pointIndex > 200) {
|
||||
throw new Error("Point index must be between 1 and 200.");
|
||||
}
|
||||
|
||||
return pointIndex;
|
||||
}
|
||||
|
||||
export function parseTrendValue(rawValue, factor = 1) {
|
||||
if (rawValue === null || rawValue === undefined) {
|
||||
return { raw: null, numeric: null };
|
||||
}
|
||||
|
||||
const raw = String(rawValue).trim();
|
||||
|
||||
if (!raw) {
|
||||
return { raw: null, numeric: null };
|
||||
}
|
||||
|
||||
const normalized = raw.replace(",", ".");
|
||||
const numeric = Number(normalized);
|
||||
const safeFactor = Number.isFinite(Number(factor)) ? Number(factor) : 1;
|
||||
|
||||
return {
|
||||
raw,
|
||||
numeric: Number.isFinite(numeric) ? numeric * safeFactor : null,
|
||||
};
|
||||
}
|
||||
|
||||
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
|
||||
FROM information_schema.TABLES t
|
||||
WHERE t.TABLE_SCHEMA = ?
|
||||
AND t.TABLE_NAME LIKE 'isp%'
|
||||
ORDER BY t.TABLE_NAME
|
||||
`,
|
||||
[dbName]
|
||||
);
|
||||
|
||||
const latestMap = new Map();
|
||||
|
||||
await Promise.all(
|
||||
tables.map(async (table) => {
|
||||
const safeTable = sanitizeIspName(table.tableName);
|
||||
const [rows] = await pool.query(
|
||||
`SELECT MAX(datum) AS latestTimestamp FROM \`${safeTable}\``
|
||||
);
|
||||
latestMap.set(safeTable, rows[0]?.latestTimestamp || null);
|
||||
})
|
||||
);
|
||||
|
||||
return tables.map((table) => ({
|
||||
tableName: sanitizeIspName(table.tableName),
|
||||
pointCount: Number(table.pointCount || 0),
|
||||
latestTimestamp: latestMap.get(sanitizeIspName(table.tableName)) || null,
|
||||
}));
|
||||
}
|
||||
|
||||
export function createDefaultMetadata(isp, pointCount = 200) {
|
||||
const points = {};
|
||||
|
||||
for (let index = 1; index <= pointCount; index += 1) {
|
||||
points[String(index)] = {
|
||||
alias: `Wert${index}`,
|
||||
unit: "",
|
||||
factor: 1,
|
||||
kind: "analog",
|
||||
min: 0,
|
||||
max: 100,
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
isp,
|
||||
displayName: isp.toUpperCase(),
|
||||
description: "",
|
||||
points,
|
||||
};
|
||||
}
|
||||
|
||||
export async function ensureIspMetadata(redis, isp, pointCount = 200) {
|
||||
const safeIsp = sanitizeIspName(isp);
|
||||
const key = `${META_PREFIX}${safeIsp}`;
|
||||
const existing = await redis.get(key);
|
||||
|
||||
if (existing) {
|
||||
const parsed = JSON.parse(existing);
|
||||
|
||||
for (let index = 1; index <= pointCount; index += 1) {
|
||||
if (!parsed.points?.[String(index)]) {
|
||||
parsed.points[String(index)] = createDefaultMetadata(safeIsp, pointCount).points[String(index)];
|
||||
}
|
||||
}
|
||||
|
||||
return parsed;
|
||||
}
|
||||
|
||||
const metadata = createDefaultMetadata(safeIsp, pointCount);
|
||||
await redis.set(key, JSON.stringify(metadata));
|
||||
return metadata;
|
||||
}
|
||||
|
||||
export async function updateIspMetadata(redis, isp, updater) {
|
||||
const metadata = await ensureIspMetadata(redis, isp);
|
||||
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) {
|
||||
const grouped = new Map();
|
||||
|
||||
points.forEach((point) => {
|
||||
const isp = sanitizeIspName(point.isp);
|
||||
const pointIndex = sanitizePointIndex(point.pointIndex);
|
||||
|
||||
if (!grouped.has(isp)) {
|
||||
grouped.set(isp, []);
|
||||
}
|
||||
|
||||
grouped.get(isp).push(pointIndex);
|
||||
});
|
||||
|
||||
const descriptors = [];
|
||||
|
||||
for (const [isp, pointIndexes] of grouped.entries()) {
|
||||
const metadata = await ensureIspMetadata(redis, isp);
|
||||
|
||||
pointIndexes.forEach((pointIndex) => {
|
||||
const pointMeta = metadata.points[String(pointIndex)];
|
||||
descriptors.push({
|
||||
isp,
|
||||
pointIndex,
|
||||
key: buildPointKey(isp, pointIndex),
|
||||
displayName: metadata.displayName,
|
||||
alias: pointMeta.alias,
|
||||
unit: pointMeta.unit,
|
||||
factor: pointMeta.factor,
|
||||
kind: pointMeta.kind,
|
||||
min: pointMeta.min,
|
||||
max: pointMeta.max,
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
return descriptors;
|
||||
}
|
||||
|
||||
export async function fetchTrendSeries(pool, redis, points, range) {
|
||||
const descriptors = await getPointDescriptors(redis, points);
|
||||
const grouped = descriptors.reduce((accumulator, descriptor) => {
|
||||
if (!accumulator.has(descriptor.isp)) {
|
||||
accumulator.set(descriptor.isp, []);
|
||||
}
|
||||
|
||||
accumulator.get(descriptor.isp).push(descriptor);
|
||||
return accumulator;
|
||||
}, new Map());
|
||||
|
||||
const mergedRows = new Map();
|
||||
|
||||
for (const [isp, pointDescriptors] of grouped.entries()) {
|
||||
const columns = pointDescriptors
|
||||
.map((descriptor) => `\`Wert${descriptor.pointIndex}\``)
|
||||
.join(", ");
|
||||
const limit = range.maxRows || 4000;
|
||||
|
||||
const [rows] = await pool.query(
|
||||
`SELECT datum, ${columns} FROM \`${isp}\` WHERE datum BETWEEN ? AND ? ORDER BY datum DESC LIMIT ?`,
|
||||
[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);
|
||||
|
||||
pointDescriptors.forEach((descriptor) => {
|
||||
const rawValue = row[`Wert${descriptor.pointIndex}`];
|
||||
const parsed = parseTrendValue(rawValue, descriptor.factor);
|
||||
target[descriptor.key] = parsed.numeric;
|
||||
target[`${descriptor.key}:raw`] = parsed.raw;
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
return {
|
||||
series: descriptors,
|
||||
rows: [...mergedRows.values()].sort((left, right) =>
|
||||
left.timestamp.localeCompare(right.timestamp)
|
||||
),
|
||||
};
|
||||
}
|
||||
|
||||
export async function fetchLatestValues(pool, redis, points) {
|
||||
const descriptors = await getPointDescriptors(redis, points);
|
||||
const grouped = descriptors.reduce((accumulator, descriptor) => {
|
||||
if (!accumulator.has(descriptor.isp)) {
|
||||
accumulator.set(descriptor.isp, []);
|
||||
}
|
||||
|
||||
accumulator.get(descriptor.isp).push(descriptor);
|
||||
return accumulator;
|
||||
}, new Map());
|
||||
|
||||
const values = [];
|
||||
|
||||
for (const [isp, pointDescriptors] of grouped.entries()) {
|
||||
const columns = pointDescriptors
|
||||
.map((descriptor) => `\`Wert${descriptor.pointIndex}\``)
|
||||
.join(", ");
|
||||
|
||||
const [rows] = await pool.query(
|
||||
`SELECT datum, ${columns} FROM \`${isp}\` ORDER BY datum DESC LIMIT 1`
|
||||
);
|
||||
|
||||
const row = rows[0];
|
||||
|
||||
pointDescriptors.forEach((descriptor) => {
|
||||
const parsed = parseTrendValue(row?.[`Wert${descriptor.pointIndex}`], descriptor.factor);
|
||||
values.push({
|
||||
...descriptor,
|
||||
timestamp: row?.datum ? new Date(row.datum).toISOString() : null,
|
||||
raw: parsed.raw,
|
||||
value: parsed.numeric,
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
return values;
|
||||
}
|
||||
|
||||
export function normalizeRange(input = {}) {
|
||||
const presetDays = Number(input.presetDays || 1);
|
||||
const maxRows = Number(input.maxRows || 4000);
|
||||
const now = new Date();
|
||||
const to = input.to ? new Date(input.to) : now;
|
||||
const from = input.from
|
||||
? new Date(input.from)
|
||||
: new Date(now.getTime() - presetDays * 24 * 60 * 60 * 1000);
|
||||
|
||||
return {
|
||||
from,
|
||||
to,
|
||||
maxRows: Number.isFinite(maxRows) ? Math.max(100, Math.min(maxRows, 10000)) : 4000,
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,67 @@
|
||||
const SELECTION_SET_KEY = "selection:ids";
|
||||
const SELECTION_SEQ_KEY = "selection:seq";
|
||||
|
||||
function safeParse(jsonValue) {
|
||||
return jsonValue ? JSON.parse(jsonValue) : null;
|
||||
}
|
||||
|
||||
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}`)));
|
||||
|
||||
return items
|
||||
.map((item) => safeParse(item))
|
||||
.filter(Boolean)
|
||||
.filter((selection) => selection.shared || selection.ownerId === user.sub || user.role === "admin")
|
||||
.sort((left, right) => right.updatedAt.localeCompare(left.updatedAt));
|
||||
}
|
||||
|
||||
export async function getSelection(redis, selectionId) {
|
||||
return safeParse(await redis.get(`selection:${selectionId}`));
|
||||
}
|
||||
|
||||
export async function saveSelection(redis, selection, existingSelectionId) {
|
||||
const selectionId = existingSelectionId || `sel-${await redis.incr(SELECTION_SEQ_KEY)}`;
|
||||
const now = new Date().toISOString();
|
||||
const payload = {
|
||||
...selection,
|
||||
id: selectionId,
|
||||
updatedAt: now,
|
||||
createdAt: selection.createdAt || now,
|
||||
};
|
||||
|
||||
await redis.set(`selection:${selectionId}`, JSON.stringify(payload));
|
||||
await redis.sadd(SELECTION_SET_KEY, selectionId);
|
||||
return payload;
|
||||
}
|
||||
|
||||
export async function deleteSelection(redis, selectionId) {
|
||||
await redis.del(`selection:${selectionId}`);
|
||||
await redis.srem(SELECTION_SET_KEY, selectionId);
|
||||
}
|
||||
|
||||
export async function getDashboard(redis, userId) {
|
||||
return safeParse(await redis.get(`dashboard:${userId}`)) || { widgets: [] };
|
||||
}
|
||||
|
||||
export async function saveDashboard(redis, userId, dashboard) {
|
||||
const payload = {
|
||||
widgets: dashboard.widgets || [],
|
||||
updatedAt: new Date().toISOString(),
|
||||
};
|
||||
await redis.set(`dashboard:${userId}`, JSON.stringify(payload));
|
||||
return payload;
|
||||
}
|
||||
|
||||
export async function getPreferences(redis, userId) {
|
||||
return safeParse(await redis.get(`preferences:${userId}`)) || { theme: "system" };
|
||||
}
|
||||
|
||||
export async function savePreferences(redis, userId, preferences) {
|
||||
const payload = {
|
||||
theme: preferences.theme || "system",
|
||||
updatedAt: new Date().toISOString(),
|
||||
};
|
||||
await redis.set(`preferences:${userId}`, JSON.stringify(payload));
|
||||
return payload;
|
||||
}
|
||||
@@ -0,0 +1,7 @@
|
||||
import Redis from "ioredis";
|
||||
import { config } from "./config.js";
|
||||
|
||||
export const redis = new Redis(config.redisUrl, {
|
||||
maxRetriesPerRequest: 2,
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user