API added
This commit is contained in:
@@ -0,0 +1,265 @@
|
||||
import { existsSync, readFileSync } from "node:fs";
|
||||
import mysql from "mysql2/promise";
|
||||
|
||||
const sourcesFile = process.env.SOURCES_FILE || "/app/collector-data/sources.json";
|
||||
const pollTickMs = Math.max(500, Number(process.env.SYMCON_POLL_TICK_MS || 1000));
|
||||
const maxConcurrency = Math.max(1, Math.min(16, Number(process.env.SYMCON_MAX_CONCURRENCY || 4)));
|
||||
const sourcePollState = new Map();
|
||||
const valueCache = new Map();
|
||||
let polling = false;
|
||||
|
||||
const db = mysql.createPool({
|
||||
host: process.env.DB_HOST || "mariadb",
|
||||
port: Number(process.env.DB_PORT || 3306),
|
||||
database: process.env.DB_NAME || "wago",
|
||||
user: process.env.DB_USER || "root",
|
||||
password: process.env.DB_PASSWORD || "SE3112",
|
||||
charset: "utf8mb4",
|
||||
waitForConnections: true,
|
||||
connectionLimit: Math.max(2, maxConcurrency + 1),
|
||||
});
|
||||
|
||||
function normalizeText(value) {
|
||||
return String(value ?? "").trim();
|
||||
}
|
||||
|
||||
function normalizeInterval(value) {
|
||||
const number = Number(value);
|
||||
return Number.isFinite(number) && number > 0 ? Math.round(number) : 60;
|
||||
}
|
||||
|
||||
function sanitizeTrendTable(value) {
|
||||
const table = String(value || "").toLowerCase();
|
||||
if (!/^(isp[a-z0-9_]*|trend_[a-z0-9_]+)$/.test(table)) {
|
||||
throw new Error("Ungültiger Trendtabellenname.");
|
||||
}
|
||||
return table;
|
||||
}
|
||||
|
||||
async function ensureWorkerStateTable() {
|
||||
await db.query(`
|
||||
CREATE TABLE IF NOT EXISTS \`source_worker_state\` (
|
||||
\`source_id\` VARCHAR(191) NOT NULL,
|
||||
\`worker\` VARCHAR(64) NOT NULL,
|
||||
\`last_poll_at\` TIMESTAMP NULL,
|
||||
\`last_poll_status\` VARCHAR(255) NOT NULL DEFAULT '',
|
||||
\`last_poll_error\` TEXT NULL,
|
||||
\`updated_at\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
|
||||
PRIMARY KEY (\`source_id\`)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci
|
||||
`);
|
||||
}
|
||||
|
||||
async function updateWorkerState(source, status, error = "") {
|
||||
await db.query(
|
||||
`
|
||||
INSERT INTO \`source_worker_state\` (\`source_id\`, \`worker\`, \`last_poll_at\`, \`last_poll_status\`, \`last_poll_error\`)
|
||||
VALUES (?, 'seltd-symcon-collector', CURRENT_TIMESTAMP, ?, ?)
|
||||
ON DUPLICATE KEY UPDATE
|
||||
\`worker\` = VALUES(\`worker\`),
|
||||
\`last_poll_at\` = CURRENT_TIMESTAMP,
|
||||
\`last_poll_status\` = VALUES(\`last_poll_status\`),
|
||||
\`last_poll_error\` = VALUES(\`last_poll_error\`)
|
||||
`,
|
||||
[String(source.id), String(status || ""), String(error || "")],
|
||||
);
|
||||
}
|
||||
|
||||
function readStore() {
|
||||
if (!existsSync(sourcesFile)) {
|
||||
return { sources: [], datapoints: [] };
|
||||
}
|
||||
try {
|
||||
const store = JSON.parse(readFileSync(sourcesFile, "utf8"));
|
||||
return {
|
||||
sources: Array.isArray(store.sources) ? store.sources : [],
|
||||
datapoints: Array.isArray(store.datapoints) ? store.datapoints : [],
|
||||
};
|
||||
} catch (error) {
|
||||
console.warn("Quellenkonfiguration konnte nicht gelesen werden:", error.message);
|
||||
return { sources: [], datapoints: [] };
|
||||
}
|
||||
}
|
||||
|
||||
function shouldPoll(source, now) {
|
||||
const lastPoll = sourcePollState.get(source.id) || 0;
|
||||
return now - lastPoll >= normalizeInterval(source.pollIntervalSeconds) * 1000;
|
||||
}
|
||||
|
||||
function matchesLogCondition(condition, value) {
|
||||
const text = String(condition || "").trim();
|
||||
if (!text) return true;
|
||||
const numeric = Number(value);
|
||||
const match = text.match(/^(?:value\s*)?(>=|<=|>|<|==|!=)\s*(-?\d+(?:[.,]\d+)?)$/i);
|
||||
if (!match || !Number.isFinite(numeric)) return true;
|
||||
const target = Number(match[2].replace(",", "."));
|
||||
if (match[1] === ">") return numeric > target;
|
||||
if (match[1] === ">=") return numeric >= target;
|
||||
if (match[1] === "<") return numeric < target;
|
||||
if (match[1] === "<=") return numeric <= target;
|
||||
if (match[1] === "==") return numeric === target;
|
||||
if (match[1] === "!=") return numeric !== target;
|
||||
return true;
|
||||
}
|
||||
|
||||
function buildEndpoint(source) {
|
||||
const rawHost = normalizeText(source.host);
|
||||
if (!rawHost) throw new Error("IP-Symcon Host oder URL fehlt.");
|
||||
const input = /^https?:\/\//i.test(rawHost) ? rawHost : `http://${rawHost}`;
|
||||
const url = new URL(input);
|
||||
if (!url.port) url.port = String(Number(source.port || 3777) || 3777);
|
||||
const path = url.pathname.replace(/\/+$/, "");
|
||||
url.pathname = !path || path === "/" ? "/hook/api" : (/\/hook\/api$/i.test(path) ? path : `${path}/hook/api`);
|
||||
url.search = "";
|
||||
url.hash = "";
|
||||
return url.toString();
|
||||
}
|
||||
|
||||
async function requestApi(source, payload) {
|
||||
const apiKey = normalizeText(source.options?.ipsymconApiKey || source.ipsymconApiKey);
|
||||
if (!apiKey) throw new Error("IP-Symcon API-Schlüssel fehlt.");
|
||||
const controller = new AbortController();
|
||||
const timeout = setTimeout(() => controller.abort(), 30000);
|
||||
try {
|
||||
const response = await fetch(buildEndpoint(source), {
|
||||
method: "POST",
|
||||
headers: { "Content-Type": "application/json; charset=utf-8" },
|
||||
body: JSON.stringify({ api_key: apiKey, ...payload }),
|
||||
signal: controller.signal,
|
||||
});
|
||||
const data = await response.json().catch(() => ({}));
|
||||
if (!response.ok || data.ok !== true) {
|
||||
throw new Error(data?.error?.message || data?.error || `IP-Symcon API HTTP ${response.status}`);
|
||||
}
|
||||
return data;
|
||||
} finally {
|
||||
clearTimeout(timeout);
|
||||
}
|
||||
}
|
||||
|
||||
async function readValue(source, datapoint) {
|
||||
const variableId = Number(datapoint.nodeId || datapoint.address);
|
||||
if (!Number.isInteger(variableId) || variableId <= 0) {
|
||||
throw new Error(`Ungültige Variablen-ID bei ${datapoint.alias || datapoint.name}.`);
|
||||
}
|
||||
const response = await requestApi(source, { function: "getvalue", id: variableId });
|
||||
const raw = response.value ?? response.variable?.value;
|
||||
if (typeof raw === "boolean") return raw ? 1 : 0;
|
||||
const numeric = Number(raw);
|
||||
return Number.isFinite(numeric) ? numeric * (Number(datapoint.scale || 1) || 1) : raw === null || raw === undefined ? null : String(raw);
|
||||
}
|
||||
|
||||
async function mapLimit(items, limit, task) {
|
||||
const result = [];
|
||||
let nextIndex = 0;
|
||||
const workers = Array.from({ length: Math.min(limit, items.length) }, async () => {
|
||||
while (nextIndex < items.length) {
|
||||
const item = items[nextIndex++];
|
||||
result.push(await task(item));
|
||||
}
|
||||
});
|
||||
await Promise.all(workers);
|
||||
return result;
|
||||
}
|
||||
|
||||
async function ensureNarrowTrendTable(tableName) {
|
||||
const table = sanitizeTrendTable(tableName);
|
||||
const [columns] = await db.query(
|
||||
"SELECT COLUMN_NAME AS name FROM information_schema.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ?",
|
||||
[table],
|
||||
);
|
||||
const names = new Set(columns.map((entry) => String(entry.name).toLowerCase()));
|
||||
if (names.size && Array.from(names).some((name) => /^wert\d+$/.test(name))) {
|
||||
throw new Error(`${table} ist eine Legacy-Tabelle und keine IP-Symcon-Trendtabelle.`);
|
||||
}
|
||||
if (!names.size) {
|
||||
await db.query(`
|
||||
CREATE TABLE \`${table}\` (
|
||||
\`id\` BIGINT NOT NULL AUTO_INCREMENT,
|
||||
\`datum\` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||||
\`userlevel\` INT NOT NULL DEFAULT 0,
|
||||
\`source_id\` VARCHAR(191) NOT NULL,
|
||||
\`point_index\` INT NOT NULL,
|
||||
\`alias\` VARCHAR(255) NOT NULL,
|
||||
\`value\` TEXT NULL,
|
||||
PRIMARY KEY (\`id\`),
|
||||
KEY \`idx_${table}_datum\` (\`datum\`),
|
||||
KEY \`idx_${table}_point_datum\` (\`point_index\`, \`datum\`)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci
|
||||
`);
|
||||
}
|
||||
return table;
|
||||
}
|
||||
|
||||
async function writeReadings(source, readings) {
|
||||
if (!readings.length) return 0;
|
||||
const table = await ensureNarrowTrendTable(source.tableName);
|
||||
const placeholders = readings.map(() => "(?, ?, ?, ?, ?)").join(", ");
|
||||
const values = readings.flatMap((reading) => [0, String(source.id), Number(reading.pointIndex), String(reading.alias), String(reading.value)]);
|
||||
await db.query(`INSERT INTO \`${table}\` (\`userlevel\`, \`source_id\`, \`point_index\`, \`alias\`, \`value\`) VALUES ${placeholders}`, values);
|
||||
return readings.length;
|
||||
}
|
||||
|
||||
async function pollSource(source, points) {
|
||||
const errors = [];
|
||||
const results = await mapLimit(points, maxConcurrency, async (datapoint) => {
|
||||
try {
|
||||
return { datapoint, value: await readValue(source, datapoint) };
|
||||
} catch (error) {
|
||||
errors.push(`${datapoint.alias || datapoint.address}: ${error.message}`);
|
||||
return null;
|
||||
}
|
||||
});
|
||||
|
||||
const readings = [];
|
||||
for (const result of results) {
|
||||
if (!result) continue;
|
||||
const { datapoint, value } = result;
|
||||
const key = `${source.id}:${datapoint.id}`;
|
||||
const serialized = String(value);
|
||||
const writeMode = datapoint.writeMode || source.writeMode || "interval";
|
||||
if (matchesLogCondition(datapoint.condition, value) && (writeMode !== "cov" || valueCache.get(key) !== serialized)) {
|
||||
readings.push({ pointIndex: datapoint.pointIndex, alias: datapoint.alias || datapoint.name || `Wert${datapoint.pointIndex}`, value });
|
||||
}
|
||||
valueCache.set(key, serialized);
|
||||
}
|
||||
|
||||
const written = await writeReadings(source, readings);
|
||||
const suffix = errors.length ? `, ${errors.length} Fehler` : "";
|
||||
const status = errors.length ? `teilweise ok: ${written} Wert(e) geschrieben${suffix}` : `ok: ${written} Wert(e) geschrieben`;
|
||||
await updateWorkerState(source, status, errors.slice(0, 5).join(" | "));
|
||||
console.log(`[${source.name || source.id}] ${written} Wert(e) geschrieben${suffix}`);
|
||||
if (errors.length) console.warn(`[${source.name || source.id}]`, errors.slice(0, 5).join(" | "));
|
||||
}
|
||||
|
||||
async function pollDueSources() {
|
||||
if (polling) return;
|
||||
polling = true;
|
||||
try {
|
||||
const store = readStore();
|
||||
const now = Date.now();
|
||||
const sources = store.sources.filter((source) => source.protocol === "ipsymcon" && shouldPoll(source, now));
|
||||
for (const source of sources) {
|
||||
const points = store.datapoints.filter((point) => point.sourceId === source.id && point.protocol === "ipsymcon" && point.enabled === true);
|
||||
sourcePollState.set(source.id, now);
|
||||
if (!points.length) {
|
||||
await updateWorkerState(source, "idle: keine Datenpunkte");
|
||||
continue;
|
||||
}
|
||||
try {
|
||||
await pollSource(source, points);
|
||||
} catch (error) {
|
||||
await updateWorkerState(source, "error", error.message || "Polling fehlgeschlagen.").catch(() => undefined);
|
||||
console.error(`[${source.name || source.id}] Polling fehlgeschlagen:`, error.message);
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
polling = false;
|
||||
}
|
||||
}
|
||||
|
||||
setInterval(() => pollDueSources().catch((error) => console.error("IP-Symcon Polling fehlgeschlagen:", error.message)), pollTickMs);
|
||||
ensureWorkerStateTable()
|
||||
.then(() => pollDueSources())
|
||||
.catch((error) => console.error("IP-Symcon Start-Polling fehlgeschlagen:", error.message));
|
||||
console.log(`SE LTD Symcon Collector gestartet (Intervall-Prüfung ${pollTickMs} ms, Parallelität ${maxConcurrency}).`);
|
||||
Reference in New Issue
Block a user