feat(filters, recursive-routes): enhance configuration history and live data handling
- Introduced configuration history management in Filters and Recursive Routes pages, allowing users to view and restore previous configurations. - Updated state management to handle live data loading and error states more effectively, improving user experience during data fetching. - Added new components for displaying configuration history and integrated them into existing pages. - Enhanced API interactions to support fetching and applying configuration revisions, ensuring data consistency across the application. - Updated tests to cover new functionalities and ensure reliability. Co-authored-by: Cursor <[email protected]>
This commit is contained in:
@@ -163,6 +163,22 @@ if (!(await withPgOrSkip())) {
|
||||
await dbQuery(`DELETE FROM servers WHERE name = 'pg-wipe-idempotent'`)
|
||||
}
|
||||
|
||||
{
|
||||
const mig = await dbQuery<{ id: string }>(
|
||||
`SELECT id FROM schema_migrations WHERE id = '0006_config_revisions'`,
|
||||
)
|
||||
assert.equal(mig.rows.length, 1, "0006 применена")
|
||||
|
||||
const { rows } = await dbQuery<{ column_name: string; udt_name: string }>(`
|
||||
SELECT column_name, udt_name FROM information_schema.columns
|
||||
WHERE table_schema = 'public' AND table_name = 'config_revisions'
|
||||
`)
|
||||
const by = Object.fromEntries(rows.map((r) => [r.column_name, r.udt_name]))
|
||||
assert.equal(by.payload, "jsonb")
|
||||
assert.equal(by.fingerprint, "text")
|
||||
assert.equal(by.section, "text")
|
||||
}
|
||||
|
||||
{
|
||||
const marker = await dbQuery<{ sqlite_imported_at: string | null }>(
|
||||
`SELECT sqlite_imported_at FROM data_migration WHERE id = 1`,
|
||||
|
||||
@@ -93,6 +93,19 @@ export const filterRules = pgTable("filter_rules", {
|
||||
index("idx_filter_rules_server_sort").on(t.serverId, t.sortOrder),
|
||||
])
|
||||
|
||||
export const configRevisions = pgTable("config_revisions", {
|
||||
id: text("id").primaryKey(),
|
||||
serverId: intPkRef().references(() => servers.id, { onDelete: "cascade" }),
|
||||
section: text("section", { enum: ["filters", "recursive-routes"] }).notNull(),
|
||||
source: text("source", { enum: ["apply", "rollback", "observed", "copy"] }).notNull(),
|
||||
fingerprint: text("fingerprint").notNull(),
|
||||
payload: jsonb("payload").$type<unknown[]>().notNull().default(sql`'[]'::jsonb`),
|
||||
note: text("note"),
|
||||
createdAt: ts("created_at").notNull().defaultNow(),
|
||||
}, (t) => [
|
||||
index("idx_config_revisions_server_section_created").on(t.serverId, t.section, t.createdAt),
|
||||
])
|
||||
|
||||
export const recursiveRoutes = pgTable("recursive_routes", {
|
||||
id: idIdentity().primaryKey(),
|
||||
serverId: intPkRef().references(() => servers.id, { onDelete: "cascade" }),
|
||||
@@ -733,6 +746,7 @@ export type ServerInsert = typeof servers.$inferInsert
|
||||
export type Snapshot = typeof serverSnapshots.$inferSelect
|
||||
export type SnapshotInsert = typeof serverSnapshots.$inferInsert
|
||||
export type FilterRuleRow = typeof filterRules.$inferSelect
|
||||
export type ConfigRevisionRow = typeof configRevisions.$inferSelect
|
||||
export type RecursiveRouteRow = typeof recursiveRoutes.$inferSelect
|
||||
export type TrafficSettingsRow = typeof trafficSettings.$inferSelect
|
||||
export type TrafficFlowSettingsRow = typeof trafficFlowSettings.$inferSelect
|
||||
|
||||
@@ -19,5 +19,31 @@ export function managedComment(label: string): string {
|
||||
}
|
||||
|
||||
export function managedRecursiveComment(comment?: string | null): string {
|
||||
return comment ? `${PRODUCT_NAME}:recursive ${comment}` : `${PRODUCT_NAME}:recursive`
|
||||
const stripped = stripManagedRecursiveComment(comment ?? "")
|
||||
return stripped ? `${PRODUCT_NAME}:recursive ${stripped}` : `${PRODUCT_NAME}:recursive`
|
||||
}
|
||||
|
||||
const LEGACY_RECURSIVE_PREFIX = /^recursive:\s*/i
|
||||
|
||||
export function stripManagedRecursiveComment(comment: string): string {
|
||||
const value = comment.trim()
|
||||
if (!value) return ""
|
||||
const managedPrefixes = [
|
||||
`${PRODUCT_NAME}:recursive`,
|
||||
`${LEGACY_PRODUCT_NAME}:recursive`,
|
||||
]
|
||||
for (const prefix of managedPrefixes) {
|
||||
if (value.startsWith(prefix)) return value.slice(prefix.length).trim()
|
||||
}
|
||||
if (LEGACY_RECURSIVE_PREFIX.test(value)) {
|
||||
return value.replace(LEGACY_RECURSIVE_PREFIX, "").trim()
|
||||
}
|
||||
return value
|
||||
}
|
||||
|
||||
/** Owned recursive route: MM/legacy prefix or old `recursive:` mask. */
|
||||
export function isOwnedRecursiveComment(comment: string | undefined): boolean {
|
||||
if (!comment) return false
|
||||
const value = comment.trim()
|
||||
return hasManagedRecursiveComment(value) || LEGACY_RECURSIVE_PREFIX.test(value)
|
||||
}
|
||||
|
||||
+204
-259
@@ -1,14 +1,20 @@
|
||||
import { and, asc, eq, inArray } from "drizzle-orm"
|
||||
import { and, asc, eq } from "drizzle-orm"
|
||||
import { z } from "zod"
|
||||
import type { FastifyPluginAsyncZod } from "@fastify/type-provider-zod"
|
||||
import { db } from "../db/index.js"
|
||||
import { filterRules, recursiveRoutes, servers } from "../db/schema.js"
|
||||
import { MikrotikClient } from "../services/mikrotik.js"
|
||||
import { parseDbServerId } from "../utils/server-id.js"
|
||||
import { appendEvent } from "../modules/events/service/events-service.js"
|
||||
import { managedComment } from "../managed-markers.js"
|
||||
import { planBgpInApply } from "../services/config-apply-plan.js"
|
||||
import {
|
||||
hasManagedCommentPrefix,
|
||||
managedComment,
|
||||
} from "../managed-markers.js"
|
||||
appendRevisionIfChanged,
|
||||
canonicalFilterRules,
|
||||
getRevisionById,
|
||||
listRevisions,
|
||||
type ConfigRevisionSource,
|
||||
} from "../services/config-revisions.js"
|
||||
|
||||
type ServerRow = typeof servers.$inferSelect
|
||||
|
||||
@@ -296,47 +302,6 @@ async function resolveRouteTargets(serverId: number, rule: ApiFilterRule): Promi
|
||||
return { gateway: rule.gateway, outIface: tid }
|
||||
}
|
||||
|
||||
function normalizeCommunity(c: string): string {
|
||||
return (c ?? "").trim()
|
||||
}
|
||||
|
||||
/** Одинаковый эффект на роутере при одинаковой community (blackhole vs gateway + out-interface) */
|
||||
async function ruleEffectSignature(serverId: number, r: ApiFilterRule): Promise<string> {
|
||||
if (r.action === "blackhole") return `bh:${normalizeCommunity(r.community)}`
|
||||
const { gateway, outIface } = await resolveRouteTargets(serverId, r)
|
||||
return `rt:${normalizeCommunity(r.community)}:${gateway}:${outIface}`
|
||||
}
|
||||
|
||||
export type FilterRouterCompareStatus = "synced" | "drift" | "missing"
|
||||
|
||||
async function compareDbRulesWithRouter(
|
||||
serverId: number,
|
||||
dbRules: ApiFilterRule[],
|
||||
remoteRules: ApiFilterRule[],
|
||||
): Promise<Record<string, FilterRouterCompareStatus>> {
|
||||
const remoteSigByComm = new Map<string, string>()
|
||||
for (const rr of remoteRules) {
|
||||
const c = normalizeCommunity(rr.community)
|
||||
if (!remoteSigByComm.has(c)) {
|
||||
remoteSigByComm.set(c, await ruleEffectSignature(serverId, rr))
|
||||
}
|
||||
}
|
||||
const out: Record<string, FilterRouterCompareStatus> = {}
|
||||
for (const dr of dbRules) {
|
||||
const c = normalizeCommunity(dr.community)
|
||||
const sigD = await ruleEffectSignature(serverId, dr)
|
||||
const sigR = remoteSigByComm.get(c)
|
||||
if (sigR === undefined) {
|
||||
out[c] = "missing"
|
||||
} else if (sigR !== sigD) {
|
||||
out[c] = "drift"
|
||||
} else {
|
||||
out[c] = "synced"
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
async function toRouterRuleBody(serverId: number, rules: ApiFilterRule[]): Promise<string> {
|
||||
if (rules.length === 0) return ""
|
||||
// Группируем по эффекту (action + gateway + out-interface). Communities с одним и тем же
|
||||
@@ -421,44 +386,80 @@ async function replaceDbRules(serverId: number, rules: ApiFilterRule[]) {
|
||||
)
|
||||
}
|
||||
|
||||
const filtersRoutes: FastifyPluginAsyncZod = async (app) => {
|
||||
/** Сравнение правил в БД с живым bgp-in на MikroTik (один запрос API к роутеру) */
|
||||
app.get("/filters/router-compare", async (req, reply) => {
|
||||
const q = req.query as { serverId?: string | number }
|
||||
const serverId = parseDbServerId(q.serverId)
|
||||
if (serverId === null) {
|
||||
return reply.status(400).send({ error: "serverId is required" })
|
||||
}
|
||||
async function cacheRulesetsForServer(serverId: number): Promise<ApiFilterRule[]> {
|
||||
const rows = await db
|
||||
.select()
|
||||
.from(filterRules)
|
||||
.where(eq(filterRules.serverId, serverId))
|
||||
.orderBy(asc(filterRules.sortOrder))
|
||||
return rows.map((r) => ({
|
||||
id: String(r.id),
|
||||
community: r.community,
|
||||
communityName: r.communityName ?? undefined,
|
||||
action: r.action,
|
||||
gateway: r.gateway,
|
||||
gatewayTunnelId: r.gatewayTunnelId,
|
||||
description: r.description,
|
||||
}))
|
||||
}
|
||||
|
||||
const server = (await db.select().from(servers).where(eq(servers.id, serverId)).limit(1))[0]
|
||||
if (!server) return reply.status(404).send({ error: "Server not found" })
|
||||
async function applyFiltersToServer(
|
||||
server: ServerRow,
|
||||
rules: ApiFilterRule[],
|
||||
source: ConfigRevisionSource,
|
||||
): Promise<{ pushed: number; action: string; conflictsRemoved: number }> {
|
||||
const client = MikrotikClient.fromServer(server)
|
||||
const existing = await client.get<RosFilterRule[]>("/routing/filter/rule")
|
||||
const plan = planBgpInApply(existing, rules.length)
|
||||
const managedCommentValue = managedComment(server.name || server.host)
|
||||
|
||||
try {
|
||||
const remote = await fetchServerFilters(server)
|
||||
const rows = await db
|
||||
.select()
|
||||
.from(filterRules)
|
||||
.where(eq(filterRules.serverId, serverId))
|
||||
.orderBy(asc(filterRules.sortOrder))
|
||||
if (plan.action === "patch" && plan.managedId) {
|
||||
const ruleBody = await toRouterRuleBody(server.id, rules)
|
||||
await client.patch(
|
||||
`/routing/filter/rule/${encodeURIComponent(plan.managedId)}`,
|
||||
{
|
||||
chain: "bgp-in",
|
||||
comment: managedCommentValue,
|
||||
rule: ruleBody,
|
||||
disabled: "no",
|
||||
},
|
||||
)
|
||||
} else if (plan.action === "create") {
|
||||
const ruleBody = await toRouterRuleBody(server.id, rules)
|
||||
await client.post("/routing/filter/rule/add", {
|
||||
chain: "bgp-in",
|
||||
comment: managedCommentValue,
|
||||
rule: ruleBody,
|
||||
})
|
||||
} else if (plan.action === "delete" && plan.managedId) {
|
||||
await client.delete(
|
||||
`/routing/filter/rule/${encodeURIComponent(plan.managedId)}`,
|
||||
)
|
||||
}
|
||||
|
||||
const dbRules: ApiFilterRule[] = rows.map(r => ({
|
||||
id: String(r.id),
|
||||
community: r.community,
|
||||
communityName: r.communityName ?? undefined,
|
||||
action: r.action,
|
||||
gateway: r.gateway,
|
||||
gatewayTunnelId: r.gatewayTunnelId,
|
||||
description: r.description,
|
||||
}))
|
||||
for (const id of plan.conflictIds) {
|
||||
await client.delete(`/routing/filter/rule/${encodeURIComponent(id)}`)
|
||||
}
|
||||
|
||||
const byCommunity = await compareDbRulesWithRouter(serverId, dbRules, remote.rules)
|
||||
return reply.send({ byCommunity })
|
||||
} catch (err) {
|
||||
app.log.error({ serverId, err: String(err) }, "filters router-compare failed")
|
||||
return reply.status(500).send({ error: String(err) })
|
||||
}
|
||||
await replaceDbRules(server.id, rules)
|
||||
await appendRevisionIfChanged({
|
||||
serverId: server.id,
|
||||
section: "filters",
|
||||
source,
|
||||
payload: canonicalFilterRules(rules),
|
||||
})
|
||||
|
||||
return {
|
||||
pushed: rules.length,
|
||||
action: plan.action,
|
||||
conflictsRemoved: plan.conflictIds.length,
|
||||
}
|
||||
}
|
||||
|
||||
const RevisionIdParamSchema = z.object({ id: z.string().min(1) })
|
||||
|
||||
const filtersRoutes: FastifyPluginAsyncZod = async (app) => {
|
||||
|
||||
/** GRE с роутеров: один сервер (?serverId) или все включённые (без query) — для /gre, карты сети */
|
||||
app.get("/filters/gre-tunnels", async (req, reply) => {
|
||||
const q = req.query as { serverId?: string | number }
|
||||
@@ -487,74 +488,103 @@ const filtersRoutes: FastifyPluginAsyncZod = async (app) => {
|
||||
return reply.send({ tunnels: results.flat() })
|
||||
})
|
||||
|
||||
/** Только правила фильтров из БД (без опроса MikroTik за GRE) */
|
||||
app.get("/filters/rules", async (_req, reply) => {
|
||||
/** Без serverId — cache для дашборда. С serverId — live с CHR, cache fallback. */
|
||||
app.get("/filters/rules", async (req, reply) => {
|
||||
const q = req.query as { serverId?: string | number }
|
||||
const serverId = parseDbServerId(q.serverId)
|
||||
const allServers = await db.select().from(servers).where(eq(servers.enabled, true))
|
||||
const dbRulesets = await toApiRulesets(allServers)
|
||||
|
||||
return reply.send({
|
||||
rulesets: dbRulesets,
|
||||
greTunnels: [] as LiveGreTunnel[],
|
||||
})
|
||||
})
|
||||
|
||||
app.put("/filters/rules", async (req, reply) => {
|
||||
const body = req.body as { rulesets?: Array<{ serverId: string; rules: ApiFilterRule[] }> }
|
||||
const payload = body.rulesets ?? []
|
||||
const serverIds = payload.map(r => Number.parseInt(r.serverId, 10)).filter(Number.isFinite)
|
||||
if (serverIds.length > 0) {
|
||||
await db.delete(filterRules).where(inArray(filterRules.serverId, serverIds))
|
||||
}
|
||||
for (const rs of payload) {
|
||||
const sid = Number.parseInt(rs.serverId, 10)
|
||||
if (!Number.isFinite(sid)) continue
|
||||
await replaceDbRules(sid, rs.rules ?? [])
|
||||
}
|
||||
return reply.send({ ok: true })
|
||||
})
|
||||
|
||||
app.post("/filters/sync/from-router", async (_req, reply) => {
|
||||
const body = _req.body as { serverId?: string | number } | undefined
|
||||
const rawServerId = body?.serverId
|
||||
const serverId = parseDbServerId(rawServerId)
|
||||
if (serverId === null) {
|
||||
return reply.status(400).send({ error: "serverId is required" })
|
||||
const dbRulesets = await toApiRulesets(allServers)
|
||||
return reply.send({
|
||||
rulesets: dbRulesets,
|
||||
greTunnels: [] as LiveGreTunnel[],
|
||||
live: false,
|
||||
stale: false,
|
||||
})
|
||||
}
|
||||
|
||||
const server = (await db.select().from(servers).where(eq(servers.id, serverId)).limit(1))[0]
|
||||
const server = allServers.find((s) => s.id === serverId)
|
||||
?? (await db.select().from(servers).where(eq(servers.id, serverId)).limit(1))[0]
|
||||
if (!server) return reply.status(404).send({ error: "Server not found" })
|
||||
|
||||
try {
|
||||
app.log.info({ serverId, host: server.host }, "Filters sync from router started")
|
||||
await appendEvent({
|
||||
level: "info",
|
||||
eventType: "filters.sync.from_router.started",
|
||||
sourceModule: "filters",
|
||||
title: "Синхронизация фильтров запущена",
|
||||
message: `${server.name || server.host} → БД`,
|
||||
entityType: "server",
|
||||
entityId: String(serverId),
|
||||
})
|
||||
const remote = await fetchServerFilters(server)
|
||||
await replaceDbRules(server.id, remote.rules)
|
||||
app.log.info({ serverId, totalRules: remote.rules.length }, "Filters sync from router completed")
|
||||
await appendRevisionIfChanged({
|
||||
serverId: server.id,
|
||||
section: "filters",
|
||||
source: "observed",
|
||||
payload: canonicalFilterRules(remote.rules),
|
||||
})
|
||||
const cached = await cacheRulesetsForServer(server.id)
|
||||
return reply.send({
|
||||
rulesets: [{ serverId: String(server.id), rules: cached }],
|
||||
greTunnels: remote.tunnels,
|
||||
live: true,
|
||||
stale: false,
|
||||
})
|
||||
} catch (err) {
|
||||
app.log.warn({ serverId, err: String(err) }, "filters live GET failed, serving cache")
|
||||
const cached = await cacheRulesetsForServer(server.id)
|
||||
return reply.send({
|
||||
rulesets: [{ serverId: String(server.id), rules: cached }],
|
||||
greTunnels: [] as LiveGreTunnel[],
|
||||
live: false,
|
||||
stale: true,
|
||||
error: String(err),
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
app.put("/filters/rules", async (req, reply) => {
|
||||
const body = req.body as {
|
||||
serverId?: string | number
|
||||
rules?: ApiFilterRule[]
|
||||
source?: ConfigRevisionSource
|
||||
}
|
||||
const serverId = parseDbServerId(body.serverId)
|
||||
if (serverId === null) return reply.status(400).send({ error: "serverId is required" })
|
||||
const server = (await db.select().from(servers).where(eq(servers.id, serverId)).limit(1))[0]
|
||||
if (!server) return reply.status(404).send({ error: "Server not found" })
|
||||
|
||||
const rules = body.rules ?? []
|
||||
const source: ConfigRevisionSource = body.source === "copy" ? "copy" : "apply"
|
||||
try {
|
||||
await appendEvent({
|
||||
level: "info",
|
||||
eventType: "filters.sync.from_router.done",
|
||||
eventType: "filters.apply.started",
|
||||
sourceModule: "filters",
|
||||
title: "Синхронизация фильтров завершена",
|
||||
message: `${server.name || server.host}: ${remote.rules.length} правил`,
|
||||
title: "Применение фильтров на роутер",
|
||||
message: `${server.name || server.host}: ${rules.length} правил`,
|
||||
entityType: "server",
|
||||
entityId: String(serverId),
|
||||
})
|
||||
return reply.send({ ok: true, updatedServers: 1, totalRules: remote.rules.length, serverId })
|
||||
const result = await applyFiltersToServer(server, rules, source)
|
||||
const cached = await cacheRulesetsForServer(server.id)
|
||||
await appendEvent({
|
||||
level: "info",
|
||||
eventType: "filters.apply.done",
|
||||
sourceModule: "filters",
|
||||
title: "Фильтры применены",
|
||||
message: `${server.name || server.host}: ${result.pushed} правил (${result.action})`,
|
||||
entityType: "server",
|
||||
entityId: String(serverId),
|
||||
})
|
||||
return reply.send({
|
||||
ok: true,
|
||||
serverId,
|
||||
pushedRules: result.pushed,
|
||||
action: result.action,
|
||||
rules: cached,
|
||||
})
|
||||
} catch (err) {
|
||||
app.log.error({ serverId, err: String(err) }, "Filters sync from router failed")
|
||||
app.log.error({ serverId, err: String(err) }, "filters apply failed")
|
||||
await appendEvent({
|
||||
level: "critical",
|
||||
eventType: "filters.sync.from_router.failed",
|
||||
eventType: "filters.apply.failed",
|
||||
sourceModule: "filters",
|
||||
title: "Ошибка синхронизации фильтров",
|
||||
title: "Ошибка применения фильтров",
|
||||
message: `${server.name || server.host}: ${String(err)}`,
|
||||
entityType: "server",
|
||||
entityId: String(serverId),
|
||||
@@ -563,145 +593,60 @@ const filtersRoutes: FastifyPluginAsyncZod = async (app) => {
|
||||
}
|
||||
})
|
||||
|
||||
app.post("/filters/sync/to-router", async (req, reply) => {
|
||||
app.get("/filters/revisions", async (req, reply) => {
|
||||
const q = req.query as { serverId?: string | number }
|
||||
const serverId = parseDbServerId(q.serverId)
|
||||
if (serverId === null) return reply.status(400).send({ error: "serverId is required" })
|
||||
const revisions = await listRevisions(serverId, "filters")
|
||||
return reply.send({ revisions })
|
||||
})
|
||||
|
||||
app.post("/filters/revisions/:id/restore", {
|
||||
schema: { params: RevisionIdParamSchema },
|
||||
}, async (req, reply) => {
|
||||
const { id } = req.params
|
||||
const body = req.body as { serverId?: string | number } | undefined
|
||||
const requestedServerId = parseDbServerId(body?.serverId)
|
||||
|
||||
const allServers = await db.select().from(servers).where(eq(servers.enabled, true))
|
||||
const targetServers = requestedServerId !== null
|
||||
? allServers.filter(s => s.id === requestedServerId)
|
||||
: allServers
|
||||
|
||||
if (requestedServerId !== null && targetServers.length === 0) {
|
||||
return reply.status(404).send({ error: "Server not found" })
|
||||
const rev = await getRevisionById(id)
|
||||
if (!rev) return reply.status(404).send({ error: "Revision not found" })
|
||||
if (rev.section !== "filters") return reply.status(400).send({ error: "Revision section mismatch" })
|
||||
const requested = parseDbServerId(body?.serverId)
|
||||
if (requested !== null && requested !== rev.serverId) {
|
||||
return reply.status(400).send({ error: "Revision belongs to another server" })
|
||||
}
|
||||
const server = (await db.select().from(servers).where(eq(servers.id, rev.serverId)).limit(1))[0]
|
||||
if (!server) return reply.status(404).send({ error: "Server not found" })
|
||||
|
||||
let updatedServers = 0
|
||||
let pushedRules = 0
|
||||
const errors: Array<{ serverId: number; error: string }> = []
|
||||
await appendEvent({
|
||||
level: "info",
|
||||
eventType: "filters.sync.to_router.started",
|
||||
sourceModule: "filters",
|
||||
title: "Отправка фильтров на роутеры запущена",
|
||||
message: `Целевых серверов: ${targetServers.length}`,
|
||||
payload: { requestedServerId },
|
||||
})
|
||||
|
||||
for (const server of targetServers) {
|
||||
try {
|
||||
app.log.info({ serverId: server.id, host: server.host }, "Filters sync to router started")
|
||||
const client = MikrotikClient.fromServer(server)
|
||||
const existing = await client.get<RosFilterRule[]>("/routing/filter/rule")
|
||||
|
||||
const isInBgpIn = (r: RosFilterRule) =>
|
||||
(r.chain ?? "").trim().toLowerCase() === "bgp-in"
|
||||
const managedCommentValue = managedComment(server.name || server.host)
|
||||
|
||||
// Уже созданное нами правило — будем PATCH'ить, чтобы сохранить ID/позицию в цепочке.
|
||||
const managedRule = existing.find(
|
||||
r => isInBgpIn(r) && hasManagedCommentPrefix(r.comment ?? ""),
|
||||
)
|
||||
|
||||
// Конфликтующие легаси-правила в bgp-in (без нашего comment, но с bgp-communities) —
|
||||
// удаляем после успешного upsert: иначе старое правило с `else { reject; }`
|
||||
// отрабатывает первым и перебивает наш upsert.
|
||||
const conflictIds = existing
|
||||
.filter(r =>
|
||||
isInBgpIn(r) &&
|
||||
!hasManagedCommentPrefix(r.comment ?? "") &&
|
||||
/bgp-communities/i.test(r.rule ?? ""),
|
||||
)
|
||||
.map(r => r[".id"])
|
||||
.filter((id): id is string => Boolean(id))
|
||||
|
||||
const rows = await db.select().from(filterRules)
|
||||
.where(and(eq(filterRules.serverId, server.id)))
|
||||
.orderBy(asc(filterRules.sortOrder))
|
||||
|
||||
const rules: ApiFilterRule[] = rows.map(r => ({
|
||||
id: String(r.id),
|
||||
community: r.community,
|
||||
communityName: r.communityName ?? undefined,
|
||||
action: r.action,
|
||||
gateway: r.gateway,
|
||||
gatewayTunnelId: r.gatewayTunnelId,
|
||||
description: r.description,
|
||||
}))
|
||||
|
||||
// Upsert: PATCH существующего managed-правила или POST /add нового.
|
||||
// Если ошибка — конфликтные правила НЕ удаляем (роутер не остаётся с пустым bgp-in).
|
||||
// Путь `/routing/filter/rule/add` обязателен: голый POST на коллекцию RouterOS REST
|
||||
// трактует как «вызов команды» и отдаёт 400 «no such command».
|
||||
// См. https://help.mikrotik.com/docs/spaces/ROS/pages/47579162/REST+API
|
||||
if (rules.length > 0) {
|
||||
const ruleBody = await toRouterRuleBody(server.id, rules)
|
||||
if (managedRule && managedRule[".id"]) {
|
||||
await client.patch(
|
||||
`/routing/filter/rule/${encodeURIComponent(managedRule[".id"])}`,
|
||||
{
|
||||
chain: "bgp-in",
|
||||
comment: managedCommentValue,
|
||||
rule: ruleBody,
|
||||
disabled: "no",
|
||||
},
|
||||
)
|
||||
app.log.info({ serverId: server.id, id: managedRule[".id"] }, "bgp-in rule updated")
|
||||
} else {
|
||||
await client.post("/routing/filter/rule/add", {
|
||||
chain: "bgp-in",
|
||||
comment: managedCommentValue,
|
||||
rule: ruleBody,
|
||||
})
|
||||
app.log.info({ serverId: server.id }, "bgp-in rule created")
|
||||
}
|
||||
pushedRules += rules.length
|
||||
} else if (managedRule && managedRule[".id"]) {
|
||||
// В БД нет правил → удаляем наш managed-rule на роутере.
|
||||
await client.delete(
|
||||
`/routing/filter/rule/${encodeURIComponent(managedRule[".id"])}`,
|
||||
)
|
||||
app.log.info({ serverId: server.id }, "bgp-in rule removed (no rules in DB)")
|
||||
}
|
||||
|
||||
for (const id of conflictIds) {
|
||||
await client.delete(`/routing/filter/rule/${encodeURIComponent(id)}`)
|
||||
}
|
||||
|
||||
updatedServers += 1
|
||||
app.log.info(
|
||||
{
|
||||
serverId: server.id,
|
||||
mode: managedRule ? "patch" : "create",
|
||||
conflictsRemoved: conflictIds.length,
|
||||
pushed: rules.length,
|
||||
},
|
||||
"Filters sync to router completed",
|
||||
)
|
||||
} catch (err) {
|
||||
app.log.error({ serverId: server.id, err: String(err) }, "filters sync to-router failed")
|
||||
errors.push({ serverId: server.id, error: String(err) })
|
||||
const raw = Array.isArray(rev.payload) ? rev.payload : []
|
||||
const rules: ApiFilterRule[] = raw.map((item, idx) => {
|
||||
const r = item as Partial<ApiFilterRule>
|
||||
return {
|
||||
id: `rev-${idx}`,
|
||||
community: r.community ?? "",
|
||||
communityName: r.communityName,
|
||||
action: r.action === "blackhole" ? "blackhole" : "route",
|
||||
gateway: r.gateway ?? "",
|
||||
gatewayTunnelId: r.gatewayTunnelId ?? "",
|
||||
description: r.description ?? "",
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
await appendEvent({
|
||||
level: errors.length === 0 ? "info" : "warning",
|
||||
eventType: errors.length === 0 ? "filters.sync.to_router.done" : "filters.sync.to_router.partial",
|
||||
sourceModule: "filters",
|
||||
title: errors.length === 0 ? "Отправка фильтров завершена" : "Отправка фильтров завершена с ошибками",
|
||||
message: `Успешно: ${updatedServers}, ошибок: ${errors.length}, правил: ${pushedRules}`,
|
||||
payload: {
|
||||
updatedServers,
|
||||
pushedRules,
|
||||
errors,
|
||||
},
|
||||
})
|
||||
return reply.send({
|
||||
ok: errors.length === 0,
|
||||
updatedServers,
|
||||
pushedRules,
|
||||
errors,
|
||||
})
|
||||
try {
|
||||
const result = await applyFiltersToServer(server, rules, "rollback")
|
||||
const cached = await cacheRulesetsForServer(server.id)
|
||||
await appendEvent({
|
||||
level: "info",
|
||||
eventType: "filters.rollback.done",
|
||||
sourceModule: "filters",
|
||||
title: "Откат фильтров",
|
||||
message: `${server.name || server.host}: ${result.pushed} правил`,
|
||||
entityType: "server",
|
||||
entityId: String(server.id),
|
||||
})
|
||||
return reply.send({ ok: true, rules: cached, pushedRules: result.pushed })
|
||||
} catch (err) {
|
||||
app.log.error({ serverId: server.id, err: String(err) }, "filters restore failed")
|
||||
return reply.status(500).send({ error: String(err) })
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -1,13 +1,23 @@
|
||||
import { asc, eq } from "drizzle-orm"
|
||||
import { z } from "zod"
|
||||
import type { FastifyPluginAsyncZod } from "@fastify/type-provider-zod"
|
||||
import { db } from "../db/index.js"
|
||||
import { recursiveRoutes, servers } from "../db/schema.js"
|
||||
import { MikrotikClient } from "../services/mikrotik.js"
|
||||
import { parseDbServerId } from "../utils/server-id.js"
|
||||
import { managedRecursiveComment } from "../managed-markers.js"
|
||||
import {
|
||||
hasManagedRecursiveComment,
|
||||
managedRecursiveComment,
|
||||
} from "../managed-markers.js"
|
||||
mapRosManagedRoutes,
|
||||
planRecursiveApply,
|
||||
userRecursiveComment,
|
||||
} from "../services/config-apply-plan.js"
|
||||
import {
|
||||
appendRevisionIfChanged,
|
||||
canonicalRecursiveRoutes,
|
||||
getRevisionById,
|
||||
listRevisions,
|
||||
type ConfigRevisionSource,
|
||||
} from "../services/config-revisions.js"
|
||||
|
||||
type ServerRow = typeof servers.$inferSelect
|
||||
|
||||
@@ -51,27 +61,6 @@ interface RecursiveRouteDto {
|
||||
disabled: boolean
|
||||
}
|
||||
|
||||
function isIpGateway(gw: string): boolean {
|
||||
return /^\d{1,3}(\.\d{1,3}){3}(?:%\S+)?$/.test(gw.trim())
|
||||
}
|
||||
|
||||
function isRecursiveRoute(r: RosRoute): boolean {
|
||||
if ((r.static ?? "false") !== "true") return false
|
||||
if ((r.dynamic ?? "false") === "true") return false
|
||||
if ((r.blackhole ?? "false") === "true") return false
|
||||
if ((r.unreachable ?? "false") === "true") return false
|
||||
if ((r.prohibit ?? "false") === "true") return false
|
||||
const dst = r["dst-address"] ?? ""
|
||||
const gw = r.gateway ?? ""
|
||||
if (!dst || !gw) return false
|
||||
return isIpGateway(gw)
|
||||
}
|
||||
|
||||
function hasRecursiveCommentMask(comment: string | undefined): boolean {
|
||||
if (!comment) return false
|
||||
return /^recursive:\s*/i.test(comment.trim())
|
||||
}
|
||||
|
||||
function splitGateway(raw: string): { ip: string; name: string } | null {
|
||||
const v = raw.trim()
|
||||
if (!v) return null
|
||||
@@ -102,6 +91,21 @@ async function mapDbRoutes(serverId: number): Promise<RecursiveRouteDto[]> {
|
||||
}))
|
||||
}
|
||||
|
||||
function mergeCachedCountry(
|
||||
live: RecursiveRouteDto[],
|
||||
cached: RecursiveRouteDto[],
|
||||
): RecursiveRouteDto[] {
|
||||
return live.map((row) => {
|
||||
if (row.country) return row
|
||||
const match = cached.find((c) =>
|
||||
c.dstAddress === row.dstAddress &&
|
||||
c.gateway === row.gateway &&
|
||||
c.distance === row.distance,
|
||||
)
|
||||
return match?.country ? { ...row, country: match.country } : row
|
||||
})
|
||||
}
|
||||
|
||||
function toRouterPayload(route: RecursiveRouteDto): Record<string, string> {
|
||||
return {
|
||||
"dst-address": route.dstAddress,
|
||||
@@ -132,7 +136,7 @@ async function replaceDbRoutes(serverId: number, routes: RecursiveRouteDto[]) {
|
||||
routingTable: r.routingTable || "main",
|
||||
checkGateway: r.checkGateway ?? "",
|
||||
country: r.country ?? "",
|
||||
comment: r.comment ?? "",
|
||||
comment: userRecursiveComment(r.comment),
|
||||
disabled: r.disabled,
|
||||
createdAt: now,
|
||||
updatedAt: now,
|
||||
@@ -140,6 +144,34 @@ async function replaceDbRoutes(serverId: number, routes: RecursiveRouteDto[]) {
|
||||
)
|
||||
}
|
||||
|
||||
async function applyRecursiveToServer(
|
||||
server: ServerRow,
|
||||
routes: RecursiveRouteDto[],
|
||||
source: ConfigRevisionSource,
|
||||
): Promise<{ pushed: number; deleted: number }> {
|
||||
const client = MikrotikClient.fromServer(server)
|
||||
const existing = await client.get<RosRoute[]>("/ip/route")
|
||||
const { deleteIds } = planRecursiveApply(existing)
|
||||
|
||||
for (const route of routes) {
|
||||
await client.post("/ip/route", toRouterPayload(route))
|
||||
}
|
||||
for (const id of deleteIds) {
|
||||
await client.delete(`/ip/route/${encodeURIComponent(id)}`)
|
||||
}
|
||||
|
||||
await replaceDbRoutes(server.id, routes)
|
||||
await appendRevisionIfChanged({
|
||||
serverId: server.id,
|
||||
section: "recursive-routes",
|
||||
source,
|
||||
payload: canonicalRecursiveRoutes(routes),
|
||||
})
|
||||
return { pushed: routes.length, deleted: deleteIds.length }
|
||||
}
|
||||
|
||||
const RevisionIdParamSchema = z.object({ id: z.string().min(1) })
|
||||
|
||||
const recursiveRoutesPlugin: FastifyPluginAsyncZod = async (app) => {
|
||||
app.get("/recursive-routes/gateways", async (req, reply) => {
|
||||
const q = req.query as { serverId?: string | number }
|
||||
@@ -175,79 +207,108 @@ const recursiveRoutesPlugin: FastifyPluginAsyncZod = async (app) => {
|
||||
const q = req.query as { serverId?: string | number }
|
||||
const serverId = parseDbServerId(q.serverId)
|
||||
if (serverId === null) return reply.status(400).send({ error: "serverId is required" })
|
||||
const server = (await db.select().from(servers).where(eq(servers.id, serverId)).limit(1))[0]
|
||||
if (!server) return reply.status(404).send({ error: "Server not found" })
|
||||
return reply.send({ routes: await mapDbRoutes(serverId) })
|
||||
})
|
||||
|
||||
app.put("/recursive-routes", async (req, reply) => {
|
||||
const body = req.body as { serverId?: string | number; routes?: RecursiveRouteDto[] }
|
||||
const serverId = parseDbServerId(body.serverId)
|
||||
if (serverId === null) return reply.status(400).send({ error: "serverId is required" })
|
||||
const server = (await db.select().from(servers).where(eq(servers.id, serverId)).limit(1))[0]
|
||||
if (!server) return reply.status(404).send({ error: "Server not found" })
|
||||
await replaceDbRoutes(serverId, body.routes ?? [])
|
||||
return reply.send({ ok: true })
|
||||
})
|
||||
|
||||
app.post("/recursive-routes/sync/from-router", async (req, reply) => {
|
||||
const body = req.body as { serverId?: string | number } | undefined
|
||||
const serverId = parseDbServerId(body?.serverId)
|
||||
if (serverId === null) return reply.status(400).send({ error: "serverId is required" })
|
||||
const server: ServerRow | undefined = (await db
|
||||
.select().from(servers)
|
||||
.where(eq(servers.id, serverId))
|
||||
.limit(1))[0]
|
||||
if (!server) return reply.status(404).send({ error: "Server not found" })
|
||||
|
||||
const cached = await mapDbRoutes(serverId)
|
||||
try {
|
||||
const client = MikrotikClient.fromServer(server)
|
||||
const rosRoutes = await client.get<RosRoute[]>("/ip/route")
|
||||
const rec = rosRoutes.filter(r =>
|
||||
isRecursiveRoute(r) && hasRecursiveCommentMask(r.comment),
|
||||
)
|
||||
const mapped: RecursiveRouteDto[] = rec.map((r, i) => ({
|
||||
id: r[".id"] ?? `ros-${i}`,
|
||||
dstAddress: r["dst-address"] ?? "",
|
||||
gateway: r.gateway ?? "",
|
||||
distance: Number.parseInt(r.distance ?? "1", 10) || 1,
|
||||
scope: r.scope ? (Number.parseInt(r.scope, 10) || null) : null,
|
||||
targetScope: r["target-scope"] ? (Number.parseInt(r["target-scope"], 10) || null) : null,
|
||||
routingTable: r["routing-table"] ?? "main",
|
||||
checkGateway: r["check-gateway"] ?? "",
|
||||
country: "",
|
||||
comment: r.comment ?? "",
|
||||
disabled: r.disabled === "true",
|
||||
}))
|
||||
await replaceDbRoutes(serverId, mapped)
|
||||
return reply.send({ ok: true, serverId, totalRoutes: mapped.length })
|
||||
const live = mergeCachedCountry(mapRosManagedRoutes(rosRoutes), cached)
|
||||
await replaceDbRoutes(serverId, live)
|
||||
await appendRevisionIfChanged({
|
||||
serverId,
|
||||
section: "recursive-routes",
|
||||
source: "observed",
|
||||
payload: canonicalRecursiveRoutes(live),
|
||||
})
|
||||
const stored = await mapDbRoutes(serverId)
|
||||
return reply.send({ routes: stored, live: true, stale: false })
|
||||
} catch (err) {
|
||||
app.log.warn({ serverId, err: String(err) }, "recursive live GET failed, serving cache")
|
||||
return reply.send({
|
||||
routes: cached,
|
||||
live: false,
|
||||
stale: true,
|
||||
error: String(err),
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
app.put("/recursive-routes", async (req, reply) => {
|
||||
const body = req.body as {
|
||||
serverId?: string | number
|
||||
routes?: RecursiveRouteDto[]
|
||||
source?: ConfigRevisionSource
|
||||
}
|
||||
const serverId = parseDbServerId(body.serverId)
|
||||
if (serverId === null) return reply.status(400).send({ error: "serverId is required" })
|
||||
const server = (await db.select().from(servers).where(eq(servers.id, serverId)).limit(1))[0]
|
||||
if (!server) return reply.status(404).send({ error: "Server not found" })
|
||||
const routes = (body.routes ?? []).map((r) => ({
|
||||
...r,
|
||||
comment: userRecursiveComment(r.comment),
|
||||
}))
|
||||
const source: ConfigRevisionSource = body.source === "copy" ? "copy" : "apply"
|
||||
try {
|
||||
const result = await applyRecursiveToServer(server, routes, source)
|
||||
const stored = await mapDbRoutes(serverId)
|
||||
return reply.send({ ok: true, routes: stored, pushedRoutes: result.pushed })
|
||||
} catch (err) {
|
||||
return reply.status(500).send({ error: String(err) })
|
||||
}
|
||||
})
|
||||
|
||||
app.post("/recursive-routes/sync/to-router", async (req, reply) => {
|
||||
const body = req.body as { serverId?: string | number } | undefined
|
||||
const serverId = parseDbServerId(body?.serverId)
|
||||
app.get("/recursive-routes/revisions", async (req, reply) => {
|
||||
const q = req.query as { serverId?: string | number }
|
||||
const serverId = parseDbServerId(q.serverId)
|
||||
if (serverId === null) return reply.status(400).send({ error: "serverId is required" })
|
||||
const server = (await db.select().from(servers).where(eq(servers.id, serverId)).limit(1))[0]
|
||||
const revisions = await listRevisions(serverId, "recursive-routes")
|
||||
return reply.send({ revisions })
|
||||
})
|
||||
|
||||
app.post("/recursive-routes/revisions/:id/restore", {
|
||||
schema: { params: RevisionIdParamSchema },
|
||||
}, async (req, reply) => {
|
||||
const { id } = req.params
|
||||
const body = req.body as { serverId?: string | number } | undefined
|
||||
const rev = await getRevisionById(id)
|
||||
if (!rev) return reply.status(404).send({ error: "Revision not found" })
|
||||
if (rev.section !== "recursive-routes") {
|
||||
return reply.status(400).send({ error: "Revision section mismatch" })
|
||||
}
|
||||
const requested = parseDbServerId(body?.serverId)
|
||||
if (requested !== null && requested !== rev.serverId) {
|
||||
return reply.status(400).send({ error: "Revision belongs to another server" })
|
||||
}
|
||||
const server = (await db.select().from(servers).where(eq(servers.id, rev.serverId)).limit(1))[0]
|
||||
if (!server) return reply.status(404).send({ error: "Server not found" })
|
||||
|
||||
const raw = Array.isArray(rev.payload) ? rev.payload : []
|
||||
const routes: RecursiveRouteDto[] = raw.map((item, idx) => {
|
||||
const r = item as Partial<RecursiveRouteDto>
|
||||
return {
|
||||
id: `rev-${idx}`,
|
||||
dstAddress: r.dstAddress ?? "",
|
||||
gateway: r.gateway ?? "",
|
||||
distance: r.distance ?? 1,
|
||||
scope: r.scope ?? null,
|
||||
targetScope: r.targetScope ?? null,
|
||||
routingTable: r.routingTable || "main",
|
||||
checkGateway: r.checkGateway ?? "",
|
||||
country: r.country ?? "",
|
||||
comment: userRecursiveComment(r.comment),
|
||||
disabled: Boolean(r.disabled),
|
||||
}
|
||||
})
|
||||
|
||||
try {
|
||||
const client = MikrotikClient.fromServer(server)
|
||||
const existing = await client.get<RosRoute[]>("/ip/route")
|
||||
const managed = existing.filter(r => hasManagedRecursiveComment(r.comment ?? ""))
|
||||
for (const r of managed) {
|
||||
if (!r[".id"]) continue
|
||||
await client.delete(`/ip/route/${encodeURIComponent(r[".id"])}`)
|
||||
}
|
||||
|
||||
const dbRows = await mapDbRoutes(serverId)
|
||||
for (const route of dbRows) {
|
||||
await client.post("/ip/route", toRouterPayload(route))
|
||||
}
|
||||
|
||||
return reply.send({ ok: true, serverId, pushedRoutes: dbRows.length })
|
||||
const result = await applyRecursiveToServer(server, routes, "rollback")
|
||||
const stored = await mapDbRoutes(server.id)
|
||||
return reply.send({ ok: true, routes: stored, pushedRoutes: result.pushed })
|
||||
} catch (err) {
|
||||
return reply.status(500).send({ error: String(err) })
|
||||
}
|
||||
|
||||
@@ -0,0 +1,96 @@
|
||||
import assert from "node:assert/strict"
|
||||
import { hasManagedCommentPrefix, isOwnedRecursiveComment, managedRecursiveComment, stripManagedRecursiveComment } from "../managed-markers.js"
|
||||
import {
|
||||
planBgpInApply,
|
||||
planRecursiveApply,
|
||||
unmanagedRouteIds,
|
||||
} from "./config-apply-plan.js"
|
||||
import {
|
||||
canonicalFilterRules,
|
||||
fingerprintPayload,
|
||||
} from "./config-revisions.js"
|
||||
import { mapRosManagedRoutes } from "./config-apply-plan.js"
|
||||
|
||||
{
|
||||
const fp1 = fingerprintPayload(canonicalFilterRules([
|
||||
{ community: "65001:100", action: "route", gateway: "10.0.0.1", gatewayTunnelId: "gre1", description: "a" },
|
||||
]))
|
||||
const fp2 = fingerprintPayload(canonicalFilterRules([
|
||||
{ community: "65001:100", action: "route", gateway: "10.0.0.1", gatewayTunnelId: "gre1", description: "a" },
|
||||
]))
|
||||
const fp3 = fingerprintPayload(canonicalFilterRules([
|
||||
{ community: "65001:100", action: "route", gateway: "10.0.0.2", gatewayTunnelId: "gre1", description: "a" },
|
||||
]))
|
||||
assert.equal(fp1, fp2)
|
||||
assert.notEqual(fp1, fp3)
|
||||
}
|
||||
|
||||
{
|
||||
const existing = [
|
||||
{ ".id": "*1", chain: "bgp-in", comment: "MikrotikManager: msk", rule: "if (true) { accept; }" },
|
||||
{ ".id": "*2", chain: "bgp-in", comment: "legacy", rule: "if (bgp-communities includes 1:1) { reject; }" },
|
||||
{ ".id": "*3", chain: "bgp-out", comment: "MikrotikManager: other", rule: "if (bgp-communities includes 1:1) { accept; }" },
|
||||
]
|
||||
const patch = planBgpInApply(existing, 3)
|
||||
assert.equal(patch.action, "patch")
|
||||
assert.equal(patch.managedId, "*1")
|
||||
assert.deepEqual(patch.conflictIds, ["*2"])
|
||||
|
||||
const create = planBgpInApply(existing.filter((r) => r[".id"] !== "*1"), 1)
|
||||
assert.equal(create.action, "create")
|
||||
assert.equal(create.managedId, undefined)
|
||||
|
||||
const del = planBgpInApply(existing, 0)
|
||||
assert.equal(del.action, "delete")
|
||||
assert.equal(del.managedId, "*1")
|
||||
|
||||
const noop = planBgpInApply([], 0)
|
||||
assert.equal(noop.action, "noop")
|
||||
}
|
||||
|
||||
{
|
||||
const routes = [
|
||||
{ ".id": "*10", comment: "MikrotikManager:recursive via de", static: "true", "dst-address": "8.8.8.8/32", gateway: "1.1.1.1" },
|
||||
{ ".id": "*11", comment: "user static", static: "true", "dst-address": "1.1.1.1/32", gateway: "9.9.9.9" },
|
||||
{ ".id": "*12", comment: "recursive: old", static: "true", "dst-address": "9.9.9.9/32", gateway: "1.1.1.1" },
|
||||
]
|
||||
const plan = planRecursiveApply(routes)
|
||||
assert.deepEqual(plan.deleteIds, ["*10", "*12"])
|
||||
assert.deepEqual(unmanagedRouteIds(routes), ["*11"])
|
||||
}
|
||||
|
||||
{
|
||||
assert.equal(stripManagedRecursiveComment("MikrotikManager:recursive via de"), "via de")
|
||||
assert.equal(stripManagedRecursiveComment("recursive: old"), "old")
|
||||
assert.equal(managedRecursiveComment("MikrotikManager:recursive via de"), "MikrotikManager:recursive via de")
|
||||
assert.equal(isOwnedRecursiveComment("MikrotikManager:recursive via de"), true)
|
||||
assert.equal(isOwnedRecursiveComment("recursive: x"), true)
|
||||
assert.equal(isOwnedRecursiveComment("user static"), false)
|
||||
assert.equal(hasManagedCommentPrefix("MikrotikManager: msk"), true)
|
||||
}
|
||||
|
||||
{
|
||||
const mapped = mapRosManagedRoutes([
|
||||
{
|
||||
".id": "*1",
|
||||
static: "true",
|
||||
"dst-address": "10.9.9.2/32",
|
||||
gateway: "1.2.3.4",
|
||||
comment: "MikrotikManager:recursive hop-de",
|
||||
distance: "1",
|
||||
},
|
||||
{
|
||||
".id": "*2",
|
||||
static: "true",
|
||||
"dst-address": "10.9.9.3/32",
|
||||
gateway: "1.2.3.4",
|
||||
comment: "not ours",
|
||||
distance: "1",
|
||||
},
|
||||
])
|
||||
assert.equal(mapped.length, 1)
|
||||
assert.equal(mapped[0]?.dstAddress, "10.9.9.2/32")
|
||||
assert.equal(mapped[0]?.comment, "hop-de")
|
||||
}
|
||||
|
||||
console.log("config-apply-plan.test.ts: ok")
|
||||
@@ -0,0 +1,140 @@
|
||||
import {
|
||||
hasManagedCommentPrefix,
|
||||
isOwnedRecursiveComment,
|
||||
stripManagedRecursiveComment,
|
||||
} from "../managed-markers.js"
|
||||
|
||||
export type RosFilterRuleLike = {
|
||||
".id"?: string
|
||||
chain?: string
|
||||
rule?: string
|
||||
comment?: string
|
||||
}
|
||||
|
||||
export type BgpInApplyAction = "patch" | "create" | "delete" | "noop"
|
||||
|
||||
export interface BgpInApplyPlan {
|
||||
action: BgpInApplyAction
|
||||
managedId?: string
|
||||
conflictIds: string[]
|
||||
}
|
||||
|
||||
function isInBgpIn(rule: RosFilterRuleLike): boolean {
|
||||
return (rule.chain ?? "").trim().toLowerCase() === "bgp-in"
|
||||
}
|
||||
|
||||
export function planBgpInApply(
|
||||
existing: RosFilterRuleLike[],
|
||||
rulesCount: number,
|
||||
): BgpInApplyPlan {
|
||||
const managed = existing.find(
|
||||
(r) => isInBgpIn(r) && hasManagedCommentPrefix(r.comment ?? ""),
|
||||
)
|
||||
const conflictIds = existing
|
||||
.filter((r) =>
|
||||
isInBgpIn(r) &&
|
||||
!hasManagedCommentPrefix(r.comment ?? "") &&
|
||||
/bgp-communities/i.test(r.rule ?? ""),
|
||||
)
|
||||
.map((r) => r[".id"])
|
||||
.filter((id): id is string => Boolean(id))
|
||||
|
||||
if (rulesCount > 0) {
|
||||
return {
|
||||
action: managed?.[".id"] ? "patch" : "create",
|
||||
managedId: managed?.[".id"],
|
||||
conflictIds,
|
||||
}
|
||||
}
|
||||
if (managed?.[".id"]) {
|
||||
return { action: "delete", managedId: managed[".id"], conflictIds }
|
||||
}
|
||||
return { action: "noop", conflictIds }
|
||||
}
|
||||
|
||||
export type RosRouteLike = {
|
||||
".id"?: string
|
||||
comment?: string
|
||||
static?: string
|
||||
dynamic?: string
|
||||
blackhole?: string
|
||||
unreachable?: string
|
||||
prohibit?: string
|
||||
"dst-address"?: string
|
||||
gateway?: string
|
||||
}
|
||||
|
||||
export function planRecursiveApply(existing: RosRouteLike[]): { deleteIds: string[] } {
|
||||
return {
|
||||
deleteIds: existing
|
||||
.filter((r) => isOwnedRecursiveComment(r.comment))
|
||||
.map((r) => r[".id"])
|
||||
.filter((id): id is string => Boolean(id)),
|
||||
}
|
||||
}
|
||||
|
||||
export function unmanagedRouteIds(existing: RosRouteLike[]): string[] {
|
||||
return existing
|
||||
.filter((r) => Boolean(r[".id"]) && !isOwnedRecursiveComment(r.comment))
|
||||
.map((r) => r[".id"] as string)
|
||||
}
|
||||
|
||||
export function userRecursiveComment(comment: string | undefined): string {
|
||||
return stripManagedRecursiveComment(comment ?? "")
|
||||
}
|
||||
|
||||
function isIpGateway(gw: string): boolean {
|
||||
return /^\d{1,3}(\.\d{1,3}){3}(?:%\S+)?$/.test(gw.trim())
|
||||
}
|
||||
|
||||
export function isManagedRecursiveRoute(r: RosRouteLike): boolean {
|
||||
if ((r.static ?? "false") !== "true") return false
|
||||
if ((r.dynamic ?? "false") === "true") return false
|
||||
if ((r.blackhole ?? "false") === "true") return false
|
||||
if ((r.unreachable ?? "false") === "true") return false
|
||||
if ((r.prohibit ?? "false") === "true") return false
|
||||
const dst = r["dst-address"] ?? ""
|
||||
const gw = r.gateway ?? ""
|
||||
if (!dst || !gw) return false
|
||||
if (!isIpGateway(gw)) return false
|
||||
return isOwnedRecursiveComment(r.comment)
|
||||
}
|
||||
|
||||
export interface MappedRecursiveRoute {
|
||||
id: string
|
||||
dstAddress: string
|
||||
gateway: string
|
||||
distance: number
|
||||
scope: number | null
|
||||
targetScope: number | null
|
||||
routingTable: string
|
||||
checkGateway: string
|
||||
country: string
|
||||
comment: string
|
||||
disabled: boolean
|
||||
}
|
||||
|
||||
export function mapRosManagedRoutes(
|
||||
rosRoutes: Array<RosRouteLike & {
|
||||
distance?: string
|
||||
scope?: string
|
||||
"target-scope"?: string
|
||||
"routing-table"?: string
|
||||
"check-gateway"?: string
|
||||
disabled?: string
|
||||
}>,
|
||||
): MappedRecursiveRoute[] {
|
||||
return rosRoutes.filter(isManagedRecursiveRoute).map((r, i) => ({
|
||||
id: r[".id"] ?? `ros-${i}`,
|
||||
dstAddress: r["dst-address"] ?? "",
|
||||
gateway: r.gateway ?? "",
|
||||
distance: Number.parseInt(r.distance ?? "1", 10) || 1,
|
||||
scope: r.scope ? (Number.parseInt(r.scope, 10) || null) : null,
|
||||
targetScope: r["target-scope"] ? (Number.parseInt(r["target-scope"], 10) || null) : null,
|
||||
routingTable: r["routing-table"] ?? "main",
|
||||
checkGateway: r["check-gateway"] ?? "",
|
||||
country: "",
|
||||
comment: userRecursiveComment(r.comment),
|
||||
disabled: r.disabled === "true",
|
||||
}))
|
||||
}
|
||||
@@ -0,0 +1,78 @@
|
||||
import assert from "node:assert/strict"
|
||||
import { withPgOrSkip } from "../test/pg.js"
|
||||
import { dbQuery } from "../db/index.js"
|
||||
import {
|
||||
appendRevisionIfChanged,
|
||||
fingerprintPayload,
|
||||
getRevisionById,
|
||||
listRevisions,
|
||||
pruneRevisions,
|
||||
} from "./config-revisions.js"
|
||||
|
||||
if (!(await withPgOrSkip())) {
|
||||
console.log("config-revisions.test.ts: skip")
|
||||
process.exit(0)
|
||||
}
|
||||
|
||||
const tag = `rev-test-${Date.now()}`
|
||||
await dbQuery(`INSERT INTO servers (name, host) VALUES ($1, '127.0.0.1')`, [tag])
|
||||
const { rows } = await dbQuery<{ id: number }>(`SELECT id FROM servers WHERE name = $1 LIMIT 1`, [tag])
|
||||
const serverId = rows[0]?.id
|
||||
assert.ok(serverId)
|
||||
|
||||
try {
|
||||
const payloadA = [{ community: "1:1", action: "route" }]
|
||||
const first = await appendRevisionIfChanged({
|
||||
serverId,
|
||||
section: "filters",
|
||||
source: "apply",
|
||||
payload: payloadA,
|
||||
})
|
||||
assert.equal(first.created, true)
|
||||
assert.equal(first.revision.source, "apply")
|
||||
|
||||
const dup = await appendRevisionIfChanged({
|
||||
serverId,
|
||||
section: "filters",
|
||||
source: "observed",
|
||||
payload: payloadA,
|
||||
})
|
||||
assert.equal(dup.created, false)
|
||||
assert.equal(dup.revision.id, first.revision.id)
|
||||
|
||||
const payloadB = [{ community: "1:2", action: "blackhole" }]
|
||||
const second = await appendRevisionIfChanged({
|
||||
serverId,
|
||||
section: "filters",
|
||||
source: "rollback",
|
||||
payload: payloadB,
|
||||
})
|
||||
assert.equal(second.created, true)
|
||||
assert.equal(second.revision.source, "rollback")
|
||||
assert.notEqual(second.revision.fingerprint, first.revision.fingerprint)
|
||||
|
||||
const listed = await listRevisions(serverId, "filters")
|
||||
assert.equal(listed.length, 2)
|
||||
assert.equal(listed[0]?.source, "rollback")
|
||||
|
||||
const stored = await getRevisionById(second.revision.id)
|
||||
assert.ok(stored)
|
||||
assert.equal(fingerprintPayload(stored.payload), second.revision.fingerprint)
|
||||
|
||||
for (let i = 0; i < 4; i++) {
|
||||
await appendRevisionIfChanged({
|
||||
serverId,
|
||||
section: "filters",
|
||||
source: "apply",
|
||||
payload: [{ community: `9:${i}`, action: "route" }],
|
||||
})
|
||||
}
|
||||
const pruned = await pruneRevisions(serverId, "filters", 3)
|
||||
assert.ok(pruned >= 1)
|
||||
const after = await listRevisions(serverId, "filters")
|
||||
assert.equal(after.length, 3)
|
||||
} finally {
|
||||
await dbQuery(`DELETE FROM servers WHERE id = $1`, [serverId])
|
||||
}
|
||||
|
||||
console.log("config-revisions.test.ts: ok")
|
||||
@@ -0,0 +1,169 @@
|
||||
import { createHash, randomUUID } from "node:crypto"
|
||||
import { and, desc, eq } from "drizzle-orm"
|
||||
import { db } from "../db/index.js"
|
||||
import { configRevisions, type ConfigRevisionRow } from "../db/schema.js"
|
||||
|
||||
export const CONFIG_REVISION_KEEP = 50
|
||||
|
||||
export type ConfigSection = "filters" | "recursive-routes"
|
||||
export type ConfigRevisionSource = "apply" | "rollback" | "observed" | "copy"
|
||||
|
||||
export interface ConfigRevisionDto {
|
||||
id: string
|
||||
serverId: string
|
||||
section: ConfigSection
|
||||
source: ConfigRevisionSource
|
||||
fingerprint: string
|
||||
createdAt: string
|
||||
note: string | null
|
||||
itemCount: number
|
||||
}
|
||||
|
||||
export function stableStringify(value: unknown): string {
|
||||
if (value === null || typeof value !== "object") return JSON.stringify(value)
|
||||
if (Array.isArray(value)) return `[${value.map(stableStringify).join(",")}]`
|
||||
const obj = value as Record<string, unknown>
|
||||
const keys = Object.keys(obj).sort()
|
||||
return `{${keys.map((k) => `${JSON.stringify(k)}:${stableStringify(obj[k])}`).join(",")}}`
|
||||
}
|
||||
|
||||
export function fingerprintPayload(payload: unknown): string {
|
||||
return createHash("sha256").update(stableStringify(payload)).digest("hex")
|
||||
}
|
||||
|
||||
export function canonicalFilterRules(
|
||||
rules: Array<{
|
||||
community?: string
|
||||
action?: string
|
||||
gateway?: string
|
||||
gatewayTunnelId?: string
|
||||
description?: string
|
||||
}>,
|
||||
): unknown[] {
|
||||
return rules.map((r) => ({
|
||||
community: (r.community ?? "").trim(),
|
||||
action: r.action === "blackhole" ? "blackhole" : "route",
|
||||
gateway: r.gateway ?? "",
|
||||
gatewayTunnelId: r.gatewayTunnelId ?? "",
|
||||
description: r.description ?? "",
|
||||
}))
|
||||
}
|
||||
|
||||
export function canonicalRecursiveRoutes(
|
||||
routes: Array<{
|
||||
dstAddress?: string
|
||||
gateway?: string
|
||||
distance?: number
|
||||
scope?: number | null
|
||||
targetScope?: number | null
|
||||
routingTable?: string
|
||||
checkGateway?: string
|
||||
comment?: string
|
||||
disabled?: boolean
|
||||
country?: string
|
||||
}>,
|
||||
): unknown[] {
|
||||
return routes.map((r) => ({
|
||||
dstAddress: (r.dstAddress ?? "").trim(),
|
||||
gateway: r.gateway ?? "",
|
||||
distance: r.distance ?? 1,
|
||||
scope: r.scope ?? null,
|
||||
targetScope: r.targetScope ?? null,
|
||||
routingTable: r.routingTable || "main",
|
||||
checkGateway: r.checkGateway ?? "",
|
||||
comment: r.comment ?? "",
|
||||
disabled: Boolean(r.disabled),
|
||||
country: r.country ?? "",
|
||||
}))
|
||||
}
|
||||
|
||||
export function toRevisionDto(row: ConfigRevisionRow): ConfigRevisionDto {
|
||||
const payload = row.payload
|
||||
const itemCount = Array.isArray(payload) ? payload.length : 0
|
||||
return {
|
||||
id: row.id,
|
||||
serverId: String(row.serverId),
|
||||
section: row.section,
|
||||
source: row.source,
|
||||
fingerprint: row.fingerprint,
|
||||
createdAt: row.createdAt,
|
||||
note: row.note ?? null,
|
||||
itemCount,
|
||||
}
|
||||
}
|
||||
|
||||
export async function listRevisions(
|
||||
serverId: number,
|
||||
section: ConfigSection,
|
||||
limit = CONFIG_REVISION_KEEP,
|
||||
): Promise<ConfigRevisionDto[]> {
|
||||
const rows = await db
|
||||
.select()
|
||||
.from(configRevisions)
|
||||
.where(and(eq(configRevisions.serverId, serverId), eq(configRevisions.section, section)))
|
||||
.orderBy(desc(configRevisions.createdAt))
|
||||
.limit(limit)
|
||||
return rows.map(toRevisionDto)
|
||||
}
|
||||
|
||||
export async function getRevisionById(id: string): Promise<ConfigRevisionRow | undefined> {
|
||||
return (await db.select().from(configRevisions).where(eq(configRevisions.id, id)).limit(1))[0]
|
||||
}
|
||||
|
||||
export async function pruneRevisions(
|
||||
serverId: number,
|
||||
section: ConfigSection,
|
||||
keep = CONFIG_REVISION_KEEP,
|
||||
): Promise<number> {
|
||||
const rows = await db
|
||||
.select({ id: configRevisions.id })
|
||||
.from(configRevisions)
|
||||
.where(and(eq(configRevisions.serverId, serverId), eq(configRevisions.section, section)))
|
||||
.orderBy(desc(configRevisions.createdAt))
|
||||
const extra = rows.slice(keep)
|
||||
if (extra.length === 0) return 0
|
||||
for (const row of extra) {
|
||||
await db.delete(configRevisions).where(eq(configRevisions.id, row.id))
|
||||
}
|
||||
return extra.length
|
||||
}
|
||||
|
||||
export async function appendRevisionIfChanged(input: {
|
||||
serverId: number
|
||||
section: ConfigSection
|
||||
source: ConfigRevisionSource
|
||||
payload: unknown
|
||||
note?: string | null
|
||||
}): Promise<{ created: boolean; revision: ConfigRevisionDto }> {
|
||||
const fingerprint = fingerprintPayload(input.payload)
|
||||
const latest = (await db
|
||||
.select()
|
||||
.from(configRevisions)
|
||||
.where(and(
|
||||
eq(configRevisions.serverId, input.serverId),
|
||||
eq(configRevisions.section, input.section),
|
||||
))
|
||||
.orderBy(desc(configRevisions.createdAt))
|
||||
.limit(1))[0]
|
||||
|
||||
if (latest?.fingerprint === fingerprint) {
|
||||
return { created: false, revision: toRevisionDto(latest) }
|
||||
}
|
||||
|
||||
const now = new Date().toISOString()
|
||||
const id = randomUUID()
|
||||
await db.insert(configRevisions).values({
|
||||
id,
|
||||
serverId: input.serverId,
|
||||
section: input.section,
|
||||
source: input.source,
|
||||
fingerprint,
|
||||
payload: Array.isArray(input.payload) ? input.payload : [],
|
||||
note: input.note ?? null,
|
||||
createdAt: now,
|
||||
})
|
||||
await pruneRevisions(input.serverId, input.section)
|
||||
const row = await getRevisionById(id)
|
||||
if (!row) throw new Error("config-revisions: insert vanished")
|
||||
return { created: true, revision: toRevisionDto(row) }
|
||||
}
|
||||
Reference in New Issue
Block a user