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 && (
+
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;
+}