chore: synchronize pending app/backend updates and repository hygiene

Includes current frontend and backend work in progress and removes generated artifacts from tracking to keep the repository clean for дальнейшая разработка.

Co-authored-by: Cursor <[email protected]>
This commit is contained in:
Denozordec
2026-05-07 12:29:04 +07:00
co-authored by Cursor
parent bdb9b72fac
commit 5f31bb47fb
81 changed files with 11976 additions and 1239 deletions
@@ -0,0 +1,71 @@
import assert from "node:assert/strict"
import { eq } from "drizzle-orm"
import { db } from "../../db/index.js"
import { alertEngineConfirmPending } from "../../db/schema.js"
import type { ApiAlertRule } from "../alerts-service.js"
import { computeStabilityReadyMap, deleteConfirmPending } from "./confirm-stability.js"
const ruleId = "test-confirm-stability-recovery-bypass"
deleteConfirmPending(ruleId)
const baseRule: ApiAlertRule = {
id: ruleId,
name: "Server transition",
type: "server",
target: "core-1",
targets: ["core-1"],
groupId: null,
condition: "перешёл в offline · восстановился",
conditions: ["перешёл в offline", "восстановился"],
severity: "critical",
enabled: true,
cooldown: "5м",
confirmStabilitySec: 120,
recoveryMode: "always",
recoveryStabilitySec: null,
chatId: "",
lastFiredAt: null,
}
const problemHit = {
ruleId,
ruleName: baseRule.name,
severity: baseRule.severity,
message: "Сервер: core-1 — перешёл в offline",
transition: "problem" as const,
payloadHash: "problem-1",
}
const recoveryHit = {
ruleId,
ruleName: baseRule.name,
severity: baseRule.severity,
message: "Сервер: core-1 — восстановился",
transition: "recovery" as const,
payloadHash: "recovery-1",
}
{
const m = computeStabilityReadyMap([baseRule], new Map([[ruleId, problemHit]]))
assert.equal(m.get(ruleId), false, "проблемный переход должен ждать confirmStabilitySec")
}
{
const m = computeStabilityReadyMap([baseRule], new Map([[ruleId, recoveryHit]]))
assert.equal(m.get(ruleId), true, "восстановление должно отправляться без задержки")
}
{
const m = computeStabilityReadyMap([baseRule], new Map([[ruleId, null]]))
assert.equal(m.get(ruleId), false, "без hit отправка не готова")
const left = db
.select()
.from(alertEngineConfirmPending)
.where(eq(alertEngineConfirmPending.ruleId, ruleId))
.limit(1)
.all()[0]
assert.equal(left, undefined, "pending должен очищаться после отсутствия hit")
}
console.log("alert-engine confirm-stability tests ok")
@@ -0,0 +1,82 @@
import { eq } from "drizzle-orm"
import { db } from "../../db/index.js"
import { alertEngineConfirmPending } from "../../db/schema.js"
import type { ApiAlertRule } from "../alerts-service.js"
import type { RuleEvalHit } from "./types.js"
export function deleteConfirmPending(ruleId: string) {
db.delete(alertEngineConfirmPending).where(eq(alertEngineConfirmPending.ruleId, ruleId)).run()
}
/**
* Обновляет состояние ожидания стабильности и возвращает, можно ли отправлять уведомление в этом тике.
* При отсутствии срабатывания (`hit == null`) ожидание сбрасывается.
*/
export function stabilityReadyForSend(
ruleId: string,
hit: RuleEvalHit | null,
delaySec: number | null | undefined,
): boolean {
if (!hit) {
deleteConfirmPending(ruleId)
return false
}
const sec = delaySec != null && delaySec > 0 ? Math.floor(delaySec) : 0
if (sec <= 0) {
deleteConfirmPending(ruleId)
return true
}
const now = Date.now()
const nowIso = new Date().toISOString()
const existing = db
.select()
.from(alertEngineConfirmPending)
.where(eq(alertEngineConfirmPending.ruleId, ruleId))
.limit(1)
.all()[0]
if (!existing || existing.payloadHash !== hit.payloadHash) {
if (existing) {
db.update(alertEngineConfirmPending)
.set({ payloadHash: hit.payloadHash, sinceAt: nowIso })
.where(eq(alertEngineConfirmPending.ruleId, ruleId))
.run()
} else {
db.insert(alertEngineConfirmPending).values({ ruleId, payloadHash: hit.payloadHash, sinceAt: nowIso }).run()
}
return false
}
const t = Date.parse(existing.sinceAt)
if (!Number.isFinite(t)) return false
return now - t >= sec * 1000
}
/** После успешной отправки — сбросить таймер стабильности для следующего цикла. */
export function clearConfirmPendingAfterSend(ruleIds: string[]) {
for (const id of ruleIds) deleteConfirmPending(id)
}
export function computeStabilityReadyMap(
rules: ApiAlertRule[],
hitByRule: Map<string, RuleEvalHit | null>,
): Map<string, boolean> {
const m = new Map<string, boolean>()
for (const rule of rules) {
if (!rule.enabled) {
deleteConfirmPending(rule.id)
m.set(rule.id, false)
continue
}
const hit = hitByRule.get(rule.id) ?? null
const delaySec =
hit?.transition === "recovery"
? rule.recoveryMode === "conditional"
? rule.recoveryStabilitySec
: 0
: rule.confirmStabilitySec
m.set(rule.id, stabilityReadyForSend(rule.id, hit, delaySec))
}
return m
}
@@ -0,0 +1,22 @@
import type { ApiAlertRule } from "../alerts-service.js"
const COOLDOWN_MS: Record<ApiAlertRule["cooldown"], number> = {
"1м": 60_000,
"5м": 5 * 60_000,
"15м": 15 * 60_000,
"1ч": 60 * 60_000,
"4ч": 4 * 60 * 60_000,
"24ч": 24 * 60 * 60_000,
}
export function cooldownToMs(cd: ApiAlertRule["cooldown"]): number {
return COOLDOWN_MS[cd] ?? COOLDOWN_MS["5м"]
}
/** Можно ли отправить снова после lastIso (ISO) */
export function isCooldownElapsed(lastFiredIso: string | null | undefined, cooldown: ApiAlertRule["cooldown"]): boolean {
if (!lastFiredIso?.trim()) return true
const t = Date.parse(lastFiredIso)
if (!Number.isFinite(t)) return true
return Date.now() - t >= cooldownToMs(cooldown)
}
@@ -0,0 +1,136 @@
import type { ApiAlertGroup, ApiAlertRule } from "../alerts-service.js"
import { computeStabilityReadyMap } from "./confirm-stability.js"
import { isCooldownElapsed } from "./cooldown.js"
import { groupShouldFire } from "./groups.js"
import { isPositiveRecoveryTelegramText } from "./telegram-emoji.js"
import { shouldAllowRecoveryByPolicy } from "./state-machine.js"
import type { AlertEngineRuleDiag, RuleEvalHit } from "./types.js"
export type EngineStateMap = Map<string, { lastFiredAt: string; lastPayloadHash: string | null }>
export type StandaloneDecision = {
kind: "rule"
rule: ApiAlertRule
hit: RuleEvalHit
}
export type GroupDecision = {
kind: "group"
group: ApiAlertGroup
members: ApiAlertRule[]
hits: RuleEvalHit[]
}
function ruleScopeKey(ruleId: string, transition?: RuleEvalHit["transition"]): string {
const suffix = transition ?? "neutral"
return `rule:${ruleId}:${suffix}`
}
function textAfterGroupAlertHeader(fullBody: string): string {
const segs = fullBody.split(/\n\n/)
if (segs.length >= 2 && /^Группа\s*«/i.test(segs[0] ?? "")) return segs.slice(1).join("\n\n").trim()
return fullBody
}
export function buildRuleDiagnostics(
rules: ApiAlertRule[],
hitByRule: Map<string, RuleEvalHit | null>,
stabilityReady: Map<string, boolean>,
state: EngineStateMap,
canSend: boolean,
): AlertEngineRuleDiag[] {
const out: AlertEngineRuleDiag[] = []
for (const rule of rules) {
if (!rule.enabled) continue
const hit = hitByRule.get(rule.id) ?? null
const inGroup = Boolean(rule.groupId)
const evalHit = Boolean(hit)
const stab = stabilityReady.get(rule.id) ?? false
const key = ruleScopeKey(rule.id, hit?.transition)
const prev = state.get(key)
const cooldownOk = isCooldownElapsed(prev?.lastFiredAt, rule.cooldown)
let blocked: AlertEngineRuleDiag["blocked"]
if (!evalHit) blocked = "no_hit"
else if (!shouldAllowRecoveryByPolicy(hit, rule.recoveryMode)) blocked = "stability"
else if (inGroup) blocked = "in_group"
else if (!stab) blocked = "stability"
else if (!cooldownOk) blocked = "cooldown"
else if (!canSend) blocked = "no_telegram"
else if (hit && isPositiveRecoveryTelegramText(hit.message) && prev?.lastPayloadHash === hit.payloadHash) blocked = "dedupe_positive"
else blocked = undefined
out.push({
ruleId: rule.id,
inGroup,
evalHit,
hitTransition: hit?.transition,
hitMessage: hit?.message,
stabilityOk: stab,
cooldownOk,
telegramOk: canSend,
blocked,
})
}
return out
}
export function computeDecisions(args: {
rules: ApiAlertRule[]
groups: ApiAlertGroup[]
hitByRule: Map<string, RuleEvalHit | null>
state: EngineStateMap
canSend: boolean
}): {
stabilityReady: Map<string, boolean>
standalone: StandaloneDecision[]
grouped: GroupDecision[]
ruleDiag: AlertEngineRuleDiag[]
} {
const { rules, groups, hitByRule, state, canSend } = args
const stabilityReady = computeStabilityReadyMap(rules, hitByRule)
const ruleDiag = buildRuleDiagnostics(rules, hitByRule, stabilityReady, state, canSend)
const standalone: StandaloneDecision[] = []
const grouped: GroupDecision[] = []
for (const rule of rules) {
if (!rule.enabled || rule.groupId) continue
const hit = hitByRule.get(rule.id)
if (!hit) continue
if (!shouldAllowRecoveryByPolicy(hit, rule.recoveryMode)) continue
if (!stabilityReady.get(rule.id)) continue
const key = ruleScopeKey(rule.id, hit.transition)
const prev = state.get(key)
if (!isCooldownElapsed(prev?.lastFiredAt, rule.cooldown)) continue
if (isPositiveRecoveryTelegramText(hit.message) && prev?.lastPayloadHash === hit.payloadHash) continue
standalone.push({ kind: "rule", rule, hit })
}
for (const g of groups) {
if (!g.enabled) continue
const members = rules.filter((r) => r.groupId === g.id && r.enabled)
if (members.length === 0) continue
const firing = members.map((r) => {
const h = hitByRule.get(r.id)
if (!shouldAllowRecoveryByPolicy(h ?? null, r.recoveryMode)) return false
return Boolean(h) && Boolean(stabilityReady.get(r.id))
})
if (!groupShouldFire(g.combineMode, firing)) continue
const hits: RuleEvalHit[] = []
for (const member of members) {
const hit = hitByRule.get(member.id)
if (!hit) continue
if (!shouldAllowRecoveryByPolicy(hit, member.recoveryMode)) continue
hits.push(hit)
}
if (hits.length === 0) continue
const key = `group:${g.id}`
const prev = state.get(key)
const cd = g.cooldownOverride ?? "5м"
if (!isCooldownElapsed(prev?.lastFiredAt, cd)) continue
const body = `Группа «${g.name}» (${g.combineMode === "any" ? "ANY" : "ALL"})\n\n${hits.map((h) => h.message).join("\n")}`
const hash = hits.map((h) => h.payloadHash).sort().join("|")
if (isPositiveRecoveryTelegramText(textAfterGroupAlertHeader(body)) && prev?.lastPayloadHash === hash) continue
grouped.push({ kind: "group", group: g, members, hits })
}
return { stabilityReady, standalone, grouped, ruleDiag }
}
@@ -0,0 +1,227 @@
import assert from "node:assert/strict"
import type { ApiAlertRule } from "../alerts-service.js"
import { evaluateRule } from "./evaluate-rule.js"
import type { SignalSnapshot } from "./types.js"
function emptySnap(over: Partial<SignalSnapshot> = {}): SignalSnapshot {
return {
sampledAt: "t",
servers: [],
probes: [],
trafficByServer: [],
greTunnels: [],
bgpSessions: [],
...over,
}
}
function rule(p: Partial<ApiAlertRule> & Pick<ApiAlertRule, "id" | "name" | "type" | "targets">): ApiAlertRule {
return {
target: p.targets[0] ?? "",
groupId: null,
condition: p.conditions?.[0] ?? p.condition ?? "",
conditions: p.conditions ?? (p.condition ? [p.condition] : []),
severity: p.severity ?? "warning",
enabled: p.enabled ?? true,
cooldown: p.cooldown ?? "5м",
confirmStabilitySec: null,
recoveryMode: "always",
recoveryStabilitySec: null,
chatId: "",
lastFiredAt: null,
...p,
} as ApiAlertRule
}
{
const r = rule({
id: "g1",
name: "GRE",
type: "gre-tunnel",
targets: ["t1 / srv"],
conditions: ["перешёл в offline"],
})
const hit = evaluateRule(
r,
emptySnap({
greTunnels: [{ targetLabel: "t1 / srv", status: "down", prevStatus: "up" }],
}),
)
assert.ok(hit)
}
{
const r = rule({
id: "g2",
name: "GRE rec",
type: "gre-tunnel",
targets: ["t1 / srv"],
conditions: ["восстановился"],
})
const hit = evaluateRule(
r,
emptySnap({
greTunnels: [{ targetLabel: "t1 / srv", status: "up", prevStatus: "down" }],
}),
)
assert.ok(hit)
}
{
const r = rule({
id: "b1",
name: "BGP",
type: "bgp-peer",
targets: ["srv: peer1"],
conditions: ["разорвал сессию"],
})
const hit = evaluateRule(
r,
emptySnap({
bgpSessions: [{ key: "srv: peer1", state: "Idle", prevState: "Established" }],
}),
)
assert.ok(hit)
}
{
const r = rule({
id: "b2",
name: "BGP up",
type: "bgp-peer",
targets: ["srv: peer1"],
conditions: ["восстановил сессию"],
})
const hit = evaluateRule(
r,
emptySnap({
bgpSessions: [{ key: "srv: peer1", state: "Established", prevState: "Active" }],
}),
)
assert.ok(hit)
}
{
const r = rule({
id: "pfx",
name: "pfx",
type: "bgp-prefix",
targets: ["любой префикс"],
conditions: ["был отозван"],
})
assert.equal(evaluateRule(r, emptySnap()), null)
}
{
const r = rule({
id: "srv-rec",
name: "rec",
type: "server",
targets: ["core"],
conditions: ["восстановился"],
})
assert.ok(
evaluateRule(
r,
emptySnap({
servers: [
{
name: "core",
status: "online",
prevStatus: "offline",
prev2Status: "offline",
sampledAt: "t1",
},
],
}),
),
"один online после offline — тоже «восстановился» (ручной REST / ресурс уже online)",
)
const hit = evaluateRule(
r,
emptySnap({
servers: [
{
name: "core",
status: "online",
prevStatus: "online",
prev2Status: "offline",
sampledAt: "t2",
},
],
}),
)
assert.ok(hit, "два подряд online после offline — одно срабатывание")
assert.equal(hit?.transition, "recovery")
}
{
const r = rule({
id: "srv-rec-offline-phrase",
name: "rec-offline-phrase",
type: "server",
targets: ["core"],
conditions: ["восстановился после offline"],
})
const hit = evaluateRule(
r,
emptySnap({
servers: [
{
name: "core",
status: "online",
prevStatus: "online",
prev2Status: "offline",
sampledAt: "t2",
},
],
}),
)
assert.ok(hit, "в тексте условия есть «offline» — всё равно восстановление по сэмплам")
}
{
const r = rule({
id: "srv-off",
name: "off",
type: "server",
targets: ["core"],
conditions: ["перешёл в offline"],
})
const edgeHit = evaluateRule(
r,
emptySnap({
servers: [
{
name: "core",
status: "offline",
prevStatus: "online",
prev2Status: "online",
sampledAt: "t",
},
],
}),
)
assert.ok(edgeHit)
assert.equal(edgeHit?.transition, "problem")
assert.equal(
evaluateRule(
r,
emptySnap({
servers: [
{
name: "core",
status: "offline",
prevStatus: "offline",
prev2Status: "online",
sampledAt: "t",
},
],
}),
),
null,
"устойчивый offline — без повторного «перешёл» на каждом тике",
)
}
console.log("alert-engine evaluate-rule tests ok")
@@ -0,0 +1,323 @@
import { createHash } from "node:crypto"
import type { ApiAlertRule } from "../alerts-service.js"
import type { RuleEvalHit, SignalSnapshot } from "./types.js"
function hashPayload(parts: string[]): string {
return createHash("sha256").update(parts.join("\0")).digest("hex").slice(0, 24)
}
function firstNumberInString(s: string): number | null {
const m = s.replace(",", ".").match(/(\d+(?:\.\d+)?)/)
if (!m) return null
const n = Number.parseFloat(m[1] ?? "")
return Number.isFinite(n) ? n : null
}
function maxSeverity(
a: ApiAlertRule["severity"],
b: ApiAlertRule["severity"],
): ApiAlertRule["severity"] {
const o = { critical: 3, warning: 2, info: 1 } as const
return o[a] >= o[b] ? a : b
}
function findServer(snap: SignalSnapshot, name: string) {
const n = name.trim().toLowerCase()
return snap.servers.find((s) => s.name.trim().toLowerCase() === n)
}
function findProbe(snap: SignalSnapshot, targetKey: string) {
const k = targetKey.trim()
return snap.probes.find((p) => p.key.trim() === k)
}
function findTraffic(snap: SignalSnapshot, serverName: string) {
const n = serverName.trim().toLowerCase()
return snap.trafficByServer.find((t) => t.serverName.trim().toLowerCase() === n)
}
function findGreTunnel(snap: SignalSnapshot, targetLabel: string) {
const n = targetLabel.trim()
return snap.greTunnels.find((g) => g.targetLabel.trim() === n)
}
function findBgpPeer(snap: SignalSnapshot, key: string) {
const n = key.trim()
return snap.bgpSessions.find((b) => b.key.trim() === n)
}
function isBgpEstablished(state: string): boolean {
return state.trim().toLowerCase() === "established"
}
/** Заголовок для Telegram: не «…недоступен: … — восстановился». */
function formatServerAlertTitle(ruleName: string, wantRecovery: boolean): string {
if (!wantRecovery) return ruleName
if (!/недоступен/i.test(ruleName)) return ruleName
const trimmed = ruleName.replace(/\s*,\s*недоступен\s*$/i, "").replace(/\s+недоступен\s*$/i, "").trim()
if (!trimmed || trimmed === ruleName) return ruleName
return `${trimmed} — снова в сети`
}
function evalServerLike(
rule: ApiAlertRule,
snap: SignalSnapshot,
cond: string,
): RuleEvalHit | null {
const c = cond.toLowerCase()
const isOtkluchilsya = /отключился|отключилась|отключилось/.test(c)
/** «Перешёл в …» — только на границе, без спама пока сэмплы остаются offline. */
const impliesOfflineEdge = /переш(ё|е)л/.test(c) && (/offline|офлайн/.test(c) || /недоступен/.test(c))
const impliesDegradedEdge = /переш(ё|е)л/.test(c) && c.includes("degraded")
const wantOffline =
c.includes("offline") || c.includes("офлайн") || c.includes("отключ")
const wantDegraded = c.includes("degraded")
/** «Восстановился» / «подключился»: двойной online в сэмплах (см. тесты). Не отсекаем по слову «offline» в фразе вроде «после offline». */
const wantRecovery =
(c.includes("восстанов") || c.includes("подключ")) && !isOtkluchilsya
const hits: string[] = []
for (const tgt of rule.targets) {
const s = findServer(snap, tgt)
if (!s) continue
if (wantRecovery) {
const p2 = s.prev2Status
if (
s.status === "online" &&
s.prevStatus === "online" &&
p2 != null &&
(p2 === "offline" || p2 === "degraded")
) {
hits.push(`${tgt}`)
} else if (
s.status === "online" &&
(s.prevStatus === "offline" || s.prevStatus === "degraded")
) {
/** Один свежий online после offline (например ресурс уже online, REST только подтвердил). */
hits.push(`${tgt}`)
}
} else if ((impliesOfflineEdge || isOtkluchilsya) && s.status === "offline") {
if (s.prevStatus === null || s.prevStatus === "online") hits.push(tgt)
} else if (impliesDegradedEdge && s.status === "degraded") {
if (s.prevStatus === null || s.prevStatus === "online") hits.push(tgt)
} else if (wantDegraded && s.status === "degraded" && !impliesDegradedEdge) {
hits.push(tgt)
} else if (
wantOffline &&
s.status === "offline" &&
!impliesOfflineEdge &&
!isOtkluchilsya
) {
hits.push(tgt)
}
}
if (hits.length === 0) return null
const title = formatServerAlertTitle(rule.name, wantRecovery)
const msg = `${title}: ${hits.join(", ")}${cond}`
return {
ruleId: rule.id,
ruleName: rule.name,
severity: rule.severity,
message: msg,
transition: wantRecovery ? "recovery" : "problem",
payloadHash: hashPayload([rule.id, "server", hits.sort().join(","), cond]),
}
}
function evalRtt(rule: ApiAlertRule, snap: SignalSnapshot, cond: string): RuleEvalHit | null {
const threshold = firstNumberInString(cond)
const wantLow = /<|ниже|меньше/i.test(cond) && !/>|выше|больше/i.test(cond)
const wantHigh = />|выше|больше/i.test(cond)
const hits: string[] = []
for (const tgt of rule.targets) {
const p = findProbe(snap, tgt)
if (!p) continue
if (p.status === "down" || p.status === "warn") {
hits.push(`${tgt} (${p.status}, RTT ${p.rttMs ?? "—"} мс)`)
continue
}
if (p.rttMs == null) continue
if (wantHigh && threshold != null && p.rttMs > threshold) hits.push(`${tgt} (${p.rttMs} мс)`)
if (wantLow && threshold != null && p.rttMs < threshold) hits.push(`${tgt} (${p.rttMs} мс)`)
}
if (hits.length === 0) return null
const msg = `${rule.name}: RTT ${cond}${hits.join("; ")}`
return {
ruleId: rule.id,
ruleName: rule.name,
severity: rule.severity,
message: msg,
transition: "neutral",
payloadHash: hashPayload([rule.id, "rtt", hits.sort().join(";"), cond]),
}
}
function evalLoss(rule: ApiAlertRule, snap: SignalSnapshot, cond: string): RuleEvalHit | null {
const threshold = firstNumberInString(cond)
if (threshold == null) return null
const hits: string[] = []
for (const tgt of rule.targets) {
const p = findProbe(snap, tgt)
if (!p || p.lossPct == null) continue
if (p.lossPct > threshold) hits.push(`${tgt} (${p.lossPct}%)`)
}
if (hits.length === 0) return null
const msg = `${rule.name}: потери ${cond}${hits.join("; ")}`
return {
ruleId: rule.id,
ruleName: rule.name,
severity: rule.severity,
message: msg,
transition: "neutral",
payloadHash: hashPayload([rule.id, "loss", hits.sort().join(";"), cond]),
}
}
function evalGreTunnel(rule: ApiAlertRule, snap: SignalSnapshot, cond: string): RuleEvalHit | null {
const lc = cond.toLowerCase()
const wantOffline = lc.includes("offline") && !lc.includes("восстанов")
const wantRecovery =
lc.includes("восстанов") && !/offline|отключ/i.test(lc)
const hits: string[] = []
for (const tgt of rule.targets) {
const g = findGreTunnel(snap, tgt)
if (!g) continue
if (wantRecovery) {
const prev = g.prevStatus
if (
g.status === "up" &&
prev != null &&
(prev === "down" || prev === "degraded")
) {
hits.push(tgt)
}
} else if (wantOffline && g.status === "down") {
hits.push(tgt)
}
}
if (hits.length === 0) return null
const msg = `${rule.name}: ${hits.join(", ")}${cond}`
return {
ruleId: rule.id,
ruleName: rule.name,
severity: rule.severity,
message: msg,
transition: wantRecovery ? "recovery" : "problem",
payloadHash: hashPayload([rule.id, "gre-tunnel", hits.sort().join(","), cond]),
}
}
function evalBgpPeer(rule: ApiAlertRule, snap: SignalSnapshot, cond: string): RuleEvalHit | null {
const lc = cond.toLowerCase()
const wantDown = lc.includes("разорвал")
const wantUp = lc.includes("восстановил")
const hits: string[] = []
for (const tgt of rule.targets) {
const b = findBgpPeer(snap, tgt)
if (!b) continue
if (wantDown && !isBgpEstablished(b.state)) hits.push(tgt)
if (
wantUp &&
isBgpEstablished(b.state) &&
b.prevState != null &&
String(b.prevState).trim() !== "" &&
!isBgpEstablished(b.prevState)
) {
hits.push(tgt)
}
}
if (hits.length === 0) return null
const msg = `${rule.name}: ${hits.join(", ")}${cond}`
return {
ruleId: rule.id,
ruleName: rule.name,
severity: rule.severity,
message: msg,
transition: wantUp ? "recovery" : "problem",
payloadHash: hashPayload([rule.id, "bgp-peer", hits.sort().join(","), cond]),
}
}
function evalTraffic(rule: ApiAlertRule, snap: SignalSnapshot, cond: string): RuleEvalHit | null {
const threshold = firstNumberInString(cond)
if (threshold == null) return null
const rxLow = /RX\s*</i.test(cond)
const txLow = /TX\s*</i.test(cond)
const hits: string[] = []
for (const tgt of rule.targets) {
const t = findTraffic(snap, tgt)
if (!t) continue
if (rxLow && t.rxMbps < threshold) hits.push(`${tgt} RX ${t.rxMbps.toFixed(2)} Мбит/с`)
if (txLow && t.txMbps < threshold) hits.push(`${tgt} TX ${t.txMbps.toFixed(2)} Мбит/с`)
}
if (hits.length === 0) return null
const msg = `${rule.name}: трафик ${cond}${hits.join("; ")}`
return {
ruleId: rule.id,
ruleName: rule.name,
severity: rule.severity,
message: msg,
transition: "neutral",
payloadHash: hashPayload([rule.id, "traffic", hits.sort().join(";"), cond]),
}
}
function evaluateRuleWithCondition(
rule: ApiAlertRule,
snap: SignalSnapshot,
cond: string,
): RuleEvalHit | null {
const c = cond.trim()
if (!c) return null
switch (rule.type) {
case "server":
case "gre-client":
return evalServerLike(rule, snap, c)
case "rtt":
return evalRtt(rule, snap, c)
case "loss":
return evalLoss(rule, snap, c)
case "traffic":
return evalTraffic(rule, snap, c)
case "gre-tunnel":
return evalGreTunnel(rule, snap, c)
case "bgp-peer":
return evalBgpPeer(rule, snap, c)
case "bgp-prefix":
return null
default:
return null
}
}
/** Одно правило: OR по targets; OR по conditions (любое из выбранных). */
export function evaluateRule(rule: ApiAlertRule, snap: SignalSnapshot): RuleEvalHit | null {
if (!rule.enabled || rule.targets.length === 0) return null
const conds =
rule.conditions?.length > 0
? rule.conditions.map((c) => c.trim()).filter(Boolean)
: [rule.condition.trim()].filter(Boolean)
if (conds.length === 0) return null
const hits: RuleEvalHit[] = []
for (const cond of conds) {
const h = evaluateRuleWithCondition(rule, snap, cond)
if (h) hits.push(h)
}
if (hits.length === 0) return null
if (hits.length === 1) return hits[0]!
const message = hits.map((h) => h.message).join("\n")
const severity = hits.reduce((a, h) => maxSeverity(a, h.severity), hits[0]!.severity)
const payloadHash = hashPayload([rule.id, "or-conds", ...hits.map((h) => h.payloadHash).sort()])
return {
ruleId: rule.id,
ruleName: rule.name,
severity,
message,
transition: hits.some((h) => h.transition === "problem")
? "problem"
: hits.some((h) => h.transition === "recovery")
? "recovery"
: "neutral",
payloadHash,
}
}
@@ -0,0 +1,17 @@
import assert from "node:assert/strict"
import { groupShouldFire } from "./groups.js"
import { cooldownToMs, isCooldownElapsed } from "./cooldown.js"
assert.equal(groupShouldFire("any", []), false)
assert.equal(groupShouldFire("any", [false, false]), false)
assert.equal(groupShouldFire("any", [false, true]), true)
assert.equal(groupShouldFire("all", [true, true]), true)
assert.equal(groupShouldFire("all", [true, false]), false)
assert.equal(groupShouldFire("all", [false, false]), false)
assert.equal(cooldownToMs("1м"), 60_000)
assert.equal(isCooldownElapsed(null, "5м"), true)
assert.equal(isCooldownElapsed(new Date(Date.now() - 6 * 60_000).toISOString(), "5м"), true)
assert.equal(isCooldownElapsed(new Date(Date.now() - 2 * 60_000).toISOString(), "5м"), false)
console.log("alert-engine groups/cooldown tests ok")
@@ -0,0 +1,9 @@
/**
* Логика группы: ANY = хотя бы один член «в огне»; ALL = все члены «в огне».
* `memberFiring` — только включённые правила группы, по порядку.
*/
export function groupShouldFire(combineMode: "any" | "all", memberFiring: boolean[]): boolean {
if (memberFiring.length === 0) return false
if (combineMode === "any") return memberFiring.some(Boolean)
return memberFiring.every(Boolean)
}
+138
View File
@@ -0,0 +1,138 @@
import { and, asc, eq, lte } from "drizzle-orm"
import { db } from "../../db/index.js"
import { alertOutbox } from "../../db/schema.js"
import { appendAlertHistory, sendTelegramAlertMessage } from "../alerts-service.js"
export type AlertOutboxPayload = {
text: string
chatId?: string
history: {
id: string
ruleId?: string | null
groupId?: string | null
ruleName: string
severity: "critical" | "warning" | "info"
message: string
firedAt: string
}
}
function safeParsePayload(raw: string): AlertOutboxPayload | null {
try {
return JSON.parse(raw) as AlertOutboxPayload
} catch {
return null
}
}
export function enqueueTelegramOutbox(input: {
id: string
dedupeKey: string
payload: AlertOutboxPayload
maxRetries?: number
}): boolean {
const existing = db
.select()
.from(alertOutbox)
.where(eq(alertOutbox.dedupeKey, input.dedupeKey))
.limit(1)
.all()[0]
if (existing && (existing.status === "pending" || existing.status === "sent")) return false
const nowIso = new Date().toISOString()
db.insert(alertOutbox)
.values({
id: input.id,
dedupeKey: input.dedupeKey,
channel: "telegram",
status: "pending",
retryCount: 0,
maxRetries: Math.max(1, Math.floor(input.maxRetries ?? 3)),
nextAttemptAt: nowIso,
payloadJson: JSON.stringify(input.payload),
})
.run()
return true
}
export function dispatchPendingOutbox(limit = 25): Promise<{ sent: number; failed: number; errors: string[] }> {
return (async () => {
const nowIso = new Date().toISOString()
const rows = db
.select()
.from(alertOutbox)
.where(and(eq(alertOutbox.status, "pending"), lte(alertOutbox.nextAttemptAt, nowIso)))
.orderBy(asc(alertOutbox.createdAt))
.limit(Math.max(1, limit))
.all()
let sent = 0
let failed = 0
const errors: string[] = []
for (const row of rows) {
const payload = safeParsePayload(row.payloadJson)
if (!payload) {
db.update(alertOutbox)
.set({
status: "failed",
lastError: "invalid payload_json",
})
.where(eq(alertOutbox.id, row.id))
.run()
failed += 1
continue
}
const send = await sendTelegramAlertMessage({ text: payload.text, chatId: payload.chatId })
if (send.ok) {
db.update(alertOutbox)
.set({
status: "sent",
sentAt: new Date().toISOString(),
lastError: null,
})
.where(eq(alertOutbox.id, row.id))
.run()
appendAlertHistory({
id: payload.history.id,
ruleId: payload.history.ruleId ?? null,
groupId: payload.history.groupId ?? null,
ruleName: payload.history.ruleName,
severity: payload.history.severity,
message: payload.history.message,
sentOk: true,
firedAt: payload.history.firedAt,
})
sent += 1
continue
}
const nextRetry = row.retryCount + 1
const exhausted = nextRetry >= row.maxRetries
db.update(alertOutbox)
.set({
retryCount: nextRetry,
status: exhausted ? "failed" : "pending",
nextAttemptAt: new Date(Date.now() + Math.min(300_000, nextRetry * 15_000)).toISOString(),
lastError: send.error,
})
.where(eq(alertOutbox.id, row.id))
.run()
if (exhausted) {
appendAlertHistory({
id: payload.history.id,
ruleId: payload.history.ruleId ?? null,
groupId: payload.history.groupId ?? null,
ruleName: payload.history.ruleName,
severity: payload.history.severity,
message: payload.history.message,
sentOk: false,
firedAt: payload.history.firedAt,
})
failed += 1
} else {
errors.push(`outbox ${row.id}: ${send.error}`)
}
}
return { sent, failed, errors }
})()
}
@@ -0,0 +1,47 @@
import { eq } from "drizzle-orm"
import { db } from "../../db/index.js"
import { alertEnginePrevLive } from "../../db/schema.js"
export type PrevLiveKind = "gre" | "bgp"
/** Предыдущие строковые состояния (GRE: up|down|degraded; BGP: state как с роутера). */
export type PrevLiveStringMap = Record<string, string>
function safeParseMap(json: string | null | undefined): PrevLiveStringMap {
if (!json || json.trim() === "") return {}
try {
const v = JSON.parse(json) as unknown
if (v == null || typeof v !== "object" || Array.isArray(v)) return {}
return v as PrevLiveStringMap
} catch {
return {}
}
}
export function loadPrevLiveMap(kind: PrevLiveKind): PrevLiveStringMap {
const row = db
.select()
.from(alertEnginePrevLive)
.where(eq(alertEnginePrevLive.kind, kind))
.limit(1)
.all()[0]
return safeParseMap(row?.payloadJson)
}
export function savePrevLiveMap(kind: PrevLiveKind, map: PrevLiveStringMap) {
const payloadJson = JSON.stringify(map)
const existing = db
.select()
.from(alertEnginePrevLive)
.where(eq(alertEnginePrevLive.kind, kind))
.limit(1)
.all()[0]
if (existing) {
db.update(alertEnginePrevLive)
.set({ payloadJson, updatedAt: new Date().toISOString() })
.where(eq(alertEnginePrevLive.kind, kind))
.run()
} else {
db.insert(alertEnginePrevLive).values({ kind, payloadJson, updatedAt: new Date().toISOString() }).run()
}
}
@@ -0,0 +1,20 @@
import type { ApiAlertRule } from "../alerts-service.js"
import { evaluateRule } from "./evaluate-rule.js"
import type { RuleEvalHit, SignalSnapshot } from "./types.js"
export function matchRules(rules: ApiAlertRule[], snap: SignalSnapshot): {
hitByRule: Map<string, RuleEvalHit | null>
errors: string[]
} {
const errors: string[] = []
const hitByRule = new Map<string, RuleEvalHit | null>()
for (const rule of rules) {
try {
hitByRule.set(rule.id, evaluateRule(rule, snap))
} catch (e) {
errors.push(`rule ${rule.id}: ${e instanceof Error ? e.message : String(e)}`)
hitByRule.set(rule.id, null)
}
}
return { hitByRule, errors }
}
@@ -0,0 +1,194 @@
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"
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)
return {
sampledAt,
rulesChecked: hasNewSources ? rules.length : 0,
standaloneFires,
groupFires,
skippedNoTelegram: !canSend,
errors,
ruleDiag: ruleDiag ?? [],
}
}
@@ -0,0 +1,6 @@
import { buildSignalSnapshotFromCollectors } from "./signals-from-collectors.js"
/** Единая точка входа для подготовки сигналов к rule matching. */
export function ingestAlertSignals() {
return buildSignalSnapshotFromCollectors()
}
@@ -0,0 +1,265 @@
import { sqliteDatabase } from "../../db/index.js"
import type {
GreBgpSnapshotRunSnapshot,
PingRunSnapshot,
ResourcesRunSnapshot,
SchedulerRunSnapshot,
ServersRestPingRunSnapshot,
TrafficRunSnapshot,
} from "../../types/scheduler-run-snapshot.js"
import type {
BgpPeerSignal,
GreTunnelSignal,
ProbeSignal,
ServerSignal,
SignalSnapshot,
TrafficServerSignal,
} from "./types.js"
function normalizeNameKey(v: string): string {
return v.trim().toLowerCase()
}
function safeParseRunSnapshot(raw: string | null): SchedulerRunSnapshot | null {
if (!raw) return null
try {
return JSON.parse(raw) as SchedulerRunSnapshot
} catch {
return null
}
}
function readLatestOkRunSnapshots(jobKey: string, limit = 32): SchedulerRunSnapshot[] {
const rows = sqliteDatabase
.prepare(
`SELECT result_json AS resultJson
FROM scheduler_runs
WHERE job_key = ? AND status = 'ok' AND result_json IS NOT NULL
ORDER BY finished_at DESC
LIMIT ?`,
)
.all(jobKey, limit) as { resultJson: string | null }[]
const out: SchedulerRunSnapshot[] = []
for (const row of rows) {
const parsed = safeParseRunSnapshot(row.resultJson)
if (!parsed || parsed.job !== jobKey) continue
out.push(parsed)
}
return out
}
type TriState = "online" | "offline" | "degraded"
function timelineToServerSignal(name: string, sampledAt: string, statuses: TriState[]): ServerSignal {
return {
name,
status: statuses[0] ?? "offline",
prevStatus: statuses[1] ?? null,
prev2Status: statuses[2] ?? null,
sampledAt,
}
}
function collectServerSignalsFromResourceRuns(runs: ResourcesRunSnapshot[]): Map<string, ServerSignal> {
const byKey = new Map<
string,
{ name: string; sampledAt: string; statuses: TriState[] }
>()
for (const run of runs) {
for (const s of run.servers ?? []) {
const name = String(s.name ?? "").trim()
if (!name) continue
const key = normalizeNameKey(name)
const rec = byKey.get(key) ?? { name, sampledAt: run.sampledAt, statuses: [] }
if (rec.statuses.length < 3) rec.statuses.push(s.status === "online" ? "online" : "offline")
if (!byKey.has(key)) rec.sampledAt = run.sampledAt
byKey.set(key, rec)
}
}
const out = new Map<string, ServerSignal>()
for (const [key, rec] of byKey) out.set(key, timelineToServerSignal(rec.name, rec.sampledAt, rec.statuses))
return out
}
function collectServerSignalsFromRestRuns(runs: ServersRestPingRunSnapshot[]): Map<string, ServerSignal> {
const byKey = new Map<
string,
{ name: string; sampledAt: string; statuses: TriState[] }
>()
for (const run of runs) {
for (const s of run.servers ?? []) {
const name = String(s.name ?? "").trim()
if (!name) continue
const key = normalizeNameKey(name)
const rec = byKey.get(key) ?? { name, sampledAt: run.sampledAt, statuses: [] }
if (rec.statuses.length < 3) rec.statuses.push(s.ok ? "online" : "offline")
if (!byKey.has(key)) rec.sampledAt = run.sampledAt
byKey.set(key, rec)
}
}
const out = new Map<string, ServerSignal>()
for (const [key, rec] of byKey) out.set(key, timelineToServerSignal(rec.name, rec.sampledAt, rec.statuses))
return out
}
function mergeResourceAndRestServer(resource: ServerSignal, rest: ServerSignal): ServerSignal {
const rUp = resource.status === "online"
const tUp = rest.status === "online"
if (tUp && !rUp) {
const rBad = resource.status === "offline" || resource.status === "degraded"
if (rBad) {
return {
name: rest.name,
status: "online",
prevStatus: "online",
prev2Status: resource.status,
sampledAt: rest.sampledAt,
}
}
return rest
}
if (rUp && !tUp) return resource
if (tUp && rUp) {
if (rest.sampledAt >= resource.sampledAt) {
if (resource.prevStatus === "offline" || resource.prevStatus === "degraded") {
return {
name: rest.name,
status: "online",
prevStatus: "online",
prev2Status: resource.prevStatus === "degraded" ? "degraded" : "offline",
sampledAt: rest.sampledAt,
}
}
if (
resource.prevStatus === "online" &&
resource.prev2Status != null &&
(resource.prev2Status === "offline" || resource.prev2Status === "degraded")
) {
return {
name: rest.name,
status: "online",
prevStatus: "online",
prev2Status: resource.prev2Status,
sampledAt: rest.sampledAt,
}
}
}
return rest.sampledAt >= resource.sampledAt ? rest : resource
}
return rest.sampledAt >= resource.sampledAt ? rest : resource
}
function collectProbeSignalsFromRuns(runs: PingRunSnapshot[]): ProbeSignal[] {
const out = new Map<string, ProbeSignal>()
for (const run of runs) {
for (const p of run.probes ?? []) {
const key = `${String(p.name ?? "").trim()}${String(p.target ?? "").trim()}`
if (!key.trim() || out.has(key)) continue
out.set(key, {
key,
status: p.status,
rttMs: p.rttMs ?? null,
lossPct: p.lossPct ?? null,
sampledAt: run.sampledAt,
})
}
}
return [...out.values()]
}
function collectTrafficSignalsFromRuns(runs: TrafficRunSnapshot[]): TrafficServerSignal[] {
const out = new Map<string, TrafficServerSignal>()
for (const run of runs) {
for (const s of run.servers ?? []) {
const name = String(s.name ?? "").trim()
if (!name) continue
const key = normalizeNameKey(name)
if (out.has(key)) continue
out.set(key, {
serverName: name,
rxMbps: Number(s.sumRxMbps ?? 0),
txMbps: Number(s.sumTxMbps ?? 0),
sampledAt: run.sampledAt,
})
}
}
return [...out.values()]
}
function collectGreSignalsFromRuns(runs: GreBgpSnapshotRunSnapshot[]): GreTunnelSignal[] {
const byTarget = new Map<string, GreTunnelSignal>()
for (const run of runs) {
const greRows = run.greTunnels ?? []
for (const row of greRows) {
const label = String(row.targetLabel ?? "").trim()
if (!label) continue
const prev = byTarget.get(label)
if (!prev) {
byTarget.set(label, {
targetLabel: label,
status: row.status,
prevStatus: null,
})
continue
}
if (prev.prevStatus == null) prev.prevStatus = row.status
}
}
return [...byTarget.values()]
}
function collectBgpSignalsFromRuns(runs: GreBgpSnapshotRunSnapshot[]): BgpPeerSignal[] {
const byKey = new Map<string, BgpPeerSignal>()
for (const run of runs) {
const peers = run.bgpPeers ?? []
for (const row of peers) {
const key = String(row.key ?? "").trim()
if (!key) continue
const prev = byKey.get(key)
if (!prev) {
byKey.set(key, { key, state: String(row.state ?? "").trim(), prevState: null })
continue
}
if (prev.prevState == null) prev.prevState = String(row.state ?? "").trim()
}
}
return [...byKey.values()]
}
/** Собирает сигналы alert_engine напрямую из результатов collector jobs (`scheduler_runs.result_json`). */
export function buildSignalSnapshotFromCollectors(): SignalSnapshot {
const sampledAt = new Date().toISOString()
const resourceRuns = readLatestOkRunSnapshots("uptime_resources", 16) as ResourcesRunSnapshot[]
const restRuns = readLatestOkRunSnapshots("servers_rest_ping", 16) as ServersRestPingRunSnapshot[]
const pingRuns = readLatestOkRunSnapshots("uptime_ping", 24) as PingRunSnapshot[]
const trafficRuns = readLatestOkRunSnapshots("traffic", 8) as TrafficRunSnapshot[]
const greBgpRuns = readLatestOkRunSnapshots("gre_bgp", 8) as GreBgpSnapshotRunSnapshot[]
const resourceByKey = collectServerSignalsFromResourceRuns(resourceRuns)
const restByKey = collectServerSignalsFromRestRuns(restRuns)
const serverKeys = new Set([...resourceByKey.keys(), ...restByKey.keys()])
const servers: ServerSignal[] = []
for (const key of serverKeys) {
const resource = resourceByKey.get(key)
const rest = restByKey.get(key)
if (!rest) {
if (resource) servers.push(resource)
continue
}
if (!resource) {
servers.push(rest)
continue
}
servers.push(mergeResourceAndRestServer(resource, rest))
}
servers.sort((a, b) => a.name.localeCompare(b.name, undefined, { sensitivity: "base" }))
return {
sampledAt,
servers,
probes: collectProbeSignalsFromRuns(pingRuns),
trafficByServer: collectTrafficSignalsFromRuns(trafficRuns),
greTunnels: collectGreSignalsFromRuns(greBgpRuns),
bgpSessions: collectBgpSignalsFromRuns(greBgpRuns),
}
}
@@ -0,0 +1,7 @@
import { buildSignalSnapshotFromCollectors } from "./signals-from-collectors.js"
import type { SignalSnapshot } from "./types.js"
/** Legacy-совместимость: единый источник сигналов — collector snapshots из scheduler_runs. */
export function buildSignalSnapshot(): SignalSnapshot {
return buildSignalSnapshotFromCollectors()
}
@@ -0,0 +1,54 @@
import { desc, eq, inArray } from "drizzle-orm"
import { db } from "../../db/index.js"
import { alertEngineCursor, schedulerRuns } from "../../db/schema.js"
const SOURCE_JOBS = [
"traffic",
"servers_rest_ping",
"uptime_resources",
"uptime_ping",
"uptime_speed",
"gre_bgp",
] as const
export function getLatestSourceFinishedAt(): string | null {
const row = db
.select({ finishedAt: schedulerRuns.finishedAt })
.from(schedulerRuns)
.where(inArray(schedulerRuns.jobKey, [...SOURCE_JOBS]))
.orderBy(desc(schedulerRuns.finishedAt))
.limit(1)
.all()[0]
return row?.finishedAt ?? null
}
export function getSourceWatermark(): string | null {
const row = db.select().from(alertEngineCursor).where(eq(alertEngineCursor.id, 1)).limit(1).all()[0]
return row?.lastSourceFinishedAt ?? null
}
export function updateSourceWatermark(lastSourceFinishedAt: string | null): void {
const existing = db
.select()
.from(alertEngineCursor)
.where(eq(alertEngineCursor.id, 1))
.limit(1)
.all()[0]
if (existing) {
db.update(alertEngineCursor)
.set({
lastSourceFinishedAt,
updatedAt: new Date().toISOString(),
})
.where(eq(alertEngineCursor.id, 1))
.run()
return
}
db.insert(alertEngineCursor)
.values({
id: 1,
lastSourceFinishedAt,
updatedAt: new Date().toISOString(),
})
.run()
}
@@ -0,0 +1,23 @@
import assert from "node:assert/strict"
import { shouldAllowRecoveryByPolicy } from "./state-machine.js"
const recoveryHit = {
ruleId: "r",
ruleName: "r",
severity: "warning" as const,
message: "восстановился",
transition: "recovery" as const,
payloadHash: "h1",
}
const problemHit = {
...recoveryHit,
transition: "problem" as const,
}
assert.equal(shouldAllowRecoveryByPolicy(recoveryHit, "always"), true)
assert.equal(shouldAllowRecoveryByPolicy(recoveryHit, "conditional"), true)
assert.equal(shouldAllowRecoveryByPolicy(recoveryHit, "never"), false)
assert.equal(shouldAllowRecoveryByPolicy(problemHit, "never"), true)
console.log("alert-engine state-machine tests ok")
@@ -0,0 +1,16 @@
import type { RuleEvalHit } from "./types.js"
export type RecoveryMode = "always" | "never" | "conditional"
export function isRecoveryTransition(hit: RuleEvalHit | null): boolean {
return hit?.transition === "recovery"
}
export function shouldAllowRecoveryByPolicy(
hit: RuleEvalHit | null,
policy: RecoveryMode,
): boolean {
if (!isRecoveryTransition(hit)) return true
if (policy === "never") return false
return true
}
@@ -0,0 +1,10 @@
import assert from "node:assert/strict"
import { isPositiveRecoveryTelegramText, pickTelegramAlertEmoji } from "./telegram-emoji.js"
assert.equal(pickTelegramAlertEmoji("Сервер: x — снова в сети: x — восстановился"), "✅")
assert.equal(pickTelegramAlertEmoji("Сервер: x недоступен: x — перешёл в offline"), "🚨")
assert.equal(isPositiveRecoveryTelegramText("A — восстановился"), true)
assert.equal(isPositiveRecoveryTelegramText("A — перешёл в offline"), false)
assert.equal(isPositiveRecoveryTelegramText("BGP: p — восстановил сессию"), true)
console.log("telegram-emoji tests ok")
@@ -0,0 +1,27 @@
/**
* Эмодзи в Telegram: восстановление / «подключился» / BGP session up — не сирена.
* Текст совпадает с форматом `evaluate-rule` (`… — условие`, несколько строк через \\n).
*/
function isRecoveryLine(line: string): boolean {
const t = line.toLowerCase()
return (
/[—\-]\s*(восстанов|подключ)/i.test(t) ||
/восстановил[аи]?\s+сессию/i.test(t) ||
/снова\s+в\s+сети/i.test(t)
)
}
/** «Хорошие» переходы: не слать дубль при том же payloadHash (см. `run-once`). */
export function isPositiveRecoveryTelegramText(fullText: string): boolean {
const lines = fullText
.split(/\n+/)
.map((l) => l.trim())
.filter(Boolean)
if (lines.length === 0) return false
if (lines.length === 1) return isRecoveryLine(lines[0]!)
return lines.every(isRecoveryLine)
}
export function pickTelegramAlertEmoji(fullText: string): string {
return isPositiveRecoveryTelegramText(fullText) ? "✅" : "🚨"
}
@@ -0,0 +1,77 @@
import type { ApiAlertRule } from "../alerts-service.js"
export type ServerSignal = {
name: string
status: "online" | "offline" | "degraded"
/** Предпоследний сэмпл (для условий «восстановился» / «подключился») */
prevStatus: "online" | "offline" | "degraded" | null
/** Третий с конца сэмпл — чтобы «восстановился» не висел истинным до второго online-сэмпла в БД. */
prev2Status: "online" | "offline" | "degraded" | null
sampledAt: string
}
export type ProbeSignal = {
/** Как в UI: `имя пробы → target` */
key: string
status: "up" | "warn" | "down"
rttMs: number | null
lossPct: number | null
sampledAt: string
}
export type TrafficServerSignal = {
serverName: string
rxMbps: number
txMbps: number
sampledAt: string
}
export type GreTunnelSignal = {
/** Как в UI: `tunnelName / serverName` */
targetLabel: string
status: "up" | "down" | "degraded"
/** Состояние на прошлом тике движка (из БД), для «восстановился». */
prevStatus: "up" | "down" | "degraded" | null
}
export type BgpPeerSignal = {
/** Как в UI: `serverName: peerName` */
key: string
state: string
/** Состояние на прошлом тике (строка state с роутера). */
prevState: string | null
}
export type SignalSnapshot = {
sampledAt: string
servers: ServerSignal[]
probes: ProbeSignal[]
trafficByServer: TrafficServerSignal[]
greTunnels: GreTunnelSignal[]
bgpSessions: BgpPeerSignal[]
}
export type RuleEvalHit = {
ruleId: string
ruleName: string
severity: ApiAlertRule["severity"]
message: string
/** Тип перехода для диагностики и логики задержки отправки. */
transition: "problem" | "recovery" | "neutral"
/** Для дедупа */
payloadHash: string
}
import type { AlertEngineRuleDiagSnapshot } from "../../types/scheduler-run-snapshot.js"
export type AlertEngineRuleDiag = AlertEngineRuleDiagSnapshot
export type AlertEngineRunResult = {
sampledAt: string
rulesChecked: number
standaloneFires: number
groupFires: number
skippedNoTelegram: boolean
errors: string[]
ruleDiag: AlertEngineRuleDiag[]
}