"""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//ban. S'abonne en retour aux commandes du master (fail2ban//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_(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_` 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//...).' ) 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//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//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_ 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 banip `) — 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, )