feat(security): give back to CrowdSec and harden the CTI budget
Gitea Actions Runner Test / test-job (push) Successful in 0s
CI / check (push) Successful in 29s
CI / tests-integration (push) Successful in 1m34s
CI / tests-unit (push) Successful in 1m36s
CI / tests-ui (push) Successful in 2m22s
CI / preflight (push) Skipped
CI / deploy (push) Successful in 1m53s
Gitea Actions Runner Test / test-job (push) Successful in 0s
CI / check (push) Successful in 29s
CI / tests-integration (push) Successful in 1m34s
CI / tests-unit (push) Successful in 1m36s
CI / tests-ui (push) Successful in 2m22s
CI / preflight (push) Skipped
CI / deploy (push) Successful in 1m53s
- Bound the in-process verdict cache (FIFO eviction at 2000 entries) so a
flood of distinct bucket-tripping IPs cannot grow it without limit.
- Record block metadata (reputation, score, behaviors, category, TTL) in
antiddos:block:meta:{ip}, surfaced as the reason in the admin block list;
unban now also clears the metadata and report locks.
- Track daily CTI enrichment usage in Redis (crowdsec:usage:{date}); warn
once at 80% and pause lookups until tomorrow at CROWDSEC_CTI_DAILY_QUOTA
(default 10000, 0 = unlimited) so a via-spread DDoS cannot burn the plan.
- Add opt-in signal push to the CrowdSec community (CAPI watcher): stable
auto-generated 48-char machine_id/password pair persisted in Redis (or via
env), one-time registration, cached JWT login, optional Console enrollment,
and POST /v3/signals with a ban decision, deduped per IP. Never throws and
reports last status to the admin panel with a verify action.
- Admin page: quota usage bar, reporting status/verify channel, and CrowdSec
block reasons in the active-blocks list.
This commit is contained in:
1 parent
ee25545b7f
commit
5e4fc9ab59
8 files changed
+1323
-37
No files matched your search
@@ -89,6 +89,25 @@ CROWDSEC_BLOCK_SCORE=4
|
|||||||
CROWDSEC_BLOCK_TTL_SECONDS=86400
|
CROWDSEC_BLOCK_TTL_SECONDS=86400
|
||||||
# Endpoint — override only for tests/staging.
|
# Endpoint — override only for tests/staging.
|
||||||
CROWDSEC_CTI_BASE_URL=https://cti.api.crowdsec.net/v2
|
CROWDSEC_CTI_BASE_URL=https://cti.api.crowdsec.net/v2
|
||||||
|
# Daily enrichment-call ceiling (freemium plan ≈ 10k/day). Once today's
|
||||||
|
# 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
|
||||||
|
|
||||||
|
# --- CROWDSEC SIGNAL PUSH (share our blocks back, optional) ---
|
||||||
|
# Opt-in: pushes blocked IPs + behaviors to the CrowdSec Central API (CAPI) so
|
||||||
|
# the community blocklist protects other members too. Set to "true" to enable.
|
||||||
|
# Requires watcher credentials — either set both CROWDSEC_REPORT_MACHINE_ID
|
||||||
|
# (48 chars, [A-Za-z0-9]) and CROWDSEC_REPORT_PASSWORD now, or leave them
|
||||||
|
# unset and let the app generate a stable pair persisted in Redis automatically.
|
||||||
|
CROWDSEC_REPORT_ENABLED=false
|
||||||
|
CROWDSEC_REPORT_MACHINE_ID=
|
||||||
|
CROWDSEC_REPORT_PASSWORD=
|
||||||
|
# Optional: attachment key from https://app.crowdsec.net → Console settings —
|
||||||
|
# links our watcher to your account so pushed signals show up there.
|
||||||
|
CROWDSEC_REPORT_ENROLL_KEY=
|
||||||
|
# Central API base — override only for tests/staging.
|
||||||
|
CROWDSEC_CAPI_BASE_URL=https://api.crowdsec.net/v3
|
||||||
|
|
||||||
# --- PATHS ---
|
# --- PATHS ---
|
||||||
BADGE_UPLOAD_DIR=./public/assets/images/badges
|
BADGE_UPLOAD_DIR=./public/assets/images/badges
|
||||||
|
|||||||
@@ -19,6 +19,10 @@ import {
|
|||||||
setLastCrowdsecVerify,
|
setLastCrowdsecVerify,
|
||||||
verifyCrowdsecConnection,
|
verifyCrowdsecConnection,
|
||||||
} from "@/lib/crowdsec-api";
|
} from "@/lib/crowdsec-api";
|
||||||
|
import {
|
||||||
|
setLastCrowdsecReport,
|
||||||
|
verifyCrowdsecReporting,
|
||||||
|
} from "@/lib/crowdsec-report";
|
||||||
import { db, WebsiteSetting } from "@/lib/db";
|
import { db, WebsiteSetting } from "@/lib/db";
|
||||||
import { logger } from "@/lib/logger";
|
import { logger } from "@/lib/logger";
|
||||||
import { PERMS } from "@/lib/permissions";
|
import { PERMS } from "@/lib/permissions";
|
||||||
@@ -223,6 +227,8 @@ export async function unbanAntiddosIp(formData: FormData): Promise<void> {
|
|||||||
if (redis) {
|
if (redis) {
|
||||||
await Promise.all([
|
await Promise.all([
|
||||||
redis.del(`antiddos:block:${ip}`),
|
redis.del(`antiddos:block:${ip}`),
|
||||||
|
redis.del(`antiddos:block:meta:${ip}`),
|
||||||
|
redis.del(`crowdsec:report:${ip}`),
|
||||||
redis.del(`antiddos:v:${ip}`),
|
redis.del(`antiddos:v:${ip}`),
|
||||||
]);
|
]);
|
||||||
}
|
}
|
||||||
@@ -272,6 +278,19 @@ export async function verifyCrowdsecConfiguration(): Promise<void> {
|
|||||||
revalidatePath("/admin/devops/antiddos");
|
revalidatePath("/admin/devops/antiddos");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Test the CrowdSec signal-push (CAPI watcher) channel. */
|
||||||
|
export async function verifyCrowdsecReportingConfiguration(): Promise<void> {
|
||||||
|
const staff = await requirePermission(PERMS.SETTINGS_VIEW);
|
||||||
|
const status = await verifyCrowdsecReporting();
|
||||||
|
await setLastCrowdsecReport(status);
|
||||||
|
logger.info("CrowdSec reporting configuration verified", {
|
||||||
|
staff: staff.username,
|
||||||
|
ok: status.ok,
|
||||||
|
message: status.message,
|
||||||
|
});
|
||||||
|
revalidatePath("/admin/devops/antiddos");
|
||||||
|
}
|
||||||
|
|
||||||
/** Test the configured Cloudflare API credentials against the zone. */
|
/** Test the configured Cloudflare API credentials against the zone. */
|
||||||
export async function verifyCloudflareConfiguration(): Promise<void> {
|
export async function verifyCloudflareConfiguration(): Promise<void> {
|
||||||
const staff = await requirePermission(PERMS.SETTINGS_VIEW);
|
const staff = await requirePermission(PERMS.SETTINGS_VIEW);
|
||||||
|
|||||||
@@ -15,6 +15,7 @@ import {
|
|||||||
unbanAntiddosIp,
|
unbanAntiddosIp,
|
||||||
verifyCloudflareConfiguration,
|
verifyCloudflareConfiguration,
|
||||||
verifyCrowdsecConfiguration,
|
verifyCrowdsecConfiguration,
|
||||||
|
verifyCrowdsecReportingConfiguration,
|
||||||
} from "@/actions/admin-antiddos";
|
} from "@/actions/admin-antiddos";
|
||||||
import { Badge } from "@/components/ui/badge";
|
import { Badge } from "@/components/ui/badge";
|
||||||
import { Button } from "@/components/ui/button";
|
import { Button } from "@/components/ui/button";
|
||||||
@@ -34,9 +35,16 @@ import {
|
|||||||
} from "@/lib/cloudflare-api";
|
} from "@/lib/cloudflare-api";
|
||||||
import {
|
import {
|
||||||
CROWDSEC_BLOCK_SOURCE,
|
CROWDSEC_BLOCK_SOURCE,
|
||||||
|
type CrowdsecBlockMeta,
|
||||||
crowdsecEnabled,
|
crowdsecEnabled,
|
||||||
|
getCrowdsecBlockMeta,
|
||||||
|
getCrowdsecQuotaUsage,
|
||||||
getLastCrowdsecVerify,
|
getLastCrowdsecVerify,
|
||||||
} from "@/lib/crowdsec-api";
|
} from "@/lib/crowdsec-api";
|
||||||
|
import {
|
||||||
|
crowdsecReportEnabled,
|
||||||
|
getLastCrowdsecReport,
|
||||||
|
} from "@/lib/crowdsec-report";
|
||||||
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";
|
||||||
@@ -76,6 +84,7 @@ export default async function AdminAntiDdosPage() {
|
|||||||
ttlMs: number;
|
ttlMs: number;
|
||||||
count: number;
|
count: number;
|
||||||
source: "gate" | "crowdsec";
|
source: "gate" | "crowdsec";
|
||||||
|
meta: CrowdsecBlockMeta | null;
|
||||||
}[] = [];
|
}[] = [];
|
||||||
let redisOk = false;
|
let redisOk = false;
|
||||||
const rateStore = redis;
|
const rateStore = redis;
|
||||||
@@ -100,15 +109,19 @@ export default async function AdminAntiDdosPage() {
|
|||||||
rateStore.pttl(key),
|
rateStore.pttl(key),
|
||||||
rateStore.get(key),
|
rateStore.get(key),
|
||||||
]);
|
]);
|
||||||
|
const ip = key.replace("antiddos:block:", "");
|
||||||
|
const source =
|
||||||
|
value === CROWDSEC_BLOCK_SOURCE
|
||||||
|
? ("crowdsec" as const)
|
||||||
|
: ("gate" as const);
|
||||||
return {
|
return {
|
||||||
ip: key.replace("antiddos:block:", ""),
|
ip,
|
||||||
ttlMs: ttlMs > 0 ? ttlMs : 0,
|
ttlMs: ttlMs > 0 ? ttlMs : 0,
|
||||||
count: violationCounts.get(key.replace("antiddos:block:", "")) ?? 0,
|
count: violationCounts.get(ip) ?? 0,
|
||||||
// The gate writes "1"; "crowdsec" marks a community-reputation block.
|
// The gate writes "1"; "crowdsec" marks a community-reputation block.
|
||||||
source:
|
source,
|
||||||
value === CROWDSEC_BLOCK_SOURCE
|
// Why CrowdSec blocked this IP, when the meta was recorded.
|
||||||
? ("crowdsec" as const)
|
meta: source === "crowdsec" ? await getCrowdsecBlockMeta(ip) : null,
|
||||||
: ("gate" as const),
|
|
||||||
};
|
};
|
||||||
}),
|
}),
|
||||||
);
|
);
|
||||||
@@ -131,6 +144,9 @@ export default async function AdminAntiDdosPage() {
|
|||||||
const lastVerify = await getLastCloudflareVerify();
|
const lastVerify = await getLastCloudflareVerify();
|
||||||
const crowdsecConfigured = crowdsecEnabled();
|
const crowdsecConfigured = crowdsecEnabled();
|
||||||
const lastCrowdsecVerify = await getLastCrowdsecVerify();
|
const lastCrowdsecVerify = await getLastCrowdsecVerify();
|
||||||
|
const crowdsecUsage = redisOk ? await getCrowdsecQuotaUsage() : null;
|
||||||
|
const reportingEnabled = await crowdsecReportEnabled();
|
||||||
|
const lastReport = await getLastCrowdsecReport();
|
||||||
|
|
||||||
return (
|
return (
|
||||||
<div className="space-y-6">
|
<div className="space-y-6">
|
||||||
@@ -195,7 +211,9 @@ export default async function AdminAntiDdosPage() {
|
|||||||
{crowdsecConfigured ? "Connected" : "Not configured"}
|
{crowdsecConfigured ? "Connected" : "Not configured"}
|
||||||
</Badge>
|
</Badge>
|
||||||
<p className="text-xs text-muted-foreground mt-1">
|
<p className="text-xs text-muted-foreground mt-1">
|
||||||
Community reputation auto-block
|
{crowdsecUsage && crowdsecUsage.quota > 0
|
||||||
|
? `${crowdsecUsage.used.toLocaleString()} / ${crowdsecUsage.quota.toLocaleString()} CTI calls today${crowdsecUsage.exhausted ? " (paused)" : ""}`
|
||||||
|
: "Community reputation auto-block"}
|
||||||
</p>
|
</p>
|
||||||
</CardContent>
|
</CardContent>
|
||||||
</Card>
|
</Card>
|
||||||
@@ -485,28 +503,44 @@ export default async function AdminAntiDdosPage() {
|
|||||||
</p>
|
</p>
|
||||||
) : (
|
) : (
|
||||||
<div className="space-y-2">
|
<div className="space-y-2">
|
||||||
{blocks.map((b) => (
|
{blocks.map((b) => {
|
||||||
<div
|
const behaviorLabel =
|
||||||
key={b.ip}
|
b.meta && b.meta.behaviors.length > 0
|
||||||
className="flex items-center justify-between gap-2 rounded-md border p-2 text-sm"
|
? b.meta.behaviors.join(", ")
|
||||||
>
|
: null;
|
||||||
<span className="font-mono">{b.ip}</span>
|
return (
|
||||||
<span className="flex items-center gap-2 text-xs text-muted-foreground">
|
<div
|
||||||
{b.source === "crowdsec" ? (
|
key={b.ip}
|
||||||
<Badge variant="default">CrowdSec</Badge>
|
className="flex flex-wrap items-center justify-between gap-2 rounded-md border p-2 text-sm"
|
||||||
) : (
|
>
|
||||||
<Badge variant="secondary">Gate</Badge>
|
<span className="font-mono">{b.ip}</span>
|
||||||
)}
|
<span className="flex items-center gap-2 text-xs text-muted-foreground">
|
||||||
TTL {seconds(b.ttlMs)} · violations {b.count}
|
{b.source === "crowdsec" ? (
|
||||||
</span>
|
<Badge variant="default">CrowdSec</Badge>
|
||||||
<form action={unbanAntiddosIp}>
|
) : (
|
||||||
<input type="hidden" name="ip" value={b.ip} />
|
<Badge variant="secondary">Gate</Badge>
|
||||||
<Button type="submit" size="sm" variant="outline">
|
)}
|
||||||
Unban
|
TTL {seconds(b.ttlMs)} · violations {b.count}
|
||||||
</Button>
|
{b.meta && (
|
||||||
</form>
|
<span
|
||||||
</div>
|
className="max-w-xs truncate"
|
||||||
))}
|
title={`${b.meta.reputation ?? "unknown"} · score ${b.meta.score} · ${b.meta.category}${behaviorLabel ? ` · ${behaviorLabel}` : ""}`}
|
||||||
|
>
|
||||||
|
{b.meta.reputation ?? "unknown"} · score{" "}
|
||||||
|
{b.meta.score} · {b.meta.category}
|
||||||
|
{behaviorLabel ? ` · ${behaviorLabel}` : ""}
|
||||||
|
</span>
|
||||||
|
)}
|
||||||
|
</span>
|
||||||
|
<form action={unbanAntiddosIp}>
|
||||||
|
<input type="hidden" name="ip" value={b.ip} />
|
||||||
|
<Button type="submit" size="sm" variant="outline">
|
||||||
|
Unban
|
||||||
|
</Button>
|
||||||
|
</form>
|
||||||
|
</div>
|
||||||
|
);
|
||||||
|
})}
|
||||||
</div>
|
</div>
|
||||||
)}
|
)}
|
||||||
</CardContent>
|
</CardContent>
|
||||||
@@ -642,9 +676,100 @@ export default async function AdminAntiDdosPage() {
|
|||||||
Verdicts are looked up lazily for IPs that already triggered a
|
Verdicts are looked up lazily for IPs that already triggered a
|
||||||
rate bucket (never on the per-request hot path), cached for an
|
rate bucket (never on the per-request hot path), cached for an
|
||||||
hour, and blocked IPs show a{" "}
|
hour, and blocked IPs show a{" "}
|
||||||
<Badge variant="default">CrowdSec</Badge> badge in the list above.
|
<Badge variant="default">CrowdSec</Badge> badge in the list above
|
||||||
|
with the community reasoning (reputation, score, behaviors).
|
||||||
</p>
|
</p>
|
||||||
)}
|
)}
|
||||||
|
|
||||||
|
{crowdsecUsage && (
|
||||||
|
<div className="rounded-md border p-3">
|
||||||
|
<p className="text-xs font-medium mb-1">
|
||||||
|
Reputation lookups today
|
||||||
|
</p>
|
||||||
|
{crowdsecUsage.quota > 0 ? (
|
||||||
|
<>
|
||||||
|
<div className="flex items-center gap-2">
|
||||||
|
<div className="h-2 flex-1 overflow-hidden rounded-full bg-muted">
|
||||||
|
<div
|
||||||
|
className="h-full rounded-full"
|
||||||
|
style={{
|
||||||
|
width: `${Math.min(100, (crowdsecUsage.used / crowdsecUsage.quota) * 100)}%`,
|
||||||
|
background: crowdsecUsage.exhausted
|
||||||
|
? "var(--color-destructive)"
|
||||||
|
: crowdsecUsage.used >= crowdsecUsage.quota * 0.8
|
||||||
|
? "var(--admin-accent)"
|
||||||
|
: "var(--color-primary)",
|
||||||
|
}}
|
||||||
|
/>
|
||||||
|
</div>
|
||||||
|
<span
|
||||||
|
className={`text-xs ${crowdsecUsage.exhausted ? "text-destructive" : "text-muted-foreground"}`}
|
||||||
|
>
|
||||||
|
{crowdsecUsage.used.toLocaleString()} /{" "}
|
||||||
|
{crowdsecUsage.quota.toLocaleString()}
|
||||||
|
</span>
|
||||||
|
</div>
|
||||||
|
<p className="text-xs text-muted-foreground mt-1">
|
||||||
|
{crowdsecUsage.exhausted
|
||||||
|
? "Quota spent for today — reputation lookups are paused until tomorrow (admin via CROWDSEC_CTI_DAILY_QUOTA)."
|
||||||
|
: "Visible in the env via CROWDSEC_CTI_DAILY_QUOTA (0 = unlimited). Lookups pause at the ceiling to protect the plan."}
|
||||||
|
</p>
|
||||||
|
</>
|
||||||
|
) : (
|
||||||
|
<p className="text-xs text-muted-foreground">
|
||||||
|
Tracking disabled (CROWDSEC_CTI_DAILY_QUOTA = 0 / unlimited).
|
||||||
|
</p>
|
||||||
|
)}
|
||||||
|
</div>
|
||||||
|
)}
|
||||||
|
|
||||||
|
<div className="rounded-md border p-3">
|
||||||
|
<p className="text-xs font-medium mb-1">Community signal push</p>
|
||||||
|
<div className="flex flex-wrap items-center gap-3">
|
||||||
|
<Badge variant={reportingEnabled ? "default" : "secondary"}>
|
||||||
|
{reportingEnabled ? "Enabled" : "Off"}
|
||||||
|
</Badge>
|
||||||
|
{!reportingEnabled && (
|
||||||
|
<p className="text-xs text-muted-foreground">
|
||||||
|
Set{" "}
|
||||||
|
<span className="font-mono">
|
||||||
|
CROWDSEC_REPORT_ENABLED=true
|
||||||
|
</span>{" "}
|
||||||
|
to share blocked IPs back into the CrowdSec community
|
||||||
|
blocklist. Watcher credentials are auto-generated and
|
||||||
|
persisted in Redis.
|
||||||
|
</p>
|
||||||
|
)}
|
||||||
|
{reportingEnabled && (
|
||||||
|
<p className="text-xs text-muted-foreground">
|
||||||
|
Blocked IPs are pushed to the Central API (deduped per IP) so
|
||||||
|
the community blocklist protects other members too.
|
||||||
|
</p>
|
||||||
|
)}
|
||||||
|
<form action={verifyCrowdsecReportingConfiguration}>
|
||||||
|
<Button
|
||||||
|
type="submit"
|
||||||
|
size="sm"
|
||||||
|
variant="outline"
|
||||||
|
disabled={!reportingEnabled}
|
||||||
|
>
|
||||||
|
Verify channel
|
||||||
|
</Button>
|
||||||
|
</form>
|
||||||
|
</div>
|
||||||
|
{lastReport && (
|
||||||
|
<p className="text-xs mt-2">
|
||||||
|
<Badge variant={lastReport.ok ? "default" : "destructive"}>
|
||||||
|
{lastReport.ok ? "Push healthy" : "Push failed"}
|
||||||
|
</Badge>
|
||||||
|
<span className="ml-2 text-muted-foreground">
|
||||||
|
{lastReport.ok
|
||||||
|
? `Last signal accepted ${new Date(lastReport.at).toLocaleString()}`
|
||||||
|
: `${lastReport.message ?? "unknown"} (${new Date(lastReport.at).toLocaleString()})`}
|
||||||
|
</span>
|
||||||
|
</p>
|
||||||
|
)}
|
||||||
|
</div>
|
||||||
</CardContent>
|
</CardContent>
|
||||||
</Card>
|
</Card>
|
||||||
|
|
||||||
|
|||||||
+25
@@ -177,6 +177,31 @@ const schema = z
|
|||||||
.int()
|
.int()
|
||||||
.positive()
|
.positive()
|
||||||
.default(86_400),
|
.default(86_400),
|
||||||
|
// Daily CTI enrichment quota guard (freemium plan ≈ 10k lookups/day).
|
||||||
|
// The gate stops consulting the API once the counter for today exceeds
|
||||||
|
// 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),
|
||||||
|
// 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.
|
||||||
|
CROWDSEC_REPORT_ENABLED: z
|
||||||
|
.string()
|
||||||
|
.optional()
|
||||||
|
.transform((value) => value === "true" || value === "1"),
|
||||||
|
// Central API (CAPI) base endpoint; overridden for tests/staging.
|
||||||
|
CROWDSEC_CAPI_BASE_URL: z
|
||||||
|
.string()
|
||||||
|
.url()
|
||||||
|
.default("https://api.crowdsec.net/v3"),
|
||||||
|
// Watcher credentials for signal push. When omitted, a stable pair is
|
||||||
|
// generated once and persisted in Redis (48-char alnum machine id,
|
||||||
|
// per the CAPI schema).
|
||||||
|
CROWDSEC_REPORT_MACHINE_ID: z.string().optional(),
|
||||||
|
CROWDSEC_REPORT_PASSWORD: z.string().optional(),
|
||||||
|
// Optional attachment key from the CrowdSec Console — links our
|
||||||
|
// watcher to your account so pushed signals show up there.
|
||||||
|
CROWDSEC_REPORT_ENROLL_KEY: z.string().optional(),
|
||||||
})
|
})
|
||||||
.superRefine((data, ctx) => {
|
.superRefine((data, ctx) => {
|
||||||
if (data.NODE_ENV !== "production") return;
|
if (data.NODE_ENV !== "production") return;
|
||||||
|
|||||||
@@ -3,7 +3,11 @@ import {
|
|||||||
type CrowdsecVerdict,
|
type CrowdsecVerdict,
|
||||||
crowdsecEnabled,
|
crowdsecEnabled,
|
||||||
getCrowdsecApiConfig,
|
getCrowdsecApiConfig,
|
||||||
|
getCrowdsecBlockMeta,
|
||||||
|
getCrowdsecQuotaUsage,
|
||||||
getLastCrowdsecVerify,
|
getLastCrowdsecVerify,
|
||||||
|
getMemoryVerdictCacheSize,
|
||||||
|
lookupCrowdsecVerdict,
|
||||||
maybeAutoBlockCrowdsec,
|
maybeAutoBlockCrowdsec,
|
||||||
resetCrowdsecCache,
|
resetCrowdsecCache,
|
||||||
setLastCrowdsecVerify,
|
setLastCrowdsecVerify,
|
||||||
@@ -34,6 +38,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,
|
||||||
pttl: async () => 60_000,
|
pttl: async () => 60_000,
|
||||||
},
|
},
|
||||||
__esModule: true,
|
__esModule: true,
|
||||||
@@ -44,6 +54,7 @@ vi.mock("@/lib/logger", () => ({
|
|||||||
info: vi.fn(),
|
info: vi.fn(),
|
||||||
warn: vi.fn(),
|
warn: vi.fn(),
|
||||||
error: vi.fn(),
|
error: vi.fn(),
|
||||||
|
debug: vi.fn(),
|
||||||
},
|
},
|
||||||
}));
|
}));
|
||||||
|
|
||||||
@@ -419,4 +430,99 @@ describe("crowdsec-api", () => {
|
|||||||
status,
|
status,
|
||||||
);
|
);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("records why it blocked an IP next to the gate key", async () => {
|
||||||
|
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
|
||||||
|
fetchMock.mockResolvedValue(jsonResponse(maliciousItem(blockIp())));
|
||||||
|
|
||||||
|
await maybeAutoBlockCrowdsec({
|
||||||
|
ip: blockIp(),
|
||||||
|
category: "api",
|
||||||
|
ttlSeconds: 86_400,
|
||||||
|
scoreThreshold: 4,
|
||||||
|
enabled: true,
|
||||||
|
});
|
||||||
|
|
||||||
|
const meta = await getCrowdsecBlockMeta(blockIp());
|
||||||
|
expect(meta).not.toBeNull();
|
||||||
|
expect(meta?.source).toBe("crowdsec");
|
||||||
|
expect(meta?.reputation).toBe("malicious");
|
||||||
|
expect(meta?.score).toBe(5);
|
||||||
|
expect(meta?.behaviors).toEqual(["http:bruteforce", "http:scan"]);
|
||||||
|
expect(meta?.category).toBe("api");
|
||||||
|
expect(meta?.ttlSeconds).toBe(86_400);
|
||||||
|
expect(meta?.blockedAt).toBeGreaterThan(0);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("stops consulting the API once today's quota is spent", async () => {
|
||||||
|
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
|
||||||
|
vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "2");
|
||||||
|
// Fresh Response per call — a consumed body must never be re-parsed.
|
||||||
|
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,
|
||||||
|
});
|
||||||
|
expect(fetchMock).toHaveBeenCalledTimes(2);
|
||||||
|
expect(state.map.get(`antiddos:block:198.51.100.2`)).toBe("crowdsec");
|
||||||
|
|
||||||
|
// Third bucket-tripping IP arrives after the quota counter hit 2.
|
||||||
|
await maybeAutoBlockCrowdsec({
|
||||||
|
ip: "198.51.100.3",
|
||||||
|
category: "api",
|
||||||
|
ttlSeconds: 600,
|
||||||
|
scoreThreshold: 4,
|
||||||
|
enabled: true,
|
||||||
|
});
|
||||||
|
expect(fetchMock).toHaveBeenCalledTimes(2);
|
||||||
|
expect(state.map.has(`antiddos:block:198.51.100.3`)).toBe(false);
|
||||||
|
|
||||||
|
const usage = await getCrowdsecQuotaUsage();
|
||||||
|
expect(usage.quota).toBe(2);
|
||||||
|
expect(usage.used).toBe(2);
|
||||||
|
expect(usage.exhausted).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("exposes today's quota usage for the admin panel", async () => {
|
||||||
|
vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "10000");
|
||||||
|
const before = await getCrowdsecQuotaUsage();
|
||||||
|
expect(before.quota).toBe(10000);
|
||||||
|
expect(before.used).toBe(0);
|
||||||
|
expect(before.exhausted).toBe(false);
|
||||||
|
expect(before.date).toMatch(/^\d{4}-\d{2}-\d{2}$/);
|
||||||
|
|
||||||
|
state.map.set(`crowdsec:usage:${before.date}`, "9876");
|
||||||
|
const after = await getCrowdsecQuotaUsage();
|
||||||
|
expect(after.used).toBe(9876);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("caps the in-process verdict cache so it cannot grow forever", async () => {
|
||||||
|
vi.stubEnv("CROWDSEC_API_KEY", "cs_key");
|
||||||
|
vi.stubEnv("CROWDSEC_CTI_DAILY_QUOTA", "0");
|
||||||
|
fetchMock.mockImplementation((url: string | URL) =>
|
||||||
|
Promise.resolve(
|
||||||
|
jsonResponse(maliciousItem(String(url).split("/").pop() ?? "ip")),
|
||||||
|
),
|
||||||
|
);
|
||||||
|
|
||||||
|
// One lookup per distinct IP (never cached before), exceeding the cap —
|
||||||
|
// the oldest entries are evicted first, so the cache stays bounded.
|
||||||
|
for (let i = 0; i < 2100; i += 1) {
|
||||||
|
await lookupCrowdsecVerdict(`198.51.100.${i}`);
|
||||||
|
}
|
||||||
|
expect(getMemoryVerdictCacheSize()).toBe(2000);
|
||||||
|
});
|
||||||
});
|
});
|
||||||
+178
-7
@@ -1,6 +1,7 @@
|
|||||||
import "server-only";
|
import "server-only";
|
||||||
|
|
||||||
import { env } from "@/env";
|
import { env } from "@/env";
|
||||||
|
import { reportCrowdsecSignal } from "@/lib/crowdsec-report";
|
||||||
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";
|
||||||
@@ -28,9 +29,15 @@ import { UNKNOWN_CLIENT_IP } from "./client-ip";
|
|||||||
* Credentials come from env only (`CROWDSEC_API_KEY`) and are never written
|
* Credentials come from env only (`CROWDSEC_API_KEY`) and are never written
|
||||||
* into the admin-visible config — same contract as the Cloudflare token.
|
* into the admin-visible config — same contract as the Cloudflare token.
|
||||||
*
|
*
|
||||||
* Note: sharing our own blocks back into the community (signal push) is not
|
* Quota guard: every enrichment call counts against a per-day Redis counter so
|
||||||
* part of this module — CrowdSec's report channel requires a full Security
|
* a spread DDoS (many distinct IPs tripping buckets) can exhaust the day's
|
||||||
* Engine / CAPI machine enrollment, not a CTI API key.
|
* freemium quota only until the configured ceiling, after which lookups pause
|
||||||
|
* until tomorrow instead of hammering a 429 wall.
|
||||||
|
*
|
||||||
|
* Sharing detections back: after a block is created this module fires the
|
||||||
|
* signal push in `@/lib/crowdsec-report` (Central API watcher login + POST
|
||||||
|
* /signals), strictly opt-in via CROWDSEC_REPORT_ENABLED and always
|
||||||
|
* fire-and-forget.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
export class CrowdsecApiError extends Error {}
|
export class CrowdsecApiError extends Error {}
|
||||||
@@ -70,6 +77,27 @@ export interface CrowdsecConnectionStatus {
|
|||||||
at: number;
|
at: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Why a CrowdSec-sourced block exists — persisted next to the block key. */
|
||||||
|
export interface CrowdsecBlockMeta {
|
||||||
|
source: typeof CROWDSEC_BLOCK_SOURCE;
|
||||||
|
category: string;
|
||||||
|
reputation: CrowdsecReputation | null;
|
||||||
|
score: number;
|
||||||
|
behaviors: string[];
|
||||||
|
ttlSeconds: number;
|
||||||
|
blockedAt: number;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Daily CTI usage counter as shown in the admin panel. */
|
||||||
|
export interface CrowdsecQuotaUsage {
|
||||||
|
/** UTC calendar day the counter belongs to (YYYY-MM-DD). */
|
||||||
|
date: string;
|
||||||
|
used: number;
|
||||||
|
/** 0 = unlimited. */
|
||||||
|
quota: number;
|
||||||
|
exhausted: boolean;
|
||||||
|
}
|
||||||
|
|
||||||
/** Value written into the shared block key so the admin UI can label the source. */
|
/** Value written into the shared block key so the admin UI can label the source. */
|
||||||
export const CROWDSEC_BLOCK_SOURCE = "crowdsec";
|
export const CROWDSEC_BLOCK_SOURCE = "crowdsec";
|
||||||
|
|
||||||
@@ -82,6 +110,13 @@ const AUTH_BACKOFF_MS = 300_000;
|
|||||||
const VERDICT_PREFIX = "crowdsec:cti:";
|
const VERDICT_PREFIX = "crowdsec:cti:";
|
||||||
const LOOKUP_LOCK_PREFIX = "crowdsec:lock:";
|
const LOOKUP_LOCK_PREFIX = "crowdsec:lock:";
|
||||||
const LAST_VERIFY_KEY = "crowdsec:last-verify";
|
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;
|
||||||
|
/** 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. */
|
||||||
|
const QUOTA_WARN_RATIO = 0.8;
|
||||||
/** Well-known, community-safe address used by the admin "verify" button. */
|
/** Well-known, community-safe address used by the admin "verify" button. */
|
||||||
const PROBE_IP = "1.1.1.1";
|
const PROBE_IP = "1.1.1.1";
|
||||||
|
|
||||||
@@ -197,6 +232,26 @@ export function verdictIsMalicious(
|
|||||||
|
|
||||||
const memoryVerdicts = new Map<string, CrowdsecVerdict>();
|
const memoryVerdicts = new Map<string, CrowdsecVerdict>();
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Insert/refresh an in-process verdict while keeping the cache bounded: Map
|
||||||
|
* iteration order is insertion order, so the oldest (leftmost) entry is
|
||||||
|
* dropped first and re-inserted entries are refreshed to the back.
|
||||||
|
*/
|
||||||
|
function rememberVerdict(verdict: CrowdsecVerdict): void {
|
||||||
|
memoryVerdicts.delete(verdict.ip);
|
||||||
|
memoryVerdicts.set(verdict.ip, verdict);
|
||||||
|
while (memoryVerdicts.size > MEMORY_VERDICT_CACHE_MAX) {
|
||||||
|
const oldest = memoryVerdicts.keys().next();
|
||||||
|
if (oldest.done) break;
|
||||||
|
memoryVerdicts.delete(oldest.value);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Test hook only — reports the bounded in-process cache size. */
|
||||||
|
export function getMemoryVerdictCacheSize(): number {
|
||||||
|
return memoryVerdicts.size;
|
||||||
|
}
|
||||||
|
|
||||||
function verdictKey(ip: string): string {
|
function verdictKey(ip: string): string {
|
||||||
return `${VERDICT_PREFIX}${ip}`;
|
return `${VERDICT_PREFIX}${ip}`;
|
||||||
}
|
}
|
||||||
@@ -211,7 +266,7 @@ async function readVerdictCache(ip: string): Promise<CrowdsecVerdict | null> {
|
|||||||
const raw = await redis.get(verdictKey(ip));
|
const raw = await redis.get(verdictKey(ip));
|
||||||
if (raw) {
|
if (raw) {
|
||||||
const parsed = JSON.parse(raw) as CrowdsecVerdict;
|
const parsed = JSON.parse(raw) as CrowdsecVerdict;
|
||||||
memoryVerdicts.set(ip, parsed);
|
rememberVerdict(parsed);
|
||||||
return parsed;
|
return parsed;
|
||||||
}
|
}
|
||||||
} catch {
|
} catch {
|
||||||
@@ -222,7 +277,7 @@ async function readVerdictCache(ip: string): Promise<CrowdsecVerdict | null> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async function writeVerdictCache(verdict: CrowdsecVerdict): Promise<void> {
|
async function writeVerdictCache(verdict: CrowdsecVerdict): Promise<void> {
|
||||||
memoryVerdicts.set(verdict.ip, verdict);
|
rememberVerdict(verdict);
|
||||||
if (redis) {
|
if (redis) {
|
||||||
try {
|
try {
|
||||||
await redis.set(
|
await redis.set(
|
||||||
@@ -256,11 +311,79 @@ async function acquireLookupLock(ip: string): Promise<boolean> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let backoffUntil = 0;
|
let backoffUntil = 0;
|
||||||
|
let quotaWarnedDate: string | null = null;
|
||||||
|
|
||||||
|
function quotaDate(): string {
|
||||||
|
return new Date().toISOString().slice(0, 10);
|
||||||
|
}
|
||||||
|
|
||||||
|
function quotaKey(date: string): string {
|
||||||
|
return `${QUOTA_PREFIX}${date}`;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Today's configured ceiling. Coerced because tests run with
|
||||||
|
* SKIP_ENV_VALIDATION (raw process.env strings, no zod defaults) while
|
||||||
|
* production gets a parsed number. 0 (or unset) = unlimited.
|
||||||
|
*/
|
||||||
|
function dailyQuota(): number {
|
||||||
|
const raw = Number(env.CROWDSEC_CTI_DAILY_QUOTA ?? 0);
|
||||||
|
return Number.isFinite(raw) && raw > 0 ? raw : 0;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Today's CTI usage against the configured daily ceiling. Counted in Redis so
|
||||||
|
* every instance shares one budget.
|
||||||
|
*/
|
||||||
|
export async function getCrowdsecQuotaUsage(): Promise<CrowdsecQuotaUsage> {
|
||||||
|
const date = quotaDate();
|
||||||
|
const quota = dailyQuota();
|
||||||
|
let used = 0;
|
||||||
|
if (redis) {
|
||||||
|
try {
|
||||||
|
used = Number((await redis.get(quotaKey(date))) ?? 0);
|
||||||
|
if (!Number.isFinite(used)) used = 0;
|
||||||
|
} catch {
|
||||||
|
// counter unavailable — report zero rather than blocking the admin
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return { date, used, quota, exhausted: quota > 0 && used >= quota };
|
||||||
|
}
|
||||||
|
|
||||||
|
async function reserveQuota(): Promise<boolean> {
|
||||||
|
const quota = dailyQuota();
|
||||||
|
if (quota <= 0) return true;
|
||||||
|
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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Community reputation verdict for an IP, from cache when possible. Returns
|
* Community reputation verdict for an IP, from cache when possible. Returns
|
||||||
* null when the API is not configured, the lookup failed, or the API is in
|
* null when the API is not configured, the lookup failed, the API is in
|
||||||
* backoff — never throws, so it is safe on the gate's hot path.
|
* backoff, or today's quota is spent — never throws, so it is safe on the
|
||||||
|
* gate's hot path.
|
||||||
*/
|
*/
|
||||||
export async function lookupCrowdsecVerdict(
|
export async function lookupCrowdsecVerdict(
|
||||||
ip: string,
|
ip: string,
|
||||||
@@ -284,6 +407,15 @@ export async function lookupCrowdsecVerdict(
|
|||||||
const raced = await readVerdictCache(ip);
|
const raced = await readVerdictCache(ip);
|
||||||
if (raced) return raced;
|
if (raced) return raced;
|
||||||
|
|
||||||
|
// Cache miss costs a paid call — reserve against today's quota first.
|
||||||
|
if (!(await reserveQuota())) {
|
||||||
|
logger.warn(
|
||||||
|
"[crowdsec-api] CTI daily quota exhausted — pausing lookups until tomorrow",
|
||||||
|
{ quota: dailyQuota() },
|
||||||
|
);
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
const response = await crowdsecRequest(`/smoke/${encodeURIComponent(ip)}`);
|
const response = await crowdsecRequest(`/smoke/${encodeURIComponent(ip)}`);
|
||||||
|
|
||||||
if (response.status === 404) {
|
if (response.status === 404) {
|
||||||
@@ -358,6 +490,27 @@ export async function maybeAutoBlockCrowdsec(input: {
|
|||||||
if (existingTtl >= ttlSeconds * 1000) return;
|
if (existingTtl >= ttlSeconds * 1000) return;
|
||||||
|
|
||||||
await redis.set(blockKey, CROWDSEC_BLOCK_SOURCE, "EX", ttlSeconds);
|
await redis.set(blockKey, CROWDSEC_BLOCK_SOURCE, "EX", ttlSeconds);
|
||||||
|
// Record why this block exists so the admin panel can surface the
|
||||||
|
// community reasoning (reputation, score, behaviors) for the IP.
|
||||||
|
const meta: CrowdsecBlockMeta = {
|
||||||
|
source: CROWDSEC_BLOCK_SOURCE,
|
||||||
|
category,
|
||||||
|
reputation: verdict.reputation,
|
||||||
|
score: verdict.score,
|
||||||
|
behaviors: verdict.behaviors,
|
||||||
|
ttlSeconds,
|
||||||
|
blockedAt: Date.now(),
|
||||||
|
};
|
||||||
|
try {
|
||||||
|
await redis.set(
|
||||||
|
`${BLOCK_META_PREFIX}${ip}`,
|
||||||
|
JSON.stringify(meta),
|
||||||
|
"EX",
|
||||||
|
ttlSeconds,
|
||||||
|
);
|
||||||
|
} catch {
|
||||||
|
// metadata is display sugar only — the block itself is already set.
|
||||||
|
}
|
||||||
// CrowdSec only records the block in the gate's own key. It never
|
// CrowdSec only records the block in the gate's own key. It never
|
||||||
// creates Cloudflare edge rules — the gate's own escalation logic is
|
// creates Cloudflare edge rules — the gate's own escalation logic is
|
||||||
// the only place that may mirror a host-level block to the edge.
|
// the only place that may mirror a host-level block to the edge.
|
||||||
@@ -372,6 +525,9 @@ export async function maybeAutoBlockCrowdsec(input: {
|
|||||||
behaviors: verdict.behaviors,
|
behaviors: verdict.behaviors,
|
||||||
},
|
},
|
||||||
);
|
);
|
||||||
|
// 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 });
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
logger.error("[crowdsec-api] Automatic IP block failed", {
|
logger.error("[crowdsec-api] Automatic IP block failed", {
|
||||||
ip,
|
ip,
|
||||||
@@ -461,9 +617,24 @@ export async function setLastCrowdsecVerify(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Why an IP is blocked by CrowdSec, when known. */
|
||||||
|
export async function getCrowdsecBlockMeta(
|
||||||
|
ip: string,
|
||||||
|
): Promise<CrowdsecBlockMeta | null> {
|
||||||
|
if (!redis) return null;
|
||||||
|
try {
|
||||||
|
const raw = await redis.get(`${BLOCK_META_PREFIX}${ip}`);
|
||||||
|
if (!raw) return null;
|
||||||
|
return JSON.parse(raw) as CrowdsecBlockMeta;
|
||||||
|
} catch {
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/** Test hook only — drop in-memory state between unit runs. */
|
/** Test hook only — drop in-memory state between unit runs. */
|
||||||
export function resetCrowdsecCache(): void {
|
export function resetCrowdsecCache(): void {
|
||||||
memoryVerdicts.clear();
|
memoryVerdicts.clear();
|
||||||
backoffUntil = 0;
|
backoffUntil = 0;
|
||||||
|
quotaWarnedDate = null;
|
||||||
lastVerifyMemory = null;
|
lastVerifyMemory = null;
|
||||||
}
|
}
|
||||||
@@ -0,0 +1,280 @@
|
|||||||
|
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||||
|
import type { CrowdsecVerdict } from "./crowdsec-api";
|
||||||
|
import {
|
||||||
|
type CrowdsecReportStatus,
|
||||||
|
crowdsecReportEnabled,
|
||||||
|
getLastCrowdsecReport,
|
||||||
|
reportCrowdsecSignal,
|
||||||
|
resetCrowdsecReportCache,
|
||||||
|
verifyCrowdsecReporting,
|
||||||
|
} from "./crowdsec-report";
|
||||||
|
|
||||||
|
// 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<string, string>() }));
|
||||||
|
|
||||||
|
vi.mock("@/lib/redis", () => ({
|
||||||
|
redis: {
|
||||||
|
get: async (key: string) => state.map.get(key) ?? null,
|
||||||
|
set: async (
|
||||||
|
key: string,
|
||||||
|
value: string,
|
||||||
|
_mode?: string,
|
||||||
|
_seconds?: number,
|
||||||
|
nx?: string,
|
||||||
|
) => {
|
||||||
|
if (nx === "NX" && state.map.has(key)) return null;
|
||||||
|
state.map.set(key, value);
|
||||||
|
return "OK";
|
||||||
|
},
|
||||||
|
del: async (...keys: string[]) => {
|
||||||
|
for (const key of keys) state.map.delete(key);
|
||||||
|
return keys.length;
|
||||||
|
},
|
||||||
|
},
|
||||||
|
__esModule: true,
|
||||||
|
}));
|
||||||
|
|
||||||
|
vi.mock("@/lib/logger", () => ({
|
||||||
|
logger: {
|
||||||
|
info: vi.fn(),
|
||||||
|
warn: vi.fn(),
|
||||||
|
error: vi.fn(),
|
||||||
|
debug: vi.fn(),
|
||||||
|
},
|
||||||
|
}));
|
||||||
|
|
||||||
|
const CAPI = "https://capi.example.test/v3";
|
||||||
|
const MACHINE = "m".repeat(48);
|
||||||
|
const PASSWORD = "Strong!1P@ssw0rdStrong!1P@ssw0rd";
|
||||||
|
|
||||||
|
function jsonResponse(body: unknown, status = 200): Response {
|
||||||
|
return new Response(JSON.stringify(body), {
|
||||||
|
status,
|
||||||
|
headers: { "content-type": "application/json" },
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
function signalInput(ip = "198.51.100.9") {
|
||||||
|
const verdict: CrowdsecVerdict = {
|
||||||
|
ip,
|
||||||
|
reputation: "malicious",
|
||||||
|
score: 5,
|
||||||
|
aggressiveness: 4,
|
||||||
|
confidence: "0.95",
|
||||||
|
behaviors: ["http:bruteforce", "http:scan"],
|
||||||
|
falsePositive: false,
|
||||||
|
checkedAt: Date.now(),
|
||||||
|
};
|
||||||
|
return {
|
||||||
|
ip,
|
||||||
|
category: "api",
|
||||||
|
ttlSeconds: 86_400,
|
||||||
|
verdict,
|
||||||
|
meta: {
|
||||||
|
source: "crowdsec" as const,
|
||||||
|
category: "api",
|
||||||
|
reputation: verdict.reputation,
|
||||||
|
score: verdict.score,
|
||||||
|
behaviors: verdict.behaviors,
|
||||||
|
ttlSeconds: 86_400,
|
||||||
|
blockedAt: Date.now(),
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
describe("crowdsec-report", () => {
|
||||||
|
let fetchMock: ReturnType<typeof vi.fn>;
|
||||||
|
|
||||||
|
function routeCapi(overrides: Record<string, number> = {}) {
|
||||||
|
const statusFor = (path: string) =>
|
||||||
|
overrides[path] ?? (path === "/signals" ? 200 : 200);
|
||||||
|
fetchMock.mockImplementation((url: string) => {
|
||||||
|
const path = String(url).replace(CAPI, "");
|
||||||
|
const status = statusFor(path);
|
||||||
|
if (status !== 200) {
|
||||||
|
return Promise.resolve(jsonResponse({ message: "boom" }, status));
|
||||||
|
}
|
||||||
|
if (path === "/watchers/login") {
|
||||||
|
return Promise.resolve(
|
||||||
|
jsonResponse({
|
||||||
|
token: "jwt-xyz",
|
||||||
|
expire: new Date(Date.now() + 3_600_000).toISOString(),
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
return Promise.resolve(jsonResponse({}));
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
beforeEach(() => {
|
||||||
|
vi.unstubAllGlobals();
|
||||||
|
vi.unstubAllEnvs();
|
||||||
|
state.map.clear();
|
||||||
|
resetCrowdsecReportCache();
|
||||||
|
fetchMock = vi.fn();
|
||||||
|
vi.stubGlobal("fetch", fetchMock);
|
||||||
|
vi.stubEnv("CROWDSEC_REPORT_ENABLED", "true");
|
||||||
|
vi.stubEnv("CROWDSEC_REPORT_MACHINE_ID", MACHINE);
|
||||||
|
vi.stubEnv("CROWDSEC_REPORT_PASSWORD", PASSWORD);
|
||||||
|
vi.stubEnv("CROWDSEC_CAPI_BASE_URL", CAPI);
|
||||||
|
});
|
||||||
|
|
||||||
|
afterEach(() => {
|
||||||
|
vi.unstubAllGlobals();
|
||||||
|
vi.unstubAllEnvs();
|
||||||
|
state.map.clear();
|
||||||
|
resetCrowdsecReportCache();
|
||||||
|
vi.restoreAllMocks();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("is enabled only when the toggle and credentials are present", async () => {
|
||||||
|
expect(await crowdsecReportEnabled()).toBe(true);
|
||||||
|
|
||||||
|
vi.stubEnv("CROWDSEC_REPORT_ENABLED", "");
|
||||||
|
resetCrowdsecReportCache();
|
||||||
|
expect(await crowdsecReportEnabled()).toBe(false);
|
||||||
|
|
||||||
|
// Machine id supplied but no password: falls back to generating a
|
||||||
|
// stable credential pair persisted in Redis.
|
||||||
|
vi.stubEnv("CROWDSEC_REPORT_ENABLED", "true");
|
||||||
|
vi.stubEnv("CROWDSEC_REPORT_PASSWORD", "");
|
||||||
|
resetCrowdsecReportCache();
|
||||||
|
expect(await crowdsecReportEnabled()).toBe(true);
|
||||||
|
const storedMachine = state.map.get("crowdsec:report:machine");
|
||||||
|
expect(storedMachine).toMatch(/^[A-Za-z0-9]{48}$/);
|
||||||
|
expect(state.map.get("crowdsec:report:pass")).toBeTruthy();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("does nothing when the channel is disabled", async () => {
|
||||||
|
vi.stubEnv("CROWDSEC_REPORT_ENABLED", "");
|
||||||
|
resetCrowdsecReportCache();
|
||||||
|
|
||||||
|
await reportCrowdsecSignal(signalInput());
|
||||||
|
expect(fetchMock).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("registers once, caches the token and pushes one signal per IP", async () => {
|
||||||
|
routeCapi();
|
||||||
|
|
||||||
|
await reportCrowdsecSignal(signalInput("198.51.100.10"));
|
||||||
|
await reportCrowdsecSignal(signalInput("198.51.100.11"));
|
||||||
|
await reportCrowdsecSignal(signalInput("198.51.100.10"));
|
||||||
|
// Let the fire-and-forget network body land.
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 20));
|
||||||
|
|
||||||
|
const urls = fetchMock.mock.calls.map((call) => String(call[0]));
|
||||||
|
expect(urls.filter((u) => u.endsWith("/watchers/register"))).toHaveLength(
|
||||||
|
1,
|
||||||
|
);
|
||||||
|
expect(urls.filter((u) => u.endsWith("/watchers/login"))).toHaveLength(1);
|
||||||
|
expect(urls.filter((u) => u.endsWith("/signals"))).toHaveLength(2);
|
||||||
|
// No enrollment requested without an attachment key.
|
||||||
|
expect(urls.some((u) => u.endsWith("/watchers/enroll"))).toBe(false);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("builds a well-formed CrowdSec signal with a ban decision", async () => {
|
||||||
|
routeCapi();
|
||||||
|
|
||||||
|
await reportCrowdsecSignal(signalInput());
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 20));
|
||||||
|
|
||||||
|
const signalsCall = fetchMock.mock.calls.find((call) =>
|
||||||
|
String(call[0]).endsWith("/signals"),
|
||||||
|
);
|
||||||
|
expect(signalsCall).toBeDefined();
|
||||||
|
if (!signalsCall) throw new Error("expected a /signals call");
|
||||||
|
const init = signalsCall[1] as {
|
||||||
|
body: string;
|
||||||
|
headers: Record<string, string>;
|
||||||
|
};
|
||||||
|
const body = JSON.parse(init.body) as Record<string, unknown>[];
|
||||||
|
expect(body).toHaveLength(1);
|
||||||
|
const signal = body[0] as {
|
||||||
|
machine_id: string;
|
||||||
|
scenario: string;
|
||||||
|
scenario_version: string;
|
||||||
|
source: { scope: string; value: string; ip: string };
|
||||||
|
decisions: {
|
||||||
|
scope: string;
|
||||||
|
type: string;
|
||||||
|
value: string;
|
||||||
|
duration: string;
|
||||||
|
}[];
|
||||||
|
context: { key: string; value: string }[];
|
||||||
|
created_at: string;
|
||||||
|
start_at: string;
|
||||||
|
stop_at: string;
|
||||||
|
};
|
||||||
|
expect(signal.machine_id).toBe(MACHINE);
|
||||||
|
expect(signal.scenario).toBe("community/anti-ddos-block");
|
||||||
|
expect(signal.scenario_version).toBe("1.0.0");
|
||||||
|
expect(signal.source).toEqual({
|
||||||
|
scope: "ip",
|
||||||
|
value: "198.51.100.9",
|
||||||
|
ip: "198.51.100.9",
|
||||||
|
});
|
||||||
|
expect(signal.decisions).toHaveLength(1);
|
||||||
|
expect(signal.decisions[0]).toMatchObject({
|
||||||
|
origin: "crowdsec",
|
||||||
|
scope: "ip",
|
||||||
|
type: "ban",
|
||||||
|
value: "198.51.100.9",
|
||||||
|
});
|
||||||
|
expect(String(signal.decisions[0].duration)).toMatch(/^24h0m0s$/);
|
||||||
|
for (const key of ["created_at", "start_at", "stop_at"] as const) {
|
||||||
|
expect(typeof signal[key]).toBe("string");
|
||||||
|
}
|
||||||
|
expect(
|
||||||
|
signal.context.find((c) => c.key === "crowdsec_reputation")?.value,
|
||||||
|
).toBe("malicious");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("records a healthy last-report state after a successful push", async () => {
|
||||||
|
routeCapi();
|
||||||
|
await reportCrowdsecSignal(signalInput());
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 20));
|
||||||
|
const last = await getLastCrowdsecReport();
|
||||||
|
expect(last?.ok).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("never throws and logs the failure when the CAPI rejects the signal", 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 expect(reportCrowdsecSignal(signalInput())).resolves.toBeUndefined();
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 20));
|
||||||
|
const last: CrowdsecReportStatus | null = await getLastCrowdsecReport();
|
||||||
|
expect(last?.ok).toBe(false);
|
||||||
|
expect(last?.message).toContain("signal push rejected");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("verifies the watcher channel end to end", async () => {
|
||||||
|
routeCapi();
|
||||||
|
const status = await verifyCrowdsecReporting();
|
||||||
|
expect(status.ok).toBe(true);
|
||||||
|
expect(String(fetchMock.mock.calls[0][0])).toContain("/watchers/register");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("reports a clear reason when verification is impossible", async () => {
|
||||||
|
vi.stubEnv("CROWDSEC_REPORT_ENABLED", "");
|
||||||
|
resetCrowdsecReportCache();
|
||||||
|
const status = await verifyCrowdsecReporting();
|
||||||
|
expect(status.ok).toBe(false);
|
||||||
|
expect(status.message).toContain("CROWDSEC_REPORT_ENABLED");
|
||||||
|
expect(fetchMock).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -0,0 +1,541 @@
|
|||||||
|
import "server-only";
|
||||||
|
|
||||||
|
import { createHash, randomBytes } from "node:crypto";
|
||||||
|
import { env } from "@/env";
|
||||||
|
import type {
|
||||||
|
CrowdsecBlockMeta,
|
||||||
|
CrowdsecConnectionStatus,
|
||||||
|
CrowdsecVerdict,
|
||||||
|
} from "@/lib/crowdsec-api";
|
||||||
|
import { logger } from "@/lib/logger";
|
||||||
|
import { redis } from "@/lib/redis";
|
||||||
|
import { UNKNOWN_CLIENT_IP } from "./client-ip";
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CrowdSec Central API (CAPI) signal push — the "give back" side of the
|
||||||
|
* anti-DDoS pipeline.
|
||||||
|
*
|
||||||
|
* When the gate blocks an IP based on the CTI community reputation, this
|
||||||
|
* module reports that detection back to CrowdSec (POST /v3/signals) so the
|
||||||
|
* community blocklist also protects every other member. Strictly opt-in
|
||||||
|
* (CROWDSEC_REPORT_ENABLED) and always fire-and-forget: a failure here never
|
||||||
|
* blocks the hot path, never throws to the caller, and records the last
|
||||||
|
* outcome for the admin panel.
|
||||||
|
*
|
||||||
|
* A plain CTI API key cannot push signals, so we act as a CAPI "watcher":
|
||||||
|
* 1. generate/load a stable 48-char alnum machine_id + password pair
|
||||||
|
* (persisted in Redis when not provided via env),
|
||||||
|
* 2. register it once (POST /v3/watchers/register),
|
||||||
|
* 3. login to obtain a JWT (POST /v3/watchers/login), cached in Redis and
|
||||||
|
* refreshed against its expiry,
|
||||||
|
* 4. optional console enrollment via attachment key (POST /v3/watchers/enroll),
|
||||||
|
* 5. push the block as a signal (POST /v3/signals), deduped per IP.
|
||||||
|
*/
|
||||||
|
|
||||||
|
const API_TIMEOUT_MS = 10_000;
|
||||||
|
const TOKEN_CACHE_KEY = "crowdsec:report:token";
|
||||||
|
const MACHINE_KEY = "crowdsec:report:machine";
|
||||||
|
const PASSWORD_KEY = "crowdsec:report:pass";
|
||||||
|
const REGISTERED_KEY = "crowdsec:report:registered";
|
||||||
|
const ENROLLED_KEY = "crowdsec:report:enrolled";
|
||||||
|
const LAST_REPORT_KEY = "crowdsec:last-report";
|
||||||
|
const REPORT_LOCK_PREFIX = "crowdsec:report:";
|
||||||
|
/** An IP is only reported once per window — the block itself already deters. */
|
||||||
|
const REPORT_DEDUPE_SECONDS = 6 * 3_600;
|
||||||
|
const SCENARIO = "community/anti-ddos-block";
|
||||||
|
const SCENARIO_VERSION = "1.0.0";
|
||||||
|
|
||||||
|
export class CrowdsecReportError extends Error {}
|
||||||
|
|
||||||
|
/** True when the reporting channel is switched on AND usable. */
|
||||||
|
export async function crowdsecReportEnabled(): Promise<boolean> {
|
||||||
|
// Production parses CROWDSEC_REPORT_ENABLED to a boolean via zod; tests
|
||||||
|
// (SKIP_ENV_VALIDATION) expose the raw env string, so accept both forms.
|
||||||
|
const flag: unknown = env.CROWDSEC_REPORT_ENABLED;
|
||||||
|
if (flag !== true && flag !== "true" && flag !== "1") return false;
|
||||||
|
return (await loadReportCredentials()) !== null;
|
||||||
|
}
|
||||||
|
|
||||||
|
function isAlnum48(value: string): boolean {
|
||||||
|
return /^[A-Za-z0-9]{48}$/.test(value);
|
||||||
|
}
|
||||||
|
|
||||||
|
function generateMachineId(): string {
|
||||||
|
// CAPI schema: exactly 48 characters, [A-Za-z0-9].
|
||||||
|
const alphabet =
|
||||||
|
"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789";
|
||||||
|
const bytes = randomBytes(48);
|
||||||
|
let id = "";
|
||||||
|
for (let i = 0; i < 48; i += 1) {
|
||||||
|
id += alphabet[bytes[i] % alphabet.length];
|
||||||
|
}
|
||||||
|
return id;
|
||||||
|
}
|
||||||
|
|
||||||
|
function generatePassword(): string {
|
||||||
|
// Deliberately generous: each class is present so common password-policy
|
||||||
|
// rules on the CAPI side are satisfied.
|
||||||
|
const upper = "ABCDEFGHIJKLMNOPQRSTUVWXYZ";
|
||||||
|
const lower = "abcdefghijklmnopqrstuvwxyz";
|
||||||
|
const digits = "0123456789";
|
||||||
|
const symbols = "!@#$%^&*()-_=+[]{};:,.?";
|
||||||
|
const charset = `${upper}${lower}${digits}${symbols}`;
|
||||||
|
const bytes = randomBytes(32);
|
||||||
|
let password = "";
|
||||||
|
for (let i = 0; i < 8; i += 1) {
|
||||||
|
// Guarantee at least one of each class.
|
||||||
|
const pool = [upper, lower, digits, symbols][i % 4];
|
||||||
|
password += pool[bytes[i] % pool.length];
|
||||||
|
}
|
||||||
|
for (let i = 8; i < 32; i += 1) {
|
||||||
|
password += charset[bytes[i] % charset.length];
|
||||||
|
}
|
||||||
|
return password;
|
||||||
|
}
|
||||||
|
|
||||||
|
interface ReportCredentials {
|
||||||
|
machineId: string;
|
||||||
|
password: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
let credentialsCache: ReportCredentials | null = null;
|
||||||
|
let credentialsMissingRedisWarned = false;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Load the configured watcher credentials, or generate a stable pair and
|
||||||
|
* persist it in Redis so restarts and other instances reuse the same identity.
|
||||||
|
*/
|
||||||
|
async function loadReportCredentials(): Promise<ReportCredentials | null> {
|
||||||
|
if (credentialsCache) return credentialsCache;
|
||||||
|
const envMachine = env.CROWDSEC_REPORT_MACHINE_ID?.trim();
|
||||||
|
const envPassword = env.CROWDSEC_REPORT_PASSWORD;
|
||||||
|
if (envMachine && envPassword) {
|
||||||
|
credentialsCache = { machineId: envMachine, password: envPassword };
|
||||||
|
return credentialsCache;
|
||||||
|
}
|
||||||
|
if (!redis) {
|
||||||
|
if (!credentialsMissingRedisWarned) {
|
||||||
|
credentialsMissingRedisWarned = true;
|
||||||
|
logger.warn(
|
||||||
|
"[crowdsec-report] Redis is required to persist auto-generated watcher credentials — set CROWDSEC_REPORT_MACHINE_ID and CROWDSEC_REPORT_PASSWORD, or REDIS_URL",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
let machineId: string | null = null;
|
||||||
|
let password: string | null = null;
|
||||||
|
const storedMachine = await redis.get(MACHINE_KEY);
|
||||||
|
const storedPassword = await redis.get(PASSWORD_KEY);
|
||||||
|
if (storedMachine && isAlnum48(storedMachine)) machineId = storedMachine;
|
||||||
|
if (storedPassword) password = storedPassword;
|
||||||
|
if (!machineId) machineId = generateMachineId();
|
||||||
|
if (!password) password = generatePassword();
|
||||||
|
await redis.set(MACHINE_KEY, machineId);
|
||||||
|
await redis.set(PASSWORD_KEY, password);
|
||||||
|
credentialsCache = { machineId, password };
|
||||||
|
return credentialsCache;
|
||||||
|
} catch {
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async function capiRequest(
|
||||||
|
path: string,
|
||||||
|
init: { method?: "POST" | "GET"; token?: string; body?: unknown } = {},
|
||||||
|
): Promise<Response> {
|
||||||
|
const controller = new AbortController();
|
||||||
|
const timer = setTimeout(() => controller.abort(), API_TIMEOUT_MS);
|
||||||
|
const headers: Record<string, string> = {
|
||||||
|
Accept: "application/json",
|
||||||
|
"Content-Type": "application/json",
|
||||||
|
};
|
||||||
|
if (init.token) headers.Authorization = `Bearer ${init.token}`;
|
||||||
|
try {
|
||||||
|
return await fetch(`${env.CROWDSEC_CAPI_BASE_URL}${path}`, {
|
||||||
|
method: init.method ?? "POST",
|
||||||
|
headers,
|
||||||
|
body: init.body === undefined ? undefined : JSON.stringify(init.body),
|
||||||
|
signal: controller.signal,
|
||||||
|
cache: "no-store",
|
||||||
|
});
|
||||||
|
} finally {
|
||||||
|
clearTimeout(timer);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async function errorDetail(response: Response): Promise<string> {
|
||||||
|
try {
|
||||||
|
const body = (await response.json()) as { message?: string };
|
||||||
|
return body.message ?? `HTTP ${response.status}`;
|
||||||
|
} catch {
|
||||||
|
return `HTTP ${response.status}`;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
interface CapToken {
|
||||||
|
raw: string;
|
||||||
|
expiresAt: number;
|
||||||
|
}
|
||||||
|
|
||||||
|
let capToken: CapToken | null = null;
|
||||||
|
let tokenPromise: Promise<string> | null = null;
|
||||||
|
|
||||||
|
function scenarioHash(): string {
|
||||||
|
return createHash("sha256")
|
||||||
|
.update(`${SCENARIO}:${SCENARIO_VERSION}`)
|
||||||
|
.digest("hex")
|
||||||
|
.slice(0, 16);
|
||||||
|
}
|
||||||
|
|
||||||
|
async function registerWatcher(
|
||||||
|
machineId: string,
|
||||||
|
password: string,
|
||||||
|
): Promise<void> {
|
||||||
|
if (!redis) return;
|
||||||
|
try {
|
||||||
|
const already = await redis.get(REGISTERED_KEY);
|
||||||
|
if (already) return;
|
||||||
|
} catch {
|
||||||
|
// proceed anyway — registering is idempotent-ish (400 = already exists)
|
||||||
|
}
|
||||||
|
const response = await capiRequest("/watchers/register", {
|
||||||
|
body: { machine_id: machineId, password },
|
||||||
|
});
|
||||||
|
if (response.ok || response.status === 400) {
|
||||||
|
try {
|
||||||
|
await redis.set(REGISTERED_KEY, "1", "EX", 30 * 24 * 3_600);
|
||||||
|
} catch {
|
||||||
|
// fine — will just attempt registration again later
|
||||||
|
}
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
throw new CrowdsecReportError(
|
||||||
|
`watcher registration failed (HTTP ${response.status}): ${await errorDetail(response)}`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Login and cache the JWT; guarded against concurrent logins. */
|
||||||
|
async function acquireCapToken(): Promise<string> {
|
||||||
|
if (capToken && capToken.expiresAt > Date.now() + 60_000) {
|
||||||
|
return capToken.raw;
|
||||||
|
}
|
||||||
|
if (tokenPromise) return tokenPromise;
|
||||||
|
const credentials = await loadReportCredentials();
|
||||||
|
if (!credentials) {
|
||||||
|
throw new CrowdsecReportError(
|
||||||
|
"CrowdSec reporting is not configured (no watcher credentials)",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
tokenPromise = (async () => {
|
||||||
|
if (capToken && capToken.expiresAt > Date.now() + 60_000) {
|
||||||
|
return capToken.raw;
|
||||||
|
}
|
||||||
|
if (redis) {
|
||||||
|
try {
|
||||||
|
const cached = await redis.get(TOKEN_CACHE_KEY);
|
||||||
|
if (cached) {
|
||||||
|
const parsed = JSON.parse(cached) as CapToken;
|
||||||
|
if (parsed.expiresAt > Date.now() + 60_000) {
|
||||||
|
capToken = parsed;
|
||||||
|
return parsed.raw;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
} catch {
|
||||||
|
// fall through to a fresh login
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// First signal push needs the watcher to exist on the CAPI.
|
||||||
|
await registerWatcher(credentials.machineId, credentials.password);
|
||||||
|
|
||||||
|
const response = await capiRequest("/watchers/login", {
|
||||||
|
body: {
|
||||||
|
machine_id: credentials.machineId,
|
||||||
|
password: credentials.password,
|
||||||
|
scenarios: [SCENARIO],
|
||||||
|
},
|
||||||
|
});
|
||||||
|
if (!response.ok) {
|
||||||
|
throw new CrowdsecReportError(
|
||||||
|
`CAPI watcher login failed (HTTP ${response.status}): ${await errorDetail(response)}`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
const body = (await response.json()) as { token?: string; expire?: string };
|
||||||
|
if (!body.token) {
|
||||||
|
throw new CrowdsecReportError("CAPI watcher login returned no token");
|
||||||
|
}
|
||||||
|
let expiresAt = Date.now() + 3_600_000;
|
||||||
|
const expire = body.expire ? Date.parse(body.expire) : NaN;
|
||||||
|
if (Number.isFinite(expire) && expire > Date.now()) {
|
||||||
|
expiresAt = expire;
|
||||||
|
}
|
||||||
|
const token: CapToken = { raw: body.token, expiresAt };
|
||||||
|
capToken = token;
|
||||||
|
if (redis) {
|
||||||
|
try {
|
||||||
|
await redis.set(TOKEN_CACHE_KEY, JSON.stringify(token), "EX", 3_600);
|
||||||
|
} catch {
|
||||||
|
// in-process view is enough
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return token.raw;
|
||||||
|
})().finally(() => {
|
||||||
|
tokenPromise = null;
|
||||||
|
});
|
||||||
|
|
||||||
|
return tokenPromise;
|
||||||
|
}
|
||||||
|
|
||||||
|
async function enrollWatcher(token: string): Promise<void> {
|
||||||
|
const attachmentKey = env.CROWDSEC_REPORT_ENROLL_KEY?.trim();
|
||||||
|
if (!attachmentKey) return;
|
||||||
|
if (!redis) return;
|
||||||
|
try {
|
||||||
|
const enrolled = await redis.get(ENROLLED_KEY);
|
||||||
|
if (enrolled) return;
|
||||||
|
} catch {
|
||||||
|
// best effort below
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
const response = await capiRequest("/watchers/enroll", {
|
||||||
|
token,
|
||||||
|
body: { attachment_key: attachmentKey, name: "atomcms-next" },
|
||||||
|
});
|
||||||
|
if (response.ok) {
|
||||||
|
try {
|
||||||
|
await redis.set(ENROLLED_KEY, "1", "EX", 30 * 24 * 3_600);
|
||||||
|
} catch {
|
||||||
|
// best effort
|
||||||
|
}
|
||||||
|
}
|
||||||
|
} catch (error) {
|
||||||
|
// Enrollment only affects Console visibility — not worth failing a push.
|
||||||
|
logger.debug("[crowdsec-report] Console enrollment skipped", {
|
||||||
|
err: error instanceof Error ? error.message : String(error),
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function formatCapiDuration(seconds: number): string {
|
||||||
|
const safe = Math.max(1, Math.floor(seconds));
|
||||||
|
const hours = Math.floor(safe / 3_600);
|
||||||
|
const minutes = Math.floor((safe % 3_600) / 60);
|
||||||
|
const rest = safe % 60;
|
||||||
|
if (hours > 0) return `${hours}h${minutes}m${rest}s`;
|
||||||
|
if (minutes > 0) return `${minutes}m${rest}s`;
|
||||||
|
return `${rest}s`;
|
||||||
|
}
|
||||||
|
|
||||||
|
async function pushSignal(input: {
|
||||||
|
credentials: ReportCredentials;
|
||||||
|
token: string;
|
||||||
|
ip: string;
|
||||||
|
category: string;
|
||||||
|
ttlSeconds: number;
|
||||||
|
verdict: CrowdsecVerdict;
|
||||||
|
meta: CrowdsecBlockMeta;
|
||||||
|
}): Promise<CrowdsecReportStatus> {
|
||||||
|
const now = new Date();
|
||||||
|
try {
|
||||||
|
await enrollWatcher(input.token);
|
||||||
|
const duration = formatCapiDuration(input.ttlSeconds);
|
||||||
|
const response = await capiRequest("/signals", {
|
||||||
|
token: input.token,
|
||||||
|
body: [
|
||||||
|
{
|
||||||
|
machine_id: input.credentials.machineId,
|
||||||
|
message: "atomcms-next anti-DDoS gate blocked a community-flagged IP",
|
||||||
|
scenario: SCENARIO,
|
||||||
|
scenario_version: SCENARIO_VERSION,
|
||||||
|
scenario_hash: scenarioHash(),
|
||||||
|
created_at: now.toISOString(),
|
||||||
|
start_at: now.toISOString(),
|
||||||
|
stop_at: now.toISOString(),
|
||||||
|
source: { scope: "ip", value: input.ip, ip: input.ip },
|
||||||
|
decisions: [
|
||||||
|
{
|
||||||
|
id: 0,
|
||||||
|
origin: "crowdsec",
|
||||||
|
scenario: SCENARIO,
|
||||||
|
scope: "ip",
|
||||||
|
type: "ban",
|
||||||
|
value: input.ip,
|
||||||
|
duration,
|
||||||
|
},
|
||||||
|
],
|
||||||
|
context: [
|
||||||
|
{ key: "crowdsec_category", value: input.category },
|
||||||
|
{
|
||||||
|
key: "crowdsec_reputation",
|
||||||
|
value: input.verdict.reputation ?? "unknown",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
key: "crowdsec_score",
|
||||||
|
value: String(input.verdict.score),
|
||||||
|
},
|
||||||
|
{
|
||||||
|
key: "crowdsec_behaviors",
|
||||||
|
value: input.verdict.behaviors.join(",") || "none",
|
||||||
|
},
|
||||||
|
],
|
||||||
|
},
|
||||||
|
],
|
||||||
|
});
|
||||||
|
if (!response.ok) {
|
||||||
|
throw new CrowdsecReportError(
|
||||||
|
`signal push rejected (HTTP ${response.status}): ${await errorDetail(response)}`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
const status: CrowdsecReportStatus = { ok: true, at: Date.now() };
|
||||||
|
await setLastCrowdsecReport(status);
|
||||||
|
logger.info("[crowdsec-report] Detection shared with the community", {
|
||||||
|
ip: input.ip,
|
||||||
|
category: input.category,
|
||||||
|
ttlSeconds: input.ttlSeconds,
|
||||||
|
score: input.verdict.score,
|
||||||
|
});
|
||||||
|
return status;
|
||||||
|
} catch (error) {
|
||||||
|
const status: CrowdsecReportStatus = {
|
||||||
|
ok: false,
|
||||||
|
at: Date.now(),
|
||||||
|
message: String(error),
|
||||||
|
};
|
||||||
|
await setLastCrowdsecReport(status);
|
||||||
|
logger.error("[crowdsec-report] Signal push failed", {
|
||||||
|
ip: input.ip,
|
||||||
|
err: error,
|
||||||
|
});
|
||||||
|
return status;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Share a CrowdSec-reputation block back into the community. Safe to call
|
||||||
|
* fire-and-forget: does nothing when reporting is off, never throws, and
|
||||||
|
* dedupes so the same IP is only reported once per window.
|
||||||
|
*/
|
||||||
|
export async function reportCrowdsecSignal(input: {
|
||||||
|
ip: string;
|
||||||
|
category: string;
|
||||||
|
ttlSeconds: number;
|
||||||
|
verdict: CrowdsecVerdict;
|
||||||
|
meta: CrowdsecBlockMeta;
|
||||||
|
}): Promise<void> {
|
||||||
|
const { ip, category, ttlSeconds, verdict, meta } = input;
|
||||||
|
if (!(await crowdsecReportEnabled())) return;
|
||||||
|
if (!ip || ip === UNKNOWN_CLIENT_IP) return;
|
||||||
|
const credentials = await loadReportCredentials();
|
||||||
|
if (!credentials) return;
|
||||||
|
|
||||||
|
try {
|
||||||
|
if (redis) {
|
||||||
|
// Cross-instance dedupe: one report per IP per window.
|
||||||
|
const acquired = await redis.set(
|
||||||
|
`${REPORT_LOCK_PREFIX}${ip}`,
|
||||||
|
"1",
|
||||||
|
"EX",
|
||||||
|
REPORT_DEDUPE_SECONDS,
|
||||||
|
"NX",
|
||||||
|
);
|
||||||
|
if (acquired !== "OK") return;
|
||||||
|
}
|
||||||
|
const token = await acquireCapToken();
|
||||||
|
await pushSignal({
|
||||||
|
credentials,
|
||||||
|
token,
|
||||||
|
ip,
|
||||||
|
category,
|
||||||
|
ttlSeconds,
|
||||||
|
verdict,
|
||||||
|
meta,
|
||||||
|
});
|
||||||
|
} catch (error) {
|
||||||
|
logger.error("[crowdsec-report] Signal push preparation failed", {
|
||||||
|
ip,
|
||||||
|
err: error,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface CrowdsecReportStatus {
|
||||||
|
ok: boolean;
|
||||||
|
message?: string;
|
||||||
|
at: number;
|
||||||
|
}
|
||||||
|
|
||||||
|
let lastReportMemory: CrowdsecReportStatus | null = null;
|
||||||
|
|
||||||
|
export async function getLastCrowdsecReport(): Promise<CrowdsecReportStatus | null> {
|
||||||
|
if (redis) {
|
||||||
|
try {
|
||||||
|
const raw = await redis.get(LAST_REPORT_KEY);
|
||||||
|
if (raw) return JSON.parse(raw) as CrowdsecReportStatus;
|
||||||
|
} catch {
|
||||||
|
// fall back to the in-process view
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return lastReportMemory;
|
||||||
|
}
|
||||||
|
|
||||||
|
export async function setLastCrowdsecReport(
|
||||||
|
status: CrowdsecReportStatus,
|
||||||
|
): Promise<void> {
|
||||||
|
lastReportMemory = status;
|
||||||
|
if (redis) {
|
||||||
|
try {
|
||||||
|
await redis.set(
|
||||||
|
LAST_REPORT_KEY,
|
||||||
|
JSON.stringify(status),
|
||||||
|
"EX",
|
||||||
|
48 * 3_600,
|
||||||
|
);
|
||||||
|
} catch {
|
||||||
|
// in-process view is enough
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Probe the full register → login path and report the outcome for the admin. */
|
||||||
|
export async function verifyCrowdsecReporting(): Promise<CrowdsecConnectionStatus> {
|
||||||
|
if (!env.CROWDSEC_REPORT_ENABLED) {
|
||||||
|
return {
|
||||||
|
ok: false,
|
||||||
|
message: "CROWDSEC_REPORT_ENABLED is not set",
|
||||||
|
at: Date.now(),
|
||||||
|
};
|
||||||
|
}
|
||||||
|
const credentials = await loadReportCredentials();
|
||||||
|
if (!credentials) {
|
||||||
|
return {
|
||||||
|
ok: false,
|
||||||
|
message:
|
||||||
|
"CrowdSec reporting is not configured (set CROWDSEC_REPORT_MACHINE_ID + CROWDSEC_REPORT_PASSWORD, or REDIS_URL)",
|
||||||
|
at: Date.now(),
|
||||||
|
};
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
await acquireCapToken();
|
||||||
|
return {
|
||||||
|
ok: true,
|
||||||
|
message: "CAPI watcher connected — signal push ready",
|
||||||
|
at: Date.now(),
|
||||||
|
};
|
||||||
|
} catch (error) {
|
||||||
|
return {
|
||||||
|
ok: false,
|
||||||
|
message: error instanceof Error ? error.message : String(error),
|
||||||
|
at: Date.now(),
|
||||||
|
};
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Test hook only. */
|
||||||
|
export function resetCrowdsecReportCache(): void {
|
||||||
|
credentialsCache = null;
|
||||||
|
capToken = null;
|
||||||
|
tokenPromise = null;
|
||||||
|
lastReportMemory = null;
|
||||||
|
}
|
||||||
Reference in new issue
Block a user