/
Krams
/
nodemon2
Обзор
Документация
Войти
/
Krams
/
nodemon2
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
backend/services/NodeService.js
1 796 строк
59 KB
Иван Новиков
init
11 июн 2026, 13:08
11 июн 2026, 13:08
900f716
Код
Авторство
О чём код?
import db from '../database.js'; import { safeJsonParse } from './DebugDump.js'; import { syncNodeStateIntervalFromCurrentStatus } from './NodeReliabilityService.js'; import { exec } from 'child_process'; import { NodeSSH } from 'node-ssh'; import sshConfig from './sshConfig.js'; import { canWriteByDedupKey, getAuditLogsPaginated, writeNodeAuditLog } from './AuditLogService.js'; import { io } from '../socketServer.js'; function getNodeFeaturesByNodeIds(nodeIds = []) { if (!Array.isArray(nodeIds) || nodeIds.length === 0) { return Promise.resolve({}); } const placeholders = nodeIds.map(() => '?').join(', '); return new Promise((resolve, reject) => { db.all( `SELECT m.node_id AS node_id, f.id AS id, f.name AS name, f.emoji AS emoji, f.color AS color, f.sort_order AS sort_order FROM node_feature_map m JOIN node_feature f ON f.id = m.feature_id WHERE m.node_id IN (${placeholders}) ORDER BY f.sort_order ASC, f.id ASC`, nodeIds, (err, rows) => { if (err) { if (String(err.message || '').includes('no such table')) { return resolve({}); } return reject(err); } const featureMap = {}; (rows || []).forEach((row) => { const nodeId = Number.parseInt(row.node_id, 10); if (!Number.isFinite(nodeId)) return; if (!featureMap[nodeId]) { featureMap[nodeId] = []; } featureMap[nodeId].push({ id: Number.parseInt(row.id, 10) || row.id, name: row.name, emoji: row.emoji, color: row.color || '#2f8f5a', sortOrder: row.sort_order, isActive: Number(row.is_active) !== 0, }); }); resolve(featureMap); } ); }); } async function applyFeaturesToNodes(nodes = []) { const nodeIds = nodes .map((node) => Number.parseInt(node?.id, 10)) .filter(Number.isFinite); const featureMap = await getNodeFeaturesByNodeIds(nodeIds); return nodes.map((node) => { const nodeId = Number.parseInt(node?.id, 10); const features = Number.isFinite(nodeId) ? (featureMap[nodeId] || []) : []; return { ...node, features, feature_ids: features.map((feature) => feature.id), }; }); } function normalizeMetricText(value) { return String(value || '').trim(); } function normalizeUnixToSec(value) { const ts = Number(value); if (!Number.isFinite(ts) || ts <= 0) return null; return ts > 1e12 ? Math.floor(ts / 1000) : Math.floor(ts); } function calcDaysInState(startedAtRaw) { const ts = normalizeUnixToSec(startedAtRaw); if (!Number.isFinite(ts) || ts <= 0) return null; const now = Math.floor(Date.now() / 1000); const delta = Math.max(0, now - ts); return Number((delta / 86400).toFixed(1)); } export async function emitNodesUpdateSafe() { try { const allNodes = await getAllNodesNested(); io.emit('nodes_update', allNodes); } catch (error) { if (String(process.env.DEBUG_SOCKET_EMIT || '').trim() === '1') { console.error('[socket] failed to emit nodes_update', error); } } } function parseReservationNames(raw = '') { const text = String(raw || '').trim(); if (!text) return []; const parts = text .split(/[,\s]+/g) .map((s) => s.trim()) .filter(Boolean); return [...new Set(parts)]; } function normalizeReservationKey(value) { const names = parseReservationNames(value); if (!names.length) return ''; return names.map((name) => name.toLowerCase()).sort().join('\0'); } function normalizeReasonKey(value) { return String(value || '').trim().toLowerCase(); } function findSinceFromDescIntervals(rowsDesc = [], normalizeFn, currentRaw, pickValue) { const targetKey = normalizeFn(currentRaw); if (!targetKey) return null; let since = null; for (const row of rowsDesc) { const rowKey = normalizeFn(pickValue(row)); if (rowKey === targetKey) { since = Number(row?.started_at) || null; } else if (since !== null) { break; } } return since; } export function computeStateDurationsFromIntervalHistory(intervalRowsDesc = [], reservation, reason) { const reservationSince = findSinceFromDescIntervals( intervalRowsDesc, normalizeReservationKey, reservation, (row) => row?.reservation ); const reasonSince = findSinceFromDescIntervals( intervalRowsDesc, normalizeReasonKey, reason, (row) => row?.reason ); return { reservation_since: reservationSince, reservation_days_in_state: calcDaysInState(reservationSince), reason_since: reasonSince, reason_days_in_state: calcDaysInState(reasonSince), }; } export async function getStateDurationsForNode(nodeId, reservation, reason) { const safeNodeId = Number.parseInt(nodeId, 10); if (!Number.isFinite(safeNodeId)) { return computeStateDurationsFromIntervalHistory([], reservation, reason); } const rows = await new Promise((resolve) => { db.all( `SELECT started_at, reason, reservation FROM node_state_interval WHERE node_id = ? ORDER BY started_at DESC, id DESC LIMIT 1000`, [safeNodeId], (err, result) => resolve(err ? [] : (result || [])) ); }); return computeStateDurationsFromIntervalHistory(rows, reservation, reason); } function debugNodeDaysLog(message, payload) { if (String(process.env.DEBUG_NODE_STATE_DAYS || '').trim() !== '1') return; try { console.log(`[node-days] ${message}`, payload); } catch { } } function parseReasonTimeToSec(value) { if (value === null || typeof value === 'undefined') return null; const raw = String(value).trim(); if (!raw) return null; if (/^\d+$/.test(raw)) { const num = Number(raw); return normalizeUnixToSec(num); } const ms = Date.parse(raw); if (!Number.isFinite(ms)) return null; return normalizeUnixToSec(ms); } async function enrichNodesWithStatus(rows = []) { let rowsWithFeatures = []; rowsWithFeatures = await applyFeaturesToNodes(rows || []); const nodeIds = rowsWithFeatures.map((n) => n?.id).filter(Number.isFinite); const commentUpdatedAtByNodeId = await new Promise((resolveLogs) => { if (nodeIds.length === 0) return resolveLogs({}); const placeholders = nodeIds.map(() => '?').join(', '); db.all( `SELECT resolved_id AS node_id, MAX(date) AS comment_updated_at FROM ( SELECT date, log, COALESCE(logical_node_id, CAST(node_id AS INTEGER)) AS resolved_id FROM node_log WHERE (logical_node_id IN (${placeholders}) OR CAST(node_id AS TEXT) IN (${placeholders})) AND log LIKE 'Комментарий изменён:%' ) t WHERE resolved_id IS NOT NULL GROUP BY resolved_id`, [...nodeIds, ...nodeIds.map((nid) => String(nid))], (logErr, logRows) => { if (logErr) return resolveLogs({}); const map = {}; (logRows || []).forEach((row) => { const id = Number(row.node_id); const ts = row.comment_updated_at; if (Number.isFinite(id) && ts !== null && ts !== undefined) { map[id] = ts; } }); resolveLogs(map); } ); }); const statusMap = await new Promise((resolveStatus, rejectStatus) => { db.all('SELECT node_id, json_data, updated_at FROM node_status', [], (err2, statusRows) => { if (err2) return rejectStatus(err2); const map = {}; (statusRows || []).forEach((row) => { try { const status = safeJsonParse(`enrichNodesWithStatus_row_${row.node_id}_before_parse`, row.json_data); status.updated_at = row.updated_at; map[row.node_id] = status; } catch { } }); resolveStatus(map); }); }); const flatNodeIds = rowsWithFeatures.map((n) => n?.id).filter(Number.isFinite); const intervalsByNodeId = await new Promise((resolveIntervals) => { if (flatNodeIds.length === 0) return resolveIntervals({}); const placeholders = flatNodeIds.map(() => '?').join(', '); db.all( `SELECT node_id, started_at, reservation, reason FROM node_state_interval WHERE node_id IN (${placeholders}) ORDER BY node_id ASC, started_at DESC, id DESC`, flatNodeIds, (intervalErr, intervalRows) => { if (intervalErr) return resolveIntervals({}); const map = {}; (intervalRows || []).forEach((row) => { const id = Number(row.node_id); if (!Number.isFinite(id)) return; if (!map[id]) map[id] = []; map[id].push(row); }); debugNodeDaysLog('interval-history-loaded', { nodeCount: flatNodeIds.length, intervalRows: Array.isArray(intervalRows) ? intervalRows.length : 0, sample: (intervalRows || []).slice(0, 5), }); resolveIntervals(map); } ); }); const nowSec = Math.floor(Date.now() / 1000); let logged = 0; return rowsWithFeatures.map((node) => { const status = statusMap[node.id]; let nodeWithStatus = { ...node, comment_updated_at: commentUpdatedAtByNodeId[node.id] ?? null, }; if (status) { const mergedStatus = { ...status }; if (Number(node.monitoring_enabled) === 0) { mergedStatus.in_queue = ''; } nodeWithStatus = { ...nodeWithStatus, ...mergedStatus }; } const currentReservation = normalizeMetricText(nodeWithStatus?.reservation); const currentReason = normalizeMetricText(nodeWithStatus?.reason); const intervalHistory = intervalsByNodeId?.[node.id] || []; const stateDurations = computeStateDurationsFromIntervalHistory( intervalHistory, currentReservation, currentReason ); nodeWithStatus = { ...nodeWithStatus, ...stateDurations, }; const reasonTimeSec = parseReasonTimeToSec(nodeWithStatus?.reason_time); if (currentReason && !nodeWithStatus.reason_since && reasonTimeSec) { nodeWithStatus.reason_since = reasonTimeSec; nodeWithStatus.reason_days_in_state = calcDaysInState(reasonTimeSec); } if (logged < 12 && (currentReservation || currentReason)) { const openInterval = intervalHistory[0] || null; logged += 1; debugNodeDaysLog('node-sample', { node_id: node.id, name: node.name, reservation: currentReservation || null, reason: currentReason || null, reason_time_raw: nodeWithStatus?.reason_time ?? null, reason_time_sec: reasonTimeSec, reservation_key: normalizeReservationKey(currentReservation), open_reservation_key: normalizeReservationKey(openInterval?.reservation), reason_key: normalizeReasonKey(currentReason), open_reason_key: normalizeReasonKey(openInterval?.reason), interval_history_rows: intervalHistory.length, reservation_since: nodeWithStatus.reservation_since, reservation_days: nodeWithStatus.reservation_days_in_state, reason_since: nodeWithStatus.reason_since, reason_days: nodeWithStatus.reason_days_in_state, now_sec: nowSec, }); } return nodeWithStatus; }); } export async function getAllNodesFlat() { return new Promise((resolve, reject) => { db.all('SELECT * FROM node', [], async (err, rows) => { if (err) return reject(err); try { const items = await enrichNodesWithStatus(rows || []); resolve({ items, total: items.length }); } catch (featureErr) { reject(featureErr); } }); }); } export async function getAllNodesNested() { return new Promise((resolve, reject) => { db.all('SELECT * FROM node', [], async (err, rows) => { if (err) return reject(err); try { const enriched = await enrichNodesWithStatus(rows || []); const data = Array.from({ length: 9 }, () => Array.from({ length: 8 }, () => Array.from({ length: 6 }, () => null) ) ); enriched.forEach((nodeWithStatus) => { const rack = nodeWithStatus.rack - 1; const shelf = nodeWithStatus.shelf - 1; const position = nodeWithStatus.position - 1; if ( rack >= 0 && rack < 9 && shelf >= 0 && shelf < 8 && position >= 0 && position < 6 ) { data[rack][shelf][position] = nodeWithStatus; } }); resolve(data); } catch (featureErr) { reject(featureErr); } }); }); } export async function updateNode(id, patch = {}) { const { mac, guid, state, comment, state_manual: stateManual, state_manual_updated_at: stateManualUpdatedAt, featureIds = [], monitoring_enabled: monitoringEnabled, binding_mac_interface: bindingMacInterface, visible_name: visibleName, } = patch; const shouldUpdateFeatures = Object.prototype.hasOwnProperty.call(patch, 'featureIds'); const normalizedFeatureIds = shouldUpdateFeatures && Array.isArray(featureIds) ? [...new Set(featureIds.map((value) => Number.parseInt(value, 10)).filter(Number.isFinite))] : []; const setParts = []; const setValues = []; if (typeof mac !== 'undefined') { setParts.push('mac = ?'); setValues.push(mac); } if (typeof guid !== 'undefined') { setParts.push('guid = ?'); setValues.push(guid); } if (typeof state !== 'undefined') { setParts.push('state = ?'); setValues.push(state); } if (typeof comment !== 'undefined') { setParts.push('comment = ?'); setValues.push(comment); } if (typeof stateManual !== 'undefined') { setParts.push('state_manual = ?'); setValues.push(Number(stateManual) === 1 || stateManual === true ? 1 : 0); } if (typeof stateManualUpdatedAt !== 'undefined') { const ts = Number.parseInt(stateManualUpdatedAt, 10); setParts.push('state_manual_updated_at = ?'); setValues.push(Number.isFinite(ts) ? ts : null); } if (typeof monitoringEnabled !== 'undefined') { setParts.push('monitoring_enabled = ?'); setValues.push(Number(monitoringEnabled) === 1 || monitoringEnabled === true ? 1 : 0); } if (typeof bindingMacInterface !== 'undefined') { const iface = String(bindingMacInterface || '').trim() || 'eth0'; setParts.push('binding_mac_interface = ?'); setValues.push(iface); } if (typeof visibleName !== 'undefined') { setParts.push('visible_name = ?'); setValues.push(visibleName === null || visibleName === '' ? null : String(visibleName)); } if (setParts.length === 0) { return Promise.reject(new Error('Nothing to update on node')); } setValues.push(id); return new Promise((resolve, reject) => { db.serialize(() => { db.run('BEGIN TRANSACTION'); db.run( `UPDATE node SET ${setParts.join(', ')} WHERE id = ?`, setValues, function (updateErr) { if (updateErr) { db.run('ROLLBACK'); return reject(updateErr); } if (this.changes === 0) { db.run('ROLLBACK'); return resolve({ changes: 0 }); } const finishCommit = () => { db.run('COMMIT', (commitErr) => { if (commitErr) { db.run('ROLLBACK'); return reject(commitErr); } resolve({ changes: 1 }); }); }; if (!shouldUpdateFeatures) { return finishCommit(); } db.run( `DELETE FROM node_feature_map WHERE node_id = ?`, [id], (deleteErr) => { if (deleteErr) { db.run('ROLLBACK'); return reject(deleteErr); } const insertFeature = (index) => { if (index >= normalizedFeatureIds.length) { return finishCommit(); } db.run( `INSERT INTO node_feature_map (node_id, feature_id) VALUES (?, ?)`, [id, normalizedFeatureIds[index]], (insertErr) => { if (insertErr) { db.run('ROLLBACK'); return reject(insertErr); } insertFeature(index + 1); } ); }; insertFeature(0); } ); } ); }); }); } export async function getNodeById(id) { return new Promise((resolve, reject) => { db.get('SELECT * FROM node WHERE id = ?', [id], async (err, row) => { if (err) return reject(err); if (!row) { return reject(new Error('Node not found')); } try { const [nodeWithFeatures] = await applyFeaturesToNodes([row]); resolve(nodeWithFeatures); } catch (featureErr) { reject(featureErr); } }); }); } export async function getNodeLogsById(id) { return new Promise((resolve, reject) => { db.all( `SELECT * FROM node_log WHERE logical_node_id = ? OR (logical_node_id IS NULL AND CAST(node_id AS TEXT) = CAST(? AS TEXT)) ORDER BY date DESC`, [id, id], (err, rows) => { if (err) return reject(err); if (!rows) { return reject(new Error('Logs not found')); } resolve(rows); }); }); } export async function getActivityLogsPaginated(page = 1, limit = 20, search = '', sortBy = 'date', sortOrder = 'desc') { const safePage = Number.isFinite(page) && page > 0 ? Math.floor(page) : 1; const safeLimit = Number.isFinite(limit) && limit > 0 ? Math.floor(limit) : 20; const offset = (safePage - 1) * safeLimit; const safeSearch = typeof search === 'string' ? search.trim() : ''; const normalizedSortBy = ['date', 'entity_name', 'worker_name'].includes(sortBy) ? sortBy : 'date'; const normalizedSortOrder = sortOrder === 'asc' ? 'ASC' : 'DESC'; const searchLike = `%${safeSearch}%`; const filterClause = safeSearch ? `WHERE ( log LIKE ? OR entity_name LIKE ? OR worker_first_name LIKE ? OR worker_last_name LIKE ? OR entity_type LIKE ? )` : ''; const sortClause = (() => { if (normalizedSortBy === 'entity_name') return `ORDER BY entity_name ${normalizedSortOrder}, sort_date DESC`; if (normalizedSortBy === 'worker_name') return `ORDER BY worker_last_name ${normalizedSortOrder}, worker_first_name ${normalizedSortOrder}, sort_date DESC`; return `ORDER BY sort_date ${normalizedSortOrder}`; })(); const countQuery = ` SELECT COUNT(*) AS total FROM ( SELECT nl.log AS log, 'node' AS entity_type, COALESCE(n.name, '') AS entity_name, COALESCE(u.first_name, '') AS worker_first_name, COALESCE(u.last_name, '') AS worker_last_name FROM node_log nl LEFT JOIN node n ON ( (nl.logical_node_id IS NOT NULL AND n.id = nl.logical_node_id) OR (nl.logical_node_id IS NULL AND CAST(n.id AS TEXT) = CAST(nl.node_id AS TEXT)) ) LEFT JOIN "user" u ON CAST(u.id AS TEXT) = CAST(nl.worker_id AS TEXT) UNION ALL SELECT bl.log AS log, 'block' AS entity_type, COALESCE(b.name, '') AS entity_name, COALESCE(u.first_name, '') AS worker_first_name, COALESCE(u.last_name, '') AS worker_last_name FROM block_log bl LEFT JOIN block b ON CAST(b.id AS TEXT) = CAST(bl.block_id AS TEXT) LEFT JOIN "user" u ON CAST(u.id AS TEXT) = CAST(bl.worker_id AS TEXT) ) merged_logs ${filterClause} `; const logsQuery = ` SELECT * FROM ( SELECT nl.id AS id, nl.date AS date, nl.log AS log, nl.worker_id AS worker_id, 'node' AS entity_type, n.id AS entity_id, COALESCE(n.name, '') AS entity_name, COALESCE(u.first_name, '') AS worker_first_name, COALESCE(u.last_name, '') AS worker_last_name, CASE WHEN CAST(nl.date AS TEXT) ~ '^[0-9]+$' THEN CAST(nl.date AS BIGINT) ELSE NULL END AS sort_date FROM node_log nl LEFT JOIN node n ON ( (nl.logical_node_id IS NOT NULL AND n.id = nl.logical_node_id) OR (nl.logical_node_id IS NULL AND CAST(n.id AS TEXT) = CAST(nl.node_id AS TEXT)) ) LEFT JOIN "user" u ON CAST(u.id AS TEXT) = CAST(nl.worker_id AS TEXT) UNION ALL SELECT bl.id AS id, bl.date AS date, bl.log AS log, bl.worker_id AS worker_id, 'block' AS entity_type, b.id AS entity_id, COALESCE(b.name, '') AS entity_name, COALESCE(u.first_name, '') AS worker_first_name, COALESCE(u.last_name, '') AS worker_last_name, CASE WHEN CAST(bl.date AS TEXT) ~ '^[0-9]+$' THEN CAST(bl.date AS BIGINT) ELSE NULL END AS sort_date FROM block_log bl LEFT JOIN block b ON CAST(b.id AS TEXT) = CAST(bl.block_id AS TEXT) LEFT JOIN "user" u ON CAST(u.id AS TEXT) = CAST(bl.worker_id AS TEXT) ) merged_logs ${filterClause} ${sortClause} LIMIT ? OFFSET ? `; return new Promise((resolve, reject) => { const filterParams = safeSearch ? [searchLike, searchLike, searchLike, searchLike, searchLike] : []; db.get(countQuery, filterParams, (countErr, countRow) => { if (countErr) return reject(countErr); db.all(logsQuery, [...filterParams, safeLimit, offset], (logsErr, rows) => { if (logsErr) return reject(logsErr); resolve({ items: rows || [], pagination: { page: safePage, limit: safeLimit, total: countRow?.total || 0, search: safeSearch, sortBy: normalizedSortBy, sortOrder: normalizedSortOrder.toLowerCase() }, }); }); }); }); } export async function getActivityLogsV2(options = {}) { return getAuditLogsPaginated(options); } export async function addNodeLog(nodeId, mac, workerId, log, options = {}) { const occurredAt = Math.floor(Date.now() / 1000); const logicalId = Number.parseInt(nodeId, 10); const logicalNodeId = Number.isFinite(logicalId) ? logicalId : null; return new Promise((resolve, reject) => { db.run( `INSERT INTO node_log (node_id, mac, worker_id, date, log, logical_node_id) VALUES (?, ?, ?, ?, ?, ?)`, [nodeId, mac, workerId, occurredAt, log, logicalNodeId], function (err) { if (err) return reject(err); resolve({ id: this.lastID, occurredAt }); } ); }).then(async (result) => { try { const node = await getNodeById(nodeId).catch(() => null); await writeNodeAuditLog({ nodeId, nodeName: node?.name || null, workerId, message: log, severity: options.severity || 'info', eventType: options.eventType || 'node.legacy_log', payload: options.payload || null, dedupKey: options.dedupKey || null, occurredAt, }); } catch (auditError) { console.error('Failed to write audit log for node_log insert:', auditError); } return { id: result.id }; }); } export async function addSystemNodeLogThrottled(nodeId, mac, message, options = {}) { const dedupKey = options.dedupKey || `system-node-${nodeId}-${message}`; const cooldownSeconds = Number.parseInt(options.cooldownSeconds, 10) || 900; const [canWriteAudit, canWriteNodeLog] = await Promise.all([ canWriteByDedupKey(dedupKey, cooldownSeconds), new Promise((resolve) => { const threshold = Math.floor(Date.now() / 1000) - cooldownSeconds; db.get( `SELECT id FROM node_log WHERE (node_id = ? OR logical_node_id = ?) AND log = ? AND date >= ? ORDER BY date DESC LIMIT 1`, [nodeId, nodeId, message, threshold], (err, row) => { if (err) return resolve(true); resolve(!row); } ); }), ]); if (!canWriteAudit || !canWriteNodeLog) { return { skipped: true }; } return addNodeLog(nodeId, mac, 6, message, { severity: options.severity || 'low', eventType: options.eventType || 'system.state_change', payload: options.payload || null, dedupKey, }); } export async function deleteNodeLog(id) { return new Promise((resolve, reject) => { db.run(`DELETE FROM node_log WHERE id = ?`, [id], function (err) { if (err) return reject(err); resolve({ changes: this.changes }); }); }); } export async function updateNodesInList(state, listId, nodeIds = []) { return new Promise((resolve, reject) => { const normalizedNodeIds = Array.isArray(nodeIds) ? nodeIds .map((id) => Number(id)) .filter((id) => Number.isInteger(id) && id > 0) : []; const useDirectNodeIds = normalizedNodeIds.length > 0; if (!useDirectNodeIds && !Number.isInteger(Number(listId))) { resolve({ changes: 0 }); return; } const targetSelectionSql = useDirectNodeIds ? `n.id IN (${normalizedNodeIds.map(() => '?').join(',')})` : 'n.id IN (SELECT node_id FROM node_list_items WHERE list_id = ?)'; const targetSelectionParams = useDirectNodeIds ? normalizedNodeIds : [listId]; db.all( `SELECT n.id, n.mac, n.state AS prev_state FROM node n WHERE ${targetSelectionSql}`, targetSelectionParams, (readErr, rows) => { if (readErr) return reject(readErr); const nowTs = Math.floor(Date.now() / 1000); const nextManual = state === 'Работает' ? 0 : 1; db.run( `UPDATE node SET state = ?, state_manual = ?, state_manual_updated_at = ? WHERE id IN ( SELECT n2.id FROM node n2 WHERE ${useDirectNodeIds ? `n2.id IN (${normalizedNodeIds.map(() => '?').join(',')})` : 'n2.id IN (SELECT node_id FROM node_list_items WHERE list_id = ?)'} )`, useDirectNodeIds ? [state, nextManual, nowTs, ...normalizedNodeIds] : [state, nextManual, nowTs, listId], async function (err) { if (err) return reject(err); const changedRows = Array.isArray(rows) ? rows.filter((row) => row.prev_state !== state) : []; for (const row of changedRows) { try { await addNodeLog(row.id, row.mac, 6, `Системное сообщение: Состояние изменено списком: '${row.prev_state || ''}' → '${state}'`, { severity: 'low', eventType: 'node.bulk_state_update', payload: { listId: useDirectNodeIds ? null : listId, nodeIds: useDirectNodeIds ? normalizedNodeIds : undefined, from: row.prev_state || '', to: state, }, }); } catch (logErr) { console.error('Failed to write bulk state update log:', logErr); } } if (this.changes > 0) { await emitNodesUpdateSafe(); } resolve({ changes: this.changes }); } ); } ); }); } export async function getNodeStatusByNodeId(nodeId) { return new Promise((resolve, reject) => { db.get('SELECT json_data FROM node_status WHERE node_id = ?', [nodeId], (err, row) => { if (err) return reject(err); if (!row) return resolve(null); try { const json = safeJsonParse(`getNodeStatusByNodeId_${nodeId}_before_parse`, row.json_data); resolve({ reservation: json.reservation || '', reason: json.reason || '', reason_time: json.reason_time || '', }); } catch (e) { resolve(null); } }); }); } export async function getNodePowerStatus(nodeName) { const ssh = new NodeSSH(); try { await ssh.connect(sshConfig); const result = await ssh.execCommand(`/home/novikovia/scripts/power_control.sh status ${nodeName}`); if (result.stderr) throw new Error(result.stderr); const status = result.stdout.trim().split(/\s+/).pop(); return status; } finally { ssh.dispose(); } } export async function getNodeByName(name) { return new Promise((resolve, reject) => { db.get('SELECT * FROM node WHERE name = ?', [name], (err, row) => { if (err) return reject(err); resolve(row); }); }); } export async function setNodePowerStatus(nodeName, action, workerId) { if (!['on', 'off'].includes(action)) throw new Error('Invalid action'); const ssh = new NodeSSH(); try { await ssh.connect(sshConfig); const result = await ssh.execCommand(`/home/novikovia/scripts/power_control.sh ${action} ${nodeName}`); if (result.stderr) throw new Error(result.stderr); const node = await getNodeByName(nodeName); if (node) { const logText = action === 'on' ? 'Узел включён' : 'Узел выключен'; await addNodeLog(node.id, node.mac, workerId, logText); } return result.stdout.trim(); } finally { ssh.dispose(); } } export async function setNodesPowerBatch(nodes, action, workerId) { if (!['on', 'off'].includes(action)) throw new Error('Invalid action'); if (!Array.isArray(nodes) || nodes.length === 0) throw new Error('No nodes specified'); const ssh = new NodeSSH(); try { await ssh.connect(sshConfig); const nodeList = nodes.join(','); const result = await ssh.execCommand(`/home/novikovia/scripts/power_control.sh ${action} ${nodeList}`); if (result.stderr) throw new Error(result.stderr); for (const nodeName of nodes) { const node = await getNodeByName(nodeName); if (node) { const logText = action === 'on' ? 'Узел включён' : 'Узел выключен'; await addNodeLog(node.id, node.mac, workerId, logText); } } return result.stdout.trim(); } finally { ssh.dispose(); } } export async function setNodesResume(nodes, workerId) { if (!Array.isArray(nodes) || nodes.length === 0) throw new Error('No nodes specified'); const ssh = new NodeSSH(); try { await ssh.connect(sshConfig); const nodeNums = nodes.map(n => n.replace(/^node/, '')); const nodeList = `node[${nodeNums.join(',')}]`; const result = await ssh.execCommand(`psh ${nodeList} resume_me.sh`); if (result.stderr) throw new Error(result.stderr); for (const nodeName of nodes) { const node = await getNodeByName(nodeName); if (node) { await addNodeLog(node.id, node.mac, workerId, 'Узел возвращен в работу'); await patchNodeStatusMeta(node.id, { reason: '', reason_time: '', }); } } await emitNodesUpdateSafe(); return result.stdout.trim(); } finally { ssh.dispose(); } } export async function setNodesReservation(action, reservation, nodes, workerId) { if (!['add', 'remove'].includes(action)) throw new Error('Invalid action'); if (!reservation) throw new Error('No reservation specified'); if (!Array.isArray(nodes) || nodes.length === 0) throw new Error('No nodes specified'); const ssh = new NodeSSH(); try { await ssh.connect(sshConfig); const nodeList = nodes.join(','); const result = await ssh.execCommand(`/home/novikovia/scripts/reservation_manager.sh ${action} ${reservation} ${nodeList}`); if (result.stderr) throw new Error(result.stderr); for (const nodeName of nodes) { const node = await getNodeByName(nodeName); if (node) { const logText = action === 'add' ? `Добавлен в резервацию ${reservation}` : `Удалён из резервации ${reservation}`; await addNodeLog(node.id, node.mac, workerId, logText); await patchNodeStatusMeta(node.id, { reservation: action === 'add' ? String(reservation).trim() : '', }); } } await emitNodesUpdateSafe(); return result.stdout.trim(); } finally { ssh.dispose(); } } export async function setNodesMonitoringBatch(nodeNames, workerId, monitoringEnabled = 1) { if (!Array.isArray(nodeNames) || nodeNames.length === 0) throw new Error('No nodes specified'); const enabled = Number(monitoringEnabled) === 1 || monitoringEnabled === true ? 1 : 0; const changed = []; const unchanged = []; const notFound = []; for (const rawName of nodeNames) { const nodeName = String(rawName || '').trim(); if (!nodeName) continue; const node = await getNodeByName(nodeName); if (!node?.id) { notFound.push(nodeName); continue; } const wasOn = Number(node.monitoring_enabled) !== 0; const willBeOn = enabled === 1; if (wasOn === willBeOn) { unchanged.push(nodeName); continue; } await updateNode(node.id, { monitoring_enabled: enabled }); if (enabled === 0) { await clearNodeInQueueMeta(node.id); } if (workerId) { await addNodeLog( node.id, node.mac, workerId, `Мониторинг (SSH/общий опрос): '${wasOn ? 'вкл' : 'выкл'}' → '${willBeOn ? 'вкл' : 'выкл'}' (пакетное действие со списка)` ); } changed.push(nodeName); } await emitNodesUpdateSafe(); const actionLabel = enabled === 1 ? 'Мониторинг включён' : 'Мониторинг выключен'; const parts = []; if (changed.length > 0) { parts.push(`${actionLabel}: ${changed.length} (${changed.join(', ')})`); } if (unchanged.length > 0) { parts.push(`Без изменений: ${unchanged.length} (${unchanged.join(', ')})`); } if (notFound.length > 0) { parts.push(`Не найдены в БД: ${notFound.join(', ')}`); } return parts.join('. ') || 'Нет узлов для обработки'; } export async function setNodesStateBatch(nodeNames, workerId, state = 'Работает') { if (!Array.isArray(nodeNames) || nodeNames.length === 0) throw new Error('No nodes specified'); const targetState = String(state || '').trim() || 'Работает'; const changed = []; const unchanged = []; const notFound = []; for (const rawName of nodeNames) { const nodeName = String(rawName || '').trim(); if (!nodeName) continue; const node = await getNodeByName(nodeName); if (!node?.id) { notFound.push(nodeName); continue; } const prevState = String(node.state || '').trim(); if (prevState === targetState) { unchanged.push(nodeName); continue; } const nowTs = Math.floor(Date.now() / 1000); await updateNode(node.id, { state: targetState, state_manual: targetState === 'Работает' ? 0 : 1, state_manual_updated_at: nowTs, }); if (workerId) { await addNodeLog( node.id, node.mac, workerId, `Состояние изменено: '${prevState}' → '${targetState}' (пакетное действие со списка)` ); } changed.push(nodeName); } await emitNodesUpdateSafe(); const parts = []; if (changed.length > 0) { parts.push(`Состояние «${targetState}»: ${changed.length} (${changed.join(', ')})`); } if (unchanged.length > 0) { parts.push(`Без изменений: ${unchanged.length} (${unchanged.join(', ')})`); } if (notFound.length > 0) { parts.push(`Не найдены в БД: ${notFound.join(', ')}`); } return parts.join('. ') || 'Нет узлов для обработки'; } export async function setNodesFullResume(nodes, workerId) { if (!Array.isArray(nodes) || nodes.length === 0) throw new Error('No nodes specified'); const ssh = new NodeSSH(); const perNode = []; try { await ssh.connect(sshConfig); const nodeNums = nodes.map((n) => String(n).replace(/^node/, '')); const nodeList = `node[${nodeNums.join(',')}]`; const resumeResult = await ssh.execCommand(`psh ${nodeList} resume_me.sh`); if (resumeResult.stderr) throw new Error(resumeResult.stderr); for (const nodeName of nodes) { const node = await getNodeByName(nodeName); if (!node) { perNode.push({ nodeName, ok: false, error: 'node not found in DB' }); continue; } await addNodeLog(node.id, node.mac, workerId, 'Узел возвращен в работу (полный возврат)'); await patchNodeStatusMeta(node.id, { reason: '', reason_time: '' }); const status = await getNodeStatusByNodeId(node.id); const reservations = parseReservationNames(status?.reservation || ''); let reservationRemoved = false; let reservationSkipped = false; let reservationName = ''; if (reservations.length === 1) { reservationName = reservations[0]; const removeResult = await ssh.execCommand( `/home/novikovia/scripts/reservation_manager.sh remove ${reservationName} ${nodeName}` ); if (removeResult.stderr) throw new Error(removeResult.stderr); await patchNodeStatusMeta(node.id, { reservation: '' }); await addNodeLog(node.id, node.mac, workerId, `Удалён из резервации ${reservationName} (полный возврат)`); reservationRemoved = true; } else if (reservations.length > 1) { reservationSkipped = true; } await updateNode(node.id, { state: 'Работает' }); perNode.push({ nodeName, ok: true, reservation: status?.reservation || '', reservationRemoved, reservationSkipped, reservationName, }); } await emitNodesUpdateSafe(); return { stdout: String(resumeResult.stdout || '').trim(), perNode, }; } finally { ssh.dispose(); } } async function touchLastManualClusterAction(nodeId) { const safeNodeId = Number.parseInt(nodeId, 10); if (!Number.isFinite(safeNodeId)) return; const ts = Math.floor(Date.now() / 1000); await new Promise((resolve, reject) => { db.run( 'UPDATE node SET last_manual_cluster_action_at = ? WHERE id = ?', [ts, safeNodeId], (err) => { if (err) return reject(err); resolve(); } ); }); } export async function clearNodeInQueueMeta(nodeId) { await patchNodeStatusMeta(nodeId, { in_queue: '' }); } export async function clearStaleInQueueForDisabledNodes() { const rows = await new Promise((resolve, reject) => { db.all( `SELECT ns.node_id, ns.json_data FROM node_status ns JOIN node n ON n.id = ns.node_id WHERE COALESCE(n.monitoring_enabled, 1) = 0`, [], (err, result) => { if (err) return reject(err); resolve(result || []); } ); }); for (const row of rows) { const nodeId = Number.parseInt(row?.node_id, 10); if (!Number.isFinite(nodeId)) continue; const parsed = row?.json_data ? safeJsonParse(`clearStaleInQueueForDisabledNodes_${nodeId}`, row.json_data) || {} : {}; if (typeof parsed.in_queue === 'string' && parsed.in_queue.trim() !== '') { await clearNodeInQueueMeta(nodeId); } } } async function patchNodeStatusMeta(nodeId, patch = {}) { const safeNodeId = Number.parseInt(nodeId, 10); if (!Number.isFinite(safeNodeId)) return; const touchesClusterSnapshot = Object.prototype.hasOwnProperty.call(patch, 'reason') || Object.prototype.hasOwnProperty.call(patch, 'reason_time') || Object.prototype.hasOwnProperty.call(patch, 'reservation'); await new Promise((resolve, reject) => { db.get('SELECT json_data FROM node_status WHERE node_id = ?', [safeNodeId], (readErr, row) => { if (readErr) return reject(readErr); const prevPayload = row?.json_data ? (safeJsonParse(`patch_node_status_${safeNodeId}`, row.json_data) || {}) : {}; const nextPayload = { ...prevPayload, ...patch }; const updatedAt = Math.floor(Date.now() / 1000); db.run( `INSERT INTO node_status (node_id, json_data, updated_at) VALUES (?, ?, ?) ON CONFLICT(node_id) DO UPDATE SET json_data=excluded.json_data, updated_at=excluded.updated_at`, [safeNodeId, JSON.stringify(nextPayload), updatedAt], async (writeErr) => { if (writeErr) return reject(writeErr); if (touchesClusterSnapshot) { try { await touchLastManualClusterAction(safeNodeId); await syncNodeStateIntervalFromCurrentStatus(safeNodeId, 'manual'); } catch (e) { console.error('[patchNodeStatusMeta] cluster snapshot sync failed', safeNodeId, e); } } resolve(); } ); }); }); } export async function getNodeMeminfo(nodeName) { console.log(nodeName) const ssh = new NodeSSH(); try { await ssh.connect(sshConfig); const execPromise = ssh.execCommand(`/home/novikovia/scripts/meminfo.sh ${nodeName}`); const timeoutPromise = new Promise((_, reject) => setTimeout(() => reject(new Error('meminfo timeout (30s)')), 30000) ); const result = await Promise.race([execPromise, timeoutPromise]); if (result.stderr) throw new Error(result.stderr); const lines = result.stdout.trim().split('\n'); let jsonStr = lines[0].startsWith('[') ? lines.join('\n') : lines.slice(1).join('\n'); return safeJsonParse(`getNodeMeminfo_${nodeName}_before_parse`, jsonStr); } finally { ssh.dispose(); } } export async function getDistinctReasonTags() { const set = new Set(); const pushReason = (raw) => { const s = String(raw ?? '').trim(); if (s) set.add(s); }; const safeDistinctQuery = (sql) => new Promise((resolve) => { db.all(sql, [], (err, rows) => { if (err) return resolve(); (rows || []).forEach((row) => pushReason(row.reason)); resolve(); }); }); await safeDistinctQuery( `SELECT DISTINCT TRIM(reason) AS reason FROM node_metric_snapshot WHERE reason IS NOT NULL AND LENGTH(TRIM(reason)) > 0`, ); await safeDistinctQuery( `SELECT DISTINCT TRIM(reason) AS reason FROM node_event WHERE reason IS NOT NULL AND LENGTH(TRIM(reason)) > 0`, ); await safeDistinctQuery( `SELECT DISTINCT TRIM(reason) AS reason FROM node_state_interval WHERE reason IS NOT NULL AND LENGTH(TRIM(reason)) > 0`, ); await new Promise((resolve) => { db.all('SELECT json_data FROM node_status', [], (err, rows) => { if (err) return resolve(); (rows || []).forEach((row) => { try { const j = safeJsonParse('distinct_reason_tags_status', row.json_data); pushReason(j?.reason); } catch { } }); resolve(); }); }); const items = [...set].sort((a, b) => a.localeCompare(b, 'ru')); return { items }; } export async function changeNodeReason(nodeName, newReason) { const ssh = new NodeSSH(); try { await ssh.connect(sshConfig); const result = await ssh.execCommand(`/home/novikovia/scripts/reason_change.sh ${nodeName} '${newReason}'`); if (result.stderr) throw new Error(result.stderr); const node = await getNodeByName(nodeName); if (node) { await patchNodeStatusMeta(node.id, { reason: String(newReason ?? ''), reason_time: '', }); await addNodeLog(node.id, node.mac, 6, `Системное сообщение: reason изменен на '${newReason}'`, { severity: 'low', eventType: 'node.reason_changed', payload: { reason: newReason, nodeName }, }); await emitNodesUpdateSafe(); } return result.stdout.trim(); } finally { ssh.dispose(); } } export async function setNodeVlan(nodeName, workerId) { const ssh = new NodeSSH(); try { await ssh.connect(sshConfig); const execPromise = ssh.execCommand(`/home/novikovia/scripts/set_vlan_node.sh ${nodeName}`); const timeoutPromise = new Promise((_, reject) => setTimeout(() => reject(new Error('set_vlan_node timeout (180s)')), 180000)); const result = await Promise.race([execPromise, timeoutPromise]); if (result && typeof result.code === 'number') { const isSuccess = (result.code === 0 || result.code === 2); if (!isSuccess) { const stderr = (result.stderr || '').trim(); const stdout = (result.stdout || '').trim(); throw new Error(`Command failed (code ${result.code}). stderr: ${stderr.slice(0,500)} | stdout: ${stdout.slice(0,500)}`); } } const node = await getNodeByName(nodeName); if (node) { await addNodeLog(node.id, node.mac, workerId, 'Настроен VLAN на узле'); } return result.stdout.trim(); } finally { ssh.dispose(); } } export async function runPuppetAgent(nodeName, workerId) { const ssh = new NodeSSH(); try { await ssh.connect(sshConfig); const execPromise = ssh.execCommand(`/home/novikovia/scripts/puppet_agent_node.sh ${nodeName}`); const timeoutPromise = new Promise((_, reject) => setTimeout(() => reject(new Error('puppet_agent_node timeout (900s)')), 900000)); const result = await Promise.race([execPromise, timeoutPromise]); if (result && typeof result.code === 'number' && result.code !== 0) { const stderr = (result.stderr || '').trim(); const stdout = (result.stdout || '').trim(); throw new Error(`Command failed (code ${result.code}). stderr: ${stderr.slice(0,500)} | stdout: ${stdout.slice(0,500)}`); } const node = await getNodeByName(nodeName); if (node) { await addNodeLog(node.id, node.mac, workerId, 'Запущен Puppet Agent на узле'); } return result.stdout.trim(); } finally { ssh.dispose(); } } export async function mountStorageNode(nodeName, workerId) { const ssh = new NodeSSH(); try { await ssh.connect(sshConfig); const execPromise = ssh.execCommand(`/home/novikovia/scripts/mount_storage_node.sh ${nodeName}`); const timeoutPromise = new Promise((_, reject) => setTimeout(() => reject(new Error('mount_storage_node timeout (180s)')), 180000)); const result = await Promise.race([execPromise, timeoutPromise]); if (result && typeof result.code === 'number' && result.code !== 0) { const stderr = (result.stderr || '').trim(); const stdout = (result.stdout || '').trim(); throw new Error(`Command failed (code ${result.code}). stderr: ${stderr.slice(0,500)} | stdout: ${stdout.slice(0,500)}`); } const node = await getNodeByName(nodeName); if (node) { await addNodeLog(node.id, node.mac, workerId, 'Смонтировано хранилище на узле'); } return result.stdout.trim(); } finally { ssh.dispose(); } } const MONITORING_SETTINGS_KEY = 'node_monitoring'; const DEFAULT_MONITORING_SETTINGS = { enabled: false, intervalMinutes: 1, disabledReason: '', }; function normalizeReservationName(value) { return String(value || '').trim(); } export async function ensureReservationMetaSchema() { return new Promise((resolve, reject) => { db.run( `CREATE TABLE IF NOT EXISTS reservation_meta ( name TEXT PRIMARY KEY, description TEXT NOT NULL DEFAULT '', updated_at INTEGER NOT NULL, updated_by INTEGER )`, [], (err) => { if (err) return reject(err); db.run( `CREATE INDEX IF NOT EXISTS idx_reservation_meta_updated_at ON reservation_meta(updated_at)`, [], (idxErr) => { if (idxErr) return reject(idxErr); resolve(); } ); } ); }); } export async function getReservationMetaMap() { await ensureReservationMetaSchema(); return new Promise((resolve) => { db.all( `SELECT rm.name AS name, rm.description AS description, rm.updated_at AS updated_at, rm.updated_by AS updated_by, u.id AS user_id, u.login AS user_login, u.first_name AS user_first_name, u.last_name AS user_last_name FROM reservation_meta rm LEFT JOIN "user" u ON u.id = rm.updated_by ORDER BY rm.name ASC`, [], (err, rows) => { if (err) return resolve(new Map()); const map = new Map(); (rows || []).forEach((row) => { const name = normalizeReservationName(row?.name); if (!name) return; map.set(name, { name, description: String(row?.description || ''), updated_at: Number(row?.updated_at) || null, updated_by: row?.updated_by === null || typeof row?.updated_by === 'undefined' ? null : Number(row.updated_by), updated_by_user: row?.user_id ? { id: Number(row.user_id), login: String(row?.user_login || ''), first_name: String(row?.user_first_name || ''), last_name: String(row?.user_last_name || ''), } : null, }); }); resolve(map); } ); }); } export async function upsertReservationMeta(nameRaw, patch = {}) { await ensureReservationMetaSchema(); const name = normalizeReservationName(nameRaw); if (!name) throw new Error('reservation name is required'); const description = typeof patch.description === 'string' ? patch.description : String(patch.description || ''); const safeDescription = String(description).trim().slice(0, 2000); const updatedByRaw = Number.parseInt(patch.updated_by ?? patch.workerId, 10); const updatedBy = Number.isFinite(updatedByRaw) ? updatedByRaw : null; const now = Math.floor(Date.now() / 1000); return new Promise((resolve, reject) => { db.run( `INSERT INTO reservation_meta (name, description, updated_at, updated_by) VALUES (?, ?, ?, ?) ON CONFLICT(name) DO UPDATE SET description = excluded.description, updated_at = excluded.updated_at, updated_by = excluded.updated_by`, [name, safeDescription, now, updatedBy], (err) => { if (err) return reject(err); resolve({ name, description: safeDescription, updated_at: now, updated_by: updatedBy, }); } ); }); } function normalizeMonitoringSettings(raw = {}) { const enabled = raw.enabled === false ? false : true; const intervalNumber = Number.parseInt(raw.intervalMinutes, 10); const intervalMinutes = Number.isFinite(intervalNumber) ? Math.min(1440, Math.max(1, intervalNumber)) : DEFAULT_MONITORING_SETTINGS.intervalMinutes; const disabledReason = String(raw.disabledReason || '').trim().slice(0, 500); return { enabled, intervalMinutes, disabledReason }; } export async function getMonitoringSettings() { return new Promise((resolve, reject) => { db.run( `CREATE TABLE IF NOT EXISTS app_settings ( key TEXT PRIMARY KEY, value TEXT NOT NULL, updated_at INTEGER NOT NULL )`, [], (createErr) => { if (createErr) return reject(createErr); db.get( 'SELECT value FROM app_settings WHERE key = ?', [MONITORING_SETTINGS_KEY], (getErr, row) => { if (getErr) return reject(getErr); if (!row || typeof row.value !== 'string') { return resolve({ ...DEFAULT_MONITORING_SETTINGS }); } try { const parsed = safeJsonParse('monitoring_settings', row.value); resolve(normalizeMonitoringSettings(parsed)); } catch { resolve({ ...DEFAULT_MONITORING_SETTINGS }); } } ); } ); }); } export async function updateMonitoringSettings(patch = {}) { const current = await getMonitoringSettings(); const next = normalizeMonitoringSettings({ ...current, ...patch, }); return new Promise((resolve, reject) => { const now = Math.floor(Date.now() / 1000); db.run( `INSERT INTO app_settings (key, value, updated_at) VALUES (?, ?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value, updated_at = excluded.updated_at`, [MONITORING_SETTINGS_KEY, JSON.stringify(next), now], (err) => { if (err) return reject(err); resolve(next); } ); }); } export async function addMonitoringLog(level = 'info', message = '') { const safeLevel = ['info', 'warn', 'error'].includes(level) ? level : 'info'; const safeMessage = String(message || '').slice(0, 4000); const now = Math.floor(Date.now() / 1000); return new Promise((resolve, reject) => { db.run( `CREATE TABLE IF NOT EXISTS monitoring_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, date INTEGER NOT NULL, level TEXT NOT NULL, message TEXT NOT NULL )`, [], (createErr) => { if (createErr) return reject(createErr); db.run( `INSERT INTO monitoring_log (date, level, message) VALUES (?, ?, ?)`, [now, safeLevel, safeMessage], function (insertErr) { if (insertErr) return reject(insertErr); resolve({ id: this.lastID, date: now, level: safeLevel, message: safeMessage }); } ); } ); }); } export async function getMonitoringLogs(limit = 50) { const safeLimit = Math.min(200, Math.max(1, Number.parseInt(limit, 10) || 50)); return new Promise((resolve, reject) => { db.run( `CREATE TABLE IF NOT EXISTS monitoring_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, date INTEGER NOT NULL, level TEXT NOT NULL, message TEXT NOT NULL )`, [], (createErr) => { if (createErr) return reject(createErr); db.all( `SELECT id, date, level, message FROM monitoring_log ORDER BY date DESC, id DESC LIMIT ?`, [safeLimit], (err, rows) => { if (err) return reject(err); resolve(rows || []); } ); } ); }); } function runStatement(sql, params = []) { return new Promise((resolve, reject) => { db.run(sql, params, function onRun(err) { if (err) return reject(err); resolve(this); }); }); } export async function ensureNodeZabbixSchema() { const safeRun = async (sql) => { try { await runStatement(sql); } catch (error) { const message = String(error?.message || ''); const pgDuplicateColumn = String(error?.code || '') === '42701'; const sqliteDuplicateColumn = message.includes('duplicate column name'); const pgAlreadyExists = message.toLowerCase().includes('already exists') || message.toLowerCase().includes('уже существует'); if (!(pgDuplicateColumn || sqliteDuplicateColumn || pgAlreadyExists)) { throw error; } } }; await safeRun(`ALTER TABLE node ADD COLUMN IF NOT EXISTS zabbix_enabled INTEGER NOT NULL DEFAULT 0`); await safeRun(`ALTER TABLE node ADD COLUMN IF NOT EXISTS zabbix_hostid TEXT`); await runStatement(`CREATE INDEX IF NOT EXISTS idx_node_zabbix_enabled ON node(zabbix_enabled)`); await runStatement(`CREATE INDEX IF NOT EXISTS idx_node_zabbix_hostid ON node(zabbix_hostid)`); } export async function ensureLogicalNodeMonitoringSchema() { const safeRun = async (sql) => { try { await runStatement(sql); } catch (error) { const message = String(error?.message || ''); const pgDuplicateColumn = String(error?.code || '') === '42701'; const sqliteDuplicateColumn = message.includes('duplicate column name'); const pgAlreadyExists = message.toLowerCase().includes('already exists') || message.toLowerCase().includes('уже существует'); if (!(pgDuplicateColumn || sqliteDuplicateColumn || pgAlreadyExists)) { throw error; } } }; await safeRun( `ALTER TABLE node ADD COLUMN IF NOT EXISTS monitoring_enabled INTEGER NOT NULL DEFAULT 1` ); await safeRun( `ALTER TABLE node ADD COLUMN IF NOT EXISTS binding_mac_interface TEXT NOT NULL DEFAULT 'eth0'` ); await safeRun(`ALTER TABLE node ADD COLUMN IF NOT EXISTS visible_name TEXT`); await runStatement(`CREATE INDEX IF NOT EXISTS idx_node_monitoring_enabled ON node(monitoring_enabled)`); try { await runStatement(`DROP VIEW IF EXISTS logical_node`); } catch (e) { console.warn('[logical_node] drop view:', e?.message || e); } await runStatement(`CREATE VIEW logical_node AS SELECT * FROM node`); } export async function ensureNodeManualStateSchema() { const safeRun = async (sql) => { try { await runStatement(sql); } catch (error) { const message = String(error?.message || ''); const pgDuplicateColumn = String(error?.code || '') === '42701'; const sqliteDuplicateColumn = message.includes('duplicate column name'); const pgAlreadyExists = message.toLowerCase().includes('already exists') || message.toLowerCase().includes('уже существует'); if (!(pgDuplicateColumn || sqliteDuplicateColumn || pgAlreadyExists)) { throw error; } } }; await safeRun(`ALTER TABLE node ADD COLUMN IF NOT EXISTS state_manual INTEGER NOT NULL DEFAULT 0`); await safeRun(`ALTER TABLE node ADD COLUMN IF NOT EXISTS state_manual_updated_at INTEGER`); await safeRun( `ALTER TABLE node ADD COLUMN IF NOT EXISTS last_manual_cluster_action_at INTEGER` ); await runStatement(`CREATE INDEX IF NOT EXISTS idx_node_state_manual ON node(state_manual)`); } export async function getZabbixEnabledNodes() { return new Promise((resolve, reject) => { db.all( `SELECT id, name, zabbix_enabled, zabbix_hostid FROM node WHERE COALESCE(zabbix_enabled, 0) = 1 ORDER BY id ASC`, [], (err, rows) => { if (err) return reject(err); resolve(rows || []); } ); }); } export async function updateNodeZabbixConfig(nodeId, patch = {}) { const safeNodeId = Number.parseInt(nodeId, 10); if (!Number.isFinite(safeNodeId)) { throw new Error('nodeId must be a number'); } const update = []; const values = []; if (typeof patch.zabbixEnabled !== 'undefined') { update.push('zabbix_enabled = ?'); values.push(patch.zabbixEnabled ? 1 : 0); } if (typeof patch.zabbixHostId !== 'undefined') { update.push('zabbix_hostid = ?'); values.push(String(patch.zabbixHostId || '').trim() || null); } if (update.length === 0) { throw new Error('Nothing to update'); } values.push(safeNodeId); await runStatement(`UPDATE node SET ${update.join(', ')} WHERE id = ?`, values); const updatedNode = await new Promise((resolve, reject) => { db.get( `SELECT id, name, zabbix_enabled, zabbix_hostid FROM node WHERE id = ?`, [safeNodeId], (err, row) => { if (err) return reject(err); if (!row) return reject(new Error('Node not found')); resolve(row); } ); }); try { await writeNodeAuditLog({ nodeId: updatedNode.id, nodeName: updatedNode.name, workerId: Number.isFinite(Number.parseInt(patch.workerId, 10)) ? Number.parseInt(patch.workerId, 10) : null, message: `Zabbix config updated: enabled=${updatedNode.zabbix_enabled}, hostid=${updatedNode.zabbix_hostid || ''}`, severity: 'low', eventType: 'node.zabbix_config_updated', payload: { zabbix_enabled: updatedNode.zabbix_enabled, zabbix_hostid: updatedNode.zabbix_hostid || null, }, }); } catch (error) { console.error('Failed to write zabbix config audit log:', error); } return updatedNode; } export async function bulkUpdateNodeZabbixEnabled(nodeIds = [], zabbixEnabled = false) { const safeIds = [...new Set( (Array.isArray(nodeIds) ? nodeIds : []) .map((value) => Number.parseInt(value, 10)) .filter(Number.isFinite) )]; if (safeIds.length === 0) { throw new Error('No node ids specified'); } const placeholders = safeIds.map(() => '?').join(', '); const values = [zabbixEnabled ? 1 : 0, ...safeIds]; const result = await runStatement( `UPDATE node SET zabbix_enabled = ? WHERE id IN (${placeholders})`, values ); return { changes: result.changes || 0 }; }