Files
MikrotikManager/backend/src/services/alert-engine/run-once.ts
T
Denozordec 11ad94f67d feat: integrate sonner for toast notifications and enhance UI feedback
Added the sonner library for toast notifications across various components, improving user feedback for actions such as saving settings, syncing rules, and handling errors. Updated the layout to include a Toaster component for consistent notification display. Refactored alert messages in the backups, gre, and filters pages to utilize the new notification system, enhancing overall user experience.
2026-05-07 20:49:35 +07:00

222 lines
8.0 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import { db } from "../../db/index.js"
import { alertEngineState } from "../../db/schema.js"
import { eq } from "drizzle-orm"
import {
getTelegramPublic,
listAlertGroups,
listAlertRules,
type ApiAlertRule,
} from "../alerts-service.js"
import { getGreBgpSnapshotCollectorState } from "../gre-bgp-snapshot-collector.js"
import { isAnyAlertSnapshotSourceJobRunning } from "../scheduler-running.js"
import { getServersRestPingCollectorState } from "../servers-rest-ping-collector.js"
import { isTrafficCollecting } from "../traffic-collector.js"
import { getUptimeCollectorsState } from "../uptime-collector.js"
import { computeDecisions } from "./decision-engine.js"
import { clearConfirmPendingAfterSend } from "./confirm-stability.js"
import { dispatchPendingOutbox, enqueueTelegramOutbox } from "./outbox.js"
import { matchRules } from "./rule-matcher.js"
import { ingestAlertSignals } from "./signal-ingestor.js"
import { getLatestSourceFinishedAt, getSourceWatermark, updateSourceWatermark } from "./source-watermark.js"
import { pickTelegramAlertEmoji } from "./telegram-emoji.js"
import type { AlertEngineRunResult, RuleEvalHit } from "./types.js"
import { appendEvent } from "../../modules/events/service/events-service.js"
const SNAPSHOT_SOURCE_WAIT_MS = 30_000
const SNAPSHOT_SOURCE_POLL_MS = 20
/**
* Не строить снимок, пока коллекторы пишут в SQLite — иначе гонка с `alert_engine` по таймеру
* (в т.ч. слот планировщика занят до первого `await` в коллекторе, когда внутренний `collecting` ещё false).
*/
async function awaitSnapshotSourcesIdle(): Promise<void> {
const t0 = Date.now()
while (Date.now() - t0 < SNAPSHOT_SOURCE_WAIT_MS) {
const busy =
isAnyAlertSnapshotSourceJobRunning() ||
getServersRestPingCollectorState().running ||
getGreBgpSnapshotCollectorState().running ||
getUptimeCollectorsState().resources ||
getUptimeCollectorsState().ping ||
isTrafficCollecting()
if (!busy) return
await new Promise<void>((r) => setTimeout(r, SNAPSHOT_SOURCE_POLL_MS))
}
}
function newHistoryId(): string {
return `ah-${Date.now()}-${Math.random().toString(36).slice(2, 10)}`
}
function newOutboxId(): string {
return `ao-${Date.now()}-${Math.random().toString(36).slice(2, 10)}`
}
function loadEngineStateMap(): Map<string, { lastFiredAt: string; lastPayloadHash: string | null }> {
const rows = db.select().from(alertEngineState).all()
const m = new Map<string, { lastFiredAt: string; lastPayloadHash: string | null }>()
for (const r of rows) {
m.set(r.scopeKey, { lastFiredAt: r.lastFiredAt, lastPayloadHash: r.lastPayloadHash ?? null })
}
return m
}
function upsertEngineState(scopeKey: string, lastFiredAt: string, lastPayloadHash: string | null) {
const existing = db.select().from(alertEngineState).where(eq(alertEngineState.scopeKey, scopeKey)).limit(1).all()[0]
if (existing) {
db.update(alertEngineState)
.set({ lastFiredAt, lastPayloadHash })
.where(eq(alertEngineState.scopeKey, scopeKey))
.run()
} else {
db.insert(alertEngineState).values({ scopeKey, lastFiredAt, lastPayloadHash }).run()
}
}
function severityRank(s: ApiAlertRule["severity"]): number {
if (s === "critical") return 3
if (s === "warning") return 2
return 1
}
function maxSeverity(a: ApiAlertRule["severity"], b: ApiAlertRule["severity"]): ApiAlertRule["severity"] {
return severityRank(a) >= severityRank(b) ? a : b
}
function ruleScopeKey(ruleId: string, transition?: RuleEvalHit["transition"]): string {
return `rule:${ruleId}:${transition ?? "neutral"}`
}
/** Один проход движка: оценка правил, группы, Telegram, история, состояние cooldown. */
export async function runAlertEngineOnce(): Promise<AlertEngineRunResult> {
const errors: string[] = []
const sampledAt = new Date().toISOString()
const latestSourceFinishedAt = getLatestSourceFinishedAt()
const prevSourceFinishedAt = getSourceWatermark()
const hasNewSources =
latestSourceFinishedAt != null &&
(prevSourceFinishedAt == null || latestSourceFinishedAt > prevSourceFinishedAt)
const rules = listAlertRules()
const groups = listAlertGroups()
const state = loadEngineStateMap()
const pub = getTelegramPublic()
const canSend = pub.tokenConfigured && Boolean(pub.chatId?.trim())
let standaloneFires = 0
let groupFires = 0
let ruleDiag: AlertEngineRunResult["ruleDiag"] = []
if (hasNewSources) {
await awaitSnapshotSourcesIdle()
const snap = ingestAlertSignals()
const matched = matchRules(rules, snap)
errors.push(...matched.errors)
const { standalone, grouped, ruleDiag: diag } = computeDecisions({
rules,
groups,
hitByRule: matched.hitByRule,
state,
canSend,
})
ruleDiag = diag
for (const d of standalone) {
if (!canSend) break
const firedAt = new Date().toISOString()
const key = ruleScopeKey(d.rule.id, d.hit.transition)
const queued = enqueueTelegramOutbox({
id: newOutboxId(),
dedupeKey: `${key}:${d.hit.payloadHash}`,
payload: {
text: `${pickTelegramAlertEmoji(d.hit.message)} ${d.hit.message}`,
chatId: d.rule.chatId?.trim() || undefined,
history: {
id: newHistoryId(),
ruleId: d.rule.id,
groupId: null,
ruleName: d.hit.ruleName,
severity: d.hit.severity,
message: d.hit.message,
firedAt,
},
},
})
if (!queued) continue
standaloneFires += 1
clearConfirmPendingAfterSend([d.rule.id])
upsertEngineState(key, firedAt, d.hit.payloadHash)
state.set(key, { lastFiredAt: firedAt, lastPayloadHash: d.hit.payloadHash })
}
for (const d of grouped) {
if (!canSend) break
const firedAt = new Date().toISOString()
const lines = d.hits.map((h) => h.message)
const sev = d.hits.reduce((a, h) => maxSeverity(a, h.severity), d.hits[0]!.severity)
const body = `Группа «${d.group.name}» (${d.group.combineMode === "any" ? "ANY" : "ALL"})\n\n${lines.join("\n")}`
const hash = d.hits.map((h) => h.payloadHash).sort().join("|")
const key = `group:${d.group.id}`
const queued = enqueueTelegramOutbox({
id: newOutboxId(),
dedupeKey: `${key}:${hash}`,
payload: {
text: `${pickTelegramAlertEmoji(body)} ${body}`,
history: {
id: newHistoryId(),
ruleId: null,
groupId: d.group.id,
ruleName: `Группа: ${d.group.name}`,
severity: sev,
message: body,
firedAt,
},
},
})
if (!queued) continue
groupFires += 1
upsertEngineState(key, firedAt, hash)
state.set(key, { lastFiredAt: firedAt, lastPayloadHash: hash })
clearConfirmPendingAfterSend(d.members.map((r) => r.id))
}
updateSourceWatermark(latestSourceFinishedAt)
}
const outbox = await dispatchPendingOutbox()
errors.push(...outbox.errors)
if (standaloneFires > 0 || groupFires > 0) {
appendEvent({
level: "info",
eventType: "alerts.engine.fired",
sourceModule: "alerts",
title: "Движок алертов обнаружил события",
message: `Правила: ${standaloneFires}, группы: ${groupFires}`,
payload: {
sampledAt,
rulesChecked: hasNewSources ? rules.length : 0,
},
})
}
if (errors.length > 0) {
appendEvent({
level: "warning",
eventType: "alerts.engine.errors",
sourceModule: "alerts",
title: "Ошибки в движке алертов",
message: errors[0] ?? "Неизвестная ошибка",
payload: {
totalErrors: errors.length,
},
})
}
return {
sampledAt,
rulesChecked: hasNewSources ? rules.length : 0,
standaloneFires,
groupFires,
skippedNoTelegram: !canSend,
errors,
ruleDiag: ruleDiag ?? [],
}
}