feat(security): ops alerts, shared backoff, atomic quota and daily stats for CrowdSec
Gitea Actions Runner Test / test-job (push) Successful in 1s
CI / check (push) Successful in 29s
CI / tests-integration (push) Successful in 1m36s
CI / tests-unit (push) Successful in 1m40s
CI / tests-ui (push) Successful in 2m28s
CI / preflight (push) Skipped
CI / deploy (push) Successful in 2m3s
Gitea Actions Runner Test / test-job (push) Successful in 1s
CI / check (push) Successful in 29s
CI / tests-integration (push) Successful in 1m36s
CI / tests-unit (push) Successful in 1m40s
CI / tests-ui (push) Successful in 2m28s
CI / preflight (push) Skipped
CI / deploy (push) Successful in 2m3s
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.
This commit is contained in:
1 parent
5e4fc9ab59
commit
301edd2c9a
9 files changed
+620
-28
No files matched your search
@@ -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
|
# counter reaches it, reputation lookups pause until tomorrow so a spread
|
||||||
# DDoS cannot silently burn the whole quota. 0 = unlimited.
|
# DDoS cannot silently burn the whole quota. 0 = unlimited.
|
||||||
CROWDSEC_CTI_DAILY_QUOTA=10000
|
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) ---
|
# --- CROWDSEC SIGNAL PUSH (share our blocks back, optional) ---
|
||||||
# Opt-in: pushes blocked IPs + behaviors to the CrowdSec Central API (CAPI) so
|
# Opt-in: pushes blocked IPs + behaviors to the CrowdSec Central API (CAPI) so
|
||||||
|
|||||||
@@ -45,6 +45,7 @@ import {
|
|||||||
crowdsecReportEnabled,
|
crowdsecReportEnabled,
|
||||||
getLastCrowdsecReport,
|
getLastCrowdsecReport,
|
||||||
} from "@/lib/crowdsec-report";
|
} from "@/lib/crowdsec-report";
|
||||||
|
import { 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";
|
||||||
@@ -147,6 +148,7 @@ export default async function AdminAntiDdosPage() {
|
|||||||
const crowdsecUsage = redisOk ? await getCrowdsecQuotaUsage() : null;
|
const crowdsecUsage = redisOk ? await getCrowdsecQuotaUsage() : null;
|
||||||
const reportingEnabled = await crowdsecReportEnabled();
|
const reportingEnabled = await crowdsecReportEnabled();
|
||||||
const lastReport = await getLastCrowdsecReport();
|
const lastReport = await getLastCrowdsecReport();
|
||||||
|
const crowdsecStats = redisOk ? await getCrowdsecStats(14) : [];
|
||||||
|
|
||||||
return (
|
return (
|
||||||
<div className="space-y-6">
|
<div className="space-y-6">
|
||||||
@@ -723,6 +725,62 @@ export default async function AdminAntiDdosPage() {
|
|||||||
</div>
|
</div>
|
||||||
)}
|
)}
|
||||||
|
|
||||||
|
{crowdsecStats.length > 0 && (
|
||||||
|
<div className="rounded-md border p-3">
|
||||||
|
<p className="text-xs font-medium mb-2">
|
||||||
|
Daily activity (last {crowdsecStats.length} days)
|
||||||
|
</p>
|
||||||
|
<div className="max-h-40 overflow-y-auto">
|
||||||
|
<table className="w-full text-xs">
|
||||||
|
<thead>
|
||||||
|
<tr className="text-left text-muted-foreground">
|
||||||
|
<th className="pb-1 pr-2 font-medium">Date</th>
|
||||||
|
<th className="pb-1 pr-2 font-medium text-right">
|
||||||
|
Lookups
|
||||||
|
</th>
|
||||||
|
<th className="pb-1 pr-2 font-medium text-right">
|
||||||
|
Blocks
|
||||||
|
</th>
|
||||||
|
<th className="pb-1 pr-2 font-medium text-right">
|
||||||
|
Reports
|
||||||
|
</th>
|
||||||
|
<th className="pb-1 font-medium text-right">Failures</th>
|
||||||
|
</tr>
|
||||||
|
</thead>
|
||||||
|
<tbody>
|
||||||
|
{crowdsecStats.map((row) => (
|
||||||
|
<tr key={row.date} className="border-t">
|
||||||
|
<td className="py-1 pr-2 text-muted-foreground">
|
||||||
|
{row.date === new Date().toISOString().slice(0, 10)
|
||||||
|
? "Today"
|
||||||
|
: row.date.slice(5)}
|
||||||
|
</td>
|
||||||
|
<td className="py-1 pr-2 text-right">
|
||||||
|
{row.lookups.toLocaleString()}
|
||||||
|
</td>
|
||||||
|
<td className="py-1 pr-2 text-right">
|
||||||
|
{row.blocks.toLocaleString()}
|
||||||
|
</td>
|
||||||
|
<td className="py-1 pr-2 text-right">
|
||||||
|
{row.reports.toLocaleString()}
|
||||||
|
</td>
|
||||||
|
<td className="py-1 text-right">
|
||||||
|
{row.reportFailures > 0 ? (
|
||||||
|
<span className="text-destructive">
|
||||||
|
{row.reportFailures.toLocaleString()}
|
||||||
|
</span>
|
||||||
|
) : (
|
||||||
|
"–"
|
||||||
|
)}
|
||||||
|
</td>
|
||||||
|
</tr>
|
||||||
|
))}
|
||||||
|
</tbody>
|
||||||
|
</table>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
)}
|
||||||
|
|
||||||
<div className="rounded-md border p-3">
|
<div className="rounded-md border p-3">
|
||||||
<p className="text-xs font-medium mb-1">Community signal push</p>
|
<p className="text-xs font-medium mb-1">Community signal push</p>
|
||||||
<div className="flex flex-wrap items-center gap-3">
|
<div className="flex flex-wrap items-center gap-3">
|
||||||
|
|||||||
@@ -182,6 +182,9 @@ const schema = z
|
|||||||
// it, so a spread DDoS can never silently burn the whole quota; 0
|
// it, so a spread DDoS can never silently burn the whole quota; 0
|
||||||
// disables the guard.
|
// disables the guard.
|
||||||
CROWDSEC_CTI_DAILY_QUOTA: z.coerce.number().int().min(0).default(10_000),
|
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
|
// Share our own detections back into the CrowdSec community blocklist
|
||||||
// (signal push over the Central API). Opt-in: flipping this on publicly
|
// (signal push over the Central API). Opt-in: flipping this on publicly
|
||||||
// shares blocked IPs + behaviors, so it defaults to off.
|
// shares blocked IPs + behaviors, so it defaults to off.
|
||||||
|
|||||||
@@ -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<false | Awaited<ReturnType<typeof sendAlert>>> {
|
||||||
|
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;
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -14,11 +14,20 @@ import {
|
|||||||
verdictIsMalicious,
|
verdictIsMalicious,
|
||||||
verifyCrowdsecConnection,
|
verifyCrowdsecConnection,
|
||||||
} from "./crowdsec-api";
|
} from "./crowdsec-api";
|
||||||
|
import { type CrowdsecDailyStat, getCrowdsecStats } from "./crowdsec-stats";
|
||||||
|
|
||||||
// Unit-test the CTI client in isolation: a deterministic in-memory Redis fake
|
// 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
|
// and a silenced logger, so fetch calls count only CrowdSec lookups. CrowdSec
|
||||||
// deliberately never touches Cloudflare, so no Cloudflare surface is stubbed.
|
// deliberately never touches Cloudflare, so no Cloudflare surface is stubbed.
|
||||||
const state = vi.hoisted(() => ({ map: new Map<string, string>() }));
|
const state = vi.hoisted(() => ({
|
||||||
|
map: new Map<string, 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: {
|
||||||
@@ -43,6 +52,11 @@ vi.mock("@/lib/redis", () => ({
|
|||||||
state.map.set(key, String(next));
|
state.map.set(key, String(next));
|
||||||
return 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,
|
expire: async () => 1,
|
||||||
pttl: async () => 60_000,
|
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 {
|
function jsonResponse(body: unknown, status = 200): Response {
|
||||||
return new Response(JSON.stringify(body), {
|
return new Response(JSON.stringify(body), {
|
||||||
status,
|
status,
|
||||||
@@ -109,6 +125,7 @@ describe("crowdsec-api", () => {
|
|||||||
vi.unstubAllGlobals();
|
vi.unstubAllGlobals();
|
||||||
vi.unstubAllEnvs();
|
vi.unstubAllEnvs();
|
||||||
state.map.clear();
|
state.map.clear();
|
||||||
|
state.sendAlert.mockReset();
|
||||||
resetCrowdsecCache();
|
resetCrowdsecCache();
|
||||||
fetchMock = vi.fn();
|
fetchMock = vi.fn();
|
||||||
vi.stubGlobal("fetch", fetchMock);
|
vi.stubGlobal("fetch", fetchMock);
|
||||||
@@ -352,6 +369,34 @@ describe("crowdsec-api", () => {
|
|||||||
expect(fetchMock).toHaveBeenCalledTimes(1);
|
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 () => {
|
it("backs off after a 429 rate limit as well", async () => {
|
||||||
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
|
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
|
||||||
fetchMock.mockResolvedValue(jsonResponse({ message: "rate limited" }, 429));
|
fetchMock.mockResolvedValue(jsonResponse({ message: "rate limited" }, 429));
|
||||||
@@ -525,4 +570,125 @@ describe("crowdsec-api", () => {
|
|||||||
}
|
}
|
||||||
expect(getMemoryVerdictCacheSize()).toBe(2000);
|
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");
|
||||||
|
});
|
||||||
});
|
});
|
||||||
+133
-25
@@ -1,7 +1,9 @@
|
|||||||
import "server-only";
|
import "server-only";
|
||||||
|
|
||||||
import { env } from "@/env";
|
import { env } from "@/env";
|
||||||
|
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 { 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";
|
||||||
@@ -113,6 +115,11 @@ const LAST_VERIFY_KEY = "crowdsec:last-verify";
|
|||||||
const BLOCK_META_PREFIX = "antiddos:block:meta:";
|
const BLOCK_META_PREFIX = "antiddos:block:meta:";
|
||||||
const QUOTA_PREFIX = "crowdsec:usage:";
|
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. */
|
||||||
|
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. */
|
/** In-process verdict cache cap so a flood of distinct IPs cannot grow it forever. */
|
||||||
const MEMORY_VERDICT_CACHE_MAX = 2_000;
|
const MEMORY_VERDICT_CACHE_MAX = 2_000;
|
||||||
/** Warn at this fraction of the daily quota, once per day. */
|
/** Warn at this fraction of the daily quota, once per day. */
|
||||||
@@ -313,6 +320,47 @@ async function acquireLookupLock(ip: string): Promise<boolean> {
|
|||||||
let backoffUntil = 0;
|
let backoffUntil = 0;
|
||||||
let quotaWarnedDate: string | null = null;
|
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<number> {
|
||||||
|
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<void> {
|
||||||
|
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 {
|
function quotaDate(): string {
|
||||||
return new Date().toISOString().slice(0, 10);
|
return new Date().toISOString().slice(0, 10);
|
||||||
}
|
}
|
||||||
@@ -350,33 +398,87 @@ export async function getCrowdsecQuotaUsage(): Promise<CrowdsecQuotaUsage> {
|
|||||||
return { date, used, quota, exhausted: quota > 0 && used >= quota };
|
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<boolean> {
|
async function reserveQuota(): Promise<boolean> {
|
||||||
const quota = dailyQuota();
|
const quota = dailyQuota();
|
||||||
if (quota <= 0) return true;
|
if (quota <= 0) return true;
|
||||||
|
if (!redis) return true; // no shared counter → unlimited best-effort
|
||||||
const date = quotaDate();
|
const date = quotaDate();
|
||||||
const usage = await getCrowdsecQuotaUsage();
|
const key = quotaKey(date);
|
||||||
if (usage.used >= quota) {
|
try {
|
||||||
// Stop consulting the API for the rest of the day: a flood of distinct
|
const used = await redis.incr(key);
|
||||||
// bucket-tripping IPs would otherwise burn every remaining call and
|
await redis.expire(key, QUOTA_KEY_TTL_SECONDS);
|
||||||
// then sit in a 429 storm anyway.
|
if (used > quota) {
|
||||||
return false;
|
// Concurrent reserves nudged us past the ceiling — give the slot
|
||||||
}
|
// back and refuse: the budget would be spent the very next call
|
||||||
if (usage.used >= quota * QUOTA_WARN_RATIO && quotaWarnedDate !== date) {
|
// anyway, so stopping here is both safe and quota-exact.
|
||||||
quotaWarnedDate = date;
|
await redis.decr(key);
|
||||||
logger.warn("[crowdsec-api] CTI daily quota nearing its limit", {
|
logger.warn(
|
||||||
used: usage.used,
|
"[crowdsec-api] CTI daily quota exhausted — pausing lookups until tomorrow",
|
||||||
quota,
|
{ quota },
|
||||||
});
|
);
|
||||||
}
|
void raiseCrowdsecAlert("quota", {
|
||||||
if (redis) {
|
type: "ddos",
|
||||||
try {
|
severity: "warning",
|
||||||
await redis.incr(quotaKey(date));
|
message: `CrowdSec reputation quota exhausted for today (${used} of ${quota} enrichment calls) — lookups are paused until tomorrow.`,
|
||||||
await redis.expire(quotaKey(date), QUOTA_KEY_TTL_SECONDS);
|
context: { used, quota, date },
|
||||||
} catch {
|
});
|
||||||
// best effort — an uncounted call is better than a failed lookup
|
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<void> {
|
||||||
|
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<CrowdsecVerdict | null> {
|
): Promise<CrowdsecVerdict | null> {
|
||||||
if (!crowdsecEnabled()) return null;
|
if (!crowdsecEnabled()) return null;
|
||||||
if (!ip || ip === UNKNOWN_CLIENT_IP) 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);
|
const cached = await readVerdictCache(ip);
|
||||||
if (cached) return cached;
|
if (cached) return cached;
|
||||||
@@ -417,6 +519,9 @@ export async function lookupCrowdsecVerdict(
|
|||||||
}
|
}
|
||||||
|
|
||||||
const response = await crowdsecRequest(`/smoke/${encodeURIComponent(ip)}`);
|
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) {
|
if (response.status === 404) {
|
||||||
// Unknown to the community — cache the negative result so a clean
|
// Unknown to the community — cache the negative result so a clean
|
||||||
@@ -426,13 +531,13 @@ export async function lookupCrowdsecVerdict(
|
|||||||
return verdict;
|
return verdict;
|
||||||
}
|
}
|
||||||
if (response.status === 403) {
|
if (response.status === 403) {
|
||||||
backoffUntil = Date.now() + AUTH_BACKOFF_MS;
|
await setBackoff(AUTH_BACKOFF_MS);
|
||||||
throw new CrowdsecApiError(
|
throw new CrowdsecApiError(
|
||||||
`CrowdSec API key rejected (HTTP 403): ${await errorDetail(response)}`,
|
`CrowdSec API key rejected (HTTP 403): ${await errorDetail(response)}`,
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
if (response.status === 429) {
|
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", {
|
logger.warn("[crowdsec-api] CTI API rate limit hit — backing off", {
|
||||||
ip,
|
ip,
|
||||||
backoffMs: RATE_LIMIT_BACKOFF_MS,
|
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
|
// The gate only ever reads its block key through shared Redis — without it
|
||||||
// there is nowhere durable to record the block.
|
// there is nowhere durable to record the block.
|
||||||
if (!redis) return;
|
if (!redis) return;
|
||||||
if (Date.now() < backoffUntil) return;
|
if (Date.now() < (await getBackoffUntil())) return;
|
||||||
|
|
||||||
try {
|
try {
|
||||||
const verdict = await lookupCrowdsecVerdict(ip);
|
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,
|
// Opt-in community signal push (CAPI), fire-and-forget: never awaited,
|
||||||
// never throws, and internally deduped per IP.
|
// never throws, and internally deduped per IP.
|
||||||
void reportCrowdsecSignal({ ip, category, ttlSeconds, verdict, meta });
|
void reportCrowdsecSignal({ ip, category, ttlSeconds, verdict, meta });
|
||||||
|
// Daily histogram + burst detection (cooldown-gated ops alert).
|
||||||
|
void bumpCrowdsecStat("blocks");
|
||||||
|
void trackBlockBurst();
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
logger.error("[crowdsec-api] Automatic IP block failed", {
|
logger.error("[crowdsec-api] Automatic IP block failed", {
|
||||||
ip,
|
ip,
|
||||||
|
|||||||
@@ -11,7 +11,10 @@ import {
|
|||||||
|
|
||||||
// The signal-push watcher is tested against a deterministic in-memory Redis
|
// 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.
|
// fake (NX lock + token cache) and a mocked fetch that routes the CAPI paths.
|
||||||
const state = vi.hoisted(() => ({ map: new Map<string, string>() }));
|
const state = vi.hoisted(() => ({
|
||||||
|
map: new Map<string, string>(),
|
||||||
|
sendAlert: vi.fn(),
|
||||||
|
}));
|
||||||
|
|
||||||
vi.mock("@/lib/redis", () => ({
|
vi.mock("@/lib/redis", () => ({
|
||||||
redis: {
|
redis: {
|
||||||
@@ -31,6 +34,12 @@ vi.mock("@/lib/redis", () => ({
|
|||||||
for (const key of keys) state.map.delete(key);
|
for (const key of keys) state.map.delete(key);
|
||||||
return keys.length;
|
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,
|
__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 CAPI = "https://capi.example.test/v3";
|
||||||
const MACHINE = "m".repeat(48);
|
const MACHINE = "m".repeat(48);
|
||||||
const PASSWORD = "Strong!1P@ssw0rdStrong!1P@ssw0rd";
|
const PASSWORD = "Strong!1P@ssw0rdStrong!1P@ssw0rd";
|
||||||
@@ -111,6 +127,7 @@ describe("crowdsec-report", () => {
|
|||||||
vi.unstubAllGlobals();
|
vi.unstubAllGlobals();
|
||||||
vi.unstubAllEnvs();
|
vi.unstubAllEnvs();
|
||||||
state.map.clear();
|
state.map.clear();
|
||||||
|
state.sendAlert.mockReset();
|
||||||
resetCrowdsecReportCache();
|
resetCrowdsecReportCache();
|
||||||
fetchMock = vi.fn();
|
fetchMock = vi.fn();
|
||||||
vi.stubGlobal("fetch", fetchMock);
|
vi.stubGlobal("fetch", fetchMock);
|
||||||
@@ -256,12 +273,56 @@ describe("crowdsec-report", () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
await expect(reportCrowdsecSignal(signalInput())).resolves.toBeUndefined();
|
await expect(reportCrowdsecSignal(signalInput())).resolves.toBeUndefined();
|
||||||
await new Promise((resolve) => setTimeout(resolve, 20));
|
await tick();
|
||||||
const last: CrowdsecReportStatus | null = await getLastCrowdsecReport();
|
const last: CrowdsecReportStatus | null = await getLastCrowdsecReport();
|
||||||
expect(last?.ok).toBe(false);
|
expect(last?.ok).toBe(false);
|
||||||
expect(last?.message).toContain("signal push rejected");
|
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 () => {
|
it("verifies the watcher channel end to end", async () => {
|
||||||
routeCapi();
|
routeCapi();
|
||||||
const status = await verifyCrowdsecReporting();
|
const status = await verifyCrowdsecReporting();
|
||||||
|
|||||||
@@ -2,11 +2,13 @@ import "server-only";
|
|||||||
|
|
||||||
import { createHash, randomBytes } from "node:crypto";
|
import { createHash, randomBytes } from "node:crypto";
|
||||||
import { env } from "@/env";
|
import { env } from "@/env";
|
||||||
|
import { raiseCrowdsecAlert } from "@/lib/crowdsec-alerts";
|
||||||
import type {
|
import type {
|
||||||
CrowdsecBlockMeta,
|
CrowdsecBlockMeta,
|
||||||
CrowdsecConnectionStatus,
|
CrowdsecConnectionStatus,
|
||||||
CrowdsecVerdict,
|
CrowdsecVerdict,
|
||||||
} from "@/lib/crowdsec-api";
|
} from "@/lib/crowdsec-api";
|
||||||
|
import { 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";
|
||||||
@@ -390,6 +392,7 @@ async function pushSignal(input: {
|
|||||||
}
|
}
|
||||||
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");
|
||||||
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,
|
||||||
@@ -404,6 +407,19 @@ async function pushSignal(input: {
|
|||||||
message: String(error),
|
message: String(error),
|
||||||
};
|
};
|
||||||
await setLastCrowdsecReport(status);
|
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", {
|
logger.error("[crowdsec-report] Signal push failed", {
|
||||||
ip: input.ip,
|
ip: input.ip,
|
||||||
err: error,
|
err: error,
|
||||||
|
|||||||
@@ -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<void> {
|
||||||
|
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<CrowdsecDailyStat[]> {
|
||||||
|
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;
|
||||||
|
}
|
||||||
Reference in new issue
Block a user