Files
MikrotikManager/backend/src/services/scheduler.ts
T
Denozordec a80caf5676
Docker images / prepare-release (push) Successful in 13s
Docker images / backend-test (push) Successful in 2m13s
Docker images / frontend-image (push) Successful in 4m31s
Docker images / updater-image (push) Successful in 48s
Docker images / backend-image (push) Successful in 3m1s
Docker images / notify-webhook (push) Skipped
Docker images / publish-release (push) Successful in 15s
feat(netflow): локальные GeoLite2-базы для стран и ASN потоков
Lookup страны/ASN при ingest теперь идёт сначала по локальным mmdb
MaxMind GeoLite2 (Country + ASN, зеркало P3TERX, без ключей), с
мгновенным синхронным ответом для IPv4/IPv6; RIPEstat остаётся
fallback до первой загрузки баз и при промахе.

- backend/src/services/traffic-flow-geoip.ts: ридеры maxmind, фасад
  resolveFlowIp (geoip-first → RIPE-кэш), статус ридеров
- geoip-update-collector.ts: conditional GET по ETag, валидация пробоем
  8.8.8.8 (US / AS15169), атомарная подмена с .prev-откатом
- джоба планировщика geoip_update (по умолчанию раз в 7 дней),
  настройки geoip_settings + миграция 0004, API /api/geoip
  (GET/PUT/update), контракты @mmapp/contracts/geoip
- engine/analytics/map-hops переведены на resolveFlowIp
- секция «GeoIP-базы» в настройках NetFlow, метки джобы на странице
  сбора данных, тесты test:geoip
2026-09-09 23:33:31 +07:00

489 lines
16 KiB
TypeScript

