From 301edd2c9addf6809f10d2168e7bab4e22417c58 Mon Sep 17 00:00:00 2001 From: openhands Date: Wed, 23 Sep 2026 14:45:35 +0200 Subject: [PATCH] feat(security): ops alerts, shared backoff, atomic quota and daily stats for CrowdSec Add an alerting/stats layer over the existing CrowdSec integration: - New crowdsec-alerts.ts: cooldown-gated ops alerts (Redis NX lock, TTL from HEALTH_ALERT_COOLDOWN_MIN) fanning out through the app's sendAlert service. Raised for daily quota exhaustion, block bursts (5-min window past CROWDSEC_ALERT_BLOCK_BURST), and signal-push failures. - New crowdsec-stats.ts: daily counters (lookups/blocks/reports/report_fail) in Redis with a 14-day reader for the admin panel. - Shared 403/429 backoff: the pause marker now lives in Redis (crowdsec:backoff-until) so every instance honours it, not just the process that hit the limit. - Atomic quota reservation: INCR-before-call with self-rollback on overshoot, so concurrent instances can never slip calls past the daily ceiling. - Admin anti-DDoS page gains a last-14-days activity table next to the quota bar. --- .env.example | 3 + src/app/admin/devops/antiddos/page.tsx | 58 +++++++++ src/env.ts | 3 + src/lib/crowdsec-alerts.ts | 66 ++++++++++ src/lib/crowdsec-api.test.ts | 168 ++++++++++++++++++++++++- src/lib/crowdsec-api.ts | 158 +++++++++++++++++++---- src/lib/crowdsec-report.test.ts | 65 +++++++++- src/lib/crowdsec-report.ts | 16 +++ src/lib/crowdsec-stats.ts | 111 ++++++++++++++++ 9 files changed, 620 insertions(+), 28 deletions(-) create mode 100644 src/lib/crowdsec-alerts.ts create mode 100644 src/lib/crowdsec-stats.ts diff --git a/.env.example b/.env.example index d2855327..f367c168 100644 --- a/.env.example +++ b/.env.example @@ -93,6 +93,9 @@ CROWDSEC_CTI_BASE_URL=https://cti.api.crowdsec.net/v2 # counter reaches it, reputation lookups pause until tomorrow so a spread # DDoS cannot silently burn the whole quota. 0 = unlimited. CROWDSEC_CTI_DAILY_QUOTA=10000 +# How many new community-reputation blocks within a 5-minute window justify an +# ops alert (quota/backoff/report alerts all use HEALTH_ALERT_COOLDOWN_MIN). +CROWDSEC_ALERT_BLOCK_BURST=10 # --- CROWDSEC SIGNAL PUSH (share our blocks back, optional) --- # Opt-in: pushes blocked IPs + behaviors to the CrowdSec Central API (CAPI) so diff --git a/src/app/admin/devops/antiddos/page.tsx b/src/app/admin/devops/antiddos/page.tsx index 7ddf1cfb..3f51cd6e 100644 --- a/src/app/admin/devops/antiddos/page.tsx +++ b/src/app/admin/devops/antiddos/page.tsx @@ -45,6 +45,7 @@ import { crowdsecReportEnabled, getLastCrowdsecReport, } from "@/lib/crowdsec-report"; +import { getCrowdsecStats } from "@/lib/crowdsec-stats"; import { db, WebsiteSetting } from "@/lib/db"; import { canAccess, getAdminContext, PERMS } from "@/lib/permissions"; import { redis } from "@/lib/redis"; @@ -147,6 +148,7 @@ export default async function AdminAntiDdosPage() { const crowdsecUsage = redisOk ? await getCrowdsecQuotaUsage() : null; const reportingEnabled = await crowdsecReportEnabled(); const lastReport = await getLastCrowdsecReport(); + const crowdsecStats = redisOk ? await getCrowdsecStats(14) : []; return (
@@ -723,6 +725,62 @@ export default async function AdminAntiDdosPage() {
)} + {crowdsecStats.length > 0 && ( +
+

+ Daily activity (last {crowdsecStats.length} days) +

+
+ + + + + + + + + + + + {crowdsecStats.map((row) => ( + + + + + + + + ))} + +
Date + Lookups + + Blocks + + Reports + Failures
+ {row.date === new Date().toISOString().slice(0, 10) + ? "Today" + : row.date.slice(5)} + + {row.lookups.toLocaleString()} + + {row.blocks.toLocaleString()} + + {row.reports.toLocaleString()} + + {row.reportFailures > 0 ? ( + + {row.reportFailures.toLocaleString()} + + ) : ( + "–" + )} +
+
+
+ )} +

