First commit

This commit is contained in:
2026-07-20 11:05:07 +02:00
commit b592ab669f
159 changed files with 10294 additions and 0 deletions
@@ -0,0 +1,439 @@
"""Client MQTT vers le broker master (Phase 2.b/4).
Relaie chaque événement fail2ban publié sur le broker local (même topic
que banevents.mqtt_listen, fail2ban/+/jail) vers le broker master en mTLS,
sous fail2ban/<MQTT_MASTER_NODE_NAME>/ban. S'abonne en retour aux commandes
du master (fail2ban/<MQTT_MASTER_NODE_NAME>/action et
fail2ban/broadcast/action) et aux mises à jour du roster (Phase 3, voir
banevents.management.commands.master_listen).
Commandes reconnues (command_handlers ci-dessous) : SYNC_BAN, BAN_ALLPORTS,
WHITELIST, NOTIFY_ONLY, RATE_LIMIT, ESCALATE. Pour en ajouter une : écrire
une méthode `execute_<nom>(self, command: dict) -> None` qui termine
toujours par un appel à `self.report_outcome(cmd, ip, status,
detail='')`, puis l'ajouter au dict `self.command_handlers` dans
`handle()`. Toute commande reçue sans handler enregistré est journalisée
sans action.
Trois formes possibles pour une commande Phase 4 choisir AVANT
d'écrire du code (voir ROADMAP.md, Phase 4, "procédure d'ajout") :
1. Réaction automatique décidée par le master (comme SYNC_BAN/ESCALATE)
un handler ici + une règle `correlation_rule_<nom>` dans
master_listen.py.
2. Décision manuelle diffusée à tous les noeuds (comme
WHITELIST/BAN_ALLPORTS/RATE_LIMIT) un handler ici + ajout aux
MANUAL_COMMANDS de publish_command.py.
3. Action master-only, jamais exécutée par un noeud (comme REPORT_ABUSE)
PAS de handler ici du tout, branchement direct dans
publish_command.py.
Ce relais est un process indépendant de mqtt_listen : les deux s'abonnent
séparément au même topic local, chacun avec sa propre connexion/identité.
"""
import ipaddress
import json
import os
import ssl
import subprocess
import threading
import uuid
from typing import Any
import paho.mqtt.client as mqtt
import redis
from asgiref.sync import async_to_sync
from channels.layers import get_channel_layer
from django.conf import settings
from django.core.management.base import BaseCommand, CommandError
from banevents.consumers import GROUP_NAME
from banevents.models import NodeRegistry
from banevents.stats import ROSTER_REDIS_DB, ROSTER_REDIS_KEY, ROSTER_TTL_SECONDS
MASTER_SYNC_JAIL = 'master-sync'
MASTER_RATELIMIT_JAIL = 'master-ratelimit'
MASTER_ESCALATE_JAIL = 'master-escalate'
# Intervalle de heartbeat (Phase 3, supervision) : assez court pour qu'une
# coupure soit détectée en quelques minutes côté dashboard (seuil "hors
# ligne" = 3x cet intervalle, cf. views.py), assez long pour rester
# négligeable en trafic MQTT.
HEARTBEAT_INTERVAL_SECONDS = 60
class Command(BaseCommand):
help = 'Relaie les événements fail2ban locaux vers le broker master et applique ses commandes (SYNC_BAN).'
def handle(self, *args: Any, **options: Any) -> None:
if not settings.MQTT_MASTER_NODE_NAME:
raise CommandError(
'MQTT_MASTER_NODE_NAME est vide. Doit être identique au nom '
'passé à scripts/generate-node-cert.sh (le CN du certificat '
"client) : c'est ce que le master utilisera pour restreindre "
'ce noeud via ACL (fail2ban/<nom>/...).'
)
self.node_topic_prefix = f'fail2ban/{settings.MQTT_MASTER_NODE_NAME}'
# Identité fail2ban locale (uuid.getnode(), posée par
# f2b_mqtt_action_banisher.py dans chaque payload) — distincte de
# MQTT_MASTER_NODE_NAME (l'alias/CN mTLS) : sert uniquement à
# reconnaître "ce SYNC_BAN vient de ce noeud", pas à s'authentifier.
self.local_node_id = str(uuid.getnode())
# Communiqué au master via le heartbeat (voir send_heartbeat) pour
# alimenter automatiquement NodeRegistry.dashboard_url côté master
# (master_listen.py::handle_heartbeat) — os.environ direct, pas un
# Django setting : DASHBOARD_DOMAIN est propre à CE noeud (vide sur
# un noeud sans tableau de bord public exposé), même lecture directe
# que enrollment.py::join_command.
dashboard_domain = os.environ.get('DASHBOARD_DOMAIN', '')
self.dashboard_url = f'https://{dashboard_domain}' if dashboard_domain else ''
self.heartbeat_timer: threading.Timer | None = None
self.command_handlers = {
'SYNC_BAN': self.execute_sync_ban,
'BAN_ALLPORTS': self.execute_ban_allports,
'WHITELIST': self.execute_whitelist,
'NOTIFY_ONLY': self.execute_notify_only,
'RATE_LIMIT': self.execute_rate_limit,
'ESCALATE': self.execute_escalate,
}
# Cache du roster (Phase 3, cf. master_listen.py) — DB dédiée
# (db=2) : db=0 déjà pris par Channels (CHANNEL_LAYERS), db=1 par
# Celery (CELERY_BROKER_URL).
self.redis_client = redis.Redis(host=settings.REDIS_HOST, port=settings.REDIS_PORT, db=ROSTER_REDIS_DB)
self.local_client = mqtt.Client(
mqtt.CallbackAPIVersion.VERSION2, client_id=f'master-relay-local-{uuid.uuid4().hex[:8]}'
)
self.local_client.username_pw_set(settings.MQTT_BROKER_USERNAME, settings.MQTT_BROKER_PASSWORD)
self.local_client.on_connect = self.on_local_connect
self.local_client.on_message = self.on_local_message
self.master_client = mqtt.Client(
mqtt.CallbackAPIVersion.VERSION2, client_id=f'master-relay-master-{uuid.uuid4().hex[:8]}'
)
# ca_certs=None : le magasin de CA système (défaut Python) sait déjà
# vérifier le certificat serveur Let's Encrypt du master. Notre CA
# privée (scripts/master-ca-init.sh) ne sert qu'à ce que LE MASTER
# vérifie NOTRE certificat client (mTLS) — elle n'a rien à voir avec
# la vérification du certificat serveur par CE client, et la fournir
# ici écrase le magasin par défaut au lieu de s'y ajouter, faisant
# échouer la validation du certificat serveur (constaté :
# SSLCertVerificationError "unable to get local issuer certificate").
self.master_client.tls_set(
certfile=settings.MQTT_MASTER_CLIENT_CERT,
keyfile=settings.MQTT_MASTER_CLIENT_KEY,
tls_version=ssl.PROTOCOL_TLSv1_2,
)
self.master_client.on_connect = self.on_master_connect
self.master_client.on_message = self.on_master_message
# Callback dédié (pas on_master_message) : le roster est un flux
# séparé des commandes (action), pas la peine de faire cohabiter
# deux formats de message dans un seul handler générique.
self.master_client.message_callback_add('fail2ban/broadcast/roster', self.on_roster_message)
self.stdout.write(f'Connexion locale à {settings.MQTT_BROKER_HOST}:{settings.MQTT_BROKER_PORT}')
self.local_client.connect(settings.MQTT_BROKER_HOST, settings.MQTT_BROKER_PORT, keepalive=60)
self.local_client.loop_start()
# connect_async + retry_first_connection : si le master est injoignable
# au démarrage (réseau, certificat pas encore en place, ...),
# loop_forever() réessaie tout seul au lieu de faire planter toute la
# commande (y compris le relais local, qui fonctionnerait pourtant
# très bien sans lui). Sans retry_first_connection, un échec de LA
# toute première tentative fait remonter l'exception au lieu d'être
# réessayé (constaté : ConnectionRefusedError non rattrapée).
self.stdout.write(f'Connexion master à {settings.MQTT_MASTER_HOST}:{settings.MQTT_MASTER_PORT}')
self.master_client.connect_async(settings.MQTT_MASTER_HOST, settings.MQTT_MASTER_PORT, keepalive=60)
self.master_client.reconnect_delay_set(min_delay=1, max_delay=30)
self.master_client.loop_forever(retry_first_connection=True)
def on_local_connect(
self,
client: mqtt.Client,
userdata: Any,
flags: mqtt.ConnectFlags,
reason_code: mqtt.ReasonCode,
properties: Any = None,
) -> None:
self.stdout.write(self.style.SUCCESS(f'Connecté au broker local ({reason_code})'))
client.subscribe(settings.MQTT_TOPIC_SUBSCRIBE)
def on_local_message(self, client: mqtt.Client, userdata: Any, message: mqtt.MQTTMessage) -> None:
try:
payload = json.loads(message.payload.decode('utf-8'))
except (json.JSONDecodeError, UnicodeDecodeError) as e:
self.stderr.write(self.style.ERROR(f'Message local invalide : {e}'))
return
# Le topic doit correspondre au CN du certificat mTLS de ce process
# (MQTT_MASTER_NODE_NAME), jamais au "node" interne du payload
# (uuid.getnode() côté fail2ban, un simple identifiant matériel) : le
# master ACL chaque connexion sur son identité de certificat, donc
# un topic qui ne correspond pas à ce CN serait rejeté.
topic = f'{self.node_topic_prefix}/ban'
self.master_client.publish(topic, json.dumps(payload), qos=1)
self.stdout.write(f'Relayé vers le master : {topic} {payload.get("action")} {payload.get("ip")}')
def on_master_connect(
self,
client: mqtt.Client,
userdata: Any,
flags: mqtt.ConnectFlags,
reason_code: mqtt.ReasonCode,
properties: Any = None,
) -> None:
self.stdout.write(self.style.SUCCESS(f'Connecté au master ({reason_code})'))
client.subscribe(f'{self.node_topic_prefix}/action')
client.subscribe('fail2ban/broadcast/action')
client.subscribe('fail2ban/broadcast/roster')
# on_master_connect peut se redéclencher à chaque reconnexion :
# annuler toute chaîne de timer précédente pour ne pas en empiler
# plusieurs en parallèle (heartbeats en double).
if self.heartbeat_timer is not None:
self.heartbeat_timer.cancel()
self.send_heartbeat()
def send_heartbeat(self) -> None:
# Le segment de topic (fail2ban/<MQTT_MASTER_NODE_NAME>/heartbeat)
# est l'alias/CN mTLS, imposé par l'ACL (pattern write
# fail2ban/%u/heartbeat) — jamais le même identifiant que
# NodeRegistry.node_id (uuid.getnode() brut, cf. commentaire plus
# haut sur ces deux identités distinctes). L'id brut voyage donc
# dans le payload, exactement comme /ban le fait déjà.
self.master_client.publish(
f'{self.node_topic_prefix}/heartbeat',
json.dumps({'node': self.local_node_id, 'dashboard_url': self.dashboard_url}),
qos=0,
)
self.heartbeat_timer = threading.Timer(HEARTBEAT_INTERVAL_SECONDS, self.send_heartbeat)
self.heartbeat_timer.daemon = True
self.heartbeat_timer.start()
def report_outcome(self, cmd: str, ip: str, status: str, detail: str = '') -> None:
"""Diffuse l'issue d'une commande reçue du master vers ses deux
destinations : accusé de réception MQTT vers le master
(fail2ban/<noeud>/ack) et notification WebSocket vers ce dashboard
(même canal que les BanEvent, cf. signals.py) un SYNC_BAN réussi
finit par apparaître indirectement via le BanEvent qu'il déclenche,
mais un échec ou une commande non gérée ne laissaient jusqu'ici
aucune trace visible hors journalctl. Point d'entrée unique pour
toute nouvelle commande : chaque execute_<nom> doit terminer par un
appel ici plutôt que de publier séparément vers les deux canaux."""
payload = {'cmd': cmd, 'ip': ip, 'status': status, 'detail': detail}
self.master_client.publish(f'{self.node_topic_prefix}/ack', json.dumps(payload), qos=1)
channel_layer = get_channel_layer()
if channel_layer is not None:
async_to_sync(channel_layer.group_send)(
GROUP_NAME, {'type': 'command.notification', 'payload': payload},
)
def on_master_message(self, client: mqtt.Client, userdata: Any, message: mqtt.MQTTMessage) -> None:
try:
command = json.loads(message.payload.decode('utf-8'))
except (json.JSONDecodeError, UnicodeDecodeError) as e:
self.stderr.write(self.style.ERROR(f'Commande master invalide : {e}'))
return
cmd = command.get('cmd')
handler = self.command_handlers.get(cmd)
if handler is not None:
handler(command)
else:
# Phase 4 restante (ESCALATE, RATE_LIMIT, REPORT_ABUSE,
# REQUEST_GEOLOCATE) : pas de handler enregistré, journalisée
# seulement.
self.stdout.write(self.style.WARNING(f'Commande reçue du master (non gérée) : {command}'))
self.report_outcome(str(cmd or '?'), '', 'unhandled')
def _banip_via_jail(self, jail: str, ip: str) -> tuple[bool, str]:
"""Bannit `ip` via la jail `jail` (`fail2ban-client set <jail>
banip <ip>`) primitive partagée par toute commande qui se
résume à "bannir cette IP dans une jail dédiée à injection
manuelle" (master-sync/SYNC_BAN, master-sync/BAN_ALLPORTS,
master-ratelimit/RATE_LIMIT, master-escalate/ESCALATE) : seuls le
nom de la jail et l'origine de la décision changent. Retourne
(succès, detail message d'erreur si échec)."""
try:
result = subprocess.run(
['sudo', settings.FAIL2BAN_CLIENT_PATH, 'set', jail, 'banip', ip],
capture_output=True, text=True, timeout=10,
)
except (subprocess.SubprocessError, OSError) as e:
return False, str(e)
if result.returncode == 0:
return True, ''
return False, result.stderr.strip() or result.stdout.strip()
def _list_active_jails(self) -> list[str]:
"""Parse la sortie de `fail2ban-client status` ("Jail list: a, b, c")
pas d'option JSON native côté fail2ban-client pour cette info."""
try:
result = subprocess.run(
['sudo', settings.FAIL2BAN_CLIENT_PATH, 'status'],
capture_output=True, text=True, timeout=10, check=True,
)
except (subprocess.SubprocessError, OSError, subprocess.CalledProcessError) as e:
self.stderr.write(self.style.ERROR(f'Impossible de lister les jails actives : {e}'))
return []
for line in result.stdout.splitlines():
if 'Jail list:' in line:
return [j.strip() for j in line.split(':', 1)[1].split(',') if j.strip()]
return []
def execute_sync_ban(self, command: dict[str, Any]) -> None:
if str(command.get('origin_node', '')) == self.local_node_id:
self.stdout.write(f'SYNC_BAN ignoré (origine = ce noeud) : {command.get("ip")}')
return
try:
ip = str(ipaddress.ip_address(command.get('ip', '')))
except (ValueError, TypeError):
self.stderr.write(self.style.ERROR(f'SYNC_BAN : IP invalide reçue : {command.get("ip")!r}'))
self.report_outcome('SYNC_BAN', str(command.get('ip', '')), 'error', 'IP invalide')
return
ok, detail = self._banip_via_jail(MASTER_SYNC_JAIL, ip)
if ok:
self.stdout.write(self.style.SUCCESS(f'SYNC_BAN appliqué : {ip} (jail {MASTER_SYNC_JAIL})'))
self.report_outcome('SYNC_BAN', ip, 'ok')
else:
self.stderr.write(self.style.ERROR(f'SYNC_BAN échec pour {ip} : {detail}'))
self.report_outcome('SYNC_BAN', ip, 'error', detail)
def execute_ban_allports(self, command: dict[str, Any]) -> None:
"""Décision ciblée explicite (pas une synchro auto : pas de check
origin_node déclenchée à la main via publish_command.py, doit
s'appliquer sur tous les noeuds qui la reçoivent, y compris
l'origine si jamais elle en porte une)."""
try:
ip = str(ipaddress.ip_address(command.get('ip', '')))
except (ValueError, TypeError):
self.stderr.write(self.style.ERROR(f'BAN_ALLPORTS : IP invalide reçue : {command.get("ip")!r}'))
self.report_outcome('BAN_ALLPORTS', str(command.get('ip', '')), 'error', 'IP invalide')
return
ok, detail = self._banip_via_jail(MASTER_SYNC_JAIL, ip)
if ok:
self.stdout.write(self.style.SUCCESS(f'BAN_ALLPORTS appliqué : {ip} (jail {MASTER_SYNC_JAIL})'))
self.report_outcome('BAN_ALLPORTS', ip, 'ok')
else:
self.stderr.write(self.style.ERROR(f'BAN_ALLPORTS échec pour {ip} : {detail}'))
self.report_outcome('BAN_ALLPORTS', ip, 'error', detail)
def execute_whitelist(self, command: dict[str, Any]) -> None:
"""Ajoute l'IP à ignoreip sur toutes les jails actives + la
débannit si elle l'était. Limitation connue : `addignoreip` est un
réglage en mémoire, non persistant perdu au prochain redémarrage
de fail2ban.service (install.sh ne touche pas jail.local donc un
`install.sh install` ne l'efface pas, mais un restart manuel du
service, oui). Persister ça proprement (fichier dédié non écrasé
par install.sh, relu au démarrage) reste à faire séparément."""
try:
ip = str(ipaddress.ip_address(command.get('ip', '')))
except (ValueError, TypeError):
self.stderr.write(self.style.ERROR(f'WHITELIST : IP invalide reçue : {command.get("ip")!r}'))
self.report_outcome('WHITELIST', str(command.get('ip', '')), 'error', 'IP invalide')
return
jails = self._list_active_jails()
if not jails:
self.report_outcome('WHITELIST', ip, 'error', 'aucune jail active trouvée')
return
for jail in jails:
subprocess.run(
['sudo', settings.FAIL2BAN_CLIENT_PATH, 'set', jail, 'addignoreip', ip],
capture_output=True, text=True, timeout=10,
)
# Échec ignoré : l'IP peut légitimement ne pas être bannie
# dans cette jail précise.
subprocess.run(
['sudo', settings.FAIL2BAN_CLIENT_PATH, 'set', jail, 'unbanip', ip],
capture_output=True, text=True, timeout=10,
)
self.stdout.write(self.style.SUCCESS(f'WHITELIST appliqué : {ip} ({len(jails)} jails)'))
self.report_outcome('WHITELIST', ip, 'ok')
def execute_notify_only(self, command: dict[str, Any]) -> None:
"""Aucune action système — la visibilité "alerte admin" demandée
pour cette commande (cf. ROADMAP.md) est déjà couverte par le
toast que report_outcome déclenche côté dashboard."""
ip = str(command.get('ip', ''))
self.stdout.write(f'NOTIFY_ONLY reçu : {ip}')
self.report_outcome('NOTIFY_ONLY', ip, 'ok')
def execute_rate_limit(self, command: dict[str, Any]) -> None:
"""Décision ciblée explicite (manuelle, via publish_command.py) :
throttle au lieu d'un blocage total, jail dédiée master-ratelimit
(banaction f2b-iptables-hashlimit)."""
try:
ip = str(ipaddress.ip_address(command.get('ip', '')))
except (ValueError, TypeError):
self.stderr.write(self.style.ERROR(f'RATE_LIMIT : IP invalide reçue : {command.get("ip")!r}'))
self.report_outcome('RATE_LIMIT', str(command.get('ip', '')), 'error', 'IP invalide')
return
ok, detail = self._banip_via_jail(MASTER_RATELIMIT_JAIL, ip)
if ok:
self.stdout.write(self.style.SUCCESS(f'RATE_LIMIT appliqué : {ip} (jail {MASTER_RATELIMIT_JAIL})'))
self.report_outcome('RATE_LIMIT', ip, 'ok')
else:
self.stderr.write(self.style.ERROR(f'RATE_LIMIT échec pour {ip} : {detail}'))
self.report_outcome('RATE_LIMIT', ip, 'error', detail)
def execute_escalate(self, command: dict[str, Any]) -> None:
"""Auto-publiée par master_listen.py (correlation_rule_escalate)
quand une IP est bannie indépendamment sur plusieurs noeuds
distincts. Contrairement à SYNC_BAN, pas de check origin_node :
doit s'appliquer aussi sur le(s) noeud(s) d'origine (leur ban
initial n'est pas permanent, master-escalate est une jail
séparée, pas de conflit à réappliquer)."""
try:
ip = str(ipaddress.ip_address(command.get('ip', '')))
except (ValueError, TypeError):
self.stderr.write(self.style.ERROR(f'ESCALATE : IP invalide reçue : {command.get("ip")!r}'))
self.report_outcome('ESCALATE', str(command.get('ip', '')), 'error', 'IP invalide')
return
ok, detail = self._banip_via_jail(MASTER_ESCALATE_JAIL, ip)
if ok:
self.stdout.write(self.style.SUCCESS(f'ESCALATE appliqué : {ip} (jail {MASTER_ESCALATE_JAIL}, permanent)'))
self.report_outcome('ESCALATE', ip, 'ok')
else:
self.stderr.write(self.style.ERROR(f'ESCALATE échec pour {ip} : {detail}'))
self.report_outcome('ESCALATE', ip, 'error', detail)
def on_roster_message(self, client: mqtt.Client, userdata: Any, message: mqtt.MQTTMessage) -> None:
"""Roster périodique publié par master_listen.py (Phase 3) : les
noeuds connus se propagent dans NodeRegistry (déjà le modèle lu par
la sidebar un client voit ainsi les autres noeuds, pas que
lui-même) ; les stats globales jail/pays (pas des lignes de table,
des agrégats transitoires) vont en cache Redis, lu par views.py
avec repli sur le calcul local si absent/périmé."""
try:
roster = json.loads(message.payload.decode('utf-8'))
except (json.JSONDecodeError, UnicodeDecodeError) as e:
self.stderr.write(self.style.ERROR(f'Roster invalide : {e}'))
return
for node in roster.get('nodes', []):
node_id = node.get('node_id')
if not node_id:
continue
# filter().update() plutôt que update_or_create() : last_seen a
# auto_now=True, qui écraserait la vraie valeur du master par
# "maintenant" si on passait par .save() (déclenché en interne
# par update_or_create) — .update() fait une requête SQL directe,
# sans repasser par la logique auto_now du champ.
fields = {
'alias': node.get('alias', ''),
'country': node.get('country', ''),
'dashboard_url': node.get('dashboard_url', ''),
}
if node.get('last_seen'):
fields['last_seen'] = node['last_seen']
if not NodeRegistry.objects.filter(node_id=node_id).update(**fields):
NodeRegistry.objects.get_or_create(node_id=node_id, defaults=fields)
self.redis_client.set(
ROSTER_REDIS_KEY,
json.dumps({'top_jails': roster.get('top_jails', []), 'top_countries': roster.get('top_countries', [])}),
ex=ROSTER_TTL_SECONDS,
)