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; let lastKnownGoodStore = { sources: [], datapoints: [] }; 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 structuredClone(lastKnownGoodStore); } try { const store = JSON.parse(readFileSync(sourcesFile, "utf8")); lastKnownGoodStore = { sources: Array.isArray(store.sources) ? store.sources : [], datapoints: Array.isArray(store.datapoints) ? store.datapoints : [], }; return structuredClone(lastKnownGoodStore); } catch (error) { console.warn("Quellenkonfiguration konnte nicht gelesen werden, letzter gültiger Stand wird verwendet:", error.message); return structuredClone(lastKnownGoodStore); } } 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}).`);