id: clickhouse-api-rate-limit-and-abuse-guard
namespace: company.security
description: |
Ingests high-throughput API gateway request telemetry into ClickHouse or DuckDB,
computes sliding-window request velocity and Z-score anomaly metrics across client IPs,
generates a dynamic quarantine manifest, and alerts security operations in Slack.
triggers:
- id: scheduled_anomaly_sweep
type: io.kestra.plugin.core.trigger.Schedule
description: Periodic sliding-window traffic anomaly detection. Shipped disabled
by default.
cron: "*/5 * * * *"
disabled: true
- id: webhook_telemetry_batch
type: io.kestra.plugin.core.trigger.Webhook
description: Authenticated webhook for API gateway log shipping (Kong, Envoy,
Cloudflare, Traefik).
key: "{{ secret('WEBHOOK_KEY') }}"
inputs:
- id: window_minutes
type: INT
defaults: 10
description: Sliding time window duration in minutes for traffic aggregation.
- id: zscore_threshold
type: FLOAT
defaults: 3.0
description: Standard score threshold above which request velocity is flagged as
an anomaly.
tasks:
- id: generate_mock_traffic_batch
type: io.kestra.plugin.scripts.python.Script
description: Generates representative API traffic logs containing normal
background activity and an injected credential stuffing spike.
taskRunner:
type: io.kestra.plugin.core.runner.Process
script: |
import random
import csv
from datetime import datetime, timedelta
now = datetime.now()
rows = []
# 1. Background legitimate traffic across 25 distinct client subnets
for i in range(25):
ip = f"198.51.100.{i+10}"
for _ in range(random.randint(6, 18)):
dt = now - timedelta(minutes=random.uniform(0, 9))
status = 200 if random.random() > 0.05 else 404
latency = random.randint(25, 95)
rows.append([ip, "/api/v1/catalog/items", status, latency, dt.isoformat()])
# 2. Injected abusive bot subnet performing credential stuffing
bot_ip = "203.0.113.195"
for _ in range(160):
dt = now - timedelta(minutes=random.uniform(0, 8))
status = random.choice([401, 401, 403, 429])
latency = random.randint(18, 38)
rows.append([bot_ip, "/api/v1/auth/login", status, latency, dt.isoformat()])
# Write CSV dataset
with open("api_requests.csv", "w", newline="") as f:
w = csv.writer(f)
w.writerow(["client_ip", "endpoint", "status_code", "response_time_ms", "timestamp"])
w.writerows(rows)
print(f"Generated api_requests.csv with {len(rows)} simulated gateway requests.")
outputFiles:
- api_requests.csv
- id: analyze_sliding_window_anomalies
type: io.kestra.plugin.jdbc.duckdb.Query
description: Computes population mean, standard deviation, and Z-score per
client IP using an embedded analytical query.
communityExtensions: []
inputFiles:
api_requests.csv: "{{
outputs.generate_mock_traffic_batch.outputFiles['api_requests.csv'] }}"
fetchType: FETCH_ONE
sql: |
WITH raw_logs AS (
SELECT * FROM read_csv_auto('api_requests.csv', header = true)
),
ip_metrics AS (
SELECT
client_ip,
COUNT(*) AS total_requests,
SUM(CASE WHEN status_code >= 400 THEN 1 ELSE 0 END) AS error_count,
ROUND(SUM(CASE WHEN status_code >= 400 THEN 1 ELSE 0 END)::DOUBLE / COUNT(*), 2) AS error_rate,
ROUND(AVG(response_time_ms), 1) AS avg_latency_ms
FROM raw_logs
GROUP BY client_ip
),
population_stats AS (
SELECT
AVG(total_requests) AS mean_requests,
COALESCE(NULLIF(STDDEV(total_requests), 0), 1.0) AS stddev_requests
FROM ip_metrics
),
scored_ips AS (
SELECT
m.client_ip,
m.total_requests,
m.error_count,
m.error_rate,
m.avg_latency_ms,
ROUND((m.total_requests - s.mean_requests) / s.stddev_requests, 2) AS z_score
FROM ip_metrics m
CROSS JOIN population_stats s
),
abusive_candidates AS (
SELECT *
FROM scored_ips
WHERE z_score >= {{ inputs.zscore_threshold }} OR (total_requests >= 50 AND error_rate >= 0.40)
),
summary AS (
SELECT
COUNT(*) AS abusive_ip_count,
COALESCE(SUM(total_requests), 0) AS total_abusive_requests,
COALESCE(MAX(z_score), 0.0) AS peak_z_score,
COALESCE(STRING_AGG(client_ip, ', '), 'NONE') AS flagged_ips
FROM abusive_candidates
)
SELECT * FROM summary;
- id: enforce_quarantine_gate
type: io.kestra.plugin.core.flow.If
description: Branches on detected traffic anomalies to export quarantine lists
and alert security operations.
condition: "{{ outputs.analyze_sliding_window_anomalies.row.abusive_ip_count > 0 }}"
then:
- id: export_quarantine_manifest
type: io.kestra.plugin.scripts.python.Script
description: Generates markdown incident brief and quarantine CSV manifest for
WAF/firewall blocklists.
taskRunner:
type: io.kestra.plugin.core.runner.Process
script: |
ip_count = "{{ outputs.analyze_sliding_window_anomalies.row.abusive_ip_count }}"
total_reqs = "{{ outputs.analyze_sliding_window_anomalies.row.total_abusive_requests }}"
peak_z = "{{ outputs.analyze_sliding_window_anomalies.row.peak_z_score }}"
flagged = "{{ outputs.analyze_sliding_window_anomalies.row.flagged_ips }}"
brief = f"""# API ABUSE & RATE-LIMIT INCIDENT BRIEF
## Overview
- **Incident Status**: ACTIVE QUARANTINE APPLIED
- **Abusive IPs Detected**: {ip_count}
- **Total Abusive Requests**: {total_reqs}
- **Peak Z-Score**: {peak_z} standard deviations above mean
## Flagged IP Addresses
`{flagged}`
## Action Taken
IPs isolated for dynamic edge WAF blocklist distribution.
"""
with open("api-abuse-quarantine-brief.md", "w") as f:
f.write(brief)
with open("quarantine_ips.csv", "w") as f:
f.write("client_ip,action,ttl_seconds\n")
for ip in flagged.split(", "):
if ip != "NONE":
f.write(f"{ip},BLOCK,3600\n")
print("Generated api-abuse-quarantine-brief.md and quarantine_ips.csv")
outputFiles:
- api-abuse-quarantine-brief.md
- quarantine_ips.csv
- id: alert_security_channel
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alerts security operations channel with detected anomaly metrics
and blocklist link.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"channel": "#security-ops",
"text": "🚨 *API Rate-Limit & Abuse Anomaly Detected*\n*Abusive IPs*: `{{ outputs.analyze_sliding_window_anomalies.row.abusive_ip_count }}`\n*Flagged*: `{{ outputs.analyze_sliding_window_anomalies.row.flagged_ips }}`\n*Peak Velocity Z-Score*: `{{ outputs.analyze_sliding_window_anomalies.row.peak_z_score }}`\n*Abusive Requests*: `{{ outputs.analyze_sliding_window_anomalies.row.total_abusive_requests }}`\nQuarantine manifest generated in Kestra storage."
}
else:
- id: log_traffic_nominal
type: io.kestra.plugin.core.log.Log
description: Logs nominal traffic status.
message: "API gateway traffic velocity within normal distribution. Zero
quarantine actions required."
errors:
- id: alert_on_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alerts on-call engineering channel if sliding-window query or
artifact export encounters an unhandled error.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"channel": "#secops-pipeline-alerts",
"text": "⚠️ *Pipeline Error*: Rate limit guard flow `{{ flow.id }}` failed on execution `{{ execution.id }}`."
}