Files

305 lines
9.7 KiB
Python

#!/usr/bin/env python3
import os
import sys
import json
import subprocess
import threading
import time
import sqlite3
import hashlib
import logging
import ipaddress
import re
# --- CONFIG ---
STATE_PATH = "/var/lib/anubis-watchdog/state.json"
DB_PATH = "/var/lib/anubis-watchdog/reports.db"
SLEEP_REFRESH = 8
DUPLICATE_TTL_SECONDS = 24 * 3600
REPORT_FLUSH_INTERVAL = 5 # seconds
INMEM_COOLDOWN = 600 # 10 minutes per IP
DOCKER_BIN = os.environ.get("DOCKER_BIN", "/run/current-system/sw/bin/docker")
ABUSEIPDB_KEY_FILE = os.environ.get("ABUSEIPDB_KEY_FILE", "/etc/anubis-watchdog/abuseipdb.key")
CURL_BIN = os.environ.get("CURL_BIN", "/run/current-system/sw/bin/curl")
IP_RE = re.compile(r'\b(?:[0-9]{1,3}\.){3}[0-9]{1,3}\b')
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
# --- DB / Directories ---
def ensure_dirs():
os.makedirs(os.path.dirname(STATE_PATH), exist_ok=True)
os.makedirs(os.path.dirname(DB_PATH), exist_ok=True)
def init_db():
conn = sqlite3.connect(DB_PATH, timeout=10, check_same_thread=False)
cur = conn.cursor()
cur.execute("""
CREATE TABLE IF NOT EXISTS reports (
ip TEXT,
logsig TEXT,
ts INTEGER,
PRIMARY KEY (ip, logsig)
)
""")
conn.commit()
return conn
# --- Deduplication ---
_recent_reports = {}
_recent_lock = threading.Lock()
def signature_for_log(container, ip, raw_log):
"""Hash IP + container + first 400 chars of log to ignore timestamps."""
m = hashlib.sha256()
core = f"{container}:{ip}"
core_part = raw_log[:400]
m.update(core.encode())
m.update(core_part.encode())
return m.hexdigest()
def already_reported(conn, ip, logsig):
now = time.time()
with _recent_lock:
if ip in _recent_reports and now - _recent_reports[ip] < INMEM_COOLDOWN:
return True
cur = conn.cursor()
cur.execute("SELECT ts FROM reports WHERE ip=? AND logsig=?", (ip, logsig))
row = cur.fetchone()
if row and (now - row[0]) < DUPLICATE_TTL_SECONDS:
return True
return False
def mark_reported(conn, ip, logsig):
now = int(time.time())
with _recent_lock:
_recent_reports[ip] = now
cur = conn.cursor()
cur.execute(
"INSERT OR REPLACE INTO reports(ip, logsig, ts) VALUES(?,?,?)",
(ip, logsig, now),
)
conn.commit()
# --- IP extraction ---
def ip_is_public(ip_str: str) -> bool:
try:
a = ipaddress.ip_address(ip_str)
if a.is_private or a.is_loopback or a.is_unspecified or a.is_multicast:
return False
if hasattr(a, "is_global") and not a.is_global:
return False
return True
except Exception:
return False
def extract_ips_from_text(line: str):
cand = IP_RE.findall(line)
return [ip for ip in cand if ip_is_public(ip)]
def pick_ip_from_json(obj: dict) -> str | None:
for k in ("x-real-ip", "x_real_ip", "x-forwarded-for", "x_forwarded_for"):
v = obj.get(k)
if isinstance(v, str) and v.strip():
if "," in v:
parts = [p.strip() for p in v.split(",") if p.strip()]
for p in parts:
if ip_is_public(p):
return p
return parts[0]
return v.strip()
return None
# --- Reporting queue ---
report_queue = []
queue_lock = threading.Lock()
def queue_report(ip, comment):
with queue_lock:
report_queue.append((ip, comment))
def flush_reports(api_key, dry_run=False):
with queue_lock:
batch = report_queue[:]
report_queue.clear()
for ip, comment in batch:
if not os.path.exists(CURL_BIN):
logging.error("curl binary not found at %s; cannot report", CURL_BIN)
continue
if dry_run:
logging.info("[dry-run] would report %s (len comment=%d)", ip, len(comment))
continue
cmd = [
CURL_BIN, "-sS", "-X", "POST", "https://api.abuseipdb.com/api/v2/report",
"-H", f"Key: {api_key}",
"-H", "Accept: application/json",
"--data-urlencode", f"ip={ip}",
"--data-urlencode", "categories=19",
"--data-urlencode", f"comment={comment}"
]
try:
p = subprocess.run(cmd, capture_output=True, text=True, timeout=30)
if p.returncode == 0:
logging.info("Reported %s to AbuseIPDB (stdout len=%d)", ip, len(p.stdout or ""))
else:
logging.warning("curl returned %s when reporting %s: %s", p.returncode, ip, (p.stderr or p.stdout or "")[:300])
except Exception as e:
logging.exception("curl failed for %s: %s", ip, e)
# --- Log line processing ---
def parse_log_line(line: str):
if "|" in line:
right = line.split("|", 1)[1].strip()
else:
right = line.strip()
try:
obj = json.loads(right)
return obj, right
except Exception:
return None, right
def process_line(container, line, conn, api_key, dry_run):
obj, raw = parse_log_line(line)
found_deny = False
if obj and isinstance(obj, dict):
cr = obj.get("check_result") or {}
if isinstance(cr, dict) and cr.get("rule") == "DENY":
found_deny = True
else:
if '"rule":"DENY"' in raw or '"rule": "DENY"' in raw:
found_deny = True
if not found_deny:
return
ips = []
if obj:
cand = pick_ip_from_json(obj)
if cand:
for p in re.split(r'[;,\s]+', cand):
try:
ips.append(str(ipaddress.ip_address(p.strip())))
except Exception:
continue
if not ips:
ips = extract_ips_from_text(raw)
ips = [ip for ip in ips if ip_is_public(ip)]
if not ips:
logging.info("DENY found but no public IP: %s", raw[:200])
return
for ip in ips:
logsig = signature_for_log(container, ip, raw)
if already_reported(conn, ip, logsig):
continue
comment = f"Anubis DENY log\nContainer: {container}\nLog: {raw[:950]}"
queue_report(ip, comment)
mark_reported(conn, ip, logsig)
# --- Container follower ---
class ContainerFollower(threading.Thread):
def __init__(self, name, docker_bin, conn, api_key, dry_run):
super().__init__(daemon=True)
self.name = name
self.docker_bin = docker_bin
self.conn = conn
self.api_key = api_key
self.dry_run = dry_run
self.proc = None
self.stopped = threading.Event()
def run(self):
while not self.stopped.is_set():
try:
# Filter DENY logs at Docker level
cmd = f'{self.docker_bin} logs -f {self.name} | grep \'"rule":"DENY"\''
self.proc = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
shell=True, text=True)
logging.info("Started DENY-only log follower: %s", self.name)
for rawline in self.proc.stdout:
if self.stopped.is_set():
break
if not rawline.strip():
continue
process_line(self.name, rawline.rstrip("\n"), self.conn, self.api_key, self.dry_run)
except Exception:
logging.exception("Follower for %s crashed; restarting in 2s", self.name)
time.sleep(2)
finally:
if self.proc and self.proc.poll() is None:
try:
self.proc.terminate()
except Exception:
pass
self.proc = None
time.sleep(1)
def stop(self):
self.stopped.set()
if self.proc and self.proc.poll() is None:
try:
self.proc.terminate()
except Exception:
pass
# --- Supervisor ---
def supervisor(dry_run=False):
ensure_dirs()
conn = init_db()
if not os.path.exists(ABUSEIPDB_KEY_FILE):
logging.warning("AbuseIPDB key file missing: %s", ABUSEIPDB_KEY_FILE)
api_key = None
else:
api_key = open(ABUSEIPDB_KEY_FILE).read().strip()
followers = {}
# reporter thread
def reporter_loop():
while True:
time.sleep(REPORT_FLUSH_INTERVAL)
flush_reports(api_key, dry_run=dry_run)
t = threading.Thread(target=reporter_loop, daemon=True)
t.start()
while True:
try:
p = subprocess.run([DOCKER_BIN, "ps", "--filter", "name=anubis", "--format", "{{.Names}}"],
capture_output=True, text=True)
if p.returncode != 0:
names = []
else:
names = [ln.strip() for ln in p.stdout.splitlines() if ln.strip()]
except Exception:
logging.exception("docker ps failed")
names = []
current = set(names)
existing = set(followers.keys())
for name in current - existing:
f = ContainerFollower(name, DOCKER_BIN, conn, api_key, dry_run)
followers[name] = f
f.start()
logging.info("Started follower for %s", name)
for name in existing - current:
logging.info("Stopping follower for %s", name)
followers[name].stop()
followers[name].join(timeout=5)
followers.pop(name, None)
time.sleep(SLEEP_REFRESH)
def main():
import argparse
parser = argparse.ArgumentParser()
parser.add_argument("--dry-run", action="store_true", help="Do not actually POST")
args = parser.parse_args()
try:
supervisor(dry_run=args.dry_run)
except KeyboardInterrupt:
logging.info("Shutting down")
sys.exit(0)
if __name__ == "__main__":
main()