Every audit_logs insert is mirrored as an RFC-5424 syslog message with a per-event severity (login failures/lockouts/access-denied → warning, security alerts → alert, config/deletes → notice, else info). A SQLAlchemy after_insert hook enqueues onto a bounded queue drained by one daemon worker (never blocks the request/flush; bursty syncs drain sequentially). Config cached 30s; TCP keeps a persistent socket with reconnect. Disabled by default → no-op until enabled. Admin-only config card (host/port/UDP-TCP/facility) with a Test button; POST /api/v1/settings/syslog/test sends a probe. No TLS yet (UDP+TCP).
226 lines
7.9 KiB
Python
226 lines
7.9 KiB
Python
"""
|
|
Forward audit-log events to an external syslog server (SIEM ingestion).
|
|
|
|
Every row inserted into `audit_logs` is mirrored as an RFC-5424 syslog message
|
|
over UDP or TCP, with a per-event severity so a SIEM can decode/rule/alert on
|
|
them (login failures, lockouts, access-denied, config changes, ...).
|
|
|
|
Design (ponytail):
|
|
- One bounded queue + one daemon worker thread. The SQLAlchemy after_insert
|
|
hook only enqueues (never blocks the request / DB flush). Bursty syncs that
|
|
write thousands of VULNERABILITY_DETECTED rows drain sequentially through the
|
|
single worker instead of spawning a thread per event.
|
|
- Config (`syslog_config` setting) is cached for 30 s so the hot path never
|
|
hits the DB. TCP keeps a persistent socket and reconnects on failure.
|
|
- Disabled by default → the hook is a cheap no-op until an admin turns it on.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import logging
|
|
import queue
|
|
import socket
|
|
import threading
|
|
import time
|
|
from datetime import datetime
|
|
from typing import Optional
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
APP_NAME = "truevuln"
|
|
_QUEUE: "queue.Queue[dict]" = queue.Queue(maxsize=10000)
|
|
_worker_started = False
|
|
_worker_lock = threading.Lock()
|
|
|
|
# Config cache
|
|
_cfg_cache: dict = {"ts": 0.0, "cfg": None}
|
|
_CFG_TTL = 30.0
|
|
|
|
# RFC-5424 severities
|
|
_SEV_EMERG, _SEV_ALERT, _SEV_CRIT, _SEV_ERR, _SEV_WARN, _SEV_NOTICE, _SEV_INFO, _SEV_DEBUG = range(8)
|
|
|
|
|
|
def _severity_for(event_type: str) -> int:
|
|
"""Map an AuditEventType name to a syslog severity so a SIEM can prioritise."""
|
|
e = (event_type or "").upper()
|
|
if "SECURITY_ALERT" in e or "PERMISSION_ESCALATION" in e:
|
|
return _SEV_ALERT
|
|
if "FAILED" in e or "DENIED" in e or "LOCK" in e:
|
|
return _SEV_WARN
|
|
if "DELETED" in e or "DEACTIVATED" in e or "DISABLED" in e or "CONFIG_CHANGE" in e:
|
|
return _SEV_NOTICE
|
|
return _SEV_INFO
|
|
|
|
|
|
def _load_config() -> Optional[dict]:
|
|
"""Cached read of the `syslog_config` setting. Returns None when disabled/
|
|
unset. Shape: {enabled, host, port, protocol('udp'|'tcp'), facility(int)}."""
|
|
now = time.monotonic()
|
|
if now - _cfg_cache["ts"] < _CFG_TTL:
|
|
return _cfg_cache["cfg"]
|
|
cfg = None
|
|
try:
|
|
from app.database import SessionLocal
|
|
from app.models.setting import Setting
|
|
db = SessionLocal()
|
|
try:
|
|
row = db.query(Setting).filter(Setting.key == "syslog_config").first()
|
|
if row and row.value:
|
|
parsed = json.loads(row.value)
|
|
if parsed.get("enabled") and parsed.get("host"):
|
|
cfg = {
|
|
"host": str(parsed["host"]).strip(),
|
|
"port": int(parsed.get("port") or 514),
|
|
"protocol": str(parsed.get("protocol") or "udp").lower(),
|
|
"facility": int(parsed.get("facility") if parsed.get("facility") is not None else 16),
|
|
}
|
|
finally:
|
|
db.close()
|
|
except Exception as e: # never let config trouble break the app
|
|
logger.debug("syslog config load failed: %s", e)
|
|
cfg = None
|
|
_cfg_cache["ts"] = now
|
|
_cfg_cache["cfg"] = cfg
|
|
return cfg
|
|
|
|
|
|
def _build_message(evt: dict, facility: int) -> bytes:
|
|
sev = _severity_for(evt.get("event_type", ""))
|
|
pri = facility * 8 + sev
|
|
ts = (evt.get("timestamp") or datetime.now()).astimezone().isoformat()
|
|
host = socket.gethostname() or "-"
|
|
msgid = (evt.get("event_type") or "AUDIT")[:32]
|
|
# MSG: human-readable + a few key=value fields for easy SIEM extraction.
|
|
parts = [evt.get("event_description") or ""]
|
|
if evt.get("user_id") is not None:
|
|
parts.append(f"user_id={evt['user_id']}")
|
|
if evt.get("resource_type"):
|
|
parts.append(f"resource={evt['resource_type']}:{evt.get('resource_id') or ''}")
|
|
if evt.get("ip_address"):
|
|
parts.append(f"src_ip={evt['ip_address']}")
|
|
msg = " ".join(p for p in parts if p)
|
|
line = f"<{pri}>1 {ts} {host} {APP_NAME} - {msgid} - {msg}"
|
|
return line.encode("utf-8", "replace")
|
|
|
|
|
|
class _Sender:
|
|
"""Holds a persistent TCP socket (reconnect on failure) or a UDP socket."""
|
|
|
|
def __init__(self):
|
|
self._sock: Optional[socket.socket] = None
|
|
self._key: tuple = ()
|
|
|
|
def _ensure(self, host: str, port: int, proto: str):
|
|
key = (host, port, proto)
|
|
if self._sock is not None and key == self._key:
|
|
return
|
|
self.close()
|
|
self._key = key
|
|
if proto == "tcp":
|
|
s = socket.create_connection((host, port), timeout=5)
|
|
else:
|
|
s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
|
self._sock = s
|
|
|
|
def send(self, data: bytes, host: str, port: int, proto: str):
|
|
self._ensure(host, port, proto)
|
|
assert self._sock is not None
|
|
if proto == "tcp":
|
|
# RFC 6587 non-transparent (LF) framing — accepted by rsyslog/syslog-ng.
|
|
self._sock.sendall(data + b"\n")
|
|
else:
|
|
self._sock.sendto(data, (host, port))
|
|
|
|
def close(self):
|
|
if self._sock is not None:
|
|
try:
|
|
self._sock.close()
|
|
except Exception:
|
|
pass
|
|
self._sock = None
|
|
self._key = ()
|
|
|
|
|
|
def _worker():
|
|
sender = _Sender()
|
|
while True:
|
|
evt = _QUEUE.get()
|
|
try:
|
|
cfg = _load_config()
|
|
if not cfg:
|
|
continue
|
|
data = _build_message(evt, cfg["facility"])
|
|
try:
|
|
sender.send(data, cfg["host"], cfg["port"], cfg["protocol"])
|
|
except Exception as e:
|
|
sender.close() # force reconnect next time
|
|
logger.debug("syslog send failed (%s:%s/%s): %s",
|
|
cfg["host"], cfg["port"], cfg["protocol"], e)
|
|
finally:
|
|
_QUEUE.task_done()
|
|
|
|
|
|
def _ensure_worker():
|
|
global _worker_started
|
|
if _worker_started:
|
|
return
|
|
with _worker_lock:
|
|
if _worker_started:
|
|
return
|
|
threading.Thread(target=_worker, name="syslog-forwarder", daemon=True).start()
|
|
_worker_started = True
|
|
|
|
|
|
def enqueue(evt: dict) -> None:
|
|
"""Best-effort: drop the event rather than block or raise if the queue is
|
|
full or the worker can't start."""
|
|
try:
|
|
_ensure_worker()
|
|
_QUEUE.put_nowait(evt)
|
|
except queue.Full:
|
|
pass
|
|
except Exception as e:
|
|
logger.debug("syslog enqueue failed: %s", e)
|
|
|
|
|
|
def send_test(cfg: dict) -> tuple[bool, str]:
|
|
"""Send a one-off test message with an explicit config (admin 'Test' button)."""
|
|
try:
|
|
facility = int(cfg.get("facility") if cfg.get("facility") is not None else 16)
|
|
evt = {"event_type": "SECURITY_ALERT", "event_description": "TrueVuln syslog test message",
|
|
"timestamp": datetime.now()}
|
|
data = _build_message(evt, facility)
|
|
s = _Sender()
|
|
try:
|
|
s.send(data, str(cfg["host"]).strip(), int(cfg.get("port") or 514),
|
|
str(cfg.get("protocol") or "udp").lower())
|
|
finally:
|
|
s.close()
|
|
return True, "sent"
|
|
except Exception as e:
|
|
return False, str(e)
|
|
|
|
|
|
def register_audit_listener() -> None:
|
|
"""Hook every AuditLog insert → enqueue. Called once at startup."""
|
|
from sqlalchemy import event
|
|
from app.models.audit_log import AuditLog
|
|
|
|
@event.listens_for(AuditLog, "after_insert")
|
|
def _after_insert(mapper, connection, target): # noqa: ARG001
|
|
# Only scalar columns here — relationships would emit SQL mid-flush.
|
|
try:
|
|
enqueue({
|
|
"event_type": target.event_type.value if getattr(target, "event_type", None) else None,
|
|
"event_description": target.event_description,
|
|
"user_id": target.user_id,
|
|
"resource_type": target.resource_type,
|
|
"resource_id": target.resource_id,
|
|
"ip_address": target.ip_address,
|
|
"timestamp": target.timestamp,
|
|
})
|
|
except Exception:
|
|
pass
|
|
|
|
logger.info("syslog: audit-log forwarder registered")
|