Community signal push

diff --git a/src/env.ts b/src/env.ts index 234cf2cc..94c3baa2 100644 --- a/src/env.ts +++ b/src/env.ts @@ -182,6 +182,9 @@ const schema = z // it, so a spread DDoS can never silently burn the whole quota; 0 // disables the guard. CROWDSEC_CTI_DAILY_QUOTA: z.coerce.number().int().min(0).default(10_000), + // How many new community-reputation blocks within a 5-minute window + // justify an ops alert (cooldown-gated via HEALTH_ALERT_COOLDOWN_MIN). + CROWDSEC_ALERT_BLOCK_BURST: z.coerce.number().int().min(1).default(10), // Share our own detections back into the CrowdSec community blocklist // (signal push over the Central API). Opt-in: flipping this on publicly // shares blocked IPs + behaviors, so it defaults to off. diff --git a/src/lib/crowdsec-alerts.ts b/src/lib/crowdsec-alerts.ts new file mode 100644 index 00000000..ba551572 --- /dev/null +++ b/src/lib/crowdsec-alerts.ts @@ -0,0 +1,66 @@ +import "server-only"; + +import { env } from "@/env"; +import { logger } from "@/lib/logger"; +import { redis } from "@/lib/redis"; +import { type SendAlertInput, sendAlert } from "@/lib/services/alert"; + +// === CrowdSec operational alerts =========================================== +// +// Thin, fire-and-forget wrapper around the app's alert service for the +// reputation pipeline. Every raise is cooldown-gated through a Redis NX lock +// (key crowdsec:alert:{key}, TTL = HEALTH_ALERT_COOLDOWN_MIN), so N instances +// and flapping conditions surface exactly one alert per window instead of +// spamming Discord/email/alert_logs. Falls back to alerting anyway when Redis +// is unreachable — a silent quota blowout is worse than one duplicate alert. + +const ALERT_PREFIX = "crowdsec:alert:"; + +function cooldownSeconds(): number { + const raw = Number(env.HEALTH_ALERT_COOLDOWN_MIN ?? 15); + return Math.ceil((Number.isFinite(raw) && raw > 0 ? raw : 15) * 60); +} + +/** + * Raise an alert unless the cooldown window is still active. Returns the + * sendAlert promise when the alert was actually raised, or false when it was + * suppressed. Never throws; the caller may `void` the result on hot paths. + */ +export async function raiseCrowdsecAlert( + key: string, + input: { + type?: string; + severity: SendAlertInput["severity"]; + message: string; + context?: SendAlertInput["context"]; + }, +): Promise>> { + if (redis) { + try { + const acquired = await redis.set( + `${ALERT_PREFIX}${key}`, + String(Date.now()), + "EX", + cooldownSeconds(), + "NX", + ); + if (acquired !== "OK") return false; + } catch { + // Cooldown bookkeeping failed — alert anyway rather than silently drop. + } + } + try { + return await sendAlert({ + type: "ddos", + severity: input.severity, + message: input.message, + context: input.context, + }); + } catch (error) { + logger.error("[crowdsec-alert] sendAlert raised an unexpected error", { + key, + err: error, + }); + return false; + } +} diff --git a/src/lib/crowdsec-api.test.ts b/src/lib/crowdsec-api.test.ts index dae6e106..358db17c 100644 --- a/src/lib/crowdsec-api.test.ts +++ b/src/lib/crowdsec-api.test.ts @@ -14,11 +14,20 @@ import { verdictIsMalicious, verifyCrowdsecConnection, } from "./crowdsec-api"; +import { type CrowdsecDailyStat, getCrowdsecStats } from "./crowdsec-stats"; // Unit-test the CTI client in isolation: a deterministic in-memory Redis fake // and a silenced logger, so fetch calls count only CrowdSec lookups. CrowdSec // deliberately never touches Cloudflare, so no Cloudflare surface is stubbed. -const state = vi.hoisted(() => ({ map: new Map() })); +const state = vi.hoisted(() => ({ + map: new Map(), + sendAlert: vi.fn(), +})); + +vi.mock("@/lib/services/alert", () => ({ + sendAlert: state.sendAlert, + ddosDetected: vi.fn(), +})); vi.mock("@/lib/redis", () => ({ redis: { @@ -43,6 +52,11 @@ vi.mock("@/lib/redis", () => ({ state.map.set(key, String(next)); return next; }, + decr: async (key: string) => { + const next = (Number(state.map.get(key)) || 0) - 1; + state.map.set(key, String(next)); + return next; + }, expire: async () => 1, pttl: async () => 60_000, }, @@ -58,6 +72,8 @@ vi.mock("@/lib/logger", () => ({ }, })); +const tick = () => new Promise((resolve) => setTimeout(resolve, 20)); + function jsonResponse(body: unknown, status = 200): Response { return new Response(JSON.stringify(body), { status, @@ -109,6 +125,7 @@ describe("crowdsec-api", () => { vi.unstubAllGlobals(); vi.unstubAllEnvs(); state.map.clear(); + state.sendAlert.mockReset(); resetCrowdsecCache(); fetchMock = vi.fn(); vi.stubGlobal("fetch", fetchMock); @@ -352,6 +369,34 @@ describe("crowdsec-api", () => { expect(fetchMock).toHaveBeenCalledTimes(1); }); + it("publishes the backoff to shared Redis so every instance respects it", async () => { + vi.stubEnv("CROWDSEC_API_KEY", "cs_key"); + fetchMock.mockResolvedValue(jsonResponse({ message: "Invalid key" }, 403)); + + await maybeAutoBlockCrowdsec({ + ip: blockIp(), + category: "api", + ttlSeconds: 600, + scoreThreshold: 4, + enabled: true, + }); + // The shared marker exists and points into the future. + const until = Number(state.map.get("crowdsec:backoff-until")); + expect(Number.isFinite(until)).toBe(true); + expect(until).toBeGreaterThan(Date.now()); + + // A fresh instance (reset in-process state) still honours the marker. + resetCrowdsecCache(); + await maybeAutoBlockCrowdsec({ + ip: "203.0.113.44", + category: "api", + ttlSeconds: 600, + scoreThreshold: 4, + enabled: true, + }); + expect(fetchMock).toHaveBeenCalledTimes(1); + }); + it("backs off after a 429 rate limit as well", async () => { vi.stubEnv("CROWDSEC_API_KEY", "cs_key"); fetchMock.mockResolvedValue(jsonResponse({ message: "rate limited" }, 429)); @@ -525,4 +570,125 @@ describe("crowdsec-api", () => { } expect(getMemoryVerdictCacheSize()).toBe(2000); }); + + it("raises an ops alert when the daily quota is exhausted", async () => { + vi.stubEnv("CROWDSEC_API_KEY", "cs_key"); + vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "1"); + fetchMock.mockImplementation(() => + Promise.resolve(jsonResponse(maliciousItem(blockIp()))), + ); + + await maybeAutoBlockCrowdsec({ + ip: blockIp(), + category: "api", + ttlSeconds: 600, + scoreThreshold: 4, + enabled: true, + }); + await maybeAutoBlockCrowdsec({ + ip: "198.51.100.2", + category: "api", + ttlSeconds: 600, + scoreThreshold: 4, + enabled: true, + }); + await tick(); + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(state.sendAlert).toHaveBeenCalledTimes(1); + const [input] = state.sendAlert.mock.calls[0]; + expect(input.type).toBe("ddos"); + expect(input.severity).toBe("warning"); + expect(input.message).toContain("quota exhausted"); + expect(input.context).toMatchObject({ quota: 1 }); + }); + + it("floods once per cooldown window when blocks burst past the threshold", async () => { + vi.stubEnv("CROWDSEC_API_KEY", "cs_key"); + vi.stubEnv("CROWDSEC_ALERT_BLOCK_BURST", "2"); + vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "0"); + fetchMock.mockImplementation((url: string | URL) => + Promise.resolve( + jsonResponse(maliciousItem(String(url).split("/").pop() ?? "ip")), + ), + ); + + await maybeAutoBlockCrowdsec({ + ip: "198.51.100.71", + category: "api", + ttlSeconds: 600, + scoreThreshold: 4, + enabled: true, + }); + await maybeAutoBlockCrowdsec({ + ip: "198.51.100.72", + category: "api", + ttlSeconds: 600, + scoreThreshold: 4, + enabled: true, + }); + // Third block in the same window: threshold crossed, but the alert is + // cooldown-gated so it still fires exactly once. + await maybeAutoBlockCrowdsec({ + ip: "198.51.100.73", + category: "api", + ttlSeconds: 600, + scoreThreshold: 4, + enabled: true, + }); + await tick(); + expect(state.sendAlert).toHaveBeenCalledTimes(1); + const [input] = state.sendAlert.mock.calls[0]; + expect(input.type).toBe("ddos"); + expect(input.context).toMatchObject({ blocks: 2, threshold: 2 }); + expect(state.map.get(`antiddos:block:198.51.100.72`)).toBe("crowdsec"); + }); + + it("tallies lookups and blocks into the daily stats histogram", async () => { + vi.stubEnv("CROWDSEC_API_KEY", "cs_key"); + vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "0"); + fetchMock.mockImplementation((url: string | URL) => { + const ip = String(url).split("/").pop() ?? "ip"; + return Promise.resolve( + jsonResponse( + ip === "198.51.100.83" ? suspiciousItem(ip, 3) : maliciousItem(ip), + ), + ); + }); + + await maybeAutoBlockCrowdsec({ + ip: "198.51.100.81", + category: "api", + ttlSeconds: 600, + scoreThreshold: 4, + enabled: true, + }); + await maybeAutoBlockCrowdsec({ + ip: "198.51.100.82", + category: "api", + ttlSeconds: 600, + scoreThreshold: 4, + enabled: true, + }); + await maybeAutoBlockCrowdsec({ + ip: "198.51.100.83", + category: "api", + ttlSeconds: 600, + scoreThreshold: 7, // suspicious/known verdicts below threshold: lookup only + enabled: true, + }); + await tick(); + + const stats = await getCrowdsecStats(1); + const today: CrowdsecDailyStat | undefined = stats.find( + (row) => row.date === new Date().toISOString().slice(0, 10), + ); + expect(today?.lookups).toBe(3); + expect(today?.blocks).toBe(2); + expect(today?.reportFailures).toBe(0); + expect( + state.map.get( + `crowdsec:stat:lookups:${new Date().toISOString().slice(0, 10)}`, + ), + ).toBe("3"); + }); }); diff --git a/src/lib/crowdsec-api.ts b/src/lib/crowdsec-api.ts index aaf60ff6..374f55c8 100644 --- a/src/lib/crowdsec-api.ts +++ b/src/lib/crowdsec-api.ts @@ -1,7 +1,9 @@ import "server-only"; import { env } from "@/env"; +import { raiseCrowdsecAlert } from "@/lib/crowdsec-alerts"; import { reportCrowdsecSignal } from "@/lib/crowdsec-report"; +import { bumpCrowdsecStat } from "@/lib/crowdsec-stats"; import { logger } from "@/lib/logger"; import { redis } from "@/lib/redis"; import { UNKNOWN_CLIENT_IP } from "./client-ip"; @@ -113,6 +115,11 @@ const LAST_VERIFY_KEY = "crowdsec:last-verify"; const BLOCK_META_PREFIX = "antiddos:block:meta:"; const QUOTA_PREFIX = "crowdsec:usage:"; const QUOTA_KEY_TTL_SECONDS = 48 * 3_600; +/** Shared 403/429 pause marker, so every instance respects the backoff. */ +const BACKOFF_KEY = "crowdsec:backoff-until"; +/** Short-window block burst counter: crowdsec:burst:{unix-5min-bucket}. */ +const BURST_PREFIX = "crowdsec:burst:"; +const BURST_WINDOW_SECONDS = 300; /** In-process verdict cache cap so a flood of distinct IPs cannot grow it forever. */ const MEMORY_VERDICT_CACHE_MAX = 2_000; /** Warn at this fraction of the daily quota, once per day. */ @@ -313,6 +320,47 @@ async function acquireLookupLock(ip: string): Promise { let backoffUntil = 0; let quotaWarnedDate: string | null = null; +/** + * Next moment (epoch ms) the CTI API may be called again — the max of the + * in-process view and the shared Redis marker so every instance respects a + * backoff discovered by any of them. Redis is only read when the local view is + * not already active, keeping the hot path cheap. + */ +async function getBackoffUntil(): Promise { + if (Date.now() < backoffUntil) return backoffUntil; + if (redis) { + try { + const raw = await redis.get(BACKOFF_KEY); + const shared = Number(raw ?? 0); + if (Number.isFinite(shared) && shared > backoffUntil) { + backoffUntil = shared; + } + } catch { + // Redis hiccup — local view is enough + } + } + return backoffUntil; +} + +async function setBackoff(ms: number): Promise { + const until = Date.now() + ms; + backoffUntil = until; + if (redis) { + try { + // EX rounds up so the marker outlives the wait it encodes, plus a + // second of slack for the read path. + await redis.set( + BACKOFF_KEY, + String(until), + "EX", + Math.ceil(ms / 1000) + 1, + ); + } catch { + // local view still protects this instance + } + } +} + function quotaDate(): string { return new Date().toISOString().slice(0, 10); } @@ -350,33 +398,87 @@ export async function getCrowdsecQuotaUsage(): Promise { return { date, used, quota, exhausted: quota > 0 && used >= quota }; } +/** + * Reserve one API call against today's quota. Atomic: the counter is INCR'd + * BEFORE the call and compared to the ceiling, so concurrent instances can + * never slip calls past the budget; a reserve that overshoots rolls itself + * back. Returns false once the budget is spent (and raises an ops alert). + */ async function reserveQuota(): Promise { const quota = dailyQuota(); if (quota <= 0) return true; + if (!redis) return true; // no shared counter → unlimited best-effort const date = quotaDate(); - const usage = await getCrowdsecQuotaUsage(); - if (usage.used >= quota) { - // Stop consulting the API for the rest of the day: a flood of distinct - // bucket-tripping IPs would otherwise burn every remaining call and - // then sit in a 429 storm anyway. - return false; - } - if (usage.used >= quota * QUOTA_WARN_RATIO && quotaWarnedDate !== date) { - quotaWarnedDate = date; - logger.warn("[crowdsec-api] CTI daily quota nearing its limit", { - used: usage.used, - quota, - }); - } - if (redis) { - try { - await redis.incr(quotaKey(date)); - await redis.expire(quotaKey(date), QUOTA_KEY_TTL_SECONDS); - } catch { - // best effort — an uncounted call is better than a failed lookup + const key = quotaKey(date); + try { + const used = await redis.incr(key); + await redis.expire(key, QUOTA_KEY_TTL_SECONDS); + if (used > quota) { + // Concurrent reserves nudged us past the ceiling — give the slot + // back and refuse: the budget would be spent the very next call + // anyway, so stopping here is both safe and quota-exact. + await redis.decr(key); + logger.warn( + "[crowdsec-api] CTI daily quota exhausted — pausing lookups until tomorrow", + { quota }, + ); + void raiseCrowdsecAlert("quota", { + type: "ddos", + severity: "warning", + message: `CrowdSec reputation quota exhausted for today (${used} of ${quota} enrichment calls) — lookups are paused until tomorrow.`, + context: { used, quota, date }, + }); + return false; } + if (used >= quota * QUOTA_WARN_RATIO && quotaWarnedDate !== date) { + quotaWarnedDate = date; + logger.warn("[crowdsec-api] CTI daily quota nearing its limit", { + used, + quota, + }); + } + return true; + } catch { + // Redis hiccup at a moment we could not count — allow the call rather + // than break the gate; the verdict cache still limits frequency. + return true; + } +} + +/** Block burst threshold from env, defensively coerced (falls back to 10). */ +function dailyBlockBurstThreshold(): number { + const raw = Number(env.CROWDSEC_ALERT_BLOCK_BURST ?? 10); + return Number.isFinite(raw) && raw > 0 ? Math.floor(raw) : 10; +} + +/** + * A burst of new community-reputation blocks is usually an automated attack + * wave. Count blocks into a rolling 5-minute bucket and alert once per + * cooldown window when they cross CROWDSEC_ALERT_BLOCK_BURST. Fire-and-forget. + */ +async function trackBlockBurst(): Promise { + if (!redis) return; + const bucket = Math.floor(Date.now() / 1000 / BURST_WINDOW_SECONDS); + const key = `${BURST_PREFIX}${bucket}`; + const threshold = dailyBlockBurstThreshold(); + try { + const count = await redis.incr(key); + await redis.expire(key, BURST_WINDOW_SECONDS * 2); + if (count >= threshold) { + void raiseCrowdsecAlert("block-burst", { + type: "ddos", + severity: "warning", + message: `CrowdSec community reputation blocked ${count} IPs in the last ${BURST_WINDOW_SECONDS / 60} minutes — likely an automated attack wave.`, + context: { + blocks: count, + windowSeconds: BURST_WINDOW_SECONDS, + threshold, + }, + }); + } + } catch { + // alert is best-effort — never break the block path } - return true; } /** @@ -390,7 +492,7 @@ export async function lookupCrowdsecVerdict( ): Promise { if (!crowdsecEnabled()) return null; if (!ip || ip === UNKNOWN_CLIENT_IP) return null; - if (Date.now() < backoffUntil) return null; + if (Date.now() < (await getBackoffUntil())) return null; const cached = await readVerdictCache(ip); if (cached) return cached; @@ -417,6 +519,9 @@ export async function lookupCrowdsecVerdict( } const response = await crowdsecRequest(`/smoke/${encodeURIComponent(ip)}`); + // The enrichment call happened — count it for the daily histogram, + // regardless of whether the verdict was positive, negative, or n/a. + void bumpCrowdsecStat("lookups"); if (response.status === 404) { // Unknown to the community — cache the negative result so a clean @@ -426,13 +531,13 @@ export async function lookupCrowdsecVerdict( return verdict; } if (response.status === 403) { - backoffUntil = Date.now() + AUTH_BACKOFF_MS; + await setBackoff(AUTH_BACKOFF_MS); throw new CrowdsecApiError( `CrowdSec API key rejected (HTTP 403): ${await errorDetail(response)}`, ); } if (response.status === 429) { - backoffUntil = Date.now() + RATE_LIMIT_BACKOFF_MS; + await setBackoff(RATE_LIMIT_BACKOFF_MS); logger.warn("[crowdsec-api] CTI API rate limit hit — backing off", { ip, backoffMs: RATE_LIMIT_BACKOFF_MS, @@ -477,7 +582,7 @@ export async function maybeAutoBlockCrowdsec(input: { // The gate only ever reads its block key through shared Redis — without it // there is nowhere durable to record the block. if (!redis) return; - if (Date.now() < backoffUntil) return; + if (Date.now() < (await getBackoffUntil())) return; try { const verdict = await lookupCrowdsecVerdict(ip); @@ -528,6 +633,9 @@ export async function maybeAutoBlockCrowdsec(input: { // Opt-in community signal push (CAPI), fire-and-forget: never awaited, // never throws, and internally deduped per IP. void reportCrowdsecSignal({ ip, category, ttlSeconds, verdict, meta }); + // Daily histogram + burst detection (cooldown-gated ops alert). + void bumpCrowdsecStat("blocks"); + void trackBlockBurst(); } catch (error) { logger.error("[crowdsec-api] Automatic IP block failed", { ip, diff --git a/src/lib/crowdsec-report.test.ts b/src/lib/crowdsec-report.test.ts index 8776b089..42231e38 100644 --- a/src/lib/crowdsec-report.test.ts +++ b/src/lib/crowdsec-report.test.ts @@ -11,7 +11,10 @@ import { // The signal-push watcher is tested against a deterministic in-memory Redis // fake (NX lock + token cache) and a mocked fetch that routes the CAPI paths. -const state = vi.hoisted(() => ({ map: new Map() })); +const state = vi.hoisted(() => ({ + map: new Map(), + sendAlert: vi.fn(), +})); vi.mock("@/lib/redis", () => ({ redis: { @@ -31,6 +34,12 @@ vi.mock("@/lib/redis", () => ({ for (const key of keys) state.map.delete(key); return keys.length; }, + incr: async (key: string) => { + const next = (Number(state.map.get(key)) || 0) + 1; + state.map.set(key, String(next)); + return next; + }, + expire: async () => 1, }, __esModule: true, })); @@ -44,6 +53,13 @@ vi.mock("@/lib/logger", () => ({ }, })); +vi.mock("@/lib/services/alert", () => ({ + sendAlert: state.sendAlert, + ddosDetected: vi.fn(), +})); + +const tick = () => new Promise((resolve) => setTimeout(resolve, 20)); + const CAPI = "https://capi.example.test/v3"; const MACHINE = "m".repeat(48); const PASSWORD = "Strong!1P@ssw0rdStrong!1P@ssw0rd"; @@ -111,6 +127,7 @@ describe("crowdsec-report", () => { vi.unstubAllGlobals(); vi.unstubAllEnvs(); state.map.clear(); + state.sendAlert.mockReset(); resetCrowdsecReportCache(); fetchMock = vi.fn(); vi.stubGlobal("fetch", fetchMock); @@ -256,12 +273,56 @@ describe("crowdsec-report", () => { }); await expect(reportCrowdsecSignal(signalInput())).resolves.toBeUndefined(); - await new Promise((resolve) => setTimeout(resolve, 20)); + await tick(); const last: CrowdsecReportStatus | null = await getLastCrowdsecReport(); expect(last?.ok).toBe(false); expect(last?.message).toContain("signal push rejected"); }); + it("counts a failed push and raises a cooldown-gated ops alert", async () => { + fetchMock.mockImplementation((url: string) => { + const path = String(url).replace(CAPI, ""); + if (path === "/watchers/login") { + return Promise.resolve( + jsonResponse({ + token: "jwt-xyz", + expire: new Date(Date.now() + 3_600_000).toISOString(), + }), + ); + } + if (path === "/signals") { + return Promise.resolve(jsonResponse({ message: "boom" }, 500)); + } + return Promise.resolve(jsonResponse({})); + }); + + await reportCrowdsecSignal(signalInput("198.51.100.20")); + await tick(); + const today = new Date().toISOString().slice(0, 10); + expect(state.map.get(`crowdsec:stat:report_fail:${today}`)).toBe("1"); + expect(state.sendAlert).toHaveBeenCalledTimes(1); + const [input] = state.sendAlert.mock.calls[0]; + expect(input.type).toBe("ddos"); + expect(input.severity).toBe("warning"); + expect(input.context).toMatchObject({ ip: "198.51.100.20" }); + + // A second failed push inside the cooldown window stays silent. + await reportCrowdsecSignal(signalInput("198.51.100.21")); + await tick(); + expect(state.sendAlert).toHaveBeenCalledTimes(1); + expect(state.map.get(`crowdsec:stat:report_fail:${today}`)).toBe("2"); + }); + + it("tallies successful pushes into the daily stats histogram", async () => { + routeCapi(); + await reportCrowdsecSignal(signalInput("198.51.100.30")); + await reportCrowdsecSignal(signalInput("198.51.100.31")); + await tick(); + const today = new Date().toISOString().slice(0, 10); + expect(state.map.get(`crowdsec:stat:reports:${today}`)).toBe("2"); + expect(state.sendAlert).not.toHaveBeenCalled(); + }); + it("verifies the watcher channel end to end", async () => { routeCapi(); const status = await verifyCrowdsecReporting(); diff --git a/src/lib/crowdsec-report.ts b/src/lib/crowdsec-report.ts index 7a6a852d..950f1fa7 100644 --- a/src/lib/crowdsec-report.ts +++ b/src/lib/crowdsec-report.ts @@ -2,11 +2,13 @@ import "server-only"; import { createHash, randomBytes } from "node:crypto"; import { env } from "@/env"; +import { raiseCrowdsecAlert } from "@/lib/crowdsec-alerts"; import type { CrowdsecBlockMeta, CrowdsecConnectionStatus, CrowdsecVerdict, } from "@/lib/crowdsec-api"; +import { bumpCrowdsecStat } from "@/lib/crowdsec-stats"; import { logger } from "@/lib/logger"; import { redis } from "@/lib/redis"; import { UNKNOWN_CLIENT_IP } from "./client-ip"; @@ -390,6 +392,7 @@ async function pushSignal(input: { } const status: CrowdsecReportStatus = { ok: true, at: Date.now() }; await setLastCrowdsecReport(status); + void bumpCrowdsecStat("reports"); logger.info("[crowdsec-report] Detection shared with the community", { ip: input.ip, category: input.category, @@ -404,6 +407,19 @@ async function pushSignal(input: { message: String(error), }; await setLastCrowdsecReport(status); + void bumpCrowdsecStat("report_fail"); + // Ops alert, cooldown-gated: a silently broken channel means the + // community never learns about the blocks we keep sharing. + void raiseCrowdsecAlert("report", { + type: "ddos", + severity: "warning", + message: + "CrowdSec signal push failed — detections are not reaching the community.", + context: { + ip: input.ip, + detail: String(error).slice(0, 300), + }, + }); logger.error("[crowdsec-report] Signal push failed", { ip: input.ip, err: error, diff --git a/src/lib/crowdsec-stats.ts b/src/lib/crowdsec-stats.ts new file mode 100644 index 00000000..5c069b88 --- /dev/null +++ b/src/lib/crowdsec-stats.ts @@ -0,0 +1,111 @@ +import "server-only"; + +import { redis } from "@/lib/redis"; + +// === CrowdSec daily counters =============================================== +// +// Small Redis counters so ops can see whether the reputation pipeline is +// actually doing anything: lookups executed, blocks created, signals pushed, +// and push failures — all bucketed per UTC calendar day +// (crowdsec:stat:{metric}:{YYYY-MM-DD}). Both the CTI client and the signal +// pusher feed them; the admin panel renders the last N days. Cheap INCRs on +// non-hot paths only, so they never tax the request path. + +export type CrowdsecStatMetric = + | "lookups" + | "blocks" + | "reports" + | "report_fail"; + +const STAT_PREFIX = "crowdsec:stat:"; +const STAT_KEY_TTL_SECONDS = 16 * 24 * 3_600; + +function statDate(): string { + return new Date().toISOString().slice(0, 10); +} + +function statKey(metric: CrowdsecStatMetric, date: string): string { + return `${STAT_PREFIX}${metric}:${date}`; +} + +/** Count one occurrence of a pipeline event for today. Best effort. */ +export async function bumpCrowdsecStat( + metric: CrowdsecStatMetric, +): Promise { + if (!redis) return; + const key = statKey(metric, statDate()); + try { + await redis.incr(key); + await redis.expire(key, STAT_KEY_TTL_SECONDS); + } catch { + // tracking is best-effort — a miss only loses a day's histogram + } +} + +export interface CrowdsecDailyStat { + /** UTC calendar day (YYYY-MM-DD). */ + date: string; + lookups: number; + blocks: number; + reports: number; + reportFailures: number; +} + +/** + * Read the per-day counters for the last `days` days (oldest first, ending + * with today). Reads the four metric keys per day via one round-trip of + * parallel GETs; never throws. + */ +export async function getCrowdsecStats( + days = 14, +): Promise { + const today = statDate(); + const rows: CrowdsecDailyStat[] = []; + if (!redis) { + // No shared store — still return a blank timeline for the UI. + for (let i = days - 1; i >= 0; i -= 1) { + const date = new Date(Date.now() - i * 86_400_000) + .toISOString() + .slice(0, 10); + rows.push({ date, lookups: 0, blocks: 0, reports: 0, reportFailures: 0 }); + } + return rows; + } + for (let i = days - 1; i >= 0; i -= 1) { + const date = new Date(Date.now() - i * 86_400_000) + .toISOString() + .slice(0, 10); + try { + const [lookups, blocks, reports, reportFailures] = await Promise.all([ + redis.get(statKey("lookups", date)), + redis.get(statKey("blocks", date)), + redis.get(statKey("reports", date)), + redis.get(statKey("report_fail", date)), + ]); + const num = (raw: string | null): number => { + const n = Number(raw ?? 0); + return Number.isFinite(n) ? n : 0; + }; + rows.push({ + date, + lookups: num(lookups), + blocks: num(blocks), + reports: num(reports), + reportFailures: num(reportFailures), + }); + } catch { + rows.push({ date, lookups: 0, blocks: 0, reports: 0, reportFailures: 0 }); + } + } + // Guard against an odd clock roll-back leaving a hole; keep chronological. + if (rows.length && rows[rows.length - 1]?.date !== today) { + rows.push({ + date: today, + lookups: 0, + blocks: 0, + reports: 0, + reportFailures: 0, + }); + } + return rows; +}