feat(security): recovery alerts, gate-block sharing, rolling-window burst and admin breakdown for CrowdSec
Gitea Actions Runner Test / test-job (push) Successful in 1s
CI / check (push) Successful in 30s
CI / tests-integration (push) Successful in 1m42s
CI / tests-unit (push) Successful in 1m50s
CI / tests-ui (push) Successful in 2m42s
CI / preflight (push) Skipped
CI / deploy (push) Successful in 2m3s

This commit is contained in:
openhands committed 2026-09-23 15:06:16 +02:00
1 parent 301edd2c9a
commit 3e1a3f92c8
8 files changed
+448 -65

No files matched your search

+71 -5
View File
@@ -45,7 +45,7 @@ import {
crowdsecReportEnabled, crowdsecReportEnabled,
getLastCrowdsecReport, getLastCrowdsecReport,
} from "@/lib/crowdsec-report"; } from "@/lib/crowdsec-report";
import { getCrowdsecStats } from "@/lib/crowdsec-stats"; import { type CrowdsecDailyStat, getCrowdsecStats } from "@/lib/crowdsec-stats";
import { db, WebsiteSetting } from "@/lib/db"; import { db, WebsiteSetting } from "@/lib/db";
import { canAccess, getAdminContext, PERMS } from "@/lib/permissions"; import { canAccess, getAdminContext, PERMS } from "@/lib/permissions";
import { redis } from "@/lib/redis"; import { redis } from "@/lib/redis";
@@ -58,6 +58,40 @@ function seconds(ttlMs: number): string {
return `${Math.floor(s / 3600)}h ${Math.floor((s % 3600) / 60)}m`; return `${Math.floor(s / 3600)}h ${Math.floor((s % 3600) / 60)}m`;
} }
function BarSparkline({ values }: { values: number[] }) {
if (values.length === 0) return null;
const max = Math.max(...values, 1);
return (
<div className="flex items-end gap-[3px] h-10" aria-hidden="true">
{values.map((v, i) => (
<div
// biome-ignore lint/suspicious/noArrayIndexKey: static timeline position is the bar's identity
key={i}
className="w-full rounded-sm bg-primary/60"
style={{
height: `${Math.max(v > 0 ? 6 : 2, (v / max) * 100)}%`,
opacity: v === 0 ? 0.15 : 0.6 + (v / max) * 0.4,
}}
/>
))}
</div>
);
}
/** Merge per-day breakdown maps (categories / reputations) into range totals. */
function mergeBreakdowns(
rows: CrowdsecDailyStat[],
kind: keyof Pick<CrowdsecDailyStat, "categories" | "reputations">,
): Record<string, number> {
const totals: Record<string, number> = {};
for (const row of rows) {
for (const [k, v] of Object.entries(row[kind])) {
totals[k] = (totals[k] ?? 0) + v;
}
}
return totals;
}
export default async function AdminAntiDdosPage() { export default async function AdminAntiDdosPage() {
const { session, permissions } = await getAdminContext(); const { session, permissions } = await getAdminContext();
if (!canAccess(permissions, PERMS.SETTINGS_VIEW, session.user.rank)) { if (!canAccess(permissions, PERMS.SETTINGS_VIEW, session.user.rank)) {
@@ -727,10 +761,42 @@ export default async function AdminAntiDdosPage() {
{crowdsecStats.length > 0 && ( {crowdsecStats.length > 0 && (
<div className="rounded-md border p-3"> <div className="rounded-md border p-3">
<p className="text-xs font-medium mb-2"> <div className="flex flex-wrap items-center justify-between gap-2">
Daily activity (last {crowdsecStats.length} days) <p className="text-xs font-medium">
</p> Daily activity (last {crowdsecStats.length} days)
<div className="max-h-40 overflow-y-auto"> </p>
<a
href="/admin/alerts"
className="text-xs text-muted-foreground underline-offset-2 hover:underline"
>
Ops alert history →
</a>
</div>
<div className="mt-2 grid gap-4 sm:grid-cols-3">
<BarSparkline values={crowdsecStats.map((row) => row.blocks)} />
<div className="col-span-2 flex flex-wrap items-center gap-1.5">
{Object.entries(
mergeBreakdowns(crowdsecStats, "categories"),
).map(([category, count]) => (
<Badge key={category} variant="secondary">
{category} · {count.toLocaleString()}
</Badge>
))}
{Object.entries(
mergeBreakdowns(crowdsecStats, "reputations"),
).map(([reputation, count]) => (
<Badge
key={reputation}
variant={
reputation === "malicious" ? "destructive" : "secondary"
}
>
{reputation} · {count.toLocaleString()}
</Badge>
))}
</div>
</div>
<div className="mt-3 max-h-40 overflow-y-auto">
<table className="w-full text-xs"> <table className="w-full text-xs">
<thead> <thead>
<tr className="text-left text-muted-foreground"> <tr className="text-left text-muted-foreground">
+104
View File
@@ -21,6 +21,7 @@ import { type CrowdsecDailyStat, getCrowdsecStats } from "./crowdsec-stats";
// deliberately never touches Cloudflare, so no Cloudflare surface is stubbed. // deliberately never touches Cloudflare, so no Cloudflare surface is stubbed.
const state = vi.hoisted(() => ({ const state = vi.hoisted(() => ({
map: new Map<string, string>(), map: new Map<string, string>(),
z: new Map<string, Array<[number, string]>>(),
sendAlert: vi.fn(), sendAlert: vi.fn(),
})); }));
@@ -59,6 +60,21 @@ vi.mock("@/lib/redis", () => ({
}, },
expire: async () => 1, expire: async () => 1,
pttl: async () => 60_000, pttl: async () => 60_000,
zadd: async (key: string, score: number, member: string) => {
const list = state.z.get(key) ?? [];
list.push([score, member]);
list.sort((a, b) => a[0] - b[0]);
state.z.set(key, list);
return 1;
},
zremrangebyscore: async (key: string, min: number, max: number) => {
const list = (state.z.get(key) ?? []).filter(
([score]) => score < min || score > max,
);
state.z.set(key, list);
return 1;
},
zcard: async (key: string) => (state.z.get(key) ?? []).length,
}, },
__esModule: true, __esModule: true,
})); }));
@@ -125,6 +141,7 @@ describe("crowdsec-api", () => {
vi.unstubAllGlobals(); vi.unstubAllGlobals();
vi.unstubAllEnvs(); vi.unstubAllEnvs();
state.map.clear(); state.map.clear();
state.z.clear();
state.sendAlert.mockReset(); state.sendAlert.mockReset();
resetCrowdsecCache(); resetCrowdsecCache();
fetchMock = vi.fn(); fetchMock = vi.fn();
@@ -418,6 +435,47 @@ describe("crowdsec-api", () => {
expect(fetchMock).toHaveBeenCalledTimes(1); expect(fetchMock).toHaveBeenCalledTimes(1);
}); });
it("raises a critical ops alert when the CTI key is rejected (403)", 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,
});
await tick();
expect(state.sendAlert).toHaveBeenCalledTimes(1);
const [input] = state.sendAlert.mock.calls[0];
expect(input.type).toBe("ddos");
expect(input.severity).toBe("critical");
expect(input.message).toContain("403");
expect(input.message).toContain("CROWDSEC_API_KEY");
expect(input.context).toMatchObject({ status: 403 });
});
it("raises a warning ops alert when the CTI rate limit is hit (429)", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse({ message: "rate limited" }, 429));
await maybeAutoBlockCrowdsec({
ip: blockIp(),
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.severity).toBe("warning");
expect(input.message).toContain("rate limited");
expect(input.context).toMatchObject({ status: 429 });
});
it("swallows API failures instead of throwing on the hot path", async () => { it("swallows API failures instead of throwing on the hot path", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key"); vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
fetchMock.mockResolvedValue(jsonResponse({ message: "boom" }, 500)); fetchMock.mockResolvedValue(jsonResponse({ message: "boom" }, 500));
@@ -602,6 +660,50 @@ describe("crowdsec-api", () => {
expect(input.context).toMatchObject({ quota: 1 }); expect(input.context).toMatchObject({ quota: 1 });
}); });
it("alerts once when quota becomes available again after exhaustion", 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,
});
// Exhaust the counter.
await maybeAutoBlockCrowdsec({
ip: "198.51.100.2",
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
// The counter is externally reset (new billing day / fresh deployment):
// the next successful reserve should call it out.
const date = new Date().toISOString().slice(0, 10);
state.map.set(`crowdsec:usage:${date}`, "0");
await maybeAutoBlockCrowdsec({
ip: "198.51.100.3",
category: "api",
ttlSeconds: 600,
scoreThreshold: 4,
enabled: true,
});
await tick();
expect(fetchMock).toHaveBeenCalledTimes(2);
expect(state.sendAlert).toHaveBeenCalledTimes(2);
const alerts = state.sendAlert.mock.calls.map(([input]) => input);
expect(alerts[0].severity).toBe("warning");
expect(alerts[1].severity).toBe("info");
expect(alerts[1].message).toContain("available again");
expect(alerts[1].context).toMatchObject({ quota: 1 });
});
it("floods once per cooldown window when blocks burst past the threshold", async () => { it("floods once per cooldown window when blocks burst past the threshold", async () => {
vi.stubEnv("CROWDSEC_API_KEY", "cs_key"); vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
vi.stubEnv("CROWDSEC_ALERT_BLOCK_BURST", "2"); vi.stubEnv("CROWDSEC_ALERT_BLOCK_BURST", "2");
@@ -685,6 +787,8 @@ describe("crowdsec-api", () => {
expect(today?.lookups).toBe(3); expect(today?.lookups).toBe(3);
expect(today?.blocks).toBe(2); expect(today?.blocks).toBe(2);
expect(today?.reportFailures).toBe(0); expect(today?.reportFailures).toBe(0);
expect(today?.categories).toMatchObject({ api: 2 });
expect(today?.reputations).toMatchObject({ malicious: 2 });
expect( expect(
state.map.get( state.map.get(
`crowdsec:stat:lookups:${new Date().toISOString().slice(0, 10)}`, `crowdsec:stat:lookups:${new Date().toISOString().slice(0, 10)}`,
+52 -11
View File
@@ -1,9 +1,13 @@
import "server-only"; import "server-only";
import { randomUUID } from "node:crypto";
import { env } from "@/env"; import { env } from "@/env";
import { raiseCrowdsecAlert } from "@/lib/crowdsec-alerts"; import { raiseCrowdsecAlert } from "@/lib/crowdsec-alerts";
import { reportCrowdsecSignal } from "@/lib/crowdsec-report"; import { reportCrowdsecSignal } from "@/lib/crowdsec-report";
import { bumpCrowdsecStat } from "@/lib/crowdsec-stats"; import {
bumpCrowdsecBreakdownStat,
bumpCrowdsecStat,
} from "@/lib/crowdsec-stats";
import { logger } from "@/lib/logger"; import { logger } from "@/lib/logger";
import { redis } from "@/lib/redis"; import { redis } from "@/lib/redis";
import { UNKNOWN_CLIENT_IP } from "./client-ip"; import { UNKNOWN_CLIENT_IP } from "./client-ip";
@@ -81,7 +85,7 @@ export interface CrowdsecConnectionStatus {
/** Why a CrowdSec-sourced block exists — persisted next to the block key. */ /** Why a CrowdSec-sourced block exists — persisted next to the block key. */
export interface CrowdsecBlockMeta { export interface CrowdsecBlockMeta {
source: typeof CROWDSEC_BLOCK_SOURCE; source: typeof CROWDSEC_BLOCK_SOURCE | "gate";
category: string; category: string;
reputation: CrowdsecReputation | null; reputation: CrowdsecReputation | null;
score: number; score: number;
@@ -117,7 +121,7 @@ const QUOTA_PREFIX = "crowdsec:usage:";
const QUOTA_KEY_TTL_SECONDS = 48 * 3_600; const QUOTA_KEY_TTL_SECONDS = 48 * 3_600;
/** Shared 403/429 pause marker, so every instance respects the backoff. */ /** Shared 403/429 pause marker, so every instance respects the backoff. */
const BACKOFF_KEY = "crowdsec:backoff-until"; const BACKOFF_KEY = "crowdsec:backoff-until";
/** Short-window block burst counter: crowdsec:burst:{unix-5min-bucket}. */ /** Short-window block burst counter: crowdsec:burst:recent (ZSET of timestamps). */
const BURST_PREFIX = "crowdsec:burst:"; const BURST_PREFIX = "crowdsec:burst:";
const BURST_WINDOW_SECONDS = 300; const BURST_WINDOW_SECONDS = 300;
/** In-process verdict cache cap so a flood of distinct IPs cannot grow it forever. */ /** In-process verdict cache cap so a flood of distinct IPs cannot grow it forever. */
@@ -319,6 +323,7 @@ async function acquireLookupLock(ip: string): Promise<boolean> {
let backoffUntil = 0; let backoffUntil = 0;
let quotaWarnedDate: string | null = null; let quotaWarnedDate: string | null = null;
let quotaExhaustedDate: string | null = null;
/** /**
* Next moment (epoch ms) the CTI API may be called again — the max of the * Next moment (epoch ms) the CTI API may be called again — the max of the
@@ -418,6 +423,7 @@ async function reserveQuota(): Promise<boolean> {
// back and refuse: the budget would be spent the very next call // back and refuse: the budget would be spent the very next call
// anyway, so stopping here is both safe and quota-exact. // anyway, so stopping here is both safe and quota-exact.
await redis.decr(key); await redis.decr(key);
quotaExhaustedDate = date;
logger.warn( logger.warn(
"[crowdsec-api] CTI daily quota exhausted — pausing lookups until tomorrow", "[crowdsec-api] CTI daily quota exhausted — pausing lookups until tomorrow",
{ quota }, { quota },
@@ -437,6 +443,17 @@ async function reserveQuota(): Promise<boolean> {
quota, quota,
}); });
} }
if (quotaExhaustedDate) {
// A reserve just succeeded after an exhaustion day (counter was
// reset or the calendar rolled over) — say so, once per cooldown.
void raiseCrowdsecAlert("quota-restored", {
type: "ddos",
severity: "info",
message: `CrowdSec reputation quota is available again (${used} of ${quota} used today) — lookups resumed.`,
context: { used, quota, date },
});
quotaExhaustedDate = null;
}
return true; return true;
} catch { } catch {
// Redis hiccup at a moment we could not count — allow the call rather // Redis hiccup at a moment we could not count — allow the call rather
@@ -452,23 +469,27 @@ function dailyBlockBurstThreshold(): number {
} }
/** /**
* A burst of new community-reputation blocks is usually an automated attack * A burst of new blocks is usually an automated attack wave. Track block
* wave. Count blocks into a rolling 5-minute bucket and alert once per * timestamps in a rolling window (Redis sorted set, 5 minutes) so a burst that
* cooldown window when they cross CROWDSEC_ALERT_BLOCK_BURST. Fire-and-forget. * straddles a bucket boundary is still counted together, and alert once per
* cooldown window when the count crosses CROWDSEC_ALERT_BLOCK_BURST.
* Fire-and-forget.
*/ */
async function trackBlockBurst(): Promise<void> { async function trackBlockBurst(): Promise<void> {
if (!redis) return; if (!redis) return;
const bucket = Math.floor(Date.now() / 1000 / BURST_WINDOW_SECONDS); const now = Date.now();
const key = `${BURST_PREFIX}${bucket}`; const key = `${BURST_PREFIX}recent`;
const threshold = dailyBlockBurstThreshold(); const threshold = dailyBlockBurstThreshold();
try { try {
const count = await redis.incr(key); await redis.zadd(key, now, randomUUID());
await redis.zremrangebyscore(key, 0, now - BURST_WINDOW_SECONDS * 1000);
const count = await redis.zcard(key);
await redis.expire(key, BURST_WINDOW_SECONDS * 2); await redis.expire(key, BURST_WINDOW_SECONDS * 2);
if (count >= threshold) { if (count >= threshold) {
void raiseCrowdsecAlert("block-burst", { void raiseCrowdsecAlert("block-burst", {
type: "ddos", type: "ddos",
severity: "warning", severity: "warning",
message: `CrowdSec community reputation blocked ${count} IPs in the last ${BURST_WINDOW_SECONDS / 60} minutes — likely an automated attack wave.`, message: `Anti-DDoS auto-block created ${count} blocks in the last ${BURST_WINDOW_SECONDS / 60} minutes — likely an automated attack wave.`,
context: { context: {
blocks: count, blocks: count,
windowSeconds: BURST_WINDOW_SECONDS, windowSeconds: BURST_WINDOW_SECONDS,
@@ -531,9 +552,18 @@ export async function lookupCrowdsecVerdict(
return verdict; return verdict;
} }
if (response.status === 403) { if (response.status === 403) {
const detail = await errorDetail(response);
await setBackoff(AUTH_BACKOFF_MS); await setBackoff(AUTH_BACKOFF_MS);
// A rejected key paralyses the whole reputation pipeline — surface
// it once (cooldown-gated) so rotating the key is an ops priority.
void raiseCrowdsecAlert("cti-auth", {
type: "ddos",
severity: "critical",
message: `CrowdSec CTI API key rejected (HTTP 403): ${detail} — reputation lookups are paused for ${Math.round(AUTH_BACKOFF_MS / 60_000)} minutes. Rotate CROWDSEC_API_KEY.`,
context: { status: 403, detail, backoffMs: AUTH_BACKOFF_MS },
});
throw new CrowdsecApiError( throw new CrowdsecApiError(
`CrowdSec API key rejected (HTTP 403): ${await errorDetail(response)}`, `CrowdSec API key rejected (HTTP 403): ${detail}`,
); );
} }
if (response.status === 429) { if (response.status === 429) {
@@ -542,6 +572,12 @@ export async function lookupCrowdsecVerdict(
ip, ip,
backoffMs: RATE_LIMIT_BACKOFF_MS, backoffMs: RATE_LIMIT_BACKOFF_MS,
}); });
void raiseCrowdsecAlert("cti-ratelimit", {
type: "ddos",
severity: "warning",
message: `CrowdSec CTI API rate limited — all instances backed off for ${Math.round(RATE_LIMIT_BACKOFF_MS / 1000)}s.`,
context: { status: 429, backoffMs: RATE_LIMIT_BACKOFF_MS },
});
return null; return null;
} }
if (!response.ok) { if (!response.ok) {
@@ -635,6 +671,10 @@ export async function maybeAutoBlockCrowdsec(input: {
void reportCrowdsecSignal({ ip, category, ttlSeconds, verdict, meta }); void reportCrowdsecSignal({ ip, category, ttlSeconds, verdict, meta });
// Daily histogram + burst detection (cooldown-gated ops alert). // Daily histogram + burst detection (cooldown-gated ops alert).
void bumpCrowdsecStat("blocks"); void bumpCrowdsecStat("blocks");
void bumpCrowdsecBreakdownStat("category", category);
if (verdict.reputation) {
void bumpCrowdsecBreakdownStat("reputation", verdict.reputation);
}
void trackBlockBurst(); void trackBlockBurst();
} catch (error) { } catch (error) {
logger.error("[crowdsec-api] Automatic IP block failed", { logger.error("[crowdsec-api] Automatic IP block failed", {
@@ -744,5 +784,6 @@ export function resetCrowdsecCache(): void {
memoryVerdicts.clear(); memoryVerdicts.clear();
backoffUntil = 0; backoffUntil = 0;
quotaWarnedDate = null; quotaWarnedDate = null;
quotaExhaustedDate = null;
lastVerifyMemory = null; lastVerifyMemory = null;
} }
+37
View File
@@ -323,6 +323,43 @@ describe("crowdsec-report", () => {
expect(state.sendAlert).not.toHaveBeenCalled(); expect(state.sendAlert).not.toHaveBeenCalled();
}); });
it("raises an info alert the first time the channel heals after failures", 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.40"));
await tick();
expect(state.sendAlert).toHaveBeenCalledTimes(1);
expect((await getLastCrowdsecReport())?.ok).toBe(false);
// Channel heals: the first success after a failure is worth a notice.
routeCapi();
await reportCrowdsecSignal(signalInput("198.51.100.41"));
await tick();
const last: CrowdsecReportStatus | null = await getLastCrowdsecReport();
expect(last?.ok).toBe(true);
expect(state.sendAlert).toHaveBeenCalledTimes(2);
const alerts = state.sendAlert.mock.calls.map(([input]) => input);
expect(alerts[0].severity).toBe("warning");
expect(alerts[1].severity).toBe("info");
expect(alerts[1].message).toContain("recovered");
expect(alerts[1].context).toMatchObject({ ip: "198.51.100.41" });
});
it("verifies the watcher channel end to end", async () => { it("verifies the watcher channel end to end", async () => {
routeCapi(); routeCapi();
const status = await verifyCrowdsecReporting(); const status = await verifyCrowdsecReporting();
+15 -1
View File
@@ -348,7 +348,9 @@ async function pushSignal(input: {
body: [ body: [
{ {
machine_id: input.credentials.machineId, machine_id: input.credentials.machineId,
message: "atomcms-next anti-DDoS gate blocked a community-flagged IP", message: input.verdict.reputation
? "atomcms-next anti-DDoS gate blocked a community-flagged IP"
: "atomcms-next anti-DDoS gate blocked a repeat rate-limit offender",
scenario: SCENARIO, scenario: SCENARIO,
scenario_version: SCENARIO_VERSION, scenario_version: SCENARIO_VERSION,
scenario_hash: scenarioHash(), scenario_hash: scenarioHash(),
@@ -390,9 +392,21 @@ async function pushSignal(input: {
`signal push rejected (HTTP ${response.status}): ${await errorDetail(response)}`, `signal push rejected (HTTP ${response.status}): ${await errorDetail(response)}`,
); );
} }
// Was the channel down before this? A first success after failures
// deserves a recovery notice (separate cooldown key from the failure).
const before = await getLastCrowdsecReport();
const status: CrowdsecReportStatus = { ok: true, at: Date.now() }; const status: CrowdsecReportStatus = { ok: true, at: Date.now() };
await setLastCrowdsecReport(status); await setLastCrowdsecReport(status);
void bumpCrowdsecStat("reports"); void bumpCrowdsecStat("reports");
if (before && !before.ok) {
void raiseCrowdsecAlert("report-recovered", {
type: "ddos",
severity: "info",
message:
"CrowdSec signal push recovered — detections are reaching the community again.",
context: { ip: input.ip },
});
}
logger.info("[crowdsec-report] Detection shared with the community", { logger.info("[crowdsec-report] Detection shared with the community", {
ip: input.ip, ip: input.ip,
category: input.category, category: input.category,
+106 -47
View File
@@ -6,7 +6,8 @@ import { redis } from "@/lib/redis";
// //
// Small Redis counters so ops can see whether the reputation pipeline is // Small Redis counters so ops can see whether the reputation pipeline is
// actually doing anything: lookups executed, blocks created, signals pushed, // actually doing anything: lookups executed, blocks created, signals pushed,
// and push failures — all bucketed per UTC calendar day // push failures, and per-block breakdowns (request category / CrowdSec
// reputation) — all bucketed per UTC calendar day
// (crowdsec:stat:{metric}:{YYYY-MM-DD}). Both the CTI client and the signal // (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 // 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. // non-hot paths only, so they never tax the request path.
@@ -17,6 +18,8 @@ export type CrowdsecStatMetric =
| "reports" | "reports"
| "report_fail"; | "report_fail";
export type CrowdsecBreakdownKind = "category" | "reputation";
const STAT_PREFIX = "crowdsec:stat:"; const STAT_PREFIX = "crowdsec:stat:";
const STAT_KEY_TTL_SECONDS = 16 * 24 * 3_600; const STAT_KEY_TTL_SECONDS = 16 * 24 * 3_600;
@@ -28,6 +31,10 @@ function statKey(metric: CrowdsecStatMetric, date: string): string {
return `${STAT_PREFIX}${metric}:${date}`; return `${STAT_PREFIX}${metric}:${date}`;
} }
function breakdownKey(kind: CrowdsecBreakdownKind, date: string): string {
return `${STAT_PREFIX}${kind}:${date}`;
}
/** Count one occurrence of a pipeline event for today. Best effort. */ /** Count one occurrence of a pipeline event for today. Best effort. */
export async function bumpCrowdsecStat( export async function bumpCrowdsecStat(
metric: CrowdsecStatMetric, metric: CrowdsecStatMetric,
@@ -42,6 +49,30 @@ export async function bumpCrowdsecStat(
} }
} }
/**
* Count one block into today's per-category / per-reputation breakdown. Stored
* as a JSON object per day and merged by the admin reader; read-modify-write
* is fine because blocks are rare and the panel is display-only.
*/
export async function bumpCrowdsecBreakdownStat(
kind: CrowdsecBreakdownKind,
value: string,
): Promise<void> {
if (!redis) return;
const key = breakdownKey(kind, statDate());
try {
const raw = await redis.get(key);
const counts: Record<string, number> = raw
? (JSON.parse(raw) as Record<string, number>)
: {};
const label = String(value).slice(0, 64);
counts[label] = (counts[label] ?? 0) + 1;
await redis.set(key, JSON.stringify(counts), "EX", STAT_KEY_TTL_SECONDS);
} catch {
// best effort — a lost breakdown entry only hides a histogram bucket
}
}
export interface CrowdsecDailyStat { export interface CrowdsecDailyStat {
/** UTC calendar day (YYYY-MM-DD). */ /** UTC calendar day (YYYY-MM-DD). */
date: string; date: string;
@@ -49,63 +80,91 @@ export interface CrowdsecDailyStat {
blocks: number; blocks: number;
reports: number; reports: number;
reportFailures: number; reportFailures: number;
/** Blocks per request category (api/pages/auth/global), today-to-date. */
categories: Record<string, number>;
/** Blocks per CrowdSec reputation (only community-sourced ones). */
reputations: Record<string, number>;
}
function blankRow(date: string): CrowdsecDailyStat {
return {
date,
lookups: 0,
blocks: 0,
reports: 0,
reportFailures: 0,
categories: {},
reputations: {},
};
}
function toCount(raw: string | null): number {
const n = Number(raw ?? 0);
return Number.isFinite(n) ? n : 0;
}
function toMap(raw: string | null): Record<string, number> {
if (!raw) return {};
try {
const parsed = JSON.parse(raw) as unknown;
if (parsed && typeof parsed === "object") {
return Object.fromEntries(
Object.entries(parsed as Record<string, unknown>)
.map(([k, v]) => [k, typeof v === "number" ? v : Number(v) || 0])
.filter(([, v]) => Number.isFinite(v)),
);
}
} catch {
// corrupt counter — treat as empty
}
return {};
} }
/** /**
* Read the per-day counters for the last `days` days (oldest first, ending * 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 * with today). Every key is fetched in one parallel burst (6 GETs per day),
* parallel GETs; never throws. * then assembled client-side; never throws.
*/ */
export async function getCrowdsecStats( export async function getCrowdsecStats(
days = 14, days = 14,
): Promise<CrowdsecDailyStat[]> { ): Promise<CrowdsecDailyStat[]> {
const today = statDate(); const dates: string[] = [];
const rows: CrowdsecDailyStat[] = []; for (let i = days - 1; i >= 0; i -= 1) {
dates.push(
new Date(Date.now() - i * 86_400_000).toISOString().slice(0, 10),
);
}
if (!redis) { if (!redis) {
// No shared store — still return a blank timeline for the UI. // No shared store — still return a blank timeline for the UI.
for (let i = days - 1; i >= 0; i -= 1) { return dates.map(blankRow);
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 store = redis;
const date = new Date(Date.now() - i * 86_400_000) return Promise.all(
.toISOString() dates.map(async (date) => {
.slice(0, 10); const [
try { lookups,
const [lookups, blocks, reports, reportFailures] = await Promise.all([ blocks,
redis.get(statKey("lookups", date)), reports,
redis.get(statKey("blocks", date)), reportFailures,
redis.get(statKey("reports", date)), categories,
redis.get(statKey("report_fail", date)), reputations,
] = await Promise.all([
store.get(statKey("lookups", date)).catch(() => null),
store.get(statKey("blocks", date)).catch(() => null),
store.get(statKey("reports", date)).catch(() => null),
store.get(statKey("report_fail", date)).catch(() => null),
store.get(breakdownKey("category", date)).catch(() => null),
store.get(breakdownKey("reputation", date)).catch(() => null),
]); ]);
const num = (raw: string | null): number => { return {
const n = Number(raw ?? 0);
return Number.isFinite(n) ? n : 0;
};
rows.push({
date, date,
lookups: num(lookups), lookups: toCount(lookups),
blocks: num(blocks), blocks: toCount(blocks),
reports: num(reports), reports: toCount(reports),
reportFailures: num(reportFailures), reportFailures: toCount(reportFailures),
}); categories: toMap(categories),
} catch { reputations: toMap(reputations),
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;
} }
+28 -1
View File
@@ -6,7 +6,16 @@ import { enforceDdosRateLimit } from "@/lib/ddos-guard";
// The gate's block escalation (and thus the CrowdSec hook) only runs when // The gate's block escalation (and thus the CrowdSec hook) only runs when
// Redis is reachable, so the integration test drives a small in-memory fake. // Redis is reachable, so the integration test drives a small in-memory fake.
const state = vi.hoisted(() => ({ map: new Map<string, string>() })); const state = vi.hoisted(() => ({
map: new Map<string, string>(),
z: new Map<string, Array<[number, string]>>(),
sendAlert: vi.fn(),
}));
vi.mock("@/lib/services/alert", () => ({
sendAlert: state.sendAlert,
ddosDetected: vi.fn(),
}));
vi.mock("@/lib/redis", () => ({ vi.mock("@/lib/redis", () => ({
redis: { redis: {
@@ -31,8 +40,24 @@ vi.mock("@/lib/redis", () => ({
state.map.set(key, String(next)); state.map.set(key, String(next));
return next; return next;
}, },
expire: async () => 1,
pexpire: async () => 1, pexpire: async () => 1,
pttl: async () => 60_000, pttl: async () => 60_000,
zadd: async (key: string, score: number, member: string) => {
const list = state.z.get(key) ?? [];
list.push([score, member]);
list.sort((a, b) => a[0] - b[0]);
state.z.set(key, list);
return 1;
},
zremrangebyscore: async (key: string, min: number, max: number) => {
const list = (state.z.get(key) ?? []).filter(
([score]) => score < min || score > max,
);
state.z.set(key, list);
return 1;
},
zcard: async (key: string) => (state.z.get(key) ?? []).length,
sadd: async (key: string, member: string) => { sadd: async (key: string, member: string) => {
const members = new Set( const members = new Set(
(state.map.get(key) ?? "").split("\u0001").filter(Boolean), (state.map.get(key) ?? "").split("\u0001").filter(Boolean),
@@ -95,6 +120,8 @@ describe("anti-DDoS automatic CrowdSec blocks", () => {
vi.unstubAllGlobals(); vi.unstubAllGlobals();
vi.unstubAllEnvs(); vi.unstubAllEnvs();
state.map.clear(); state.map.clear();
state.z.clear();
state.sendAlert.mockReset();
resetCrowdsecCache(); resetCrowdsecCache();
invalidateAntiddosConfig(); invalidateAntiddosConfig();
fetchMock = vi.fn(); fetchMock = vi.fn();
+35
View File
@@ -8,6 +8,11 @@ import { resolveClientIp } from "@/lib/client-ip";
import { isCloudflareProxied } from "@/lib/cloudflare"; import { isCloudflareProxied } from "@/lib/cloudflare";
import { maybeAutoBlockCloudflare } from "@/lib/cloudflare-api"; import { maybeAutoBlockCloudflare } from "@/lib/cloudflare-api";
import { maybeAutoBlockCrowdsec } from "@/lib/crowdsec-api"; import { maybeAutoBlockCrowdsec } from "@/lib/crowdsec-api";
import { reportCrowdsecSignal } from "@/lib/crowdsec-report";
import {
bumpCrowdsecBreakdownStat,
bumpCrowdsecStat,
} from "@/lib/crowdsec-stats";
import { classifyDdos, isSuspiciousPath } from "@/lib/ddos"; import { classifyDdos, isSuspiciousPath } from "@/lib/ddos";
import { rateLimit } from "@/lib/rate-limit"; import { rateLimit } from "@/lib/rate-limit";
import { redis } from "@/lib/redis"; import { redis } from "@/lib/redis";
@@ -113,6 +118,36 @@ export async function enforceDdosRateLimit(
const ttl = blockTtlForViolations(violations, config.blockTiers); const ttl = blockTtlForViolations(violations, config.blockTiers);
if (violations >= config.maxViolations) { if (violations >= config.maxViolations) {
await redis.set(blockKey, "1", "EX", ttl); await redis.set(blockKey, "1", "EX", ttl);
// Share the block in the daily activity histogram + breakdown.
void bumpCrowdsecStat("blocks");
void bumpCrowdsecBreakdownStat("category", category);
// Opt-in: also push our OWN detection (not just CrowdSec-
// reputation blocks) into the community blocklist via the same
// CAPI channel. Fire-and-forget; deduped per IP internally.
void reportCrowdsecSignal({
ip,
category,
ttlSeconds: ttl,
verdict: {
ip,
reputation: null,
score: 0,
aggressiveness: 0,
confidence: null,
behaviors: [`gate:${category}`],
falsePositive: false,
checkedAt: Date.now(),
},
meta: {
source: "gate",
category,
reputation: null,
score: 0,
behaviors: [`gate:${category}`],
ttlSeconds: ttl,
blockedAt: Date.now(),
},
});
// Mirror the host-level block to the Cloudflare edge (IP Access // Mirror the host-level block to the Cloudflare edge (IP Access
// Rules) so a repeat offender is shed before it reaches the // Rules) so a repeat offender is shed before it reaches the
// origin. Only when this request demonstrably transited // origin. Only when this request demonstrably transited