/**
* In-process планировщик: отдельный интервал на job (`refreshScheduler`).
*
* **Приоритеты при конкуренции** (для дальнейшего per-server mutex; сейчас зафиксировано в дизайне):
* 1. `uptime_ping` — выше (короткий критичный сигнал).
* 2. `traffic`, `uptime_resources`, `servers_rest_ping` — средний (REST `/system/identity` по каталогу).
* 3. `uptime_speed` — ниже (BW-test, тяжёлый); локи узлов — `withBtestNodeLocks` в speed-сервисе.
*
* **Оповещения (`alert_engine`):** читают PostgreSQL после джоб сбора (см. `buildSignalSnapshot`), в т.ч.
* `gre_bgp` → snapshot в `scheduler_runs.result_json` (`greTunnels` / `bgpPeers`). После успешного завершения джоб
* `traffic`, `uptime_*`, `servers_rest_ping`, `gre_bgp` планируется **дополнительный** прогон движка
* (debounce), см. [`alert-collector-hooks.ts`](./alert-collector-hooks.ts); любой другой писатель сэмплов
* для снимка оповещений тоже должен вызывать `scheduleAlertEngineAfterDataCollectors()` после коммита.
*
* **Кастомные job без произвольного кода (MVP):** пер-тайминг ping в `uptime_probes.interval_sec`;
* таблица `scheduler_job_instances` (kind + config JSON) — при необходимости следующий этап.
*/
import { desc, eq, lt } from "drizzle-orm"
import { db } from "../db/index.js"
import { events, schedulerRuns } from "../db/schema.js"
import type { SchedulerRunSnapshot } from "../types/scheduler-run-snapshot.js"
import { SCHEDULER_RUN_SNAPSHOT_VERSION } from "../types/scheduler-run-snapshot.js"
import {
collectServersRestPingOnce,
getServersApiPingSettings,
} from "./servers-rest-ping-collector.js"
import { collectTrafficOnce, getTrafficSettings } from "./traffic-collector.js"
import {
collectPingProbesOnce,
collectResourceSamplesOnce,
getSettings as getUptimeSettings,
} from "./uptime-collector.js"
import { runScheduledSpeedProbesOnce } from "./uptime-speed-service.js"
import { collectGreBgpSnapshotOnce } from "./gre-bgp-snapshot-collector.js"
import { runAlertEngineOnce } from "./alert-engine/run-once.js"
import {
clearAlertCollectorHooksTimer,
scheduleAlertEngineAfterDataCollectors,
wireAlertEngineRunner,
} from "./alert-collector-hooks.js"
import {
collectInternetPathSnapshotOnce,
getInternetPathSettings,
} from "./internet-path-collector.js"
import { collectCertificatesRenewOnce } from "./certificate-renew-collector.js"
import { getCertificateRenewSettings } from "./certificates-service.js"
import { collectScheduledBackupsOnce } from "./backup-scheduler-collector.js"
import { getBackupScheduleSettings } from "./backup-service.js"
import { collectGeoipUpdateOnce } from "./geoip-update-collector.js"
import { getGeoipSettings } from "./geoip-settings.js"
import {
endSchedulerJob,
isSchedulerJobRunning,
tryBeginSchedulerJob,
} from "./scheduler-running.js"
import { appendEvent } from "../modules/events/service/events-service.js"
export const JOB_KEYS = [
"traffic",
"servers_rest_ping",
"uptime_resources",
"uptime_ping",
"uptime_speed",
"internet_path",
"gre_bgp",
"certificates_renew",
"backups",
"geoip_update",
"alert_engine",
] as const
export type SchedulerJobKey = (typeof JOB_KEYS)[number]
const timers = new Map<string, ReturnType<typeof setInterval>>()
const RUN_LOG_RETENTION_MS = 30 * 24 * 60 * 60 * 1000
const QUIET_SCHEDULER_OK_JOBS = new Set<SchedulerJobKey>([
"traffic",
"servers_rest_ping",
"uptime_resources",
"uptime_ping",
"uptime_speed",
"gre_bgp",
"alert_engine",
])
export function shouldAppendSchedulerOkEvent(jobKey: SchedulerJobKey): boolean {
return !QUIET_SCHEDULER_OK_JOBS.has(jobKey)
}
function newRunId(): string {
return `sch-${Date.now()}-${Math.random().toString(36).slice(2, 10)}`
}
async function appendSchedulerRun(row: {
jobKey: string
startedAt: string
finishedAt: string
status: "ok" | "error"
error: string | null
durationMs: number
result?: SchedulerRunSnapshot | null
}) {
const cutoff = new Date(Date.now() - RUN_LOG_RETENTION_MS).toISOString()
await db.transaction(async (tx) => {
await tx.insert(schedulerRuns).values({
id: newRunId(),
jobKey: row.jobKey,
startedAt: row.startedAt,
finishedAt: row.finishedAt,
status: row.status,
error: row.error,
durationMs: row.durationMs,
resultJson: row.result ?? null,
})
await tx.delete(schedulerRuns).where(lt(schedulerRuns.finishedAt, cutoff))
await tx.delete(events).where(lt(events.createdAt, cutoff))
})
}
async function runSchedulerJobBody(jobKey: SchedulerJobKey): Promise<void> {
const startedAt = Date.now()
const startedIso = new Date().toISOString()
let snapshot: SchedulerRunSnapshot | undefined
try {
switch (jobKey) {
case "traffic":
snapshot = await collectTrafficOnce()
break
case "uptime_resources":
snapshot = await collectResourceSamplesOnce()
break
case "uptime_ping":
snapshot = await collectPingProbesOnce()
break
case "uptime_speed":
snapshot = await runScheduledSpeedProbesOnce()
break
case "servers_rest_ping":
snapshot = await collectServersRestPingOnce()
break
case "gre_bgp":
snapshot = await collectGreBgpSnapshotOnce()
break
case "internet_path":
snapshot = await collectInternetPathSnapshotOnce()
break
case "certificates_renew":
snapshot = await collectCertificatesRenewOnce()
break
case "backups":
snapshot = await collectScheduledBackupsOnce()
break
case "geoip_update":
snapshot = await collectGeoipUpdateOnce()
break
case "alert_engine": {
const r = await runAlertEngineOnce()
snapshot = {
v: SCHEDULER_RUN_SNAPSHOT_VERSION,
job: "alert_engine",
sampledAt: r.sampledAt,
rulesChecked: r.rulesChecked,
standaloneFires: r.standaloneFires,
groupFires: r.groupFires,
skippedNoTelegram: r.skippedNoTelegram,
...(r.ruleDiag.length ? { ruleDiag: r.ruleDiag } : {}),
...(r.errors.length ? { errors: r.errors } : {}),
}
break
}
default:
throw new Error(`Unknown job: ${jobKey}`)
}
const finishedIso = new Date().toISOString()
await appendSchedulerRun({
jobKey,
startedAt: startedIso,
finishedAt: finishedIso,
status: "ok",
error: null,
durationMs: Date.now() - startedAt,
result: snapshot ?? null,
})
if (shouldAppendSchedulerOkEvent(jobKey)) {
await appendEvent({
level: "info",
eventType: "scheduler.job.ok",
sourceModule: "scheduler",
title: "Задача планировщика завершена",
message: `${jobKey}: выполнено за ${Date.now() - startedAt} мс`,
entityType: "job",
entityId: jobKey,
payload: {
startedAt: startedIso,
finishedAt: finishedIso,
},
})
}
if (
jobKey === "traffic" ||
jobKey === "servers_rest_ping" ||
jobKey === "uptime_resources" ||
jobKey === "uptime_ping" ||
jobKey === "uptime_speed" ||
jobKey === "internet_path" ||
jobKey === "gre_bgp"
) {
scheduleAlertEngineAfterDataCollectors()
}
} catch (e) {
const msg = e instanceof Error ? e.message : String(e)
const finishedIso = new Date().toISOString()
await appendSchedulerRun({
jobKey,
startedAt: startedIso,
finishedAt: finishedIso,
status: "error",
error: msg,
durationMs: Date.now() - startedAt,
result: snapshot ?? null,
})
await appendEvent({
level: "critical",
eventType: "scheduler.job.failed",
sourceModule: "scheduler",
title: "Ошибка задачи планировщика",
message: `${jobKey}: ${msg}`,
entityType: "job",
entityId: jobKey,
payload: {
startedAt: startedIso,
finishedAt: finishedIso,
},
})
throw e
}
}
/** Фоновый тик: пропуск, если предыдущий прогон ещё идёт. */
export async function executeSchedulerJob(jobKey: SchedulerJobKey): Promise<void> {
if (!tryBeginSchedulerJob(jobKey)) return
try {
await runSchedulerJobBody(jobKey)
} catch {
/* залогировано в runSchedulerJobBody */
} finally {
endSchedulerJob(jobKey)
}
}
/** Ручной запуск: 409, если job уже выполняется. */
export async function runSchedulerJobNow(jobKey: string): Promise<void> {
if (!JOB_KEYS.includes(jobKey as SchedulerJobKey)) {
throw new Error(`Unknown job key: ${jobKey}`)
}
if (!tryBeginSchedulerJob(jobKey)) {
const err = new Error("Job already running")
;(err as Error & { statusCode?: number }).statusCode = 409
throw err
}
try {
await runSchedulerJobBody(jobKey as SchedulerJobKey)
} finally {
endSchedulerJob(jobKey)
}
}
function clearAllTimers() {
clearAlertCollectorHooksTimer()
for (const t of timers.values()) clearInterval(t)
timers.clear()
}
/** Пересоздать интервалы после смены настроек. */
export async function refreshScheduler(): Promise<void> {
clearAllTimers()
const traffic = await getTrafficSettings()
if (traffic.enabled) {
const ms = Math.max(5_000, traffic.intervalSec * 1000)
void executeSchedulerJob("traffic").catch(() => {})
timers.set(
"traffic",
setInterval(() => {
void executeSchedulerJob("traffic").catch(() => {})
}, ms),
)
}
const apiPing = await getServersApiPingSettings()
const internetPath = await getInternetPathSettings()
if (apiPing.enabled) {
const apiMs = Math.max(10_000, apiPing.intervalSec * 1000)
void executeSchedulerJob("servers_rest_ping").catch(() => {})
timers.set(
"servers_rest_ping",
setInterval(() => {
void executeSchedulerJob("servers_rest_ping").catch(() => {})
}, apiMs),
)
}
const uptime = await getUptimeSettings()
const resOn = uptime.resourcesEnabled ?? uptime.enabled
const pingOn = uptime.pingEnabled ?? uptime.enabled
const spdOn = uptime.speedEnabled ?? uptime.enabled
if (resOn) {
const resMs = Math.max(5_000, uptime.intervalSec * 1000)
void executeSchedulerJob("uptime_resources").catch(() => {})
timers.set(
"uptime_resources",
setInterval(() => {
void executeSchedulerJob("uptime_resources").catch(() => {})
}, resMs),
)
}
if (pingOn) {
const pingMs = Math.max(1_000, uptime.probeIntervalSec * 1000)
void executeSchedulerJob("uptime_ping").catch(() => {})
timers.set(
"uptime_ping",
setInterval(() => {
void executeSchedulerJob("uptime_ping").catch(() => {})
}, pingMs),
)
}
if (spdOn) {
const spdMs = Math.max(10_000, uptime.speedIntervalSec * 1000)
void executeSchedulerJob("uptime_speed").catch(() => {})
timers.set(
"uptime_speed",
setInterval(() => {
void executeSchedulerJob("uptime_speed").catch(() => {})
}, spdMs),
)
}
if (internetPath.enabled) {
const internetPathMs = Math.max(30_000, internetPath.intervalSec * 1000)
void executeSchedulerJob("internet_path").catch(() => {})
timers.set(
"internet_path",
setInterval(() => {
void executeSchedulerJob("internet_path").catch(() => {})
}, internetPathMs),
)
}
const greBgpMs = 30_000
void executeSchedulerJob("gre_bgp").catch(() => {})
timers.set(
"gre_bgp",
setInterval(() => {
void executeSchedulerJob("gre_bgp").catch(() => {})
}, greBgpMs),
)
const certRenew = await getCertificateRenewSettings()
if (certRenew.enabled) {
const certRenewMs = Math.max(300_000, certRenew.intervalSec * 1000)
void executeSchedulerJob("certificates_renew").catch(() => {})
timers.set(
"certificates_renew",
setInterval(() => {
void executeSchedulerJob("certificates_renew").catch(() => {})
}, certRenewMs),
)
}
const backupSchedule = await getBackupScheduleSettings()
if (backupSchedule.enabled) {
const backupMs = 60_000
void executeSchedulerJob("backups").catch(() => {})
timers.set(
"backups",
setInterval(() => {
void executeSchedulerJob("backups").catch(() => {})
}, backupMs),
)
}
const geoip = await getGeoipSettings()
if (geoip.enabled) {
const geoipMs = Math.max(6 * 3600_000, geoip.updateIntervalSec * 1000)
void executeSchedulerJob("geoip_update").catch(() => {})
timers.set(
"geoip_update",
setInterval(() => {
void executeSchedulerJob("geoip_update").catch(() => {})
}, geoipMs),
)
}
const alertMs = 20_000
void executeSchedulerJob("alert_engine").catch(() => {})
timers.set(
"alert_engine",
setInterval(() => {
void executeSchedulerJob("alert_engine").catch(() => {})
}, alertMs),
)
}
export function stopScheduler(): void {
clearAllTimers()
}
export function isJobRunning(jobKey: string): boolean {
return isSchedulerJobRunning(jobKey)
}
async function lastRunForJob(jobKey: string) {
return (await db.select().from(schedulerRuns)
.where(eq(schedulerRuns.jobKey, jobKey))
.orderBy(desc(schedulerRuns.finishedAt))
.limit(1))[0]
}
export async function getSchedulerStatus() {
const traffic = await getTrafficSettings()
const uptime = await getUptimeSettings()
const apiPing = await getServersApiPingSettings()
const internetPath = await getInternetPathSettings()
const certRenew = await getCertificateRenewSettings()
const backupSchedule = await getBackupScheduleSettings()
const geoip = await getGeoipSettings()
const resOn = uptime.resourcesEnabled ?? uptime.enabled
const pingOn = uptime.pingEnabled ?? uptime.enabled
const spdOn = uptime.speedEnabled ?? uptime.enabled
const jobMeta: Record<SchedulerJobKey, { enabled: boolean; intervalSec: number }> = {
traffic: { enabled: traffic.enabled, intervalSec: traffic.intervalSec },
servers_rest_ping: { enabled: apiPing.enabled, intervalSec: apiPing.intervalSec },
uptime_resources: { enabled: resOn, intervalSec: uptime.intervalSec },
uptime_ping: { enabled: pingOn, intervalSec: uptime.probeIntervalSec },
uptime_speed: { enabled: spdOn, intervalSec: uptime.speedIntervalSec },
internet_path: { enabled: internetPath.enabled, intervalSec: internetPath.intervalSec },
gre_bgp: { enabled: true, intervalSec: 30 },
certificates_renew: { enabled: certRenew.enabled, intervalSec: certRenew.intervalSec },
backups: { enabled: backupSchedule.enabled, intervalSec: 60 },
geoip_update: { enabled: geoip.enabled, intervalSec: geoip.updateIntervalSec },
alert_engine: { enabled: true, intervalSec: 20 },
}
return {
jobs: await Promise.all(JOB_KEYS.map(async (jobKey) => {
const last = await lastRunForJob(jobKey)
const m = jobMeta[jobKey]
return {
jobKey,
enabled: m.enabled,
intervalSec: m.intervalSec,
running: isSchedulerJobRunning(jobKey),
lastFinishedAt: last?.finishedAt ?? null,
lastStatus: last?.status ?? null,
lastDurationMs: last?.durationMs ?? null,
lastError: last?.error ?? null,
}
})),
}
}
export async function listSchedulerRuns(opts: { jobKey?: string; limit: number; offset: number }) {
const limit = Math.min(200, Math.max(1, opts.limit))
const offset = Math.max(0, opts.offset)
if (opts.jobKey) {
return {
runs: await db.select().from(schedulerRuns)
.where(eq(schedulerRuns.jobKey, opts.jobKey))
.orderBy(desc(schedulerRuns.finishedAt))
.limit(limit)
.offset(offset),
}
}
return {
runs: await db.select().from(schedulerRuns)
.orderBy(desc(schedulerRuns.finishedAt))
.limit(limit)
.offset(offset),
}
}
wireAlertEngineRunner(() => executeSchedulerJob("alert_engine"))