/
Krams
/
nodemon2
Обзор
Документация
Войти
/
Krams
/
nodemon2
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
backend/services/NodeStatusSyncService.js
1 058 строк
42 KB
Иван Новиков
init
11 июн 2026, 13:08
11 июн 2026, 13:08
900f716
Код
Авторство
О чём код?
import { NodeSSH } from 'node-ssh'; import { safeJsonParse } from './DebugDump.js'; import db from '../database.js'; import sshConfig from './sshConfig.js'; import { io } from '../socketServer.js'; import { getAllNodesNested } from './NodeService.js'; import { getAllBlocksGrouped } from './BlockService.js'; import { addMonitoringLog, addSystemNodeLogThrottled, clearNodeInQueueMeta, getZabbixEnabledNodes, setNodesReservation } from './NodeService.js'; import { applyPhysicalBindingFromMetrics } from './PhysicalNodeBindingService.js'; import { recordNodeReliabilitySnapshot } from './NodeReliabilityService.js'; import { fetchZabbixItemByNodesAndKey, fetchZabbixItemHistoryByNodesAndKey, fetchZabbixStatusesByNodes } from './ZabbixExternalService.js'; import { closeRelayConnection, executeOnNodeViaRelay, executeOnRelay } from '../collector/relayExecutor.js'; import { runAnalyzerOnce } from './AnalyzerService.js'; import { upsertNodeAnalysisLatest } from './AnalyzerStorageService.js'; const ssh = new NodeSSH(); const NO_METRICS_AUTORIP_AFTER_SECONDS = 1 * 60; async function fetchNodesViaSSH() { await ssh.connect(sshConfig); const statusNodeHost = process.env.STATUS_NODE_HOST || 'node241'; const onlyHostname = String(process.env.STATUS_ONLY_HOSTNAME || '').trim() === '1'; if (onlyHostname) { const precheckCmd = `ssh ${statusNodeHost} "hostname"`; const precheckRes = await ssh.execCommand(precheckCmd); if (precheckRes.stderr) throw new Error(precheckRes.stderr); console.log(`[SYNC] status node hostname (${statusNodeHost}):`, (precheckRes.stdout || '').trim()); return []; } const precheckEnabled = String(process.env.STATUS_NODE_PRECHECK || '').trim() === '1'; if (precheckEnabled) { const precheckCmd = `ssh ${statusNodeHost} "hostname"`; const precheckRes = await ssh.execCommand(precheckCmd); if (precheckRes.stderr) throw new Error(precheckRes.stderr); console.log(`[SYNC] status node precheck hostname (${statusNodeHost}):`, (precheckRes.stdout || '').trim()); } const remoteScriptPath = '/home/novikovia/scripts/check_all_nodes.sh'; const cmd = `ssh ${statusNodeHost} "bash ${remoteScriptPath}"`; console.log(`[SYNC] status collection cmd (runs on head): ${cmd}`); const startedAt = Date.now(); const result = await ssh.execCommand(cmd); const elapsedMs = Date.now() - startedAt; const stdoutLen = (result.stdout || '').length; const stderrLen = (result.stderr || '').length; console.log(`[SYNC] status collection finished in ${elapsedMs}ms; code=${result.code}; stdout=${stdoutLen}B; stderr=${stderrLen}B`); if (stderrLen) { console.log('[SYNC] stderr head 500 chars:', (result.stderr || '').slice(0, 500)); } if (typeof result.code === 'number' && result.code !== 0) { throw new Error(`status collection failed (code ${result.code}). stderr: ${(result.stderr || '').slice(0, 500)}`); } try { console.log('[SYNC] stdout length:', (result.stdout || '').length); if (result.stdout) { console.log('[SYNC] stdout head 500 chars:', result.stdout); } const nodes = safeJsonParse('sync_stdout_before_parse', result.stdout); console.log(`[SYNC] parsed nodes: ${Array.isArray(nodes) ? nodes.length : 'not an array'}`); if (!Array.isArray(nodes)) { throw new Error('Parsed result is not an array'); } if (Array.isArray(nodes)) { nodes.forEach((n, i) => { const id = n && (n.node_id || n.name || n.id || `index_${i}`); console.log(`[SYNC][node ${i}]`, id); }); } return nodes; } catch (e) { console.error('[SYNC] Ошибка парсинга JSON. Полный stdout (обрезан до 2000 символов):'); console.error((result.stdout || '')); throw new Error('Ошибка парсинга JSON: ' + e.message); } } async function mergeZabbixDataForTestNodes(nodes = []) { const zabbixEnabledNodes = await getZabbixEnabledNodes(); if (!Array.isArray(zabbixEnabledNodes) || zabbixEnabledNodes.length === 0) { return { mergedNodes: nodes, matchedCount: 0, configuredCount: 0, details: [], }; } const zabbixStatusesByName = await fetchZabbixStatusesByNodes(zabbixEnabledNodes); const details = []; const mergedNodes = (Array.isArray(nodes) ? nodes : []).map((node) => { const nodeName = String(node?.node_id || '').trim(); const zabbix = zabbixStatusesByName.get(nodeName); if (!zabbix) return node; return { ...node, zabbix_hostid: zabbix.hostid, zabbix_available: zabbix.available, zabbix_status: zabbix.status, zabbix_error: zabbix.error, zabbix_tags: zabbix.tags, zabbix_checked_at: Math.floor(Date.now() / 1000), }; }); zabbixEnabledNodes.forEach((node) => { const nodeName = String(node?.name || '').trim(); const matched = zabbixStatusesByName.get(nodeName); details.push({ node_name: nodeName, requested_hostid: String(node?.zabbix_hostid || '').trim() || null, matched_hostid: matched?.hostid || null, available: matched?.available ?? null, status: matched?.status ?? null, error: matched?.error || (matched ? '' : 'Host not found in Zabbix response'), }); }); return { mergedNodes, matchedCount: zabbixStatusesByName.size, configuredCount: zabbixEnabledNodes.length, details, }; } async function collectZabbixTestDetails(options = {}) { const zabbixEnabledNodes = await getZabbixEnabledNodes(); if (!Array.isArray(zabbixEnabledNodes) || zabbixEnabledNodes.length === 0) { return { configuredCount: 0, matchedCount: 0, details: [], }; } const itemKey = String(options?.itemKey || '').trim(); const zabbixStatusesByName = await fetchZabbixStatusesByNodes(zabbixEnabledNodes); const zabbixItemsByName = itemKey ? await fetchZabbixItemByNodesAndKey(zabbixEnabledNodes, itemKey) : new Map(); const zabbixItemHistoryByName = itemKey ? await fetchZabbixItemHistoryByNodesAndKey(zabbixEnabledNodes, itemKey, { windowMinutes: 15 }) : new Map(); const details = zabbixEnabledNodes.map((node) => { const nodeName = String(node?.name || '').trim(); const matched = zabbixStatusesByName.get(nodeName); const item = zabbixItemsByName.get(nodeName) || null; const itemHistory = zabbixItemHistoryByName.get(nodeName) || null; return { node_name: nodeName, requested_hostid: String(node?.zabbix_hostid || '').trim() || null, matched_hostid: matched?.hostid || null, available: matched?.available ?? null, status: matched?.status ?? null, error: matched?.error || (matched ? '' : 'Host not found in Zabbix response'), item_key_requested: itemKey || null, item: item || null, item_history_last_15m: itemHistory, }; }); return { configuredCount: zabbixEnabledNodes.length, matchedCount: zabbixStatusesByName.size, details, }; } async function upsertNode(node, options = {}) { const bypassMonitoringGate = options.bypassMonitoringGate === true; return new Promise((resolve, reject) => { db.get( `SELECT id, COALESCE(monitoring_enabled, 1) AS monitoring_enabled, COALESCE(last_manual_cluster_action_at, 0) AS last_manual_cluster_action_at FROM node WHERE name = ?`, [node.node_id], (err, row) => { if (err) return reject(err); if (!row) return resolve(); if (!bypassMonitoringGate && Number(row.monitoring_enabled) === 0) return resolve(); const dbNodeId = row.id; const lastManualClusterActionAt = Number.parseInt(row.last_manual_cluster_action_at, 10); db.get('SELECT json_data FROM node_status WHERE node_id = ?', [dbNodeId], async (prevErr, prevRow) => { if (prevErr) return reject(prevErr); try { await applyPhysicalBindingFromMetrics({ dbNodeId, hostname: node.node_id, metrics: node, }); } catch (bindErr) { console.warn('[SYNC] physical binding:', bindErr?.message || bindErr); } const updatedAt = Math.floor(Date.now() / 1000); let prevPayload = {}; try { if (prevRow?.json_data) { prevPayload = safeJsonParse(`upsert_prev_${dbNodeId}`, prevRow.json_data) || {}; } } catch { prevPayload = {}; } const hasNoMetricsError = typeof node?.error === 'string' && node.error.includes('no metrics received'); const prevLastSuccessAt = Number.parseInt(prevPayload?.last_success_at, 10); const prevNoMetricsSince = Number.parseInt(prevPayload?.no_metrics_since, 10); let enrichedNode = { ...node, last_sync_at: updatedAt, last_success_at: hasNoMetricsError ? (Number.isFinite(prevLastSuccessAt) ? prevLastSuccessAt : null) : updatedAt, no_metrics_since: hasNoMetricsError ? (Number.isFinite(prevNoMetricsSince) ? prevNoMetricsSince : updatedAt) : null, }; const sampledRaw = Number(node?.sampled_at); const sampledAtSec = Number.isFinite(sampledRaw) && sampledRaw > 0 ? sampledRaw > 1e12 ? Math.floor(sampledRaw / 1000) : Math.floor(sampledRaw) : null; const lastManualSec = Number.isFinite(lastManualClusterActionAt) ? lastManualClusterActionAt : 0; if ( sampledAtSec !== null && lastManualSec > 0 && sampledAtSec < lastManualSec ) { enrichedNode = { ...enrichedNode, reason: typeof prevPayload.reason === 'string' ? prevPayload.reason : String(enrichedNode.reason ?? ''), reason_time: typeof prevPayload.reason_time === 'string' ? prevPayload.reason_time : String(enrichedNode.reason_time ?? ''), reservation: typeof prevPayload.reservation === 'string' ? prevPayload.reservation : String(enrichedNode.reservation ?? ''), }; } const jsonData = JSON.stringify(enrichedNode); 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`, [dbNodeId, jsonData, updatedAt], (err2) => { if (err2) return reject(err2); resolve(); } ); }); }); }); } async function getNodeStatusPayloadByName(nodeName) { const safeNodeName = String(nodeName || '').trim(); if (!safeNodeName) { throw new Error('nodeName is required'); } await ssh.connect(sshConfig); const metricsResult = await ssh.execCommand(`timeout 8 ssh -o ConnectTimeout=2 ${safeNodeName} "bash /home/novikovia/scripts/status_script.sh"`); if (typeof metricsResult.code === 'number' && metricsResult.code !== 0) { throw new Error(`Не удалось получить метрики с ${safeNodeName}`); } const metricsRaw = String(metricsResult.stdout || '').trim(); return parseManualSyncMetricsPayload(safeNodeName, metricsRaw, { inQueueLoader: async () => { const inQueueResult = await ssh.execCommand(`squeue -h -w ${safeNodeName} | awk 'NR==1{print $4}'`); return String(inQueueResult.stdout || '').trim(); }, }); } async function getNodeStatusPayloadByNameViaCollector(nodeName) { const safeNodeName = String(nodeName || '').trim(); if (!safeNodeName) { throw new Error('nodeName is required'); } const statusScriptPath = String( process.env.COLLECTOR_STATUS_SCRIPT_PATH || '/home/novikovia/scripts/status_script.sh' ).trim(); const metricsResult = await executeOnNodeViaRelay(safeNodeName, `bash ${statusScriptPath}`); if (!metricsResult?.success) { throw new Error(`Не удалось получить метрики с ${safeNodeName}`); } const metricsRaw = String(metricsResult.stdout || '').trim(); return parseManualSyncMetricsPayload(safeNodeName, metricsRaw, { inQueueLoader: async () => { const inQueueResult = await executeOnRelay(`squeue -h -w ${safeNodeName} | awk 'NR==1{print $4}'`); return inQueueResult?.success ? String(inQueueResult.stdout || '').trim() : ''; }, }); } async function parseManualSyncMetricsPayload(safeNodeName, metricsRaw, options = {}) { if (!metricsRaw) { throw new Error(`Узел ${safeNodeName} не вернул метрики`); } const lines = metricsRaw .split('\n') .map((line) => String(line || '').replace(/\u001b\[[0-9;]*m/g, '').trim()) .filter(Boolean); const parseCandidates = [ metricsRaw, lines.join(''), lines.slice(0, -1).join(''), `${lines.slice(0, -1).join('')}}`, ].filter(Boolean); let metrics = null; let parseError = null; for (let i = 0; i < parseCandidates.length; i += 1) { try { metrics = safeJsonParse(`manual_sync_${safeNodeName}_variant_${i}`, parseCandidates[i]); if (metrics && typeof metrics === 'object') { break; } } catch (error) { parseError = error; } } if (!metrics || typeof metrics !== 'object') { throw new Error( `Не удалось разобрать JSON от ${safeNodeName}: ${parseError?.message || 'unknown parse error'}` ); } const inQueueLoader = typeof options.inQueueLoader === 'function' ? options.inQueueLoader : null; const inQueue = inQueueLoader ? await inQueueLoader() : ''; return { ...metrics, node_id: safeNodeName, in_queue: inQueue, }; } async function getNodeStatusMetaByNodeId(nodeId) { return new Promise((resolve) => { db.get('SELECT json_data FROM node_status WHERE node_id = ?', [nodeId], (err, row) => { if (err || !row?.json_data) { return resolve({ reservation: '', reason: '', reason_time: '' }); } try { const parsed = safeJsonParse(`node_status_meta_${nodeId}_before_parse`, row.json_data); resolve({ reservation: String(parsed?.reservation || ''), reason: String(parsed?.reason || ''), reason_time: String(parsed?.reason_time || ''), }); } catch { resolve({ reservation: '', reason: '', reason_time: '' }); } }); }); } async function analyzeAndSetNodeStates() { const conditions = [ { name: 'rip_reserved', check: (data) => typeof data.reservation === 'string' && data.reservation.trim() === 'RIP', state: 'Отключен от опроса', comment: 'Узел исключен из автоопроса (RIP)', }, { name: 'system_node', check: (data) => data.reason === 'SCC_SYSTEM', state: 'Системный', comment: 'Системный узел', }, { name: 'linpack_in_progress', check: (data) => typeof data.reason === 'string' && /SCC_LINPACK_\d+C/.test(data.reason), state: 'Линпак', comment: 'Узел в линпаке', }, { name: 'error_present', check: (data) => typeof data.error === 'string' && data.error.length > 0, state: 'Нужно обслужить', comment: 'Узел недоступен по SSH или не работает ls /home', }, { name: 'linpack_done', check: (data) => typeof data.reason === 'string' && data.reason.includes('SCC_LINPACK_DONE'), state: 'Ожидает возврата', comment: 'Узел прошел линпак', }, { name: 'reason_not_responding', check: (data) => typeof data.reason === 'string' && data.reason.includes('Not responding'), state: 'Нужно обслужить', comment: 'Завис (not responding)', }, { name: 'reason_scc_error', check: (data) => typeof data.reason === 'string' && data.reason.includes('SCC_ERROR'), state: 'Нужно обслужить', comment: 'Метка узла содержит SCC_ERROR', }, { name: 'reason_rebooted', check: (data) => typeof data.reason === 'string' && data.reason.includes('Node unexpectedly rebooted'), state: 'Нужно обслужить', comment: 'Перезагрузился', }, { name: 'reservation_exists', check: (data) => typeof data.reservation === 'string' && data.reservation.trim() !== '', state: 'Нужно обслужить', comment: 'Узел находится в резервации', }, { name: 'nonstandard_total_mem', check: (data) => { if (typeof data.total_mem_mb !== 'number' && typeof data.total_mem_mb !== 'string') return false; const str = String(data.total_mem_mb).trim(); return !(str.startsWith('24') || str.startsWith('32') || str.startsWith('48')); }, state: 'Нужно обслужить', comment: 'Оперативная память: нестандартное количество', }, { name: 'skipped', check: (data) => data.skipped === true, state: 'Неисправен', comment: 'Узел пропущен скриптом', }, { name: 'ib_state_down', check: (data) => data.ib_state === 'Down', state: 'Нужно обслужить', comment: 'IB-Link не работает', }, { name: 'mtu_not_standard', check: (data) => typeof data.ib0_mtu_correct !== 'undefined' && data.ib0_mtu_correct !== true, state: 'Нужно обслужить', comment: 'MTU отличается от стандартного (65520)', }, { name: 'scc_linpack_run', check: (data) => typeof data.reason === 'string' && data.reason.includes('SCC_LINPACK_RUN'), state: 'Линпак', comment: 'Узел запускает/заканчивает линпак', }, { name: 'node_in_queue', check: (data) => typeof data.in_queue === 'string' && data.in_queue !== "", state: 'Работает', comment: 'Узел в расчетах' } ]; return new Promise((resolve, reject) => { db.all('SELECT node_id, json_data FROM node_status', [], (err, statusRows) => { if (err) return reject(err); db.all( 'SELECT id, state, COALESCE(monitoring_enabled, 1) AS monitoring_enabled, COALESCE(state_manual, 0) AS state_manual FROM node', [], (err2, nodeRows) => { if (err2) return reject(err2); const nodeStateMap = {}; const nodeMonitoringMap = {}; const nodeStateManualMap = {}; nodeRows.forEach((row) => { nodeStateMap[row.id] = row.state; nodeMonitoringMap[row.id] = Number(row.monitoring_enabled) !== 0 ? 1 : 0; nodeStateManualMap[row.id] = Number(row.state_manual) === 1 ? 1 : 0; }); const nodesToReserve = []; const nodesToDisablePolling = []; const updates = statusRows.map(async (row) => { if (nodeMonitoringMap[row.node_id] === 0) return Promise.resolve(); const currentState = nodeStateMap[row.node_id]; const isManualState = nodeStateManualMap[row.node_id] === 1; if (isManualState) return Promise.resolve(); let state = 'Работает'; let comment = ''; let matched = false; let problemComments = []; let problemMatched = false; let shouldLog = false; let prevState = currentState; let prevComment = null; let mac = null; await new Promise((res, rej) => { db.get('SELECT mac, system_comment FROM node WHERE id = ?', [row.node_id], (err, nodeRow) => { if (!err && nodeRow) { mac = nodeRow.mac; prevComment = nodeRow.system_comment; } res(); }); }); try { if (!row.json_data || typeof row.json_data !== 'string' || row.json_data.trim() === '') { console.error('Empty or invalid json_data for node_id:', row.node_id, 'json_data:', row.json_data); return Promise.resolve(); } const data = safeJsonParse(`analyze_row_${row.node_id}_before_parse`, row.json_data); const isLinpack = typeof data.reason === 'string' && (data.reason.includes('SCC_LINPACK_RUN') || /SCC_LINPACK_\d+C/.test(data.reason)); if (isLinpack) { const nodeName = data.node_id; const nodeRow = await new Promise((resolve) => { db.get('SELECT id, mac, system_comment FROM node WHERE name = ?', [nodeName], (err, row) => { resolve(row); }); }); if (!nodeRow) return Promise.resolve(); const dbNodeId = nodeRow.id; const mac = nodeRow.mac; const prevComment = nodeRow.system_comment; const ssh = new NodeSSH(); try { await ssh.connect(sshConfig); const result = await ssh.execCommand(`/home/novikovia/scripts/power_control.sh status ${nodeName}`); const powerStatus = result.stdout.trim().split(':').pop().trim(); if (powerStatus === 'off') { await new Promise((resolve) => { db.run("DELETE FROM linpack WHERE node_id = ? AND end_time = ''", [dbNodeId], function (err) { resolve(); }); }); state = 'Нужно обслужить'; comment = 'Выключился в линпаке'; matched = true; shouldLog = (prevState !== state || prevComment !== comment); try { const resetResult = await ssh.execCommand(`/home/novikovia/scripts/reason_change.sh ${nodeName} SCC_ERROR_OFF`); } catch (e) { console.error('Ошибка при сбросе reason через reason_change.sh:', nodeName, e); } return new Promise((res, rej) => { db.run('UPDATE node SET state = ?, system_comment = ? WHERE id = ?', [state, comment, dbNodeId], async function (err3) { if (err3) return rej(err3); const suppressInQueueSpam = (state === 'Работает' && comment === 'Узел в расчетах'); if (shouldLog && mac && !suppressInQueueSpam) { try { await addSystemNodeLogThrottled(dbNodeId, mac, `Системное сообщение: ${comment}`, { dedupKey: `sync:${dbNodeId}:${state}:${comment}`, cooldownSeconds: 1800, eventType: 'sync.state_transition', severity: state === 'Нужно обслужить' ? 'medium' : 'low', payload: { state, comment, condition: 'linpack_power_check' }, }); } catch (e) {} } res(); }); }); } } catch (e) { console.error('Ошибка при проверке питания для линпак-узла:', nodeName, e); } finally { ssh.dispose(); } } if (prevState === 'Ожидает ремонта' || prevState === 'Готов к установке') { return Promise.resolve(); } { const reservationExists = typeof data.reservation === 'string' && data.reservation.trim() !== ''; if (prevState === 'Ожидает возврата' && reservationExists) { return Promise.resolve(); } } for (const cond of conditions) { if (cond.check(data)) { if (cond.state !== 'Нужно обслужить') { state = cond.state; comment = cond.comment || ''; matched = true; shouldLog = (prevState !== state || prevComment !== comment); break; } if (cond.name === 'skipped') { state = cond.state; comment = cond.comment || ''; matched = true; shouldLog = (prevState !== state || prevComment !== comment); break; } if (cond.name === 'reason_not_responding') { state = cond.state; comment = cond.comment || ''; matched = true; shouldLog = (prevState !== state || prevComment !== comment); break; } if (cond.name === 'error_present') { state = cond.state; comment = cond.comment || ''; matched = true; shouldLog = (prevState !== state || prevComment !== comment); break; } problemComments.push(cond.comment || ''); problemMatched = true; } } if (problemMatched && !matched) { state = 'Нужно обслужить'; comment = problemComments.join('; '); matched = true; shouldLog = (prevState !== state || prevComment !== comment); } if (state === 'Нужно обслужить' && (!data.reservation || data.reservation.trim() === '')) { const nodeName = data.node_id; const hasNoMetricsError = typeof data.error === 'string' && data.error.includes('no metrics received'); if (hasNoMetricsError) { const nowTs = Math.floor(Date.now() / 1000); const lastSuccessAt = Number.parseInt(data.last_success_at, 10); const noMetricsSince = Number.parseInt(data.no_metrics_since, 10); const silenceSeconds = Number.isFinite(lastSuccessAt) ? (nowTs - lastSuccessAt) : (Number.isFinite(noMetricsSince) ? (nowTs - noMetricsSince) : 0); if (silenceSeconds >= NO_METRICS_AUTORIP_AFTER_SECONDS) { nodesToDisablePolling.push({ nodeId: row.node_id, nodeName, mac, reason: 'Длительное отсутствие метрик: SSH-опрос автоматически отключён', }); } } else { nodesToReserve.push(nodeName); } } if (typeof data.reason === 'string' && data.reason.includes('SCC_LINPACK_RUN')) { openLinpackIfNeeded(row.node_id); } else if (typeof data.reason === 'string' && data.reason.includes('SCC_LINPACK_DONE')) { closeLinpackIfNeeded(row.node_id); } else if (currentState === 'Линпак' && state !== 'Линпак') { closeLinpackIfNeeded(row.node_id); } } catch (e) { console.error('JSON parse error for node_id:', row.node_id, 'json_data:', row.json_data, e); } if (!matched && currentState && currentState !== 'Работает' && currentState !== 'Нужно обслужить' && currentState !== 'Неисправен') { return Promise.resolve(); } return new Promise((res, rej) => { db.run('UPDATE node SET state = ?, system_comment = ? WHERE id = ?', [state, comment, row.node_id], async function (err3) { if (err3) return rej(err3); const suppressInQueueSpam = (state === 'Работает' && comment === 'Узел в расчетах'); if (shouldLog && mac && !suppressInQueueSpam) { try { await addSystemNodeLogThrottled(row.node_id, mac, `Системное сообщение: ${comment}`, { dedupKey: `sync:${row.node_id}:${state}:${comment}`, cooldownSeconds: 1800, eventType: 'sync.state_transition', severity: state === 'Нужно обслужить' ? 'medium' : 'low', payload: { state, comment, condition: 'analyze_conditions' }, }); } catch (e) { console.error('Ошибка при записи системного лога:', e); } } res(); }); }); }); Promise.all(updates).then(async () => { if (nodesToDisablePolling.length > 0) { const seenNodeIds = new Set(); const uniqueNodes = nodesToDisablePolling.filter((item) => { if (!item || !Number.isFinite(Number(item.nodeId))) return false; if (seenNodeIds.has(item.nodeId)) return false; seenNodeIds.add(item.nodeId); return true; }); for (const item of uniqueNodes) { try { await new Promise((res, rej) => { db.run( `UPDATE node SET monitoring_enabled = 0, state = 'Отключен от опроса', system_comment = ? WHERE id = ?`, [item.reason, item.nodeId], function onDisable(err) { if (err) return rej(err); res(); } ); }); await clearNodeInQueueMeta(item.nodeId); if (item.mac) { await addSystemNodeLogThrottled(item.nodeId, item.mac, `Системное сообщение: ${item.reason}`, { dedupKey: `sync:disable-ssh-polling:${item.nodeId}`, cooldownSeconds: 3600, eventType: 'sync.ssh_polling_disabled', severity: 'medium', payload: { nodeName: item.nodeName, reason: 'no_metrics_timeout' }, }); } } catch (e) { console.error('[AUTO-SSH-DISABLE] Ошибка авто-отключения опроса:', item?.nodeName || item?.nodeId, e); } } const disabledNames = uniqueNodes.map((item) => item.nodeName).filter(Boolean); if (disabledNames.length > 0) { console.log(`[AUTO-SSH-DISABLE] Nodes ${disabledNames.join(', ')} polling disabled (monitoring_enabled=0)`); } } if (nodesToReserve.length > 0) { try { await setNodesReservation('add', 'RESERVATION', nodesToReserve, 6); console.log(`[AUTO-RESERVATION] Nodes ${nodesToReserve.join(', ')} added to RESERVATION`); } catch (e) { console.error('[AUTO-RESERVATION] Ошибка при добавлении в резервацию:', e); } } resolve(); }).catch(reject); }); }); }); } function openLinpackIfNeeded(nodeId) { return new Promise((resolve, reject) => { db.get("SELECT * FROM linpack WHERE node_id = ? AND end_time = ''", [nodeId], (err, row) => { if (err) return reject(err); if (!row) { const now = new Date().toISOString(); db.run("INSERT INTO linpack (node_id, start_time, end_time) VALUES (?, ?, '')", [nodeId, now], function (err2) { if (err2) return reject(err2); resolve(); }); } else { resolve(); } }); }); } function closeLinpackIfNeeded(nodeId) { return new Promise((resolve, reject) => { db.get("SELECT * FROM linpack WHERE node_id = ? AND end_time = ''", [nodeId], (err, row) => { if (err) return reject(err); if (row) { const now = new Date().toISOString(); db.run('UPDATE linpack SET end_time = ? WHERE id = ?', [now, row.id], function (err2) { if (err2) return reject(err2); resolve(); }); } else { resolve(); } }); }); } function cleanupStaleLinpackRecords() { return new Promise((resolve, reject) => { db.run( `DELETE FROM linpack WHERE end_time = '' AND start_time <= ?`, [new Date(Date.now() - (3 * 60 * 60 * 1000)).toISOString()], function (err) { if (err) return reject(err); resolve(); } ); }); } async function recordReliabilitySnapshotsFromCurrentStatus() { return new Promise((resolve, reject) => { db.all( `SELECT n.id AS node_id, ns.updated_at AS source_updated_at, ns.json_data AS json_data, COALESCE(n.state, 'Работает') AS current_state FROM node n LEFT JOIN node_status ns ON ns.node_id = n.id`, [], (err, rows) => { if (err) return reject(err); resolve(rows || []); } ); }).then(async (rows) => { const ts = Math.floor(Date.now() / 1000); for (const row of rows) { let parsed = {}; try { if (row.json_data) { parsed = safeJsonParse(`record_reliability_${row.node_id}`, row.json_data); } } catch (e) { parsed = {}; } const errorText = typeof parsed.error === 'string' ? parsed.error : ''; await recordNodeReliabilitySnapshot(row.node_id, { ts, state: row.current_state || 'Работает', reason: typeof parsed.reason === 'string' ? parsed.reason : '', reservation: typeof parsed.reservation === 'string' ? parsed.reservation : '', hasError: Boolean(errorText), errorText, ibState: typeof parsed.ib_state === 'string' ? parsed.ib_state : '', ib0MtuCorrect: typeof parsed.ib0_mtu_correct === 'boolean' ? parsed.ib0_mtu_correct : null, inQueue: typeof parsed.in_queue === 'string' ? parsed.in_queue : '', skipped: typeof parsed.skipped === 'boolean' ? parsed.skipped : null, totalMemMb: parsed.total_mem_mb, sourceUpdatedAt: row.source_updated_at, }); } }); } export async function ingestCollectorNodes(nodes = [], options = {}) { const bypassMonitoringGate = options?.bypassMonitoringGate === true; const normalized = Array.isArray(nodes) ? nodes : []; await Promise.all( normalized.map((node) => upsertNode(node, { bypassMonitoringGate })) ); await analyzeAndSetNodeStates(); await recordReliabilitySnapshotsFromCurrentStatus(); const allNodes = await getAllNodesNested(); io.emit('nodes_update', allNodes); const allBlocks = await getAllBlocksGrouped(); io.emit('blocks_update', allBlocks); const analyzerShadowEnabled = String(process.env.ANALYZER_SHADOW || '').trim() === '1'; if (analyzerShadowEnabled) { try { const nodeNames = [...new Set( normalized .map((n) => String(n?.node_id || '').trim()) .filter(Boolean) )]; const nodeIds = await new Promise((resolve) => { if (nodeNames.length === 0) return resolve([]); const placeholders = nodeNames.map(() => '?').join(', '); db.all( `SELECT id FROM node WHERE name IN (${placeholders})`, nodeNames, (err, rows) => { if (err) return resolve([]); resolve((rows || []).map((r) => Number(r?.id)).filter(Number.isFinite)); } ); }); const analysis = await runAnalyzerOnce({ nodeIds, limit: 5000 }); const results = Array.isArray(analysis?.results) ? analysis.results : []; const currentById = await new Promise((resolve) => { if (nodeIds.length === 0) return resolve(new Map()); const placeholders = nodeIds.map(() => '?').join(', '); db.all( `SELECT id, state, system_comment FROM node WHERE id IN (${placeholders})`, nodeIds, (err, rows) => { const map = new Map(); if (!err) { (rows || []).forEach((row) => { map.set(Number(row.id), { state: String(row?.state || ''), system_comment: String(row?.system_comment || ''), }); }); } resolve(map); } ); }); const mismatches = []; for (const r of results) { const id = Number(r?.node_id); if (!Number.isFinite(id)) continue; const current = currentById.get(id) || { state: '', system_comment: '' }; const recState = String(r?.recommended_state || '').trim(); const curState = String(current.state || '').trim(); const recComment = String(r?.recommended_comment || '').trim(); const curComment = String(current.system_comment || '').trim(); if (recState && curState && recState !== curState) { mismatches.push({ node_id: id, node_name: r?.node_name, current_state: curState, recommended_state: recState, current_comment: curComment, recommended_comment: recComment, confidence_pct: r?.confidence_pct, matched_rules: r?.matched_rules, }); } } const persistLatest = String(process.env.ANALYZER_PERSIST_LATEST || '').trim() === '1'; if (persistLatest) { for (const r of results) { const id = Number(r?.node_id); if (!Number.isFinite(id) || id <= 0) continue; await upsertNodeAnalysisLatest(id, r).catch(() => {}); } } if (mismatches.length > 0) { console.log( `[ANALYZER][shadow] mismatches=${mismatches.length}/${results.length}; sample=`, mismatches.slice(0, 5) ); } else { console.log(`[ANALYZER][shadow] ok results=${results.length}`); } } catch (e) { console.warn('[ANALYZER][shadow] failed:', e?.message || e); } } return { ingested: normalized.length, }; } export async function syncNodeStatus(options = {}) { const manualZabbix = options?.manualZabbix === true; const result = { totalNodesFromSsh: 0, zabbixConfiguredNodes: 0, zabbixMatchedNodes: 0, zabbixDetails: [], }; try { await cleanupStaleLinpackRecords(); const nodes = await fetchNodesViaSSH(); result.totalNodesFromSsh = Array.isArray(nodes) ? nodes.length : 0; let nodesToPersist = nodes; if (manualZabbix) { const mergeResult = await mergeZabbixDataForTestNodes(nodes); nodesToPersist = mergeResult.mergedNodes; result.zabbixConfiguredNodes = mergeResult.configuredCount; result.zabbixMatchedNodes = mergeResult.matchedCount; result.zabbixDetails = mergeResult.details; await addMonitoringLog( 'info', `[ZABBIX][manual] configured=${mergeResult.configuredCount}, matched=${mergeResult.matchedCount}` ); } await Promise.all(nodesToPersist.map((n) => upsertNode(n))); await analyzeAndSetNodeStates(); await recordReliabilitySnapshotsFromCurrentStatus(); console.log('Node statuses updated successfully.'); const allNodes = await getAllNodesNested(); console.log(`[Socket.IO] Отправка nodes_update, количество узлов: ${allNodes.flat(Infinity).length}`); io.emit('nodes_update', allNodes); const allBlocks = await getAllBlocksGrouped(); console.log(`[Socket.IO] Отправка blocks_update, количество блоков: ${allBlocks.length}`); io.emit('blocks_update', allBlocks); return result; } catch (e) { console.error('Ошибка при получении данных:', e); if (manualZabbix) { const message = e instanceof Error ? e.message : String(e || 'unknown error'); await addMonitoringLog('error', `[ZABBIX][manual] sync failed: ${message}`); throw e; } return result; } finally { ssh.dispose(); } } export async function runZabbixTestSyncOnce(options = {}) { await addMonitoringLog('info', '[ZABBIX][manual] started'); const result = await collectZabbixTestDetails(options); await addMonitoringLog('info', '[ZABBIX][manual] completed'); return { ok: true, zabbixConfiguredNodes: result.configuredCount, zabbixMatchedNodes: result.matchedCount, zabbixDetails: result.details, }; } export async function syncSingleNodeStatus(nodeName) { const safeNodeName = String(nodeName || '').trim(); if (!safeNodeName) { throw new Error('nodeName is required'); } const monitoringDriver = String(process.env.MONITORING_DRIVER || 'collector').trim().toLowerCase(); const useCollectorPath = monitoringDriver !== 'legacy'; try { const nodeRow = await new Promise((resolve, reject) => { db.get('SELECT id FROM node WHERE name = ?', [safeNodeName], (err, row) => { if (err) return reject(err); resolve(row || null); }); }); if (!nodeRow?.id) { throw new Error('Node not found'); } const payload = useCollectorPath ? await getNodeStatusPayloadByNameViaCollector(safeNodeName) : await getNodeStatusPayloadByName(safeNodeName); const prevMeta = await getNodeStatusMetaByNodeId(nodeRow.id); const mergedPayload = { ...payload, reservation: payload?.reservation || prevMeta.reservation || '', reason: payload?.reason || prevMeta.reason || '', reason_time: payload?.reason_time || prevMeta.reason_time || '', }; await upsertNode(mergedPayload, { bypassMonitoringGate: true }); const allNodes = await getAllNodesNested(); io.emit('nodes_update', allNodes); const allBlocks = await getAllBlocksGrouped(); io.emit('blocks_update', allBlocks); return { ok: true, node: safeNodeName, }; } finally { if (useCollectorPath) { await closeRelayConnection().catch(() => {}); } else { ssh.dispose(); } } }