# Fichier: python_cheats/cheatsheets/rabbitmq.txt
# RabbitMQ - Guide Complet pour Grands Débutants
# Comprendre RabbitMQ de A à Z : Architecture, Fonctionnement, Pratique

"""
[OK] INTRODUCTION : POURQUOI CE GUIDE ?

# Tu débutes avec les systèmes de messaging et RabbitMQ ?
# Ce guide est fait pour TOI !

# Objectif :
# Te montrer EXACTEMENT comment RabbitMQ fonctionne "sous le capot"
# Comprendre ce qui se passe quand tu publies un message
# Découvrir TOUS les composants et leur rôle précis
# Maîtriser les patterns de messaging
# Aucun concept ne sera laissé dans le flou !

# Plan du guide :
# 1. Qu'est-ce qu'un Message Broker et pourquoi RabbitMQ ?
# 2. Concepts fondamentaux (Producteur, Consommateur, Queue, Exchange)
# 3. Installation et fichiers créés (détail complet)
# 4. Architecture interne de RabbitMQ
# 5. Du message à la livraison : voyage complet
# 6. Types d'Exchanges et Routing
# 7. Patterns de Messaging (Point-to-Point, Pub/Sub, RPC, etc.)
# 8. Persistence, Durabilité et Fiabilité
# 9. Clustering et Haute Disponibilité
# 10. Performance et Optimisations
# 11. Sécurité et Monitoring
# 12. Premiers pas pratiques
# 13. Cas d'usage réels et bonnes pratiques
# 14. Comparaison avec Kafka, Redis, ActiveMQ
"""


# [OK] PARTIE 1 : QU'EST-CE QU'UN MESSAGE BROKER ET POURQUOI RABBITMQ ?

"""
┌────────────────────────────────────────────────────────────────────────┐
│               MESSAGE BROKER : UNE NÉCESSITÉ DANS LE MONDE DISTRIBUÉ   │
└────────────────────────────────────────────────────────────────────────┘

CONTEXTE HISTORIQUE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

AVANT LES MESSAGE BROKERS : COMMUNICATION DIRECTE
-> Service A appelle directement Service B via HTTP/REST
-> Couplage fort entre services
-> Synchrone : A attend la réponse de B
-> Pas de résilience : si B est down, A échoue

PROBLÈMES APPARUS AVEC LES MICROSERVICES :
[X] Couplage fort
   -> Service A doit connaître l'adresse de B
   -> Si B change, A doit être modifié
   
[X] Synchrone = Bloquant
   -> A attend que B réponde
   -> Cascades de timeouts
   -> Performance dégradée

[X] Pas de résilience
   -> Si B est down, les requêtes de A échouent
   -> Pas de retry automatique
   -> Perte de messages

[X] Scalabilité limitée
   -> Difficile d'ajouter des instances
   -> Load balancing complexe
   -> Pas de buffer pour les pics de charge

[X] Complexité du code
   -> Gestion des timeouts
   -> Retry logic partout
   -> Circuit breakers nécessaires


NAISSANCE DES MESSAGE BROKERS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Un MESSAGE BROKER = Intermédiaire entre services

Service A -> [Message Broker] -> Service B

AVANTAGES :
[OK] Découplage : A ne connaît pas B
[OK] Asynchrone : A n'attend pas B
[OK] Résilience : Buffer si B est down
[OK] Scalabilité : Multiple instances de B
[OK] Simplicité : Le broker gère la complexité


PRINCIPAUX MESSAGE BROKERS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. RabbitMQ (2007)
   -> Basé sur AMQP 0-9-1
   -> Erlang
   -> Routing complexe
   -> Haute fiabilité

2. Apache Kafka (2011)
   -> Log distribué
   -> Très haute performance
   -> Streaming de données
   -> Pas un broker traditionnel

3. Redis Pub/Sub
   -> Très rapide
   -> En mémoire
   -> Pas de persistence par défaut
   -> Simple

4. ActiveMQ
   -> Java
   -> JMS compliant
   -> Mature

5. AWS SQS / SNS
   -> Cloud managé
   -> Pay-per-use


ANALOGIE : SYSTÈME POSTAL
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

SANS MESSAGE BROKER (Livraison directe) :
┌────────────────────────────────────────────────────────────────────┐
│ [POSTBOX] LIVRAISON DIRECTE EN MAIN PROPRE                                │
│                                                                    │
│ Expéditeur A veut envoyer un colis à Destinataire B               │
│                                                                    │
│ Processus :                                                        │
│ -> A doit connaître l'adresse exacte de B                          │
│ -> A doit se déplacer jusqu'à B                                    │
│ -> A attend que B soit présent pour recevoir                       │
│ -> Si B n'est pas là, A doit revenir                               │
│                                                                    │
│ Problèmes :                                                        │
│ [X] A perd du temps à attendre                                     │
│ [X] Si B déménage, A doit trouver la nouvelle adresse             │
│ [X] A ne peut rien faire d'autre pendant la livraison              │
│ [X] Si B est absent, le colis est perdu                            │
└────────────────────────────────────────────────────────────────────┘


AVEC MESSAGE BROKER (Poste) :
┌────────────────────────────────────────────────────────────────────┐
│ [EMAIL] SYSTÈME POSTAL                                                  │
│                                                                    │
│ Expéditeur A -> [Bureau de Poste] -> Destinataire B                 │
│                                                                    │
│ Processus :                                                        │
│ -> A dépose le colis au bureau de poste                            │
│ -> A part immédiatement (asynchrone)                               │
│ -> La poste stocke le colis                                        │
│ -> La poste livre à B quand B est disponible                       │
│ -> Si B est absent, la poste réessaye                              │
│                                                                    │
│ Avantages :                                                        │
│ [OK] A n'a pas besoin de connaître l'adresse exacte de B           │
│ [OK] A ne perd pas de temps                                         │
│ [OK] Si B déménage, il suffit d'informer la poste                   │
│ [OK] Le colis est stocké en sécurité                                │
│ [OK] Livraison garantie même si B est temporairement absent         │
│ [OK] Plusieurs postiers peuvent livrer en parallèle (scalabilité)   │
└────────────────────────────────────────────────────────────────────┘


QU'EST-CE QUE RABBITMQ ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RabbitMQ = "Rabbit" (lapin) + "MQ" (Message Queue)
-> Le lapin est connu pour sa rapidité et sa reproduction (scalabilité)

HISTORIQUE :
-> Créé en 2007 par Rabbit Technologies
-> Racheté par SpringSource (2010), puis VMware, puis Pivotal
-> Open Source (Mozilla Public License)
-> Écrit en Erlang (langage conçu pour la téléphonie, ultra-résilient)
-> Version actuelle : 3.13.x (décembre 2024)

CARACTÉRISTIQUES PRINCIPALES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] AMQP 0-9-1 (Advanced Message Queuing Protocol)
   -> Standard ouvert de messaging
   -> Interopérabilité entre langages et plateformes
   
[OK] Routing flexible
   -> Direct, Fanout, Topic, Headers exchanges
   -> Routing complexe basé sur des règles
   
[OK] Fiabilité
   -> Acknowledgments (ACKs)
   -> Persistence des messages
   -> Clustering pour haute disponibilité
   
[OK] Plugins riches
   -> Management UI (interface web)
   -> MQTT, STOMP (autres protocoles)
   -> Federation, Shovel (inter-cluster)
   
[OK] Multi-protocole
   -> AMQP 0-9-1 (principal)
   -> AMQP 1.0 (via plugin)
   -> MQTT (IoT)
   -> STOMP (WebSocket)
   -> HTTP (REST API)

[OK] Écosystème mature
   -> Clients officiels pour tous les langages
   -> Documentation exhaustive
   -> Large communauté
   -> Production-ready


QUI UTILISE RABBITMQ ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Reddit (messaging backend)
-> T-Mobile (télécommunications)
-> Runtastic (fitness tracking)
-> 9GAG (social media)
-> Bloomberg (finance)
-> Mozilla (Firefox telemetry)
-> Klarna (payments)
-> Accenture (enterprise)


RABBITMQ vs AUTRES MESSAGE BROKERS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

┌──────────────────┬─────────────┬────────────┬──────────────┬─────────┐
│  Caractéristique │  RabbitMQ   │   Kafka    │  Redis Pub/Sub│ ActiveMQ│
├──────────────────┼─────────────┼────────────┼──────────────┼─────────┤
│ Protocole        │ AMQP 0-9-1  │ Propriétaire│ Redis Protocol│ JMS    │
│ Langage          │ Erlang      │ Scala/Java │ C            │ Java    │
│ Performance      │ 10-50K msg/s│ 1M+ msg/s  │ 100K+ msg/s  │ 10K/s   │
│ Persistence      │ Oui (opt-in)│ Oui (natif)│ Non par défaut│ Oui    │
│ Routing          │ Très flexible│ Basique   │ Basique      │ Flexible│
│ Use case         │ Microservices│ Streaming │ Cache/Pub-Sub│ JEE apps│
│ Complexité       │ Moyenne     │ Haute      │ Faible       │ Moyenne │
│ Ordre garanti    │ Oui (queue) │ Oui (part.)│ Non          │ Oui     │
│ Replay messages  │ Non (limité)│ Oui        │ Non          │ Non     │
│ Latence          │ ~1-10 ms    │ ~10-50 ms  │ <1 ms        │ ~5-20ms │
│ Clustering       │ Oui         │ Natif      │ Oui (Sentinel)│ Oui    │
└──────────────────┴─────────────┴────────────┴──────────────┴─────────┘


QUAND UTILISER RABBITMQ ?
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] UTILISER RABBITMQ SI :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[OK] Architecture microservices
[OK] Besoin de routing complexe
[OK] Messages < 1 MB
[OK] Ordre des messages important (par queue)
[OK] Fiabilité critique (ACKs, persistence)
[OK] Découplage services
[OK] Task queues (workers)
[OK] Event-driven architecture
[OK] RPC (Remote Procedure Call)
[OK] Publish/Subscribe patterns

Exemples :
-> E-commerce : commande -> paiement -> expédition -> notification
-> Traitement d'images : upload -> resize -> watermark -> CDN
-> Notifications : événement -> email + SMS + push
-> Task queues : jobs longs en arrière-plan


[X] NE PAS UTILISER RABBITMQ SI :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[X] Streaming de données massif (utiliser Kafka)
[X] Messages > 1 MB (utiliser S3 + référence)
[X] Replay des messages essentiel (utiliser Kafka)
[X] Analytics temps réel (utiliser Kafka)
[X] Latence ultra-faible (<1ms) (utiliser Redis)
[X] Simple cache (utiliser Redis)
[X] Communication synchrone simple (utiliser HTTP/gRPC)


RABBITMQ vs KAFKA : DIFFÉRENCES FONDAMENTALES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RABBITMQ = SMART BROKER, DUMB CONSUMER
-> Le broker fait le routing intelligent
-> Le consommateur reçoit juste les messages
-> Messages supprimés après consommation (par défaut)
-> Push model : broker pousse vers consommateur

KAFKA = DUMB BROKER, SMART CONSUMER
-> Le broker stocke tout dans un log
-> Le consommateur décide quoi lire et où
-> Messages conservés (rétention configurable)
-> Pull model : consommateur tire du broker

ANALOGIE :
RabbitMQ = Serveur de restaurant qui apporte les plats
Kafka = Buffet où vous vous servez vous-même
"""


# [OK] PARTIE 2 : CONCEPTS FONDAMENTAUX

"""
┌────────────────────────────────────────────────────────────────────────┐
│              CONCEPTS DE BASE RABBITMQ                                 │
└────────────────────────────────────────────────────────────────────────┘

ARCHITECTURE SIMPLIFIÉE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

┌──────────────┐       ┌──────────────────────────────┐       ┌──────────────┐
│              │       │      RABBITMQ BROKER         │       │              │
│  PRODUCER    │──────[BLACK_RIGHT-POINTING_TRIANGLE]│                              │──────[BLACK_RIGHT-POINTING_TRIANGLE]│  CONSUMER    │
│ (Producteur) │       │  ┌─────────┐   ┌─────────┐  │       │(Consommateur)│
│              │       │  │EXCHANGE │──[BLACK_RIGHT-POINTING_TRIANGLE]│ QUEUE   │  │       │              │
│  Publie des  │       │  │         │   │         │  │       │  Consomme    │
│  messages    │       │  └─────────┘   └─────────┘  │       │  messages    │
│              │       │       ^             │        │       │              │
│              │       │       └─[BINDING]───┘        │       │              │
└──────────────┘       └──────────────────────────────┘       └──────────────┘


COMPOSANTS PRINCIPAUX :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. PRODUCER (Producteur)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RÔLE :
-> Application qui ENVOIE des messages au broker
-> Ne connaît PAS les queues (généralement)
-> Envoie à un EXCHANGE

ANALOGIE :
-> L'expéditeur qui dépose une lettre à la poste
-> Il spécifie l'adresse (routing key) mais ne livre pas lui-même

EXEMPLE :
# Python (pika)
import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Publier un message
channel.basic_publish(
    exchange='',              # Exchange par défaut
    routing_key='hello',      # Nom de la queue
    body='Hello World!'       # Message
)

print("Message envoyé [OK]")
connection.close()


PROPRIÉTÉS D'UN MESSAGE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
-> Body (corps) : le contenu du message
-> Properties (propriétés) : métadonnées

Properties importantes :
- content_type : "application/json", "text/plain", etc.
- content_encoding : "utf-8", "gzip", etc.
- delivery_mode : 1 (transient) ou 2 (persistent)
- priority : 0-9 (priorité du message)
- correlation_id : pour RPC (associer requête/réponse)
- reply_to : queue pour la réponse (RPC)
- expiration : TTL en millisecondes
- message_id : identifiant unique
- timestamp : horodatage
- type : type de message (custom)
- user_id : utilisateur émetteur
- app_id : application émettrice
- headers : métadonnées custom (dict)

EXEMPLE AVEC PROPERTIES :
channel.basic_publish(
    exchange='orders',
    routing_key='order.created',
    body=json.dumps({
        'order_id': 12345,
        'customer_id': 67890,
        'total': 150.00
    }),
    properties=pika.BasicProperties(
        content_type='application/json',
        delivery_mode=2,          # Persistent
        priority=5,
        correlation_id=str(uuid.uuid4()),
        timestamp=int(time.time()),
        headers={'source': 'web-app', 'version': '1.0'}
    )
)


2. EXCHANGE (Échangeur)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RÔLE :
-> Point d'entrée des messages
-> ROUTE les messages vers les queues
-> Basé sur des RULES (bindings + routing keys)

ANALOGIE :
-> Le centre de tri postal
-> Reçoit les lettres et décide vers quel bureau local les envoyer

TYPES D'EXCHANGES (détail dans Partie 6) :
1. Direct : routage exact par routing key
2. Fanout : broadcast à toutes les queues liées
3. Topic : routage par pattern (wildcards)
4. Headers : routage par headers du message

DÉCLARATION D'UN EXCHANGE :
channel.exchange_declare(
    exchange='orders',        # Nom de l'exchange
    exchange_type='topic',    # Type : direct, fanout, topic, headers
    durable=True,            # Survit au redémarrage
    auto_delete=False,       # Ne se supprime pas automatiquement
    internal=False,          # Accessible de l'extérieur
    arguments={}             # Arguments additionnels
)


DEFAULT EXCHANGE :
-> Exchange sans nom ('')
-> Type direct
-> Routing key = nom de la queue
-> Toutes les queues y sont liées automatiquement

EXEMPLE :
channel.basic_publish(
    exchange='',           # Default exchange
    routing_key='hello',   # Nom de la queue directement
    body='Hello!'
)


3. BINDING (Liaison)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RÔLE :
-> LIE un exchange à une queue
-> Définit les RÈGLES de routage

ANALOGIE :
-> Les règles de tri du centre postal
-> "Toutes les lettres avec code postal 75xxx vont à Paris"

DÉCLARATION D'UN BINDING :
channel.queue_bind(
    queue='orders_processing',      # Queue destination
    exchange='orders',              # Exchange source
    routing_key='order.created'     # Routing key (règle)
)

PLUSIEURS BINDINGS POSSIBLES :
# Queue orders_processing reçoit plusieurs types de messages
channel.queue_bind(queue='orders_processing', exchange='orders', routing_key='order.created')
channel.queue_bind(queue='orders_processing', exchange='orders', routing_key='order.updated')
channel.queue_bind(queue='orders_processing', exchange='orders', routing_key='order.paid')


4. QUEUE (File d'attente)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RÔLE :
-> STOCKE les messages
-> Ordre FIFO (First In First Out) par défaut
-> Un message ne peut être dans qu'UNE SEULE queue

ANALOGIE :
-> La boîte aux lettres du destinataire
-> Les lettres attendent d'être lues

DÉCLARATION D'UNE QUEUE :
channel.queue_declare(
    queue='hello',              # Nom de la queue
    durable=True,              # Survit au redémarrage du broker
    exclusive=False,           # Accessible par plusieurs connexions
    auto_delete=False,         # Ne se supprime pas quand vide
    arguments={
        'x-message-ttl': 60000,           # TTL des messages (60s)
        'x-max-length': 10000,            # Nombre max de messages
        'x-max-length-bytes': 1048576,    # Taille max (1 MB)
        'x-overflow': 'drop-head',        # drop-head, reject-publish
        'x-dead-letter-exchange': 'dlx',  # Exchange pour messages morts
        'x-dead-letter-routing-key': 'dead',
        'x-max-priority': 10,             # Support des priorités
        'x-queue-mode': 'lazy'            # Mode lazy (disque)
    }
)


PROPRIÉTÉS DES QUEUES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

durable = True :
-> Queue survit au redémarrage de RabbitMQ
-> Metadata persisté sur disque
-> MAIS messages peuvent être perdus (voir delivery_mode)

exclusive = True :
-> Queue accessible uniquement par la connexion qui l'a créée
-> Supprimée automatiquement à la fermeture de la connexion
-> Utile pour RPC (queues temporaires)

auto_delete = True :
-> Queue supprimée quand plus aucun consommateur
-> Utile pour des workers temporaires

x-message-ttl (Time To Live) :
-> Durée de vie des messages en millisecondes
-> Messages expirés sont supprimés ou DLQ

x-max-length :
-> Nombre maximum de messages dans la queue
-> Si dépassé : x-overflow définit l'action

x-overflow :
-> drop-head : supprime les plus anciens messages
-> reject-publish : refuse les nouveaux messages
-> reject-publish-dlx : envoie vers DLX

x-dead-letter-exchange (DLX) :
-> Exchange vers lequel envoyer les messages "morts"
-> Messages morts = rejetés, expirés, ou TTL dépassé
-> Permet de traiter les erreurs

x-max-priority :
-> Active le système de priorité (0-255, max 10 recommandé)
-> Messages avec priorité haute traités en premier

x-queue-mode = lazy :
-> Messages stockés sur disque immédiatement
-> Économise RAM
-> Plus lent mais supporte plus de messages


QUEUE ANONYME (EXCLUSIVE) :
result = channel.queue_declare(queue='', exclusive=True)
queue_name = result.method.queue  # Nom auto-généré : amq.gen-Xx0zY...

-> Utile pour RPC ou Pub/Sub temporaire


5. CONSUMER (Consommateur)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RÔLE :
-> Application qui REÇOIT des messages d'une queue
-> Traite les messages
-> Envoie un ACK (acknowledgment) ou NACK

ANALOGIE :
-> Le destinataire qui ouvre sa boîte aux lettres
-> Lit les lettres et confirme réception

EXEMPLE BASIQUE :
def callback(ch, method, properties, body):
    print(f"Reçu : {body}")
    # Traitement du message
    ch.basic_ack(delivery_tag=method.delivery_tag)  # ACK

channel.basic_consume(
    queue='hello',
    on_message_callback=callback,
    auto_ack=False  # ACK manuel (recommandé)
)

print('En attente de messages...')
channel.start_consuming()


PREFETCH COUNT (QoS - Quality of Service) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

channel.basic_qos(prefetch_count=1)

-> Nombre de messages non-ACKés qu'un consommateur peut avoir
-> prefetch_count=1 : un seul message à la fois (fair dispatch)
-> prefetch_count=10 : max 10 messages non-ACKés

POURQUOI ?
Sans prefetch, RabbitMQ envoie tous les messages possibles au consommateur
-> Si un worker est lent, il accumule beaucoup de messages
-> Round-robin devient unfair

Avec prefetch_count=1 :
-> RabbitMQ envoie 1 message
-> Attend l'ACK
-> Envoie le suivant au worker libre
-> Fair dispatch [OK]


ACKNOWLEDGMENTS (ACK / NACK) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

AUTO-ACK (auto_ack=True) :
-> Message considéré livré dès qu'envoyé au consommateur
-> Dangereux : si crash, message perdu [X]

MANUAL ACK (auto_ack=False, recommandé) :
-> Consommateur doit explicitement ACK
-> Si crash avant ACK, RabbitMQ redeliver [OK]

ACK :
ch.basic_ack(delivery_tag=method.delivery_tag)
-> "J'ai bien traité le message"

NACK (Negative ACK) :
ch.basic_nack(
    delivery_tag=method.delivery_tag,
    requeue=True    # True = remettre en queue, False = supprimer
)
-> "Je n'ai pas pu traiter, réessaye ou supprime"

REJECT :
ch.basic_reject(
    delivery_tag=method.delivery_tag,
    requeue=False
)
-> Comme NACK mais un seul message à la fois


EXEMPLE COMPLET :
def callback(ch, method, properties, body):
    try:
        print(f"Traitement de : {body}")
        # Traitement (peut échouer)
        process_message(body)
        
        # Succès -> ACK
        ch.basic_ack(delivery_tag=method.delivery_tag)
        
    except Exception as e:
        print(f"Erreur : {e}")
        # Échec -> NACK (requeue)
        ch.basic_nack(
            delivery_tag=method.delivery_tag,
            requeue=True  # Réessayer
        )


6. CONNECTION & CHANNEL
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONNECTION :
-> Connexion TCP au broker
-> Coûteuse à créer (handshake, auth)
-> Une par application généralement

CHANNEL :
-> Canal virtuel dans une connexion
-> Léger (multiplexing)
-> Isolation : chaque thread/worker a son channel

ANALOGIE :
-> Connection = Câble téléphonique
-> Channel = Appel téléphonique sur ce câble
-> Un câble peut porter plusieurs appels

EXEMPLE :
# Créer une connexion
connection = pika.BlockingConnection(
    pika.ConnectionParameters(
        host='localhost',
        port=5672,
        virtual_host='/',
        credentials=pika.PlainCredentials('guest', 'guest'),
        heartbeat=600,        # Heartbeat toutes les 10 minutes
        blocked_connection_timeout=300
    )
)

# Créer des channels
channel1 = connection.channel()  # Pour producer
channel2 = connection.channel()  # Pour consumer

# Fermer proprement
channel1.close()
channel2.close()
connection.close()


HEARTBEATS :
-> RabbitMQ envoie un "ping" régulier au client
-> Détecte les connexions mortes (réseau, crash)
-> Configurable (défaut : 60 secondes)


7. VIRTUAL HOST (vhost)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RÔLE :
-> Isolation logique dans un broker
-> Comme des "databases" séparées
-> Exchanges, queues, bindings isolés

ANALOGIE :
-> Appartements dans un immeuble
-> Chaque locataire (app) a son espace

UTILISATION :
-> Multi-tenancy (plusieurs clients)
-> Environnements (dev, staging, prod)
-> Isolation sécurité

VHOST PAR DÉFAUT : /

CRÉER UN VHOST :
rabbitmqctl add_vhost my_vhost
rabbitmqctl set_permissions -p my_vhost guest ".*" ".*" ".*"

CONNEXION :
connection = pika.BlockingConnection(
    pika.ConnectionParameters(
        host='localhost',
        virtual_host='my_vhost'
    )
)


RÉSUMÉ DU FLUX COMPLET :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. PRODUCER publie un message vers un EXCHANGE
   └─ Avec une ROUTING KEY
   
2. EXCHANGE examine la ROUTING KEY
   └─ Consulte ses BINDINGS
   
3. EXCHANGE route le message vers les QUEUES liées
   └─ Basé sur les règles de binding
   
4. MESSAGE stocké dans la QUEUE
   └─ En mémoire (rapide) ou disque (durable)
   
5. CONSUMER reçoit le message
   └─ Traite le message
   
6. CONSUMER envoie un ACK
   └─ RabbitMQ supprime le message de la queue


EXEMPLE COMPLET MINIMAL :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# producer.py
import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Déclarer la queue
channel.queue_declare(queue='hello', durable=True)

# Publier
channel.basic_publish(
    exchange='',
    routing_key='hello',
    body='Hello World!',
    properties=pika.BasicProperties(delivery_mode=2)  # Persistent
)

print("Message envoyé [OK]")
connection.close()


# consumer.py
import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Déclarer la queue (idempotent)
channel.queue_declare(queue='hello', durable=True)

def callback(ch, method, properties, body):
    print(f"Reçu : {body.decode()}")
    ch.basic_ack(delivery_tag=method.delivery_tag)

# QoS
channel.basic_qos(prefetch_count=1)

# Consommer
channel.basic_consume(queue='hello', on_message_callback=callback)

print('En attente...')
channel.start_consuming()
"""


# [OK] PARTIE 3 : INSTALLATION ET FICHIERS CRÉÉS

"""
┌────────────────────────────────────────────────────────────────────────┐
│              INSTALLATION RABBITMQ (LINUX/MAC/WINDOWS)                 │
└────────────────────────────────────────────────────────────────────────┘

MÉTHODES D'INSTALLATION :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
1. Package manager (apt, yum, brew, chocolatey)
2. Binaires officiels
3. Docker (recommandé pour développement)
4. CloudAMQP (cloud managé, gratuit pour démarrer)


PRÉ-REQUIS : ERLANG
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RabbitMQ est écrit en Erlang/OTP
-> Erlang/OTP doit être installé AVANT RabbitMQ

VERSION REQUISE :
-> RabbitMQ 3.13.x requiert Erlang/OTP 25.x ou 26.x


INSTALLATION UBUNTU/DEBIAN :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# 1. Installer Erlang
curl -1sLf 'https://dl.cloudsmith.io/public/rabbitmq/rabbitmq-erlang/setup.deb.sh' | \
sudo -E bash

sudo apt-get install -y erlang-base \
                        erlang-asn1 erlang-crypto erlang-eldap erlang-ftp erlang-inets \
                        erlang-mnesia erlang-os-mon erlang-parsetools erlang-public-key \
                        erlang-runtime-tools erlang-snmp erlang-ssl \
                        erlang-syntax-tools erlang-tftp erlang-tools erlang-xmerl

# 2. Ajouter le dépôt RabbitMQ
curl -1sLf 'https://dl.cloudsmith.io/public/rabbitmq/rabbitmq-server/setup.deb.sh' | \
sudo -E bash

# 3. Installer RabbitMQ
sudo apt-get install -y rabbitmq-server

# 4. Démarrer RabbitMQ
sudo systemctl start rabbitmq-server
sudo systemctl enable rabbitmq-server

# 5. Vérifier le statut
sudo systemctl status rabbitmq-server

[ALARM_CLOCK] Durée d'installation : 5-10 minutes


INSTALLATION macOS (Homebrew) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Installer RabbitMQ (inclut Erlang)
brew install rabbitmq

# Ajouter au PATH
export PATH=$PATH:/usr/local/sbin

# Démarrer RabbitMQ
brew services start rabbitmq

# Ou manuellement
rabbitmq-server


INSTALLATION WINDOWS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. Télécharger et installer Erlang/OTP
   -> https://www.erlang.org/downloads
   
2. Télécharger et installer RabbitMQ
   -> https://www.rabbitmq.com/install-windows.html
   
3. RabbitMQ s'installe comme service Windows
   -> Démarrage automatique au boot

Répertoire par défaut :
C:\Program Files\RabbitMQ Server\


INSTALLATION DOCKER (Recommandé pour développement) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Image officielle avec Management Plugin
docker run -d --name rabbitmq \
  -p 5672:5672 \
  -p 15672:15672 \
  -e RABBITMQ_DEFAULT_USER=admin \
  -e RABBITMQ_DEFAULT_PASS=password \
  rabbitmq:3.13-management

PORTS :
-> 5672 : AMQP 0-9-1 (protocole principal)
-> 15672 : Management UI (http://localhost:15672)

ACCÈS WEB :
-> URL : http://localhost:15672
-> User : admin
-> Pass : password

AVEC VOLUME (persistence) :
docker run -d --name rabbitmq \
  -p 5672:5672 \
  -p 15672:15672 \
  -v rabbitmq_data:/var/lib/rabbitmq \
  -e RABBITMQ_DEFAULT_USER=admin \
  -e RABBITMQ_DEFAULT_PASS=password \
  rabbitmq:3.13-management

Avantages Docker :
[OK] Installation propre (pas de pollution système)
[OK] Versions multiples possibles
[OK] Facile à détruire et recréer
[OK] Même environnement sur tous les OS


ACTIVER LE MANAGEMENT PLUGIN :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

sudo rabbitmq-plugins enable rabbitmq_management

-> Interface web de gestion
-> REST API
-> CLI amélioré

ACCÈS :
-> http://localhost:15672
-> User défaut : guest / guest (localhost only)


STRUCTURE DES FICHIERS RABBITMQ :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ARBORESCENCE COMPLÈTE (Linux) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

/var/lib/rabbitmq/                    <- Répertoire de données
├── mnesia/                           <- Base Mnesia (metadata)
│   └── rabbit@hostname/
│       ├── DECISION_TAB.LOG
│       ├── LATEST.LOG
│       ├── cluster_nodes.config
│       ├── nodes_running_at_shutdown
│       ├── rabbit_durable_exchange.DCD
│       ├── rabbit_durable_queue.DCD
│       ├── rabbit_durable_route.DCD
│       ├── rabbit_runtime_parameters.DCD
│       ├── rabbit_serial
│       ├── rabbit_user.DCD
│       ├── rabbit_user_permission.DCD
│       ├── rabbit_vhost.DCD
│       ├── schema.DAT
│       ├── msg_store_persistent/     <- Messages persistants
│       │   ├── 0.rdq
│       │   ├── 1.rdq
│       │   └── ...
│       └── msg_store_transient/      <- Messages non-persistants
│           ├── 0.rdq
│           └── 1.rdq
│
├── .erlang.cookie                    <- Cookie pour clustering
└── enabled_plugins                   <- Plugins activés

/var/log/rabbitmq/                    <- Logs
├── rabbit@hostname.log               <- Log principal
├── rabbit@hostname_upgrade.log       <- Log upgrades
└── log/

/etc/rabbitmq/                        <- Configuration
├── rabbitmq.conf                     <- Config principale (nouveau format)
├── advanced.config                   <- Config avancée (Erlang)
└── enabled_plugins                   <- Symlink


DÉTAIL DE CHAQUE TYPE DE FICHIER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. MNESIA DATABASE (metadata)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

/var/lib/rabbitmq/mnesia/rabbit@hostname/

MNESIA :
-> Base de données distribuée Erlang
-> Stocke les METADATA (pas les messages)
-> Queues, exchanges, bindings, users, permissions, vhosts

FICHIERS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

schema.DAT :
-> Schéma de la base Mnesia
-> Structure des tables

rabbit_durable_exchange.DCD :
-> Définitions des exchanges durables
-> Nom, type, durabilité, arguments

rabbit_durable_queue.DCD :
-> Définitions des queues durables
-> Nom, durabilité, auto-delete, arguments

rabbit_durable_route.DCD :
-> Bindings durables
-> Exchange -> Queue mappings

rabbit_user.DCD :
-> Utilisateurs RabbitMQ
-> Username, password hash, tags

rabbit_user_permission.DCD :
-> Permissions des utilisateurs
-> Configure, write, read permissions par vhost

rabbit_vhost.DCD :
-> Virtual hosts

rabbit_runtime_parameters.DCD :
-> Paramètres runtime (policies, etc.)

LATEST.LOG :
-> Log des transactions Mnesia
-> Permet la récupération après crash

cluster_nodes.config :
-> Nœuds du cluster

nodes_running_at_shutdown :
-> Nœuds en cours d'exécution au dernier arrêt


2. MESSAGE STORE (stockage des messages)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

msg_store_persistent/ :
-> Messages avec delivery_mode=2 (persistent)
-> Survit au redémarrage

msg_store_transient/ :
-> Messages avec delivery_mode=1 (transient)
-> Perdu au redémarrage

FORMAT .rdq (RabbitMQ Data Queue) :
-> Fichiers de segment
-> Taille par défaut : 16 MB par segment
-> Anciens segments supprimés quand messages consommés

STRUCTURE INTERNE :
-> Index : position des messages
-> Journal : log des opérations
-> Segments : données des messages


3. .erlang.cookie (CRITIQUE pour clustering)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

/var/lib/rabbitmq/.erlang.cookie

RÔLE :
-> Secret partagé pour authentification inter-nœuds
-> Tous les nœuds d'un cluster doivent avoir le MÊME cookie

CONTENU :
Une chaîne aléatoire de 20 caractères

EXEMPLE :
XPKWBDHQWFHZJXZTMQIG

SÉCURITÉ :
-> Permissions 400 (lecture seule par owner)
-> Ne JAMAIS partager publiquement
-> Changer en production

CHANGER LE COOKIE :
sudo rabbitmqctl stop_app
echo "NEW_SECRET_COOKIE" > /var/lib/rabbitmq/.erlang.cookie
chmod 400 /var/lib/rabbitmq/.erlang.cookie
sudo rabbitmqctl start_app


4. FICHIER DE CONFIGURATION : /etc/rabbitmq/rabbitmq.conf
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Format : sysctl (nouveau format depuis 3.7.0)

EXEMPLE COMPLET :

# Network & listeners
listeners.tcp.default = 5672
management.tcp.port = 15672

# Memory & disk
vm_memory_high_watermark.relative = 0.6
disk_free_limit.absolute = 2GB

# Default user
default_user = admin
default_pass = secretpassword
default_vhost = /
default_permissions.configure = .*
default_permissions.read = .*
default_permissions.write = .*

# Logging
log.file.level = info
log.console = true
log.console.level = info

# Message store
queue_index_embed_msgs_below = 4096

# Clustering
cluster_partition_handling = autoheal
cluster_name = my-cluster

# Heartbeats & timeouts
heartbeat = 60
frame_max = 131072
channel_max = 2047

# Connection limits
num_acceptors.tcp = 10
tcp_listen_options.backlog = 128
tcp_listen_options.nodelay = true

# TLS/SSL
listeners.ssl.default = 5671
ssl_options.cacertfile = /path/to/ca_certificate.pem
ssl_options.certfile = /path/to/server_certificate.pem
ssl_options.keyfile = /path/to/server_key.pem
ssl_options.verify = verify_peer
ssl_options.fail_if_no_peer_cert = true


5. FICHIER DE LOG : /var/log/rabbitmq/rabbit@hostname.log
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONTENU TYPE :

2024-12-08 10:30:00.123 [info] <0.222.0> Server startup complete; 4 plugins started.
 * rabbitmq_management
 * rabbitmq_management_agent
 * rabbitmq_web_dispatch
 * amqp_client

2024-12-08 10:30:05.456 [info] <0.567.0> accepting AMQP connection <0.567.0> (192.168.1.10:54321 -> 192.168.1.5:5672)

2024-12-08 10:30:05.789 [info] <0.567.0> connection <0.567.0> (192.168.1.10:54321 -> 192.168.1.5:5672): user 'guest' authenticated and granted access to vhost '/'

2024-12-08 10:35:12.345 [warning] <0.890.0> closing AMQP connection <0.890.0> (192.168.1.10:54321 -> 192.168.1.5:5672, vhost: '/', user: 'guest'):
client unexpectedly closed TCP connection

2024-12-08 10:40:00.678 [error] <0.234.0> Error on AMQP connection <0.234.0> (192.168.1.11:54322 -> 192.168.1.5:5672, vhost: '/', user: 'admin', state: running), channel 1:
operation none caused a channel exception frame_error: "invalid field {consumer_tag,<<>>}"

NIVEAUX DE LOG :
debug, info, warning, error, critical

ROTATION :
-> Automatique par RabbitMQ
-> Nouveaux fichiers quand taille ou date dépassée


6. enabled_plugins
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

/etc/rabbitmq/enabled_plugins

FORMAT : Liste Erlang

EXEMPLE :
[rabbitmq_management,rabbitmq_prometheus,rabbitmq_shovel].

PLUGINS UTILES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

rabbitmq_management :
-> Interface web de gestion
-> REST API
-> CLI amélioré

rabbitmq_shovel :
-> Déplacer des messages entre brokers
-> Utile pour migration, réplication

rabbitmq_federation :
-> Fédération de brokers
-> Partage d'exchanges entre brokers distants

rabbitmq_prometheus :
-> Export des métriques Prometheus
-> Monitoring

rabbitmq_mqtt :
-> Support protocole MQTT (IoT)

rabbitmq_stomp :
-> Support protocole STOMP
-> WebSocket

rabbitmq_auth_backend_ldap :
-> Authentification via LDAP

rabbitmq_delayed_message_exchange :
-> Messages avec délai (delayed exchange)
-> Très utile !

ACTIVER UN PLUGIN :
sudo rabbitmq-plugins enable rabbitmq_management

DÉSACTIVER :
sudo rabbitmq-plugins disable rabbitmq_management

LISTER :
sudo rabbitmq-plugins list


COMMANDES UTILES (rabbitmqctl) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Statut du serveur
sudo rabbitmqctl status

# Lister les queues
sudo rabbitmqctl list_queues name messages consumers

# Lister les exchanges
sudo rabbitmqctl list_exchanges name type durable

# Lister les bindings
sudo rabbitmqctl list_bindings

# Lister les connexions
sudo rabbitmqctl list_connections name peer_host peer_port state

# Lister les channels
sudo rabbitmqctl list_channels connection number consumer_count

# Créer un utilisateur
sudo rabbitmqctl add_user myuser mypassword
sudo rabbitmqctl set_user_tags myuser administrator
sudo rabbitmqctl set_permissions -p / myuser ".*" ".*" ".*"

# Supprimer un utilisateur
sudo rabbitmqctl delete_user myuser

# Créer un vhost
sudo rabbitmqctl add_vhost my_vhost

# Supprimer un vhost
sudo rabbitmqctl delete_vhost my_vhost

# Purger une queue
sudo rabbitmqctl purge_queue queue_name

# Arrêter RabbitMQ
sudo rabbitmqctl stop

# Redémarrer
sudo systemctl restart rabbitmq-server

# Environnement
sudo rabbitmqctl environment


OUTILS DE DIAGNOSTIC :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Health check
sudo rabbitmqctl node_health_check

# Rapport complet
sudo rabbitmqctl report > rabbitmq_report.txt

# Alarms
sudo rabbitmqctl list_alarms
"""


# [OK] PARTIE 4 : ARCHITECTURE INTERNE DE RABBITMQ

"""
┌────────────────────────────────────────────────────────────────────────┐
│                   ARCHITECTURE RABBITMQ                                │
└────────────────────────────────────────────────────────────────────────┘

VUE D'ENSEMBLE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

┌─────────────────────────────────────────────────────────────────────┐
│                       RABBITMQ BROKER (Nœud Erlang)                 │
│                                                                     │
│  ┌────────────────────────────────────────────────────────────────┐ │
│  │                   COUCHE RÉSEAU (TCP/TLS)                      │ │
│  │  - AMQP 0-9-1 (port 5672)                                      │ │
│  │  - AMQPS (port 5671)                                           │ │
│  │  - HTTP Management (port 15672)                                │ │
│  │  - Clustering (port 25672)                                     │ │
│  └────────────────────────────────────────────────────────────────┘ │
│                              v                                      │
│  ┌────────────────────────────────────────────────────────────────┐ │
│  │              COUCHE PROTOCOLE (AMQP 0-9-1)                     │ │
│  │  - Parsing des frames                                          │ │
│  │  - Connection handling                                         │ │
│  │  - Channel multiplexing                                        │ │
│  │  - Flow control                                                │ │
│  └────────────────────────────────────────────────────────────────┘ │
│                              v                                      │
│  ┌────────────────────────────────────────────────────────────────┐ │
│  │                   COUCHE ROUTING                               │ │
│  │                                                                │ │
│  │  ┌──────────────┐  ┌──────────────┐  ┌──────────────┐        │ │
│  │  │  EXCHANGES   │  │   BINDINGS   │  │   QUEUES     │        │ │
│  │  │              │  │              │  │              │        │ │
│  │  │ - Direct     │  │ - Rules      │  │ - Storage    │        │ │
│  │  │ - Fanout     │──│ - Routing    │──│ - Ordering   │        │ │
│  │  │ - Topic      │  │   keys       │  │ - Consumers  │        │ │
│  │  │ - Headers    │  │              │  │              │        │ │
│  │  └──────────────┘  └──────────────┘  └──────────────┘        │ │
│  │                                                                │ │
│  └────────────────────────────────────────────────────────────────┘ │
│                              v                                      │
│  ┌────────────────────────────────────────────────────────────────┐ │
│  │              COUCHE STOCKAGE                                   │ │
│  │                                                                │ │
│  │  ┌────────────┐  ┌────────────┐  ┌────────────┐              │ │
│  │  │  MNESIA    │  │ Message    │  │   Index    │              │ │
│  │  │ (Metadata) │  │   Store    │  │  (Queue)   │              │ │
│  │  │            │  │            │  │            │              │ │
│  │  │ - Queues   │  │ - Persist. │  │ - Position │              │ │
│  │  │ - Exchanges│  │ - Transient│  │ - Ordering │              │ │
│  │  │ - Users    │  │ - Segments │  │            │              │ │
│  │  └────────────┘  └────────────┘  └────────────┘              │ │
│  │                                                                │ │
│  └────────────────────────────────────────────────────────────────┘ │
│                              v                                      │
│  ┌────────────────────────────────────────────────────────────────┐ │
│  │                   FICHIERS SUR DISQUE                          │ │
│  │  - Mnesia tables (.DCD)                                        │ │
│  │  - Message store (.rdq)                                        │ │
│  │  - Logs (.log)                                                 │ │
│  └────────────────────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────────────┘


DÉTAIL DES COUCHES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. COUCHE RÉSEAU (Acceptors)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PROCESSUS ACCEPTOR :
-> Pool de processus Erlang qui acceptent les connexions TCP
-> Par défaut : 10 acceptors

HANDSHAKE TCP :
1. Client -> SYN -> RabbitMQ
2. RabbitMQ -> SYN-ACK -> Client
3. Client -> ACK -> RabbitMQ
4. Connexion établie

TLS/SSL (si activé) :
-> Handshake TLS après TCP
-> Échange de certificats
-> Établissement de la session chiffrée

HEARTBEATS :
-> Ping/Pong régulier entre client et serveur
-> Détection de connexions mortes
-> Configurable (défaut : 60 secondes)

FLUX :
Client établit connexion -> Acceptor -> Nouveau processus Erlang créé
                                      -> Gère cette connexion


2. COUCHE PROTOCOLE AMQP
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

AMQP 0-9-1 FRAME FORMAT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

FRAME STRUCTURE :
┌────────┬────────┬─────────┬─────────────────┬────────┐
│  Type  │Channel │  Size   │     Payload     │  End   │
│ (1 byte)│(2 bytes│(4 bytes)│  (Size bytes)   │(1 byte)│
└────────┴────────┴─────────┴─────────────────┴────────┘

TYPES DE FRAMES :
1. METHOD : commande AMQP (publish, consume, ack, etc.)
2. HEADER : propriétés du message (content-type, etc.)
3. BODY : contenu du message (peut être fragmenté)
4. HEARTBEAT : keepalive

EXEMPLE : basic.publish
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Frame 1 (METHOD) :
{
  "type": "method",
  "channel": 1,
  "method": {
    "class": "basic",
    "name": "publish",
    "exchange": "orders",
    "routing_key": "order.created"
  }
}

Frame 2 (HEADER) :
{
  "type": "header",
  "channel": 1,
  "properties": {
    "content_type": "application/json",
    "delivery_mode": 2,
    "priority": 5
  },
  "body_size": 1024
}

Frame 3 (BODY) :
{
  "type": "body",
  "channel": 1,
  "body": "...message content..."
}

CHANNEL MULTIPLEXING :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Une connexion TCP peut avoir plusieurs channels
-> Chaque channel = flux logique indépendant
-> Channel ID dans chaque frame

AVANTAGES :
[OK] Économie de connexions TCP
[OK] Isolation des opérations
[OK] Parallélisme

LIMITE :
-> Max 65535 channels par connexion (2^16 - 1)

PROCESSUS ERLANG :
-> 1 processus Erlang par channel
-> Léger : Erlang peut gérer millions de processus


FLOW CONTROL :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Si RabbitMQ est surchargé :
-> Envoie connection.blocked à tous les clients
-> Clients arrêtent d'envoyer
-> Quand ressources disponibles : connection.unblocked

RAISONS DE BLOCAGE :
-> Mémoire haute (vm_memory_high_watermark dépassé)
-> Disque plein (disk_free_limit atteint)


3. COUCHE ROUTING (Cœur de RabbitMQ)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

EXCHANGE PROCESS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Chaque exchange = processus Erlang supervisé

TÂCHES :
1. Recevoir le message du channel
2. Consulter la table de bindings (Mnesia)
3. Déterminer les queues cibles
4. Router le message vers ces queues

ALGORITHME DE ROUTAGE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DIRECT :
-> Hash map : routing_key -> [queue1, queue2, ...]
-> Lookup O(1)

TOPIC :
-> Trie-based matching
-> Wildcards (* et #)
-> Plus lent que direct mais toujours rapide

FANOUT :
-> Broadcast à toutes les queues liées
-> Pas de routing key utilisée

HEADERS :
-> Match sur les headers du message
-> x-match: all (tous les headers) ou any (au moins un)
-> Le plus lent

ALTERNATE EXCHANGE :
-> Si aucun binding trouvé, envoyer vers alternate exchange
-> Utile pour gérer les messages "perdus"


QUEUE PROCESS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Chaque queue = processus Erlang supervisé

TÂCHES :
1. Recevoir messages de l'exchange
2. Stocker dans l'index (ordre)
3. Stocker le body dans le message store
4. Distribuer aux consommateurs
5. Gérer les ACKs

QUEUE STATES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RUNNING :
-> Normal, traite les messages

IDLE :
-> Pas de messages, pas de consommateurs
-> Économie de ressources

FLOW :
-> En backpressure (ralentit la production)

BLOCKED :
-> Bloquée (mémoire ou disque plein)

SYNCHRONIZING (mirrored queue) :
-> Synchronisation avec les mirrors

MODES DE QUEUE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DEFAULT (classic) :
-> Messages en mémoire sauf si pression
-> Rapide pour petites queues

LAZY :
-> Messages sur disque immédiatement
-> Économise RAM
-> Supporte plus de messages
-> Légèrement plus lent

QUORUM :
-> Réplication Raft (3+ nœuds)
-> Haute disponibilité
-> Perte de messages impossible
-> Performance similaire à lazy


4. COUCHE STOCKAGE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

MNESIA (Metadata Store) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RÔLE :
-> Stocke les DEFINITIONS (pas les messages)
-> Queues, exchanges, bindings, users, vhosts, permissions

TYPE :
-> Base distribuée Erlang
-> Réplication entre nœuds du cluster
-> Transactions ACID

TABLES PRINCIPALES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
rabbit_durable_exchange : exchanges durables
rabbit_durable_queue : queues durables
rabbit_durable_route : bindings durables
rabbit_user : utilisateurs
rabbit_user_permission : permissions
rabbit_vhost : virtual hosts
rabbit_runtime_parameters : policies, limits


MESSAGE STORE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DEUX STORES :
1. msg_store_persistent : delivery_mode=2
2. msg_store_transient : delivery_mode=1

STRUCTURE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

SEGMENTS :
-> Fichiers .rdq (RabbitMQ Data Queue)
-> Taille : 16 MB par défaut
-> Append-only (ajout à la fin)

INDEX :
-> Map : message_id -> (segment, offset)
-> En mémoire pour les messages récents
-> Sur disque pour les anciens

GARBAGE COLLECTION :
-> Ancien segment supprimé quand tous ses messages sont ACKés
-> Compaction périodique

EMBEDDING :
-> Petits messages (<= 4096 bytes par défaut) :
   -> Stockés directement dans l'index de la queue
   -> Pas de passage par le message store
   -> Plus rapide [OK]


QUEUE INDEX :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RÔLE :
-> Maintient l'ORDRE des messages dans la queue
-> Pointeurs vers le message store

STRUCTURE :
┌────────────────────────────────────────────┐
│  Queue Index (par queue)                   │
│                                            │
│  Message 1 -> (msg_store_id, segment, pos) │
│  Message 2 -> (msg_store_id, segment, pos) │
│  Message 3 -> (embedded_message_body)      │
│  Message 4 -> (msg_store_id, segment, pos) │
│  ...                                       │
└────────────────────────────────────────────┘

FICHIERS :
-> Un fichier par segment
-> Format journal (append-only)


MEMORY MANAGEMENT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

vm_memory_high_watermark :
-> Pourcentage de RAM disponible
-> Par défaut : 0.4 (40%)
-> Si dépassé : BLOCAGE des producteurs

PAGING :
-> Si mémoire haute, messages écrits sur disque
-> Queue entre en mode "flow"
-> Producteurs ralentis

ALARM :
-> Si watermark dépassé : alarm déclenchée
-> connection.blocked envoyé aux clients
-> Tous les publishers stoppés


DISK MANAGEMENT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

disk_free_limit :
-> Espace disque minimum requis
-> Par défaut : 50 MB
-> Si atteint : ALARM + BLOCAGE

RECOMMANDATION PRODUCTION :
disk_free_limit.absolute = 2GB
-> ou
disk_free_limit.relative = 1.5
   (1.5x la RAM totale)


SUPERVISION TREE (Erlang OTP) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RabbitMQ utilise OTP (Open Telecom Platform)
-> Framework de supervision Erlang
-> Tolérance aux pannes

STRUCTURE :
┌────────────────────────────────────────────┐
│        rabbit_sup (root supervisor)        │
│                    │                       │
│       ┌────────────┴──────────────┐        │
│       │                           │        │
│  rabbit_tcp_      rabbit_queue_   ...     │
│   sup               sup                   │
│   │                 │                     │
│   ├─ Acceptor 1     ├─ Queue 1           │
│   ├─ Acceptor 2     ├─ Queue 2           │
│   ├─ Acceptor 3     ├─ Queue 3           │
│   ...               ...                   │
└────────────────────────────────────────────┘

RESTART STRATEGY :
-> Si un processus crashe, superviseur le redémarre
-> one_for_one : redémarre seulement le crashé
-> one_for_all : redémarre tous
-> rest_for_one : redémarre le crashé + suivants

HAUTE RÉSILIENCE :
-> Erlang a été conçu pour téléphonie (99.9999999% uptime)
-> RabbitMQ hérite de cette résilience
"""


# [OK] PARTIE 5 : VOYAGE COMPLET D'UN MESSAGE

"""
┌────────────────────────────────────────────────────────────────────────┐
│     SCÉNARIO : PUBLISH + CONSUME (E-COMMERCE ORDER)                    │
└────────────────────────────────────────────────────────────────────────┘

CONTEXTE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Application e-commerce à Dakar
RabbitMQ à Paris (serveur dédié)
Distance : ~4500 km
Latence réseau : ~150ms (aller-retour)

Scénario :
-> Client commande un produit sur le site web
-> Service Web publie un message "order.created"
-> Service Payment consomme et traite

TOPOLOGIE :
Exchange : orders (type: topic)
Queue : payment_processing
Binding : orders -> payment_processing (routing key: order.*)


ÉTAPE 0 : SETUP & CONNEXION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Connexion (Dakar -> Paris)
connection = pika.BlockingConnection(
    pika.ConnectionParameters(
        host='rabbitmq.example.com',
        port=5672,
        virtual_host='ecommerce',
        credentials=pika.PlainCredentials('web-app', 'secret'),
        heartbeat=600,
        blocked_connection_timeout=300
    )
)

# Créer un channel
channel = connection.channel()

# Déclarer l'exchange
channel.exchange_declare(
    exchange='orders',
    exchange_type='topic',
    durable=True
)

# Déclarer la queue
channel.queue_declare(
    queue='payment_processing',
    durable=True,
    arguments={
        'x-max-priority': 10,
        'x-message-ttl': 300000  # 5 minutes
    }
)

# Créer le binding
channel.queue_bind(
    exchange='orders',
    queue='payment_processing',
    routing_key='order.*'
)

TEMPS TOTAL SETUP : ~500ms (connexion TCP + handshakes)


ÉTAPE 1 : PUBLISH DU MESSAGE (Dakar -> Paris)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CODE (Dakar) :
order = {
    'order_id': 12345,
    'customer_id': 67890,
    'items': [
        {'product_id': 1, 'quantity': 2, 'price': 50.00},
        {'product_id': 2, 'quantity': 1, 'price': 30.00}
    ],
    'total': 130.00,
    'currency': 'EUR',
    'created_at': '2024-12-08T10:30:00Z'
}

channel.basic_publish(
    exchange='orders',
    routing_key='order.created',
    body=json.dumps(order),
    properties=pika.BasicProperties(
        content_type='application/json',
        delivery_mode=2,          # Persistent
        priority=5,
        correlation_id=str(uuid.uuid4()),
        headers={
            'source': 'web-app',
            'version': '1.0'
        }
    ),
    mandatory=True  # Erreur si pas de queue cible
)

PROCESSUS DÉTAILLÉ :

1. PRÉPARATION CLIENT (Dakar, local) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) Sérialisation JSON :
   order dict -> JSON string
   [TEMPS] ~1 ms
   
   JSON size : ~250 bytes

b) Construction des frames AMQP :
   
   Frame 1 (METHOD - basic.publish) :
   ┌────────┬────────┬─────────┬───────────────────────┬────────┐
   │  0x01  │  0001  │  0x0050 │ class=60, method=40   │  0xCE  │
   │ METHOD │Channel │  Size   │ exchange="orders"     │  END   │
   │        │   1    │  80 b   │ routing_key="order.   │  BYTE  │
   │        │        │         │  created"             │        │
   └────────┴────────┴─────────┴───────────────────────┴────────┘
   
   Frame 2 (HEADER - properties) :
   ┌────────┬────────┬─────────┬───────────────────────┬────────┐
   │  0x02  │  0001  │  0x0080 │ content-type, deliv.  │  0xCE  │
   │ HEADER │Channel │  128 b  │ mode, priority, etc.  │  END   │
   │        │   1    │         │ body_size=250         │        │
   └────────┴────────┴─────────┴───────────────────────┴────────┘
   
   Frame 3 (BODY - message content) :
   ┌────────┬────────┬─────────┬───────────────────────┬────────┐
   │  0x03  │  0001  │  0x00FA │ {"order_id": 12345... │  0xCE  │
   │  BODY  │Channel │  250 b  │ ...}                  │  END   │
   │        │   1    │         │                       │        │
   └────────┴────────┴─────────┴───────────────────────┴────────┘
   
   Total frames : 3
   Total bytes : ~460 bytes (80 + 128 + 250 + overhead)
   [TEMPS] ~1 ms

2. ENVOI RÉSEAU (Dakar -> Paris) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) Écriture dans le socket TCP :
   -> Client -> OS -> Network stack
   [TEMPS] <1 ms
   
b) Traversée du réseau :
   -> Dakar -> FAI Sénégal -> Câble sous-marin -> Europe -> Paris
   [TEMPS] ~75 ms (one-way)

3. RÉCEPTION PAR RABBITMQ (Paris) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) Acceptor process reçoit les frames :
   -> Lecture du socket TCP
   -> Désérialisation des frames AMQP
   [TEMPS] ~1 ms
   
b) Validation :
   -> Vérifier le format AMQP
   -> Vérifier que l'exchange existe
   -> Vérifier les permissions (publish)
   [TEMPS] <1 ms

4. ROUTING PAR L'EXCHANGE (Paris) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) Exchange process 'orders' (type: topic) :
   
   Routing key : "order.created"
   
   Consultation de la table de bindings (Mnesia) :
   ┌─────────────┬──────────────────────┬────────────────────┐
   │  Exchange   │    Routing Pattern   │       Queue        │
   ├─────────────┼──────────────────────┼────────────────────┤
   │  orders     │     order.*          │ payment_processing │
   │  orders     │     order.*          │ shipping_queue     │
   │  orders     │     order.created    │ analytics_queue    │
   └─────────────┴──────────────────────┴────────────────────┘
   
   Match :
   -> "order.*" match "order.created" [OK]
      -> payment_processing [OK]
      -> shipping_queue [OK]
   -> "order.created" exact match [OK]
      -> analytics_queue [OK]
   
   Queues cibles : 3 queues
   [TEMPS] ~2 ms (topic matching)

b) Duplication du message :
   -> Message copié 3 fois (une par queue)
   -> Pointeur vers le message body (pas de copie du body)
   [TEMPS] <1 ms

5. STOCKAGE DANS LES QUEUES (Paris) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

QUEUE : payment_processing

a) Message Store (delivery_mode=2, persistent) :
   
   -> Écriture dans msg_store_persistent :
   
   Segment file : /var/lib/rabbitmq/mnesia/.../msg_store_persistent/12.rdq
   
   Message entry :
   {
     msg_id: <<binary_id>>,
     properties: {...},
     body: <<json_bytes>>,
     size: 250
   }
   
   -> Append au fichier (O_APPEND)
   [TEMPS] ~5 ms (écriture disque SSD)
   
   -> Flush obligatoire (fsync) pour persistance :
   [TEMPS] ~10 ms (selon I/O load)

b) Queue Index (payment_processing) :
   
   -> Ajout de l'entrée dans l'index de la queue :
   
   {
     msg_seq_id: 12345,
     msg_store_ref: <<msg_id>>,
     properties: {...},
     is_persistent: true
   }
   
   -> Écriture dans le queue index file
   [TEMPS] ~3 ms
   
c) Queue process state update :
   -> Incrément du compteur de messages
   -> Notification aux consommateurs (si présents)
   [TEMPS] <1 ms

TOTAL STOCKAGE (par queue) : ~19 ms

PARALLÈLE :
-> Les 3 queues écrivent en parallèle
-> Temps total ≈ temps d'une queue ≈ 19 ms

6. CONFIRMATION AU PRODUCER (Paris -> Dakar) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RabbitMQ envoie basic.ack (si Publisher Confirms activé)

Frame :
┌────────┬────────┬─────────┬───────────────────────┬────────┐
│  0x01  │  0001  │  0x0020 │ class=60, method=80   │  0xCE  │
│ METHOD │Channel │  32 b   │ delivery_tag=1        │  END   │
│        │   1    │         │ multiple=false        │        │
└────────┴────────┴─────────┴───────────────────────┴────────┘

-> Envoi réseau (Paris -> Dakar)
[TEMPS] ~75 ms (one-way)

-> Réception par le client
[TEMPS] <1 ms

RÉCAPITULATIF TEMPS PUBLISH :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Préparation client :            2 ms
Envoi (Dakar -> Paris) :        75 ms
Réception RabbitMQ :            1 ms
Validation :                    1 ms
Routing (exchange) :            2 ms
Stockage (queues) :            19 ms
Confirmation (Paris -> Dakar) : 75 ms
──────────────────────────────────────
TOTAL :                       ~175 ms

SANS Publisher Confirms : ~100 ms
(pas d'attente de confirmation)


ÉTAPE 2 : CONSUME DU MESSAGE (Paris)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CODE (Service Payment, Paris) :

def process_payment(ch, method, properties, body):
    try:
        order = json.loads(body)
        print(f"Traitement paiement : {order['order_id']}")
        
        # Traiter le paiement (appel API Stripe, etc.)
        result = payment_gateway.charge(
            amount=order['total'],
            currency=order['currency'],
            customer_id=order['customer_id']
        )
        
        if result.success:
            print(f"Paiement réussi [OK]")
            # ACK
            ch.basic_ack(delivery_tag=method.delivery_tag)
        else:
            print(f"Paiement échoué [X]")
            # NACK + requeue
            ch.basic_nack(
                delivery_tag=method.delivery_tag,
                requeue=True
            )
            
    except Exception as e:
        print(f"Erreur : {e}")
        # NACK + requeue
        ch.basic_nack(
            delivery_tag=method.delivery_tag,
            requeue=True
        )

# Setup consumer
channel.basic_qos(prefetch_count=1)

channel.basic_consume(
    queue='payment_processing',
    on_message_callback=process_payment,
    auto_ack=False
)

print('Payment service en attente...')
channel.start_consuming()


PROCESSUS DÉTAILLÉ :

1. CONSUMER ENREGISTREMENT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

basic_consume() envoie :
┌────────────────────────────────────────┐
│  basic.consume (METHOD frame)          │
│  queue = "payment_processing"          │
│  consumer_tag = "ctag-12345"           │
│  no_ack = False (manual ack)           │
│  exclusive = False                     │
│  arguments = {}                        │
└────────────────────────────────────────┘

RabbitMQ répond :
┌────────────────────────────────────────┐
│  basic.consume-ok (METHOD frame)       │
│  consumer_tag = "ctag-12345"           │
└────────────────────────────────────────┘

-> Consumer enregistré dans la queue
[TEMPS] ~2 ms

2. QUEUE DÉTECTE LE CONSOMMATEUR :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Queue process (payment_processing) :
-> Voit qu'un consumer est prêt
-> Vérifie prefetch_count (1)
-> Envoie le premier message

3. LIVRAISON DU MESSAGE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) Queue lit le message depuis l'index :
   -> Récupère msg_store_ref
   [TEMPS] <1 ms (en mémoire si récent)
   
b) Lecture du body depuis Message Store :
   -> Lookup dans msg_store_persistent
   -> Lecture du segment file
   [TEMPS] ~5 ms (si cache miss, SSD)
   [TEMPS] <1 ms (si cache hit)
   
c) Construction des frames de livraison :
   
   Frame 1 (METHOD - basic.deliver) :
   ┌────────────────────────────────────────┐
   │  basic.deliver                         │
   │  consumer_tag = "ctag-12345"           │
   │  delivery_tag = 1                      │
   │  redelivered = false                   │
   │  exchange = "orders"                   │
   │  routing_key = "order.created"         │
   └────────────────────────────────────────┘
   
   Frame 2 (HEADER) :
   -> Properties (content-type, delivery_mode, etc.)
   
   Frame 3 (BODY) :
   -> Message content (JSON)
   
   [TEMPS] ~1 ms

d) Envoi au consumer (même serveur Paris) :
   -> Écriture dans le socket TCP
   [TEMPS] <1 ms (local)
   
e) Consumer reçoit et décode :
   -> Lecture des 3 frames
   -> Désérialisation JSON
   -> Appel du callback
   [TEMPS] ~2 ms

TOTAL LIVRAISON : ~10 ms (cache hit)
                  ~16 ms (cache miss)

4. TRAITEMENT PAR LE CONSUMER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

process_payment() exécuté :
-> Désérialisation JSON : ~1 ms
-> Appel API Stripe : ~200 ms
-> Logique business : ~10 ms
[TEMPS] TOTAL : ~211 ms

5. ACKNOWLEDGMENT (ACK) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ch.basic_ack(delivery_tag=1)

Frame envoyé :
┌────────────────────────────────────────┐
│  basic.ack (METHOD frame)              │
│  delivery_tag = 1                      │
│  multiple = false                      │
└────────────────────────────────────────┘

-> Envoi au RabbitMQ
[TEMPS] <1 ms (local)

6. RABBITMQ TRAITE L'ACK :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

a) Queue process reçoit l'ACK :
   -> Marque le message comme ACKé
   -> Retire de la liste des messages non-ACKés
   [TEMPS] <1 ms

b) Suppression du message :
   
   -> Queue index : retrait de l'entrée
   [TEMPS] <1 ms
   
   -> Message Store : marque pour suppression
   -> Garbage collection plus tard (pas immédiat)
   [TEMPS] asynchrone
   
c) Queue est prête pour le message suivant :
   -> Envoie le prochain message au consumer
   [TEMPS] <1 ms

RÉCAPITULATIF TEMPS CONSUME :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Livraison (queue -> consumer) :  10 ms
Traitement (callback) :        211 ms
ACK (consumer -> queue) :         1 ms
Traitement ACK (queue) :         2 ms
──────────────────────────────────────────
TOTAL :                        ~224 ms


SCÉNARIO : ÉCHEC & RETRY
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

SI LE PAIEMENT ÉCHOUE :

ch.basic_nack(delivery_tag=1, requeue=True)

-> Message remis en queue (FRONT de la queue)
-> Renvoyé au consumer (ou un autre)
-> Retry immédiat

PROBLÈME : BOUCLE INFINIE
Si le message échoue toujours (erreur permanente) :
-> Retry infini [X]

SOLUTIONS :
1. Dead Letter Exchange (DLX)
2. Compteur de retries dans le message
3. TTL sur la queue


SOLUTION : DEAD LETTER EXCHANGE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Queue principale avec DLX
channel.queue_declare(
    queue='payment_processing',
    durable=True,
    arguments={
        'x-dead-letter-exchange': 'dlx',
        'x-dead-letter-routing-key': 'payment.failed',
        'x-message-ttl': 60000  # 60 secondes
    }
)

# Dead Letter Exchange
channel.exchange_declare(
    exchange='dlx',
    exchange_type='direct',
    durable=True
)

# Dead Letter Queue
channel.queue_declare(
    queue='payment_failed',
    durable=True
)

channel.queue_bind(
    exchange='dlx',
    queue='payment_failed',
    routing_key='payment.failed'
)

COMPORTEMENT :
-> Si message rejeté (NACK requeue=False) : envoyé vers DLX
-> Si message expire (TTL) : envoyé vers DLX
-> Queue pleine (x-max-length) : anciens messages vers DLX

TRAITEMENT DES MESSAGES MORTS :
-> Consumer sur payment_failed
-> Logging, alerting
-> Retry manuel après investigation


PERFORMANCE : DÉBIT MAXIMAL
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

FACTEURS LIMITANTS :

1. RÉSEAU :
   -> Latence : 150ms (Dakar-Paris)
   -> Bande passante : typiquement non-limitante

2. DISK I/O (si persistent) :
   -> SSD : ~10K IOPS
   -> HDD : ~200 IOPS

3. SERIALIZATION / DESERIALIZATION :
   -> JSON parsing
   -> AMQP framing

OPTIMISATIONS :

[OK] BATCHING :
-> Publier plusieurs messages en un batch
-> Réduit overhead réseau

# Publisher
with channel.connection.channel() as ch:
    for msg in messages:
        ch.basic_publish(...)
    # Confirm en batch

[OK] PUBLISHER CONFIRMS (batch) :
channel.confirm_delivery()  # Mode confirm

# Publier plusieurs
for msg in messages:
    channel.basic_publish(...)

# Attendre les confirms de tous
channel.wait_for_confirms()

[OK] CONSUMER PREFETCH :
channel.basic_qos(prefetch_count=10)

-> Consumer reçoit 10 messages d'avance
-> Pipeline parallèle

[OK] MULTIPLE CONSUMERS :
-> Scale horizontalement
-> Fair dispatch automatique

[OK] NO PERSISTENCE (si acceptable) :
delivery_mode=1

-> Pas de write disque
-> ~10x plus rapide
-> Mais perte possible

DÉBIT TYPIQUE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Transient (RAM) : 10K-50K msg/s par nœud
Persistent (SSD) : 5K-20K msg/s par nœud
Persistent (HDD) : 500-2K msg/s par nœud

(Messages ~1KB, single queue)
"""


# [OK] PARTIE 6 : TYPES D'EXCHANGES ET ROUTING

"""
┌────────────────────────────────────────────────────────────────────────┐
│              TYPES D'EXCHANGES - ROUTING DÉTAILLÉ                      │
└────────────────────────────────────────────────────────────────────────┘

Les Exchanges sont le CŒUR du routing dans RabbitMQ
-> 4 types principaux : Direct, Fanout, Topic, Headers
-> Chaque type a un algorithme de routing différent


1. DIRECT EXCHANGE (Routage exact)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> Routing par ÉGALITÉ EXACTE de la routing key
-> Message routé vers les queues dont binding key = routing key

ALGORITHME :
if message.routing_key == binding.routing_key:
    deliver_to(binding.queue)

ANALOGIE :
-> Système postal avec codes postaux
-> Lettre avec code 75001 va UNIQUEMENT à Paris 1er

DÉCLARATION :
channel.exchange_declare(
    exchange='direct_logs',
    exchange_type='direct',
    durable=True
)

EXEMPLE : SYSTÈME DE LOGS PAR NIVEAU
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

TOPOLOGIE :
┌────────────────────────────────────────────────────────────┐
│                                                            │
│         EXCHANGE: direct_logs (type: direct)               │
│                          │                                 │
│          ┌───────────────┼───────────────┐                │
│          │               │               │                │
│     "error"          "warning"        "info"              │
│          │               │               │                │
│    ┌─────[BLACK_DOWN-POINTING_TRIANGLE]─────┐   ┌────[BLACK_DOWN-POINTING_TRIANGLE]────┐   ┌─────[BLACK_DOWN-POINTING_TRIANGLE]─────┐          │
│    │  errors   │   │ warnings│   │   infos   │          │
│    │  queue    │   │  queue  │   │   queue   │          │
│    └───────────┘   └─────────┘   └───────────┘          │
│                                                            │
└────────────────────────────────────────────────────────────┘

SETUP :
# Déclarer l'exchange
channel.exchange_declare(
    exchange='direct_logs',
    exchange_type='direct',
    durable=True
)

# Déclarer les queues
for severity in ['error', 'warning', 'info']:
    channel.queue_declare(
        queue=f'{severity}s',
        durable=True
    )
    
    # Binding avec routing key = severity
    channel.queue_bind(
        exchange='direct_logs',
        queue=f'{severity}s',
        routing_key=severity
    )

PUBLISHER :
# Publier un log ERROR
channel.basic_publish(
    exchange='direct_logs',
    routing_key='error',  # Routing key
    body='Database connection failed'
)
-> Va UNIQUEMENT dans la queue "errors"

# Publier un log INFO
channel.basic_publish(
    exchange='direct_logs',
    routing_key='info',
    body='User logged in'
)
-> Va UNIQUEMENT dans la queue "infos"

CONSUMER (errors seulement) :
channel.basic_consume(
    queue='errors',
    on_message_callback=handle_error
)
-> Reçoit SEULEMENT les messages "error"

MULTIPLE BINDINGS :
Une queue peut avoir plusieurs bindings sur le même exchange

# Queue "all_logs" reçoit TOUS les logs
channel.queue_declare(queue='all_logs', durable=True)

for severity in ['error', 'warning', 'info']:
    channel.queue_bind(
        exchange='direct_logs',
        queue='all_logs',
        routing_key=severity
    )

-> all_logs reçoit error + warning + info


CAS D'USAGE DIRECT EXCHANGE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[OK] Logs par niveau de gravité
[OK] Task routing par type (image, video, pdf)
[OK] RPC (response routing par correlation_id)
[OK] Tout besoin de routing 1:1 exact


2. FANOUT EXCHANGE (Broadcast)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> IGNORE la routing key
-> Broadcast à TOUTES les queues liées
-> Copie du message vers chaque queue

ALGORITHME :
for binding in exchange.bindings:
    deliver_to(binding.queue)

ANALOGIE :
-> Haut-parleur dans un stade
-> Tout le monde entend le même message

DÉCLARATION :
channel.exchange_declare(
    exchange='notifications',
    exchange_type='fanout',
    durable=True
)

EXEMPLE : SYSTÈME DE NOTIFICATIONS MULTI-CANAL
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

TOPOLOGIE :
┌────────────────────────────────────────────────────────────┐
│                                                            │
│      EXCHANGE: notifications (type: fanout)                │
│                          │                                 │
│          ┌───────────────┼───────────────┐                │
│          │               │               │                │
│     (broadcast)      (broadcast)    (broadcast)           │
│          │               │               │                │
│    ┌─────[BLACK_DOWN-POINTING_TRIANGLE]─────┐   ┌────[BLACK_DOWN-POINTING_TRIANGLE]────┐   ┌─────[BLACK_DOWN-POINTING_TRIANGLE]─────┐          │
│    │   email   │   │   sms   │   │   push    │          │
│    │   queue   │   │  queue  │   │   queue   │          │
│    └───────────┘   └─────────┘   └───────────┘          │
│          │               │               │                │
│    [Email       [SMS          [Push                       │
│     Service]     Service]      Service]                   │
│                                                            │
└────────────────────────────────────────────────────────────┘

SETUP :
# Exchange
channel.exchange_declare(
    exchange='notifications',
    exchange_type='fanout',
    durable=True
)

# Queues
for channel_type in ['email', 'sms', 'push']:
    queue_name = f'{channel_type}_queue'
    channel.queue_declare(queue=queue_name, durable=True)
    
    # Binding (pas de routing key nécessaire)
    channel.queue_bind(
        exchange='notifications',
        queue=queue_name,
        routing_key=''  # Ignoré de toute façon
    )

PUBLISHER :
notification = {
    'user_id': 12345,
    'title': 'Order Shipped',
    'message': 'Your order #67890 has been shipped',
    'timestamp': '2024-12-08T10:30:00Z'
}

channel.basic_publish(
    exchange='notifications',
    routing_key='',  # Ignoré !
    body=json.dumps(notification)
)

-> Message copié vers :
   - email_queue
   - sms_queue
   - push_queue

CONSUMERS :
# Email Service
def send_email(ch, method, properties, body):
    notif = json.loads(body)
    send_email_to_user(notif['user_id'], notif['message'])
    ch.basic_ack(delivery_tag=method.delivery_tag)

channel.basic_consume(queue='email_queue', on_message_callback=send_email)

# SMS Service
def send_sms(ch, method, properties, body):
    notif = json.loads(body)
    send_sms_to_user(notif['user_id'], notif['message'])
    ch.basic_ack(delivery_tag=method.delivery_tag)

channel.basic_consume(queue='sms_queue', on_message_callback=send_sms)

# Push Service
def send_push(ch, method, properties, body):
    notif = json.loads(body)
    send_push_to_user(notif['user_id'], notif['message'])
    ch.basic_ack(delivery_tag=method.delivery_tag)

channel.basic_consume(queue='push_queue', on_message_callback=send_push)

AVANTAGES :
[OK] Découplage total
[OK] Ajouter un nouveau canal = créer une queue + bind
[OK] Pas de modification du publisher
[OK] Chaque service scale indépendamment

CAS D'USAGE FANOUT EXCHANGE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[OK] Notifications multi-canal (email, SMS, push)
[OK] Real-time updates (WebSocket broadcast)
[OK] Cache invalidation (tous les nœuds)
[OK] Event broadcasting (microservices)
[OK] Metrics/monitoring (plusieurs systèmes)


3. TOPIC EXCHANGE (Pattern matching)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> Routing par PATTERN avec wildcards
-> Routing key = mots séparés par des points
-> Wildcards : * (un mot) et # (zéro ou plusieurs mots)

ALGORITHME :
if pattern_match(message.routing_key, binding.pattern):
    deliver_to(binding.queue)

ANALOGIE :
-> Système de catégories hiérarchiques
-> Filtrage par tags multiples

WILDCARDS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
* (étoile) : correspond à EXACTEMENT un mot
# (dièse) : correspond à ZÉRO ou PLUSIEURS mots

EXEMPLES DE PATTERNS :
"*.orange.*"     match   "quick.orange.rabbit"
                 match   "lazy.orange.fox"
                 NO match "quick.orange.male.rabbit" (3 mots, pas 2)

"*.*.rabbit"     match   "quick.orange.rabbit"
                 match   "lazy.brown.rabbit"
                 NO match "quick.rabbit" (seulement 2 mots)

"lazy.#"         match   "lazy"
                 match   "lazy.orange"
                 match   "lazy.orange.male.rabbit"
                 
"#.rabbit"       match   "rabbit"
                 match   "orange.rabbit"
                 match   "quick.orange.rabbit"
                 
"#"              match   TOUT (équivalent fanout)

"lazy.orange.*"  match   "lazy.orange.rabbit"
                 match   "lazy.orange.fox"
                 NO match "lazy.orange"

DÉCLARATION :
channel.exchange_declare(
    exchange='logs',
    exchange_type='topic',
    durable=True
)

EXEMPLE : SYSTÈME DE LOGS AVANCÉ
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

FORMAT DE ROUTING KEY :
<facility>.<severity>.<source>

Exemples :
- auth.error.api
- auth.info.web
- database.warning.mysql
- cache.error.redis

TOPOLOGIE :
┌────────────────────────────────────────────────────────────┐
│                                                            │
│         EXCHANGE: logs (type: topic)                       │
│                          │                                 │
│          ┌───────────────┼────────────────┐               │
│          │               │                │               │
│   "*.error.*"      "auth.#"        "*.*.api"              │
│          │               │                │               │
│    ┌─────[BLACK_DOWN-POINTING_TRIANGLE]─────┐   ┌────[BLACK_DOWN-POINTING_TRIANGLE]────┐    ┌──────[BLACK_DOWN-POINTING_TRIANGLE]──────┐       │
│    │  critical │   │   auth  │    │     api     │       │
│    │   queue   │   │  queue  │    │    queue    │       │
│    └───────────┘   └─────────┘    └─────────────┘       │
│                                                            │
└────────────────────────────────────────────────────────────┘

SETUP :
# Exchange
channel.exchange_declare(
    exchange='logs',
    exchange_type='topic',
    durable=True
)

# Queue pour tous les errors (toute facilité, toute source)
channel.queue_declare(queue='critical', durable=True)
channel.queue_bind(
    exchange='logs',
    queue='critical',
    routing_key='*.error.*'
)

# Queue pour tous les logs auth (tout niveau)
channel.queue_declare(queue='auth_logs', durable=True)
channel.queue_bind(
    exchange='logs',
    queue='auth_logs',
    routing_key='auth.#'
)

# Queue pour tous les logs de l'API (toute facilité, tout niveau)
channel.queue_declare(queue='api_logs', durable=True)
channel.queue_bind(
    exchange='logs',
    queue='api_logs',
    routing_key='*.*.api'
)

PUBLISHER :
# Log 1
channel.basic_publish(
    exchange='logs',
    routing_key='auth.error.api',
    body='Failed login attempt'
)
-> Routé vers :
   - critical (match *.error.*)
   - auth_logs (match auth.#)
   - api_logs (match *.*.api)

# Log 2
channel.basic_publish(
    exchange='logs',
    routing_key='database.warning.mysql',
    body='Slow query detected'
)
-> Routé vers : AUCUNE queue (aucun pattern match)

# Log 3
channel.basic_publish(
    exchange='logs',
    routing_key='auth.info.web',
    body='User logged in'
)
-> Routé vers :
   - auth_logs (match auth.#)

PATTERNS AVANCÉS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Tous les errors critiques de l'API
routing_key='*.error.api'

# Tous les logs auth (info, warning, error)
routing_key='auth.#'

# Tous les logs venant de services externes
routing_key='#.external'

# Errors et warnings seulement
routing_key='*.error.*'  # Binding 1
routing_key='*.warning.*'  # Binding 2 (même queue)

EXEMPLE : E-COMMERCE EVENTS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

FORMAT :
<entity>.<action>.<status>

Events possibles :
- order.created.pending
- order.updated.paid
- order.canceled.refunded
- order.shipped.completed
- product.created.draft
- product.updated.published
- user.registered.active

BINDINGS :
# Service Payment : tous les events de commande payés
routing_key='order.*.paid'

# Service Shipping : commandes à expédier
routing_key='order.updated.paid'
routing_key='order.*.ready_to_ship'

# Service Analytics : TOUS les events
routing_key='#'

# Service Notification : events importants
routing_key='order.created.*'
routing_key='order.shipped.*'
routing_key='order.canceled.*'

CAS D'USAGE TOPIC EXCHANGE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[OK] Logs hiérarchiques (facility.severity.source)
[OK] Events avec catégories (entity.action.status)
[OK] Multi-tenant routing (tenant_id.resource.action)
[OK] Géolocalisation (country.region.city)
[OK] Tout routing complexe avec patterns


4. HEADERS EXCHANGE (Routing par headers)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> Routing basé sur les HEADERS du message
-> Ignore la routing key
-> Match : all (tous les headers) ou any (au moins un)

ALGORITHME :
headers_match = check_headers(message.headers, binding.arguments)
if headers_match:
    deliver_to(binding.queue)

ANALOGIE :
-> Filtre de recherche avec plusieurs critères
-> AND (all) ou OR (any)

DÉCLARATION :
channel.exchange_declare(
    exchange='images',
    exchange_type='headers',
    durable=True
)

EXEMPLE : TRAITEMENT D'IMAGES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

TOPOLOGIE :
┌────────────────────────────────────────────────────────────┐
│                                                            │
│         EXCHANGE: images (type: headers)                   │
│                          │                                 │
│          ┌───────────────┼────────────────┐               │
│          │               │                │               │
│   {format: jpg,     {priority: high}  {format: png,       │
│    x-match: all}                       size: large,       │
│          │                             x-match: all}      │
│    ┌─────[BLACK_DOWN-POINTING_TRIANGLE]─────┐   ┌────[BLACK_DOWN-POINTING_TRIANGLE]────┐    ┌──────[BLACK_DOWN-POINTING_TRIANGLE]──────┐       │
│    │jpeg_proc  │   │ priority│    │  png_large  │       │
│    │   queue   │   │  queue  │    │    queue    │       │
│    └───────────┘   └─────────┘    └─────────────┘       │
│                                                            │
└────────────────────────────────────────────────────────────┘

SETUP :
# Exchange
channel.exchange_declare(
    exchange='images',
    exchange_type='headers',
    durable=True
)

# Queue 1 : Traiter SEULEMENT les JPEG
channel.queue_declare(queue='jpeg_processing', durable=True)
channel.queue_bind(
    exchange='images',
    queue='jpeg_processing',
    routing_key='',  # Ignoré
    arguments={
        'x-match': 'all',  # TOUS les headers doivent matcher
        'format': 'jpg'
    }
)

# Queue 2 : Traiter haute priorité (format quelconque)
channel.queue_declare(queue='priority_queue', durable=True)
channel.queue_bind(
    exchange='images',
    queue='priority_queue',
    routing_key='',
    arguments={
        'x-match': 'any',  # AU MOINS UN header doit matcher
        'priority': 'high',
        'urgent': 'true'
    }
)

# Queue 3 : PNG grandes images
channel.queue_declare(queue='png_large', durable=True)
channel.queue_bind(
    exchange='images',
    queue='png_large',
    routing_key='',
    arguments={
        'x-match': 'all',
        'format': 'png',
        'size': 'large'
    }
)

PUBLISHER :
# Image 1 : JPEG haute priorité
channel.basic_publish(
    exchange='images',
    routing_key='',
    body=image_data,
    properties=pika.BasicProperties(
        headers={
            'format': 'jpg',
            'size': 'medium',
            'priority': 'high'
        }
    )
)
-> Routé vers :
   - jpeg_processing (match format=jpg)
   - priority_queue (match priority=high)

# Image 2 : PNG large
channel.basic_publish(
    exchange='images',
    routing_key='',
    body=image_data,
    properties=pika.BasicProperties(
        headers={
            'format': 'png',
            'size': 'large',
            'priority': 'normal'
        }
    )
)
-> Routé vers :
   - png_large (match format=png AND size=large)

# Image 3 : PNG medium
channel.basic_publish(
    exchange='images',
    routing_key='',
    body=image_data,
    properties=pika.BasicProperties(
        headers={
            'format': 'png',
            'size': 'medium'
        }
    )
)
-> Routé vers : AUCUNE queue

x-match: all vs any :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

x-match: 'all' :
-> TOUS les headers du binding doivent être présents dans le message
-> Message peut avoir des headers supplémentaires

Binding : {'format': 'jpg', 'priority': 'high'}
Message : {'format': 'jpg', 'priority': 'high', 'size': 'large'}
-> MATCH [OK] (tous présents, + extra OK)

Message : {'format': 'jpg'}
-> NO MATCH [X] (priority manquant)

x-match: 'any' :
-> AU MOINS UN header du binding doit matcher

Binding : {'format': 'jpg', 'priority': 'high'}
Message : {'format': 'png', 'priority': 'high'}
-> MATCH [OK] (priority match)

Message : {'size': 'large'}
-> NO MATCH [X] (aucun header du binding présent)

NOTES IMPORTANTES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[ATTENTION] Headers exchange est le PLUS LENT
-> Doit parser et comparer tous les headers
-> Pas de hash lookup comme direct
-> Pas de trie comme topic

[ATTENTION] Moins utilisé en pratique
-> Topic exchange souvent suffisant
-> Headers utile pour cas très spécifiques

CAS D'USAGE HEADERS EXCHANGE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
[OK] Routing très complexe (multiples critères)
[OK] Filtrage dynamique
[OK] Traitement conditionnel (format, taille, type)
[OK] Quand routing key insuffisant


DEFAULT EXCHANGE (Exchange anonyme)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

NOM : '' (chaîne vide)
TYPE : direct
CARACTÉRISTIQUES :
-> Pré-créé par RabbitMQ
-> Ne peut pas être supprimé
-> Toutes les queues y sont automatiquement liées
-> Binding key = nom de la queue

UTILISATION :
channel.basic_publish(
    exchange='',         # Default exchange
    routing_key='my_queue',  # Nom de la queue
    body='Hello'
)

-> Équivalent à :
   Direct exchange avec binding key = nom de la queue

QUAND UTILISER :
-> Point-to-point simple
-> Pas besoin de routing complexe
-> Prototypage rapide


ALTERNATE EXCHANGE (Exchange de secours)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> Si message ne match aucun binding
-> Envoyé vers l'alternate exchange
-> Évite la perte de messages

SETUP :
# Créer l'alternate exchange (fanout généralement)
channel.exchange_declare(
    exchange='unrouted',
    exchange_type='fanout',
    durable=True
)

# Queue pour messages non-routés
channel.queue_declare(queue='unrouted_messages', durable=True)
channel.queue_bind(
    exchange='unrouted',
    queue='unrouted_messages'
)

# Exchange principal avec alternate
channel.exchange_declare(
    exchange='orders',
    exchange_type='topic',
    durable=True,
    arguments={
        'alternate-exchange': 'unrouted'
    }
)

COMPORTEMENT :
# Message avec routing key qui ne match rien
channel.basic_publish(
    exchange='orders',
    routing_key='unknown.action',
    body='...'
)
-> Aucun binding ne match
-> Envoyé vers 'unrouted' exchange
-> Arrive dans 'unrouted_messages' queue

UTILISATION :
-> Debugging (voir les messages perdus)
-> Fallback processing
-> Audit trail


COMPARAISON DES EXCHANGES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

┌─────────────┬──────────┬───────────┬──────────┬────────────────┐
│    Type     │Routing Key│ Wildcards │Performance│  Cas d'usage   │
├─────────────┼──────────┼───────────┼──────────┼────────────────┤
│   Direct    │   Oui    │    Non    │ Très bon │ 1:1 exact      │
│   Fanout    │   Non    │    Non    │ Excellent│ Broadcast      │
│   Topic     │   Oui    │    Oui    │   Bon    │ Patterns       │
│   Headers   │   Non    │    Non    │  Moyen   │ Critères multi │
│   Default   │   Oui    │    Non    │ Très bon │ Simple P2P     │
└─────────────┴──────────┴───────────┴──────────┴────────────────┘


BONNES PRATIQUES DE ROUTING :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] UTILISER LE BON TYPE :
-> Direct : 99% des cas simples
-> Topic : cas avec catégories/hiérarchie
-> Fanout : broadcast authentique
-> Headers : seulement si vraiment nécessaire

[OK] NAMING CONVENTIONS (Topic) :
-> Format cohérent : <entity>.<action>.<status>
-> Lowercase
-> Max 3-4 niveaux
-> Éviter les caractères spéciaux

[OK] DECLARER AVANT UTILISATION :
-> Producer ET Consumer déclarent exchanges/queues
-> Déclarations idempotentes (pas d'erreur si existe)

[OK] DURABLE = TRUE (production) :
-> Exchange survive au redémarrage
-> Queue survive (si durable=True)
-> Messages persistent (si delivery_mode=2)

[OK] ALTERNATE EXCHANGE :
-> Toujours configurer en production
-> Évite la perte silencieuse de messages

[OK] ÉVITER TROP DE BINDINGS :
-> 1000+ bindings = ralentissement
-> Simplifier la topologie si possible
"""


# [OK] PARTIE 7 : PATTERNS DE MESSAGING

"""
┌────────────────────────────────────────────────────────────────────────┐
│              PATTERNS DE MESSAGING CLASSIQUES                          │
└────────────────────────────────────────────────────────────────────────┘

RabbitMQ supporte tous les patterns de messaging
-> Point-to-Point
-> Publish/Subscribe
-> Request/Reply (RPC)
-> Work Queues
-> Routing
-> Topics


1. POINT-TO-POINT (Queue simple)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> 1 message -> 1 seul consommateur
-> Messages consommés une seule fois
-> Order garanti (FIFO)

TOPOLOGIE :
┌────────────┐      ┌────────┐      ┌────────────┐
│  Producer  │─────[BLACK_RIGHT-POINTING_TRIANGLE]│ Queue  │─────[BLACK_RIGHT-POINTING_TRIANGLE]│  Consumer  │
└────────────┘      └────────┘      └────────────┘

UTILISATION :
-> Default exchange
-> Routing key = nom de la queue

CODE :
# Producer
channel.queue_declare(queue='tasks', durable=True)

channel.basic_publish(
    exchange='',
    routing_key='tasks',
    body='Process this task'
)

# Consumer
def callback(ch, method, properties, body):
    print(f"Processing: {body}")
    # Work...
    ch.basic_ack(delivery_tag=method.delivery_tag)

channel.basic_consume(queue='tasks', on_message_callback=callback)
channel.start_consuming()

CAS D'USAGE :
-> Task queue
-> Command queue
-> Job processing


2. WORK QUEUES (Competing Consumers)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> 1 queue -> N consommateurs
-> Distribution round-robin (avec prefetch)
-> Load balancing automatique

TOPOLOGIE :
                    ┌─────────────┐
              ┌────[BLACK_RIGHT-POINTING_TRIANGLE]│  Consumer 1 │
              │     └─────────────┘
┌──────────┐  │     ┌─────────────┐
│ Producer │─[BLACK_RIGHT-POINTING_TRIANGLE]│Queue├────[BLACK_RIGHT-POINTING_TRIANGLE]│  Consumer 2 │
└──────────┘  │     └─────────────┘
              │     ┌─────────────┐
              └────[BLACK_RIGHT-POINTING_TRIANGLE]│  Consumer 3 │
                    └─────────────┘

FAIR DISPATCH :
-> Sans prefetch : round-robin simple
-> Avec prefetch=1 : dispatch vers worker libre

CODE :
# Producer (identique)
for i in range(100):
    channel.basic_publish(
        exchange='',
        routing_key='tasks',
        body=f'Task {i}'
    )

# Consumer (multiple instances)
channel.basic_qos(prefetch_count=1)  # Fair dispatch

def callback(ch, method, properties, body):
    print(f"[Worker {worker_id}] Processing: {body}")
    time.sleep(5)  # Simule travail long
    ch.basic_ack(delivery_tag=method.delivery_tag)

channel.basic_consume(queue='tasks', on_message_callback=callback)
channel.start_consuming()

# Lancer plusieurs workers :
# Terminal 1 : python worker.py
# Terminal 2 : python worker.py
# Terminal 3 : python worker.py

AVANTAGES :
[OK] Scalabilité horizontale facile
[OK] Résilience (si un worker crashe, messages redistribués)
[OK] Load balancing automatique

CAS D'USAGE :
-> Traitement d'images
-> Envoi d'emails en masse
-> Data processing
-> Video encoding


3. PUBLISH/SUBSCRIBE (Fan-out)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> 1 message -> TOUS les abonnés
-> Chaque abonné reçoit une copie
-> Abonnés indépendants

TOPOLOGIE :
                  ┌────────┐      ┌──────────┐
            ┌────[BLACK_RIGHT-POINTING_TRIANGLE]│ Queue1 │─────[BLACK_RIGHT-POINTING_TRIANGLE]│Consumer 1│
            │     └────────┘      └──────────┘
┌─────────┐ │     ┌────────┐      ┌──────────┐
│Producer │─┼────[BLACK_RIGHT-POINTING_TRIANGLE]│Exchange│─────[BLACK_RIGHT-POINTING_TRIANGLE]│ Queue2 │─[BLACK_RIGHT-POINTING_TRIANGLE]│Consumer 2│
└─────────┘ │     └────────┘      └──────────┘
            │     ┌────────┐      ┌──────────┐
            └────[BLACK_RIGHT-POINTING_TRIANGLE]│ Queue3 │─────[BLACK_RIGHT-POINTING_TRIANGLE]│Consumer 3│
                  └────────┘      └──────────┘

CODE :
# Publisher
channel.exchange_declare(
    exchange='events',
    exchange_type='fanout',
    durable=True
)

channel.basic_publish(
    exchange='events',
    routing_key='',
    body=json.dumps({'event': 'user_registered', 'user_id': 123})
)

# Subscriber 1 (Email)
channel.exchange_declare(exchange='events', exchange_type='fanout')

result = channel.queue_declare(queue='', exclusive=True)
queue_name = result.method.queue

channel.queue_bind(exchange='events', queue=queue_name)

def callback(ch, method, properties, body):
    event = json.loads(body)
    send_welcome_email(event['user_id'])
    ch.basic_ack(delivery_tag=method.delivery_tag)

channel.basic_consume(queue=queue_name, on_message_callback=callback)
channel.start_consuming()

# Subscriber 2 (Analytics) - code identique, traitement différent
# Subscriber 3 (CRM) - code identique, traitement différent

EXCLUSIVE QUEUES :
-> Queues temporaires pour chaque subscriber
-> Supprimées quand subscriber déconnecté
-> Pas de persistence nécessaire

CAS D'USAGE :
-> Event broadcasting
-> Real-time notifications
-> Cache invalidation
-> Monitoring alerts


4. ROUTING (Direct Exchange)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> Messages routés par clé exacte
-> Subscriber filtre les messages voulus

TOPOLOGIE :
                        ┌────────┐   ┌──────────┐
              ┌─"error"─│ Queue1 │──[BLACK_RIGHT-POINTING_TRIANGLE]│Consumer 1│
              │         └────────┘   └──────────┘
┌─────────┐   │   ┌──────────┐
│Producer │───┼──[BLACK_RIGHT-POINTING_TRIANGLE]│ Exchange │
└─────────┘   │   └──────────┘
              │         ┌────────┐   ┌──────────┐
              └─"info"─[BLACK_RIGHT-POINTING_TRIANGLE]│ Queue2 │──[BLACK_RIGHT-POINTING_TRIANGLE]│Consumer 2│
                        └────────┘   └──────────┘

CODE :
# Publisher
channel.exchange_declare(exchange='logs', exchange_type='direct')

severity = 'error'  # ou 'info', 'warning'
channel.basic_publish(
    exchange='logs',
    routing_key=severity,
    body=f'Log message with severity {severity}'
)

# Subscriber (errors seulement)
channel.queue_declare(queue='error_logs', durable=True)
channel.queue_bind(
    exchange='logs',
    queue='error_logs',
    routing_key='error'
)

channel.basic_consume(queue='error_logs', on_message_callback=callback)
channel.start_consuming()

CAS D'USAGE :
-> Log routing par sévérité
-> Task routing par type
-> Message filtering


5. TOPICS (Pattern-based routing)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> Routing par patterns avec wildcards
-> Très flexible

CODE :
# Publisher
channel.exchange_declare(exchange='events', exchange_type='topic')

channel.basic_publish(
    exchange='events',
    routing_key='order.created.paid',
    body=json.dumps({'order_id': 123})
)

# Subscriber (tous les orders)
channel.queue_bind(
    exchange='events',
    queue='order_processor',
    routing_key='order.#'
)

# Subscriber (seulement orders paid)
channel.queue_bind(
    exchange='events',
    queue='payment_processor',
    routing_key='order.*.paid'
)

CAS D'USAGE :
-> Event-driven architecture
-> Microservices communication
-> Multi-tenant systems


6. REQUEST/REPLY (RPC)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> Client envoie request -> queue
-> Server traite et répond -> reply queue
-> Synchrone (client attend réponse)

TOPOLOGIE :
┌────────┐  request   ┌──────────┐  request  ┌────────┐
│ Client │───────────[BLACK_RIGHT-POINTING_TRIANGLE]│rpc_queue │──────────[BLACK_RIGHT-POINTING_TRIANGLE]│ Server │
│        │            └──────────┘            │        │
│        │  response  ┌──────────┐  response │        │
│        │[BLACK_LEFT-POINTING_TRIANGLE]───────────│reply_queue[BLACK_LEFT-POINTING_TRIANGLE]───────────│        │
└────────┘            └──────────┘            └────────┘

MÉCANISME :
1. Client génère une reply queue exclusive
2. Client envoie request avec properties :
   - reply_to : nom de la reply queue
   - correlation_id : ID unique pour matcher requête/réponse
3. Server traite et répond à reply_to
4. Client reçoit réponse et match par correlation_id

CODE COMPLET :
# RPC Server
def fib(n):
    if n == 0:
        return 0
    elif n == 1:
        return 1
    else:
        return fib(n-1) + fib(n-2)

def on_request(ch, method, properties, body):
    n = int(body)
    print(f"Computing fib({n})")
    
    response = fib(n)
    
    ch.basic_publish(
        exchange='',
        routing_key=properties.reply_to,
        properties=pika.BasicProperties(
            correlation_id=properties.correlation_id
        ),
        body=str(response)
    )
    
    ch.basic_ack(delivery_tag=method.delivery_tag)

channel.queue_declare(queue='rpc_queue', durable=True)
channel.basic_qos(prefetch_count=1)
channel.basic_consume(queue='rpc_queue', on_message_callback=on_request)

print("RPC Server en attente...")
channel.start_consuming()


# RPC Client
import uuid

class FibonacciRpcClient:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
        
        # Créer reply queue exclusive
        result = self.channel.queue_declare(queue='', exclusive=True)
        self.callback_queue = result.method.queue
        
        self.channel.basic_consume(
            queue=self.callback_queue,
            on_message_callback=self.on_response,
            auto_ack=True
        )
        
        self.response = None
        self.corr_id = None
    
    def on_response(self, ch, method, properties, body):
        if self.corr_id == properties.correlation_id:
            self.response = body
    
    def call(self, n):
        self.response = None
        self.corr_id = str(uuid.uuid4())
        
        self.channel.basic_publish(
            exchange='',
            routing_key='rpc_queue',
            properties=pika.BasicProperties(
                reply_to=self.callback_queue,
                correlation_id=self.corr_id
            ),
            body=str(n)
        )
        
        # Attendre la réponse
        while self.response is None:
            self.connection.process_data_events()
        
        return int(self.response)

# Utilisation
client = FibonacciRpcClient()

print(f"fib(30) = {client.call(30)}")

NOTES IMPORTANTES :
[ATTENTION] RPC = Synchrone = Bloquant
-> Client attend la réponse
-> Si server slow/down, client bloqué

[ATTENTION] Pas vraiment "async messaging"
-> Utiliser seulement si nécessaire
-> Préférer async patterns si possible

CAS D'USAGE :
-> Microservices sync calls
-> Calculat

ions distribuées
-> Data validation services


7. PRIORITY QUEUES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> Messages avec priorité
-> Haute priorité traités en premier

CODE :
# Créer queue avec priorités
channel.queue_declare(
    queue='priority_tasks',
    durable=True,
    arguments={'x-max-priority': 10}
)

# Publier avec priorité
channel.basic_publish(
    exchange='',
    routing_key='priority_tasks',
    body='High priority task',
    properties=pika.BasicProperties(
        delivery_mode=2,
        priority=9  # 0-10
    )
)

channel.basic_publish(
    exchange='',
    routing_key='priority_tasks',
    body='Low priority task',
    properties=pika.BasicProperties(
        delivery_mode=2,
        priority=1
    )
)

COMPORTEMENT :
-> Messages triés par priorité dans la queue
-> Priorité 9 traité avant priorité 1

CAS D'USAGE :
-> VIP customers
-> Urgent tasks
-> Critical alerts


8. DELAYED MESSAGES (via Plugin)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> Message livré après un délai
-> Nécessite plugin rabbitmq_delayed_message_exchange

INSTALLATION :
rabbitmq-plugins enable rabbitmq_delayed_message_exchange

CODE :
# Créer delayed exchange
channel.exchange_declare(
    exchange='delayed',
    exchange_type='x-delayed-message',
    durable=True,
    arguments={'x-delayed-type': 'direct'}
)

channel.queue_declare(queue='tasks', durable=True)
channel.queue_bind(exchange='delayed', queue='tasks', routing_key='task')

# Publier avec délai
channel.basic_publish(
    exchange='delayed',
    routing_key='task',
    body='Execute this in 10 seconds',
    properties=pika.BasicProperties(
        headers={'x-delay': 10000}  # Délai en ms
    )
)

CAS D'USAGE :
-> Scheduled tasks
-> Reminder notifications
-> Retry avec backoff
"""

# [OK] PARTIE 8 : PERSISTENCE, DURABILITÉ ET FIABILITÉ

"""
┌────────────────────────────────────────────────────────────────────────┐
│              GARANTIR LA FIABILITÉ DES MESSAGES                        │
└────────────────────────────────────────────────────────────────────────┘

RabbitMQ offre plusieurs niveaux de garantie
-> De "best effort" à "exactly once"
-> Trade-off performance vs fiabilité


NIVEAUX DE GARANTIE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. AT-MOST-ONCE (0 ou 1 fois)
   -> Message peut être perdu
   -> Auto-ack, pas de persistence
   -> Plus rapide

2. AT-LEAST-ONCE (1 ou N fois)
   -> Message jamais perdu
   -> Peut être dupliqué
   -> Manual ack, persistence

3. EXACTLY-ONCE (1 fois exactement)
   -> Message ni perdu ni dupliqué
   -> Nécessite idempotence côté consumer
   -> Plus complexe


COMPOSANTS DE LA DURABILITÉ :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Pour une fiabilité COMPLÈTE, tout doit être durable :

1. EXCHANGE DURABLE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

channel.exchange_declare(
    exchange='orders',
    exchange_type='topic',
    durable=True  # [OK] Survit au redémarrage
)

-> Métadonnées stockées dans Mnesia
-> Exchange recréé automatiquement après crash

2. QUEUE DURABLE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

channel.queue_declare(
    queue='payment_processing',
    durable=True  # [OK] Survit au redémarrage
)

-> Métadonnées stockées dans Mnesia
-> Queue recréée automatiquement

ATTENTION :
-> durable=True ne rend PAS les messages durables !
-> Seulement la DÉFINITION de la queue

3. MESSAGES PERSISTENT
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

channel.basic_publish(
    exchange='orders',
    routing_key='order.created',
    body='...',
    properties=pika.BasicProperties(
        delivery_mode=2  # [OK] Persistent
    )
)

delivery_mode=1 : Transient (en RAM)
delivery_mode=2 : Persistent (sur disque)

STOCKAGE :
-> delivery_mode=1 : msg_store_transient/ (perdu au redémarrage)
-> delivery_mode=2 : msg_store_persistent/ (survit)

4. BINDINGS DURABLE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

-> Automatiquement durable si exchange ET queue sont durables
-> Stocké dans rabbit_durable_route.DCD


CONFIGURATION COMPLÈTE POUR DURABILITÉ MAXIMALE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Exchange durable
channel.exchange_declare(
    exchange='orders',
    exchange_type='topic',
    durable=True
)

# Queue durable
channel.queue_declare(
    queue='payment_processing',
    durable=True
)

# Binding (automatiquement durable)
channel.queue_bind(
    exchange='orders',
    queue='payment_processing',
    routing_key='order.#'
)

# Publisher
channel.basic_publish(
    exchange='orders',
    routing_key='order.created',
    body=json.dumps(order),
    properties=pika.BasicProperties(
        delivery_mode=2  # Persistent
    )
)

# Consumer avec manual ACK
def callback(ch, method, properties, body):
    try:
        process_order(body)
        ch.basic_ack(delivery_tag=method.delivery_tag)
    except Exception as e:
        ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)

channel.basic_qos(prefetch_count=1)
channel.basic_consume(
    queue='payment_processing',
    on_message_callback=callback,
    auto_ack=False  # Manual ACK
)


PUBLISHER CONFIRMS (Garantie de livraison)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PROBLÈME :
Sans confirmation, le publisher ne sait pas si :
-> Message bien reçu par RabbitMQ
-> Message bien écrit sur disque
-> Message bien routé vers une queue

SOLUTION : Publisher Confirms

ACTIVATION :
channel.confirm_delivery()

MÉTHODE 1 : SYNCHRONE (Bloquant)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

try:
    channel.basic_publish(
        exchange='orders',
        routing_key='order.created',
        body=json.dumps(order),
        properties=pika.BasicProperties(delivery_mode=2),
        mandatory=True  # Erreur si pas de queue
    )
    
    # Attendre confirmation
    channel.wait_for_confirms()
    print("Message confirmé [OK]")
    
except pika.exceptions.UnroutableError:
    print("Message non-routable [X]")
except pika.exceptions.NackError:
    print("Message refusé par le broker [X]")

-> Bloque jusqu'à confirmation
-> Simple mais LENT (1 publish à la fois)

MÉTHODE 2 : ASYNCHRONE (Callbacks)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

acked = []
nacked = []

def on_delivery_confirmation(method_frame):
    if method_frame.method.NAME == 'Basic.Ack':
        acked.append(method_frame.method.delivery_tag)
        print(f"Message {method_frame.method.delivery_tag} confirmé [OK]")
    elif method_frame.method.NAME == 'Basic.Nack':
        nacked.append(method_frame.method.delivery_tag)
        print(f"Message {method_frame.method.delivery_tag} refusé [X]")

channel.confirm_delivery(callback=on_delivery_confirmation)

# Publier plusieurs messages sans attendre
for i in range(100):
    channel.basic_publish(
        exchange='orders',
        routing_key='order.created',
        body=f'Order {i}',
        properties=pika.BasicProperties(delivery_mode=2)
    )

print(f"Confirmés : {len(acked)}, Refusés : {len(nacked)}")

-> Non-bloquant
-> Plus rapide (pipeline)

MÉTHODE 3 : BATCH (Best practice)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

channel.confirm_delivery()

messages = [...]

for msg in messages:
    channel.basic_publish(...)

# Attendre tous les confirms
channel.wait_for_confirms()

-> Pipeline de N messages
-> Attend la confirmation de tous
-> Bon compromis performance/fiabilité


MANDATORY FLAG (Détection messages non-routés)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PROBLÈME :
Message publié vers exchange, mais aucun binding ne match
-> Message PERDU silencieusement

SOLUTION :
channel.basic_publish(
    exchange='orders',
    routing_key='unknown.routing.key',
    body='...',
    mandatory=True  # [OK] Erreur si non-routable
)

-> Si aucun binding : basic.return
-> Exception UnroutableError

CALLBACK POUR BASIC.RETURN :
def on_return(channel, method, properties, body):
    print(f"Message non-routable : {method.routing_key}")
    # Logger, alerter, retry, etc.

channel.add_on_return_callback(on_return)


CONSUMER ACKNOWLEDGMENTS (Garantie de traitement)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

AUTO-ACK (Dangereux [X])
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

channel.basic_consume(
    queue='tasks',
    on_message_callback=callback,
    auto_ack=True  # [X] Dangereux
)

-> Message ACKé dès qu'envoyé au consumer
-> Si consumer crash AVANT traitement : message PERDU

MANUAL ACK (Recommandé [OK])
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

def callback(ch, method, properties, body):
    try:
        # Traitement
        process(body)
        
        # Succès -> ACK
        ch.basic_ack(delivery_tag=method.delivery_tag)
        
    except Exception as e:
        # Échec -> NACK + requeue
        ch.basic_nack(
            delivery_tag=method.delivery_tag,
            requeue=True
        )

channel.basic_consume(
    queue='tasks',
    on_message_callback=callback,
    auto_ack=False  # [OK] Manual ACK
)

-> ACK seulement si traitement réussi
-> Si consumer crash : message redelivered

REQUEUE vs NO REQUEUE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

requeue=True :
-> Message remis en queue (front)
-> Réessayé immédiatement
-> Risque de boucle infinie si erreur permanente

requeue=False :
-> Message supprimé (ou envoyé vers DLX si configuré)
-> Pas de retry
-> Utile pour erreurs permanentes

STRATÉGIE RECOMMANDÉE :
1. Erreurs temporaires (network timeout) : requeue=True
2. Erreurs permanentes (invalid data) : requeue=False + DLX
3. Compteur de retries dans message headers


PREFETCH COUNT (QoS)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

channel.basic_qos(prefetch_count=1)

-> Nombre max de messages non-ACKés qu'un consumer peut avoir
-> prefetch_count=1 : fair dispatch (recommandé)
-> prefetch_count=10 : plus de throughput mais unfair

IMPACT :

SANS PREFETCH (unfair) :
Consumer 1 : [fast] reçoit messages 1,3,5,7,9... (traite vite)
Consumer 2 : [slow] reçoit messages 2,4,6,8... (traite lentement)
-> Consumer 2 devient bottleneck

AVEC PREFETCH=1 (fair) :
Consumer 1 : [fast] reçoit message, traite, ACK, reçoit suivant
Consumer 2 : [slow] reçoit message, traite (plus long), ACK, reçoit suivant
-> Consumer 1 traite plus de messages automatiquement


DEAD LETTER EXCHANGE (DLX)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
Messages "morts" envoyés vers un exchange spécial
-> Permet de traiter les erreurs
-> Évite la perte de messages

MESSAGES "MORTS" :
1. Rejetés (NACK requeue=False)
2. TTL expiré
3. Queue pleine (x-max-length dépassé)

CONFIGURATION :
# DLX et queue
channel.exchange_declare(exchange='dlx', exchange_type='direct')
channel.queue_declare(queue='failed_orders', durable=True)
channel.queue_bind(exchange='dlx', queue='failed_orders', routing_key='failed')

# Queue principale avec DLX
channel.queue_declare(
    queue='orders',
    durable=True,
    arguments={
        'x-dead-letter-exchange': 'dlx',
        'x-dead-letter-routing-key': 'failed',
        'x-message-ttl': 300000  # 5 minutes
    }
)

COMPORTEMENT :
# Message rejeté
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
-> Envoyé vers DLX avec routing key 'failed'

# Message expiré (TTL)
-> Automatiquement envoyé vers DLX après 5 minutes

HEADERS AJOUTÉS :
-> x-death : historique des morts
-> x-first-death-exchange : premier exchange
-> x-first-death-queue : première queue
-> x-first-death-reason : raison (rejected, expired, maxlen)

CONSUMER DLX :
def handle_failed(ch, method, properties, body):
    print(f"Message failed: {body}")
    
    # Logger
    logger.error(f"Failed message: {body}, headers: {properties.headers}")
    
    # Alerter
    send_alert_to_ops_team()
    
    # ACK pour retirer de la DLQ
    ch.basic_ack(delivery_tag=method.delivery_tag)

channel.basic_consume(queue='failed_orders', on_message_callback=handle_failed)


RETRY PATTERN AVEC DLX
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
Retry avec backoff exponentiel

TOPOLOGIE :
orders_queue -> (fail) -> retry_queue (TTL) -> orders_queue
                              v (max retries)
                         dead_queue

CODE :
# Queue principale
channel.queue_declare(
    queue='orders',
    durable=True,
    arguments={
        'x-dead-letter-exchange': 'retry'
    }
)

# Retry exchange
channel.exchange_declare(exchange='retry', exchange_type='direct')

# Retry queue avec TTL
channel.queue_declare(
    queue='orders_retry',
    durable=True,
    arguments={
        'x-message-ttl': 5000,  # 5 secondes
        'x-dead-letter-exchange': '',  # Default exchange
        'x-dead-letter-routing-key': 'orders'  # Retour vers orders
    }
)

channel.queue_bind(exchange='retry', queue='orders_retry', routing_key='retry')

# Dead queue (max retries atteint)
channel.queue_declare(queue='orders_dead', durable=True)

CONSUMER :
def callback(ch, method, properties, body):
    # Vérifier nombre de retries
    if properties.headers is None:
        properties.headers = {}
    
    retry_count = properties.headers.get('x-retry-count', 0)
    
    if retry_count >= 3:
        # Max retries atteint -> dead queue
        ch.basic_publish(
            exchange='',
            routing_key='orders_dead',
            body=body,
            properties=properties
        )
        ch.basic_ack(delivery_tag=method.delivery_tag)
        return
    
    try:
        process_order(body)
        ch.basic_ack(delivery_tag=method.delivery_tag)
        
    except Exception as e:
        # Incrémenter retry count
        properties.headers['x-retry-count'] = retry_count + 1
        
        # Envoyer vers retry queue
        ch.basic_publish(
            exchange='retry',
            routing_key='retry',
            body=body,
            properties=properties
        )
        
        ch.basic_ack(delivery_tag=method.delivery_tag)


QUORUM QUEUES (Haute disponibilité)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> Réplication via Raft consensus
-> Pas de perte de messages possible
-> Alternative aux mirrored queues (deprecated)

CARACTÉRISTIQUES :
[OK] Réplication synchrone (quorum)
[OK] Pas de perte de données
[OK] Élection automatique du leader
[OK] Performance similaire à lazy queues

CRÉATION :
channel.queue_declare(
    queue='critical_orders',
    durable=True,
    arguments={
        'x-queue-type': 'quorum'
    }
)

RÉPLICATION :
-> Minimum 3 nœuds recommandé
-> Quorum = (N/2) + 1
-> 3 nœuds : quorum = 2
-> 5 nœuds : quorum = 3

WRITE :
-> Leader écrit
-> Attend ACK de (N/2) + 1 nœuds
-> Plus lent que classic queue
-> Mais garanti durable

READ :
-> Toujours depuis le leader
-> Pas de stale reads

QUAND UTILISER :
[OK] Messages critiques (paiements, commandes)
[OK] Clustering
[OK] Zéro perte acceptable

[X] Ne pas utiliser pour :
-> Messages temporaires
-> Performance critique
-> Single node


TRANSACTIONS (AMQP)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> Grouper plusieurs publishes
-> Commit ou Rollback atomique

[ATTENTION] TRÈS LENT : ~250x plus lent que normal
-> Utiliser Publisher Confirms à la place

CODE :
# Activer mode transactionnel
channel.tx_select()

try:
    # Publier plusieurs messages
    channel.basic_publish(exchange='', routing_key='q1', body='msg1')
    channel.basic_publish(exchange='', routing_key='q2', body='msg2')
    channel.basic_publish(exchange='', routing_key='q3', body='msg3')
    
    # Valider
    channel.tx_commit()
    print("Transaction committed [OK]")
    
except Exception as e:
    # Annuler
    channel.tx_rollback()
    print("Transaction rolled back [X]")

RECOMMANDATION :
[X] Ne pas utiliser transactions AMQP
[OK] Utiliser Publisher Confirms (100x plus rapide)


IDEMPOTENCE (Garantie Exactly-Once)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PROBLÈME :
At-Least-Once = messages peuvent être dupliqués
-> Consumer peut recevoir le même message 2 fois
-> Dû à : retry, network issues, redelivery

SOLUTION : IDEMPOTENCE côté consumer

MÉTHODE 1 : ID UNIQUE + CACHE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

import redis

redis_client = redis.Redis()

def callback(ch, method, properties, body):
    message_id = properties.message_id
    
    # Vérifier si déjà traité
    if redis_client.exists(f'processed:{message_id}'):
        print(f"Message {message_id} déjà traité, skip")
        ch.basic_ack(delivery_tag=method.delivery_tag)
        return
    
    # Traiter
    result = process_message(body)
    
    # Marquer comme traité (avec TTL)
    redis_client.setex(f'processed:{message_id}', 3600, '1')
    
    ch.basic_ack(delivery_tag=method.delivery_tag)

MÉTHODE 2 : OPERATION IDEMPOTENTE PAR NATURE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# [X] Non-idempotent
UPDATE users SET balance = balance + 100 WHERE id = 123

# [OK] Idempotent
UPDATE transactions SET status = 'completed' WHERE id = 456

# [OK] Idempotent (upsert)
INSERT INTO ... ON CONFLICT DO UPDATE ...


RÉSUMÉ : CHECKLIST FIABILITÉ MAXIMALE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PUBLISHER :
[ ] Exchange durable=True
[ ] Queue durable=True
[ ] delivery_mode=2 (persistent)
[ ] Publisher Confirms activé
[ ] mandatory=True
[ ] Retry logic

CONSUMER :
[ ] auto_ack=False (manual ACK)
[ ] prefetch_count=1
[ ] Try/catch avec NACK
[ ] Idempotence

QUEUE :
[ ] durable=True
[ ] x-dead-letter-exchange configuré
[ ] Quorum queue si critique

MONITORING :
[ ] Alertes sur DLQ
[ ] Métriques de redelivery
[ ] Logs des erreurs
"""


# [OK] PARTIE 9 : CLUSTERING ET HAUTE DISPONIBILITÉ

"""
┌────────────────────────────────────────────────────────────────────────┐
│              CLUSTERING RABBITMQ                                       │
└────────────────────────────────────────────────────────────────────────┘

RabbitMQ peut être clusterisé pour :
-> Haute disponibilité
-> Scalabilité
-> Résilience


ARCHITECTURE CLUSTER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

┌────────────────────────────────────────────────────────────────────┐
│                        RABBITMQ CLUSTER                            │
│                                                                    │
│   ┌──────────────┐    ┌──────────────┐    ┌──────────────┐      │
│   │   NODE 1     │────│   NODE 2     │────│   NODE 3     │      │
│   │  rabbit@srv1 │    │  rabbit@srv2 │    │  rabbit@srv3 │      │
│   │              │    │              │    │              │      │
│   │ - Mnesia     │    │ - Mnesia     │    │ - Mnesia     │      │
│   │ - Queues     │    │ - Queues     │    │ - Queues     │      │
│   │ - Exchanges  │    │ - Exchanges  │    │ - Exchanges  │      │
│   └──────────────┘    └──────────────┘    └──────────────┘      │
│          │                    │                    │              │
│          └────────────────────┴────────────────────┘              │
│                    Erlang Distribution                            │
│                   (port 25672 par défaut)                         │
└────────────────────────────────────────────────────────────────────┘

CARACTÉRISTIQUES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

METADATA RÉPLIQUÉE :
[OK] Exchanges
[OK] Queues (définitions)
[OK] Bindings
[OK] Users
[OK] Virtual hosts
[OK] Permissions
[OK] Policies

DONNÉES NON-RÉPLIQUÉES (Classic queues) :
[X] Messages dans les queues
[X] Contenu des queues

-> Queue existe sur UN SEUL nœud (master)
-> Si ce nœud crash : queue et messages perdus [X]

SOLUTION : Quorum Queues ou Mirrored Queues


PRÉ-REQUIS CLUSTERING :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. MÊME VERSION RABBITMQ
   -> Tous les nœuds même version exacte

2. MÊME ERLANG COOKIE
   -> Fichier .erlang.cookie identique

3. RÉSOLUTION DNS / HOSTS
   -> Chaque nœud doit résoudre les autres

4. PORTS OUVERTS
   -> 5672 : AMQP
   -> 25672 : Erlang distribution (inter-node)
   -> 15672 : Management UI

5. MÊME PLUGINS
   -> Mêmes plugins activés sur tous les nœuds


CRÉATION D'UN CLUSTER (3 nœuds)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ÉTAPE 1 : Préparer les nœuds
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Sur chaque serveur

# Installer RabbitMQ
sudo apt-get install rabbitmq-server

# Définir le même cookie
echo "SHARED_SECRET_COOKIE" | sudo tee /var/lib/rabbitmq/.erlang.cookie
sudo chmod 400 /var/lib/rabbitmq/.erlang.cookie
sudo chown rabbitmq:rabbitmq /var/lib/rabbitmq/.erlang.cookie

# Redémarrer
sudo systemctl restart rabbitmq-server

ÉTAPE 2 : Configurer /etc/hosts
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Sur chaque serveur
192.168.1.10  rabbit1
192.168.1.11  rabbit2
192.168.1.12  rabbit3

ÉTAPE 3 : Joindre le cluster
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Sur rabbit1 (master, rien à faire)

# Sur rabbit2
sudo rabbitmqctl stop_app
sudo rabbitmqctl reset
sudo rabbitmqctl join_cluster rabbit@rabbit1
sudo rabbitmqctl start_app

# Sur rabbit3
sudo rabbitmqctl stop_app
sudo rabbitmqctl reset
sudo rabbitmqctl join_cluster rabbit@rabbit1
sudo rabbitmqctl start_app

ÉTAPE 4 : Vérifier le cluster
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

sudo rabbitmqctl cluster_status

Cluster status of node rabbit@rabbit1 ...
Basics

Cluster name: rabbit@rabbit1

Disk Nodes

rabbit@rabbit1
rabbit@rabbit2
rabbit@rabbit3

Running Nodes

rabbit@rabbit1
rabbit@rabbit2
rabbit@rabbit3

Versions

rabbit@rabbit1: RabbitMQ 3.13.0 on Erlang 26.1
rabbit@rabbit2: RabbitMQ 3.13.0 on Erlang 26.1
rabbit@rabbit3: RabbitMQ 3.13.0 on Erlang 26.1


NODE TYPES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

DISK NODE :
-> Stocke metadata sur disque (Mnesia)
-> Au moins 1 disk node nécessaire
-> Recommandé : majorité en disk nodes

RAM NODE :
-> Metadata en RAM uniquement
-> Plus rapide
-> Perd metadata au crash (resync au redémarrage)

CRÉER RAM NODE :
sudo rabbitmqctl join_cluster --ram rabbit@rabbit1

CHANGER TYPE :
sudo rabbitmqctl change_cluster_node_type ram
sudo rabbitmqctl change_cluster_node_type disc

RECOMMANDATION :
-> 3 nœuds : 3 disk nodes
-> 5 nœuds : 3 disk + 2 ram (ou 5 disk)


QUORUM QUEUES (Réplication Raft)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PRINCIPE :
-> Réplication basée sur Raft consensus
-> Haute disponibilité native
-> Pas de perte de données

CRÉATION :
channel.queue_declare(
    queue='critical_orders',
    durable=True,
    arguments={
        'x-queue-type': 'quorum',
        'x-quorum-initial-group-size': 3  # 3 replicas
    }
)

FONCTIONNEMENT :
-> 1 leader, N-1 followers
-> Write : leader écrit, attend quorum
-> Read : toujours depuis leader
-> Élection automatique si leader crash

QUORUM :
-> 3 nœuds : quorum = 2
-> 5 nœuds : quorum = 3
-> Tolère (N-1)/2 pannes

EXEMPLE :
3 nœuds : tolère 1 panne
5 nœuds : tolère 2 pannes

AVANTAGES :
[OK] Pas de perte de données
[OK] Failover automatique
[OK] Pas de split-brain

INCONVÉNIENTS :
[X] Plus lent que classic queue
[X] Nécessite au moins 3 nœuds


MIRRORED QUEUES (Deprecated, legacy)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[ATTENTION] DEPRECATED depuis RabbitMQ 3.8
-> Utiliser Quorum Queues à la place

PRINCIPE :
-> Queue répliquée sur N nœuds
-> 1 master, N-1 mirrors

CONFIGURATION (via Policy) :
sudo rabbitmqctl set_policy ha-all "^orders\." \
  '{"ha-mode":"exactly","ha-params":3,"ha-sync-mode":"automatic"}' \
  --apply-to queues

-> Toutes les queues commençant par "orders." sont mirrored
-> 3 replicas
-> Synchronisation automatique

PROBLÈMES :
[X] Split-brain possible
[X] Performance dégradée
[X] Complexité

-> UTILISER QUORUM QUEUES


LOAD BALANCING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CLIENT SE CONNECTE À :
-> Un seul nœud à la fois
-> Mais peut se connecter à n'importe quel nœud

MÉTHODE 1 : LISTE DE NŒUDS (Driver)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Python
servers = [
    'rabbit1:5672',
    'rabbit2:5672',
    'rabbit3:5672'
]

for server in servers:
    try:
        connection = pika.BlockingConnection(
            pika.ConnectionParameters(server)
        )
        break
    except:
        continue

-> Client essaye chaque nœud jusqu'à connexion
-> Failover manuel

MÉTHODE 2 : HAProxy
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# /etc/haproxy/haproxy.cfg

listen rabbitmq
    bind *:5672
    mode tcp
    balance roundrobin
    
    option tcplog
    option tcp-check
    
    server rabbit1 192.168.1.10:5672 check inter 5s rise 2 fall 3
    server rabbit2 192.168.1.11:5672 check inter 5s rise 2 fall 3
    server rabbit3 192.168.1.12:5672 check inter 5s rise 2 fall 3

-> Client se connecte à HAProxy
-> HAProxy distribue vers nœuds sains

MÉTHODE 3 : DNS Round-Robin
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# DNS
rabbitmq.example.com  A  192.168.1.10
rabbitmq.example.com  A  192.168.1.11
rabbitmq.example.com  A  192.168.1.12

-> DNS retourne IPs en round-robin
-> Simple mais pas de health check


NETWORK PARTITIONS (Split-Brain)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PROBLÈME :
-> Coupure réseau entre nœuds
-> Cluster divisé en 2+ groupes
-> Chaque groupe pense être le cluster complet

COMPORTEMENT PAR DÉFAUT :
-> RabbitMQ détecte la partition
-> Met une alarme
-> Continue de fonctionner (dangereux!)

MODES DE GESTION :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

Configuration dans rabbitmq.conf :
cluster_partition_handling = MODE

MODES :
1. ignore (défaut) :
   -> Ne fait rien
   -> Dangereux [X]

2. pause_minority :
   -> Groupe minoritaire se pause
   -> Groupe majoritaire continue
   -> Recommandé [OK]

3. autoheal :
   -> Choisit un nœud "gagnant"
   -> Redémarre les autres
   -> Perte de données possible

4. pause_if_all_down :
   -> Se pause si tous les autres nœuds sont down

RECOMMANDATION :
cluster_partition_handling = pause_minority

-> Nécessite nombre impair de nœuds (3, 5, 7)


MAINTENANCE ET UPGRADES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ROLLING UPGRADE (sans downtime) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Nœud par nœud :

# Nœud 1
sudo rabbitmqctl stop_app
sudo apt-get update && sudo apt-get upgrade rabbitmq-server
sudo rabbitmqctl start_app

# Attendre que nœud 1 soit synchronisé

# Nœud 2
sudo rabbitmqctl stop_app
sudo apt-get update && sudo apt-get upgrade rabbitmq-server
sudo rabbitmqctl start_app

# Nœud 3
...

COMPATIBILITÉ :
-> N et N+1 compatibles
-> 3.12 peut cluster avec 3.13
-> Ne pas sauter de version majeure


RETIRER UN NŒUD :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Sur le nœud à retirer
sudo rabbitmqctl stop_app

# Sur un autre nœud
sudo rabbitmqctl forget_cluster_node rabbit@nodename

-> Métadonnées supprimées
-> Nœud retiré du cluster
"""


# [OK] PARTIE 10 : PERFORMANCE ET OPTIMISATIONS

"""
┌────────────────────────────────────────────────────────────────────────┐
│              PERFORMANCE RABBITMQ                                      │
└────────────────────────────────────────────────────────────────────────┘

FACTEURS IMPACTANT LA PERFORMANCE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. PERSISTENCE (delivery_mode=2)
   -> Écriture disque
   -> ~10x plus lent que transient

2. ACKNOWLEDGMENTS
   -> Manual ACK plus lent que auto-ack
   -> Mais nécessaire pour fiabilité

3. PREFETCH COUNT
   -> Impacte throughput
   -> prefetch=1 : fair mais moins de débit
   -> prefetch=100 : plus de débit mais unfair

4. MESSAGE SIZE
   -> Messages > 128KB : ralentissement
   -> Recommandation : < 1MB

5. QUEUE DEPTH
   -> Queue avec millions de messages : ralentissement
   -> Garder queues courtes

6. CONNECTIONS / CHANNELS
   -> Trop de connexions : overhead
   -> Réutiliser les connexions
   -> Pool de connexions

7. CLUSTER
   -> Communication inter-nœuds
   -> Quorum queues plus lentes que classic


BENCHMARKS TYPIQUES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

TRANSIENT MESSAGES (delivery_mode=1) :
-> 50K-100K msg/s (single queue, single node)
-> Messages 1KB

PERSISTENT MESSAGES (delivery_mode=2, SSD) :
-> 10K-20K msg/s (single queue, single node)

PERSISTENT + ACK :
-> 5K-15K msg/s

AVEC CLUSTERING :
-> Linear scaling pour publishing
-> Queues limitées au nœud master


OPTIMISATIONS CÔTÉ PUBLISHER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. BATCHING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# [X] Lent (1 par 1)
for msg in messages:
    channel.basic_publish(...)
    channel.wait_for_confirms()  # Attente à chaque fois

# [OK] Rapide (batch)
for msg in messages:
    channel.basic_publish(...)

channel.wait_for_confirms()  # Attente une fois

-> ~100x plus rapide

2. CONNECTION POOLING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Pool de connexions
from queue import Queue

connection_pool = Queue(maxsize=10)

for _ in range(10):
    conn = pika.BlockingConnection(...)
    connection_pool.put(conn)

# Utilisation
conn = connection_pool.get()
channel = conn.channel()
channel.basic_publish(...)
connection_pool.put(conn)  # Retour au pool

3. TRANSIENT SI POSSIBLE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Si perte acceptable
properties=pika.BasicProperties(delivery_mode=1)

-> ~10x plus rapide

4. COMPRESSION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

import gzip

body = json.dumps(large_message)
compressed = gzip.compress(body.encode())

channel.basic_publish(
    exchange='',
    routing_key='queue',
    body=compressed,
    properties=pika.BasicProperties(
        content_encoding='gzip'
    )
)

-> Réduit taille réseau
-> Trade-off CPU vs network


OPTIMISATIONS CÔTÉ CONSUMER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. PREFETCH OPTIMAL
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Si traitement rapide (<10ms)
channel.basic_qos(prefetch_count=100)

# Si traitement lent (>1s)
channel.basic_qos(prefetch_count=1)

-> Expérimenter pour trouver l'optimal

2. MULTIPLE CONSUMERS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Lancer N workers
for i in range(10):
    subprocess.Popen(['python', 'worker.py'])

-> Scaling horizontal
-> RabbitMQ fait le load balancing

3. ASYNC PROCESSING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

import asyncio
import aio_pika

async def consume():
    connection = await aio_pika.connect_robust("amqp://localhost/")
    channel = await connection.channel()
    await channel.set_qos(prefetch_count=10)
    
    queue = await channel.declare_queue('tasks')
    
    async with queue.iterator() as queue_iter:
        async for message in queue_iter:
            async with message.process():
                await process_message(message.body)

asyncio.run(consume())

-> Non-bloquant
-> Plus de throughput


OPTIMISATIONS QUEUE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. LAZY QUEUES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

channel.queue_declare(
    queue='big_queue',
    durable=True,
    arguments={'x-queue-mode': 'lazy'}
)

-> Messages sur disque immédiatement
-> Économise RAM
-> Supporte millions de messages
-> Légèrement plus lent mais stable

2. MESSAGE TTL
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

arguments={'x-message-ttl': 60000}  # 60 secondes

-> Auto-suppression des vieux messages
-> Garde queues courtes

3. MAX LENGTH
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

arguments={
    'x-max-length': 10000,
    'x-overflow': 'drop-head'  # Supprimer les plus anciens
}

-> Limite le nombre de messages
-> Évite saturation


OPTIMISATIONS SYSTÈME :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. FILE DESCRIPTORS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# /etc/systemd/system/rabbitmq-server.service.d/limits.conf
[Service]
LimitNOFILE=65536

sudo systemctl daemon-reload
sudo systemctl restart rabbitmq-server

-> Augmente limite de connexions

2. VM MEMORY
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# rabbitmq.conf
vm_memory_high_watermark.relative = 0.6

-> Utilise 60% de la RAM
-> Au-delà : flow control activé

3. DISK SPACE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

disk_free_limit.absolute = 2GB

-> Alarme si < 2GB

4. SSD vs HDD
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

-> SSD : ~10-20K msg/s persistent
-> HDD : ~500-2K msg/s persistent
-> SSD fortement recommandé


MONITORING PERFORMANCE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Via Management API
curl -u guest:guest http://localhost:15672/api/queues

# Métriques importantes :
- message_stats.publish_details.rate : msg/s published
- message_stats.deliver_details.rate : msg/s delivered
- messages : nombre de messages en attente
- consumers : nombre de consumers
- memory : RAM utilisée

# Via rabbitmqctl
sudo rabbitmqctl list_queues name messages consumers memory

# Prometheus
-> Plugin rabbitmq_prometheus
-> Export métriques vers Prometheus
"""

# [OK] PARTIE 11 : SÉCURITÉ ET MONITORING

"""
┌────────────────────────────────────────────────────────────────────────┐
│              SÉCURITÉ RABBITMQ                                         │
└────────────────────────────────────────────────────────────────────────┘

NIVEAUX DE SÉCURITÉ :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. AUTHENTIFICATION (Qui es-tu ?)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

UTILISATEUR PAR DÉFAUT :
-> Username : guest
-> Password : guest
-> [ATTENTION] Accès UNIQUEMENT depuis localhost (sécurité)

CRÉER UN UTILISATEUR :
sudo rabbitmqctl add_user myuser mypassword

# Définir comme administrateur
sudo rabbitmqctl set_user_tags myuser administrator

# Ou monitoring only
sudo rabbitmqctl set_user_tags myuser monitoring

# Ou management
sudo rabbitmqctl set_user_tags myuser management

TAGS DISPONIBLES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
administrator : Tous les droits
monitoring    : Vue lecture seule (Management UI)
management    : Gestion utilisateurs, vhosts, policies
policymaker   : Gestion policies uniquement
None          : Aucun accès Management UI

CHANGER MOT DE PASSE :
sudo rabbitmqctl change_password myuser newpassword

SUPPRIMER UTILISATEUR :
sudo rabbitmqctl delete_user myuser

LISTER UTILISATEURS :
sudo rabbitmqctl list_users

Listing users ...
user    tags
guest   [administrator]
myuser  [administrator]


2. AUTORISATION (Que peux-tu faire ?)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PERMISSIONS PAR VHOST :
-> Configure : créer/supprimer exchanges, queues
-> Write : publier des messages
-> Read : consommer des messages

FORMAT :
sudo rabbitmqctl set_permissions -p <vhost> <user> <configure> <write> <read>

EXEMPLES :
# Tous les droits sur vhost "/"
sudo rabbitmqctl set_permissions -p / myuser ".*" ".*" ".*"

# Lecture seule
sudo rabbitmqctl set_permissions -p / readonly_user "" "" ".*"

# Écriture seule (publisher)
sudo rabbitmqctl set_permissions -p / publisher_user "" ".*" ""

# Queues spécifiques (regex)
sudo rabbitmqctl set_permissions -p / limited_user "^orders.*" "^orders.*" "^orders.*"

LISTER PERMISSIONS :
sudo rabbitmqctl list_permissions -p /

Listing permissions for vhost "/" ...
user    configure  write  read
guest   .*         .*     .*
myuser  .*         .*     .*

SUPPRIMER PERMISSIONS :
sudo rabbitmqctl clear_permissions -p / myuser


3. VIRTUAL HOSTS (Isolation)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CRÉER VHOST :
sudo rabbitmqctl add_vhost production
sudo rabbitmqctl add_vhost staging
sudo rabbitmqctl add_vhost development

LISTER VHOSTS :
sudo rabbitmqctl list_vhosts

Listing vhosts ...
name
/
production
staging
development

PERMISSIONS SUR VHOST :
sudo rabbitmqctl set_permissions -p production prod_user ".*" ".*" ".*"
sudo rabbitmqctl set_permissions -p staging staging_user ".*" ".*" ".*"

SUPPRIMER VHOST :
sudo rabbitmqctl delete_vhost staging
# [ATTENTION] Supprime TOUT (exchanges, queues, messages)

CONNEXION À UN VHOST :
connection = pika.BlockingConnection(
    pika.ConnectionParameters(
        host='localhost',
        virtual_host='production',
        credentials=pika.PlainCredentials('prod_user', 'password')
    )
)


4. TLS/SSL (Chiffrement)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

GÉNÉRER CERTIFICATS (auto-signés pour test) :
# CA
openssl req -x509 -newkey rsa:4096 -days 365 -nodes \
  -keyout ca-key.pem -out ca-cert.pem \
  -subj "/CN=MyCA"

# Server certificate
openssl req -newkey rsa:4096 -nodes \
  -keyout server-key.pem -out server-req.pem \
  -subj "/CN=rabbitmq.example.com"

openssl x509 -req -in server-req.pem -days 365 \
  -CA ca-cert.pem -CAkey ca-key.pem -CAcreateserial \
  -out server-cert.pem

# Client certificate (optionnel)
openssl req -newkey rsa:4096 -nodes \
  -keyout client-key.pem -out client-req.pem \
  -subj "/CN=client"

openssl x509 -req -in client-req.pem -days 365 \
  -CA ca-cert.pem -CAkey ca-key.pem -CAcreateserial \
  -out client-cert.pem

CONFIGURATION RABBITMQ :
# rabbitmq.conf
listeners.ssl.default = 5671

ssl_options.cacertfile = /path/to/ca-cert.pem
ssl_options.certfile   = /path/to/server-cert.pem
ssl_options.keyfile    = /path/to/server-key.pem
ssl_options.verify     = verify_peer
ssl_options.fail_if_no_peer_cert = true

# Redémarrer
sudo systemctl restart rabbitmq-server

CLIENT AVEC TLS :
import ssl

context = ssl.create_default_context(cafile='ca-cert.pem')
context.load_cert_chain('client-cert.pem', 'client-key.pem')

parameters = pika.ConnectionParameters(
    host='rabbitmq.example.com',
    port=5671,
    ssl_options=pika.SSLOptions(context),
    credentials=pika.PlainCredentials('user', 'password')
)

connection = pika.BlockingConnection(parameters)


5. AUTRES MÉCANISMES D'AUTHENTIFICATION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

LDAP (Enterprise) :
# Installer plugin
sudo rabbitmq-plugins enable rabbitmq_auth_backend_ldap

# rabbitmq.conf
auth_backends.1 = ldap
auth_backends.2 = internal

auth_ldap.servers.1 = ldap.example.com
auth_ldap.user_dn_pattern = cn=${username},ou=users,dc=example,dc=com

OAUTH 2.0 / JWT :
# Via plugin rabbitmq_auth_backend_oauth2

EXTERNAL (certificat x.509) :
auth_mechanisms.1 = PLAIN
auth_mechanisms.2 = EXTERNAL


BONNES PRATIQUES SÉCURITÉ :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] DÉSACTIVER guest EN PRODUCTION :
sudo rabbitmqctl delete_user guest

[OK] UTILISER TLS :
-> Toujours en production
-> Évite man-in-the-middle

[OK] PRINCIPE DU MOINDRE PRIVILÈGE :
-> Chaque app a son propre user
-> Permissions minimales

[OK] FIREWALL :
-> Port 5672/5671 : seulement apps
-> Port 15672 : seulement admins
-> Port 25672 : seulement nœuds cluster

[OK] VHOSTS SÉPARÉS :
-> Par environnement (prod/staging/dev)
-> Par tenant (multi-tenant)

[OK] ROTATION CREDENTIALS :
-> Changer mots de passe régulièrement
-> Utiliser secrets management (Vault)

[OK] AUDIT LOGS :
-> Activer logging des connexions
-> Monitoring des accès


MONITORING RABBITMQ :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. MANAGEMENT UI
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ACCÈS :
http://localhost:15672
-> Username : guest
-> Password : guest

SECTIONS :
- Overview : Vue d'ensemble (msg/s, connexions)
- Connections : Connexions actives
- Channels : Channels ouverts
- Exchanges : Liste des exchanges
- Queues : Liste des queues (détails, graphes)
- Admin : Users, vhosts, policies

MÉTRIQUES CLÉS :
-> Message rates (publish, deliver, ack)
-> Queue depth (nombre de messages)
-> Consumer count
-> Connection count
-> Memory usage
-> Disk space


2. MANAGEMENT API (HTTP REST)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

EXEMPLES :
# Overview
curl -u guest:guest http://localhost:15672/api/overview

# Liste des queues
curl -u guest:guest http://localhost:15672/api/queues

# Détails d'une queue
curl -u guest:guest http://localhost:15672/api/queues/%2F/myqueue

# Purger une queue
curl -u guest:guest -X DELETE \
  http://localhost:15672/api/queues/%2F/myqueue/contents

# Publier un message
curl -u guest:guest -H "content-type:application/json" \
  -X POST http://localhost:15672/api/exchanges/%2F/amq.default/publish \
  -d '{"properties":{},"routing_key":"myqueue","payload":"hello","payload_encoding":"string"}'

PYTHON CLIENT :
import requests

class RabbitMQMonitor:
    def __init__(self, host='localhost', user='guest', password='guest'):
        self.base_url = f'http://{host}:15672/api'
        self.auth = (user, password)
    
    def get_queues(self):
        response = requests.get(f'{self.base_url}/queues', auth=self.auth)
        return response.json()
    
    def get_queue_depth(self, vhost, queue_name):
        url = f'{self.base_url}/queues/{vhost}/{queue_name}'
        response = requests.get(url, auth=self.auth)
        data = response.json()
        return data.get('messages', 0)
    
    def get_connections(self):
        response = requests.get(f'{self.base_url}/connections', auth=self.auth)
        return response.json()

monitor = RabbitMQMonitor()
queues = monitor.get_queues()

for queue in queues:
    print(f"Queue: {queue['name']}, Messages: {queue['messages']}")


3. PROMETHEUS METRICS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

INSTALLATION :
sudo rabbitmq-plugins enable rabbitmq_prometheus

CONFIGURATION :
# rabbitmq.conf
prometheus.tcp.port = 15692

ACCÈS METRICS :
curl http://localhost:15692/metrics

MÉTRIQUES EXPORTÉES :
# Queues
rabbitmq_queue_messages{queue="orders"} 1523
rabbitmq_queue_messages_ready{queue="orders"} 1500
rabbitmq_queue_messages_unacknowledged{queue="orders"} 23
rabbitmq_queue_consumers{queue="orders"} 3

# Memory
rabbitmq_process_resident_memory_bytes 1.5e+08

# Connections
rabbitmq_connections 42

# Message rates
rabbitmq_queue_messages_published_total{queue="orders"} 150000

PROMETHEUS CONFIG :
# prometheus.yml
scrape_configs:
  - job_name: 'rabbitmq'
    static_configs:
      - targets: ['localhost:15692']

GRAFANA DASHBOARD :
-> Dashboard ID : 10991 (RabbitMQ-Overview)
-> Import depuis grafana.com


4. ALERTING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ALERTES CRITIQUES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Queue depth trop élevé
alert: RabbitMQQueueDepthHigh
expr: rabbitmq_queue_messages{queue="critical"} > 10000
for: 5m
annotations:
  summary: "Queue depth élevé sur {{ $labels.queue }}"

# Pas de consumers
alert: RabbitMQNoConsumers
expr: rabbitmq_queue_consumers{queue="orders"} == 0
for: 5m
annotations:
  summary: "Aucun consumer sur queue {{ $labels.queue }}"

# Memory haute
alert: RabbitMQMemoryHigh
expr: rabbitmq_process_resident_memory_bytes > 2e+09
annotations:
  summary: "Memory RabbitMQ > 2GB"

# Disk plein
alert: RabbitMQDiskSpaceLow
expr: rabbitmq_disk_space_available_bytes < 2e+09
annotations:
  summary: "Espace disque < 2GB"

# Messages en erreur (DLQ)
alert: RabbitMQDeadLetterQueueGrowing
expr: increase(rabbitmq_queue_messages{queue="dlq"}[5m]) > 100
annotations:
  summary: "DLQ augmente rapidement"

# Connexions fermées anormalement
alert: RabbitMQConnectionsClosed
expr: increase(rabbitmq_connections_closed_total[5m]) > 50
annotations:
  summary: "Beaucoup de connexions fermées"


5. LOGS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

LOGS IMPORTANTS À MONITORER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Connexions
[info] accepting AMQP connection
[warning] closing AMQP connection (client unexpectedly closed TCP connection)

# Memory
[warning] memory resource limit alarm set on node rabbit@hostname

# Disk
[error] disk resource limit alarm set on node rabbit@hostname

# Cluster
[error] node rabbit@node2 down

CENTRALISER LOGS :
-> ELK Stack (Elasticsearch, Logstash, Kibana)
-> Splunk
-> Graylog

EXEMPLE FILEBEAT :
# filebeat.yml
filebeat.inputs:
  - type: log
    paths:
      - /var/log/rabbitmq/rabbit@*.log
    fields:
      service: rabbitmq

output.elasticsearch:
  hosts: ["localhost:9200"]


6. HEALTH CHECKS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RABBITMQCTL :
sudo rabbitmqctl node_health_check

Node rabbit@hostname health check passed

SCRIPT PYTHON :
import pika
import sys

def health_check():
    try:
        # Tester connexion
        connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost', socket_timeout=5)
        )
        channel = connection.channel()
        
        # Tester publish
        channel.basic_publish(
            exchange='',
            routing_key='health_check',
            body='ping'
        )
        
        connection.close()
        print("[OK] RabbitMQ is healthy")
        return 0
        
    except Exception as e:
        print(f"[X] RabbitMQ is unhealthy: {e}")
        return 1

sys.exit(health_check())

KUBERNETES LIVENESS PROBE :
livenessProbe:
  exec:
    command:
      - rabbitmq-diagnostics
      - ping
  initialDelaySeconds: 30
  periodSeconds: 10

readinessProbe:
  exec:
    command:
      - rabbitmq-diagnostics
      - check_port_connectivity
  initialDelaySeconds: 10
  periodSeconds: 5
"""


# [OK] PARTIE 12 : PREMIERS PAS PRATIQUES

"""
┌────────────────────────────────────────────────────────────────────────┐
│              TUTORIEL COMPLET ÉTAPE PAR ÉTAPE                          │
└────────────────────────────────────────────────────────────────────────┘

PROJET : SYSTÈME DE TRAITEMENT DE COMMANDES E-COMMERCE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ARCHITECTURE :
Web API -> [RabbitMQ] -> Payment Service
                    -> Inventory Service
                    -> Notification Service
                    -> Analytics Service


ÉTAPE 1 : INSTALLATION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Docker (le plus simple)
docker run -d --name rabbitmq \
  -p 5672:5672 \
  -p 15672:15672 \
  rabbitmq:3.13-management

# Attendre démarrage
sleep 10

# Vérifier
curl -u guest:guest http://localhost:15672/api/overview


ÉTAPE 2 : STRUCTURE DU PROJET
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ecommerce-rabbitmq/
├── config.py           # Configuration commune
├── web_api.py          # API qui publie les commandes
├── payment_service.py  # Service paiement
├── inventory_service.py # Service inventaire
├── notification_service.py # Service notification
├── analytics_service.py # Service analytics
└── requirements.txt


# requirements.txt
pika==1.3.2
flask==3.0.0


ÉTAPE 3 : CONFIGURATION COMMUNE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# config.py
import pika
import json

# Configuration RabbitMQ
RABBITMQ_HOST = 'localhost'
RABBITMQ_PORT = 5672
RABBITMQ_USER = 'guest'
RABBITMQ_PASS = 'guest'

# Exchanges
ORDERS_EXCHANGE = 'orders'
NOTIFICATIONS_EXCHANGE = 'notifications'

# Queues
PAYMENT_QUEUE = 'payment_processing'
INVENTORY_QUEUE = 'inventory_update'
NOTIFICATION_QUEUE = 'notifications_all'
ANALYTICS_QUEUE = 'analytics_events'

def get_connection():
    """Créer une connexion RabbitMQ"""
    credentials = pika.PlainCredentials(RABBITMQ_USER, RABBITMQ_PASS)
    parameters = pika.ConnectionParameters(
        host=RABBITMQ_HOST,
        port=RABBITMQ_PORT,
        credentials=credentials,
        heartbeat=600
    )
    return pika.BlockingConnection(parameters)

def setup_topology():
    """Créer la topologie (exchanges, queues, bindings)"""
    connection = get_connection()
    channel = connection.channel()
    
    # Exchange pour les commandes (topic)
    channel.exchange_declare(
        exchange=ORDERS_EXCHANGE,
        exchange_type='topic',
        durable=True
    )
    
    # Exchange pour les notifications (fanout)
    channel.exchange_declare(
        exchange=NOTIFICATIONS_EXCHANGE,
        exchange_type='fanout',
        durable=True
    )
    
    # Queue paiement
    channel.queue_declare(
        queue=PAYMENT_QUEUE,
        durable=True,
        arguments={
            'x-dead-letter-exchange': 'dlx',
            'x-message-ttl': 300000  # 5 minutes
        }
    )
    channel.queue_bind(
        exchange=ORDERS_EXCHANGE,
        queue=PAYMENT_QUEUE,
        routing_key='order.created'
    )
    
    # Queue inventaire
    channel.queue_declare(queue=INVENTORY_QUEUE, durable=True)
    channel.queue_bind(
        exchange=ORDERS_EXCHANGE,
        queue=INVENTORY_QUEUE,
        routing_key='order.paid'
    )
    
    # Queue notifications
    channel.queue_declare(queue=NOTIFICATION_QUEUE, durable=True)
    channel.queue_bind(
        exchange=NOTIFICATIONS_EXCHANGE,
        queue=NOTIFICATION_QUEUE
    )
    
    # Queue analytics
    channel.queue_declare(queue=ANALYTICS_QUEUE, durable=True)
    channel.queue_bind(
        exchange=ORDERS_EXCHANGE,
        queue=ANALYTICS_QUEUE,
        routing_key='order.#'  # Tous les events
    )
    
    # Dead Letter Exchange
    channel.exchange_declare(exchange='dlx', exchange_type='direct', durable=True)
    channel.queue_declare(queue='dead_letters', durable=True)
    channel.queue_bind(exchange='dlx', queue='dead_letters', routing_key='')
    
    connection.close()
    print("[OK] Topologie créée")

if __name__ == '__main__':
    setup_topology()


ÉTAPE 4 : WEB API (Publisher)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# web_api.py
from flask import Flask, request, jsonify
import pika
import json
import uuid
from datetime import datetime
from config import get_connection, ORDERS_EXCHANGE

app = Flask(__name__)

class OrderPublisher:
    def __init__(self):
        self.connection = get_connection()
        self.channel = self.connection.channel()
        self.channel.confirm_delivery()  # Publisher confirms
    
    def publish_order(self, order_data, routing_key):
        """Publier une commande"""
        try:
            message = json.dumps(order_data)
            
            self.channel.basic_publish(
                exchange=ORDERS_EXCHANGE,
                routing_key=routing_key,
                body=message,
                properties=pika.BasicProperties(
                    delivery_mode=2,  # Persistent
                    content_type='application/json',
                    message_id=str(uuid.uuid4()),
                    timestamp=int(datetime.now().timestamp()),
                    headers={'source': 'web-api'}
                ),
                mandatory=True
            )
            
            print(f"[OK] Published order {order_data['order_id']}")
            return True
            
        except pika.exceptions.UnroutableError:
            print("[X] Message non-routable")
            return False
        except Exception as e:
            print(f"[X] Erreur publication: {e}")
            return False
    
    def close(self):
        self.connection.close()

# Instance globale
publisher = OrderPublisher()

@app.route('/orders', methods=['POST'])
def create_order():
    """
    API endpoint pour créer une commande
    
    POST /orders
    {
        "customer_id": 12345,
        "items": [
            {"product_id": 1, "quantity": 2, "price": 50.00}
        ]
    }
    """
    data = request.json
    
    # Créer l'objet order
    order = {
        'order_id': str(uuid.uuid4()),
        'customer_id': data['customer_id'],
        'items': data['items'],
        'total': sum(item['price'] * item['quantity'] for item in data['items']),
        'status': 'created',
        'created_at': datetime.now().isoformat()
    }
    
    # Publier vers RabbitMQ
    success = publisher.publish_order(order, 'order.created')
    
    if success:
        return jsonify({
            'success': True,
            'order_id': order['order_id'],
            'message': 'Order created and queued for processing'
        }), 201
    else:
        return jsonify({
            'success': False,
            'message': 'Failed to queue order'
        }), 500

@app.route('/health', methods=['GET'])
def health():
    return jsonify({'status': 'healthy'}), 200

if __name__ == '__main__':
    app.run(port=5000, debug=True)


ÉTAPE 5 : PAYMENT SERVICE (Consumer)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# payment_service.py
import pika
import json
import time
import random
from config import get_connection, PAYMENT_QUEUE, ORDERS_EXCHANGE

def process_payment(order):
    """Simuler traitement paiement"""
    print(f"[CARTE] Processing payment for order {order['order_id']}")
    
    # Simuler appel API Stripe
    time.sleep(2)
    
    # Succès aléatoire (90%)
    success = random.random() < 0.9
    
    if success:
        print(f"[OK] Payment successful for order {order['order_id']}")
        return True
    else:
        print(f"[X] Payment failed for order {order['order_id']}")
        return False

def callback(ch, method, properties, body):
    """Callback pour chaque message reçu"""
    try:
        # Parser le message
        order = json.loads(body)
        print(f"\n[PACKAGE] Received order: {order['order_id']}")
        
        # Traiter le paiement
        success = process_payment(order)
        
        if success:
            # Paiement réussi -> publier event "order.paid"
            order['status'] = 'paid'
            ch.basic_publish(
                exchange=ORDERS_EXCHANGE,
                routing_key='order.paid',
                body=json.dumps(order),
                properties=pika.BasicProperties(delivery_mode=2)
            )
            
            # ACK
            ch.basic_ack(delivery_tag=method.delivery_tag)
            print(f"[OK] Order {order['order_id']} marked as paid")
            
        else:
            # Paiement échoué -> NACK (retry)
            print(f"[ATTENTION]  Payment failed, will retry...")
            ch.basic_nack(
                delivery_tag=method.delivery_tag,
                requeue=True
            )
    
    except Exception as e:
        print(f"[X] Error processing order: {e}")
        ch.basic_nack(
            delivery_tag=method.delivery_tag,
            requeue=True
        )

def main():
    """Démarrer le consumer"""
    connection = get_connection()
    channel = connection.channel()
    
    # QoS : 1 message à la fois
    channel.basic_qos(prefetch_count=1)
    
    # Consommer
    channel.basic_consume(
        queue=PAYMENT_QUEUE,
        on_message_callback=callback,
        auto_ack=False
    )
    
    print('[CARTE] Payment Service démarré. En attente de commandes...')
    channel.start_consuming()

if __name__ == '__main__':
    main()


ÉTAPE 6 : INVENTORY SERVICE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# inventory_service.py
import pika
import json
from config import get_connection, INVENTORY_QUEUE, NOTIFICATIONS_EXCHANGE

# Inventaire simulé
inventory = {
    1: {'name': 'Product A', 'stock': 100},
    2: {'name': 'Product B', 'stock': 50},
    3: {'name': 'Product C', 'stock': 200}
}

def update_inventory(order):
    """Mettre à jour l'inventaire"""
    print(f"[PACKAGE] Updating inventory for order {order['order_id']}")
    
    for item in order['items']:
        product_id = item['product_id']
        quantity = item['quantity']
        
        if product_id in inventory:
            inventory[product_id]['stock'] -= quantity
            print(f"  - {inventory[product_id]['name']}: "
                  f"{inventory[product_id]['stock']} remaining")

def callback(ch, method, properties, body):
    try:
        order = json.loads(body)
        print(f"\n[PACKAGE] Received paid order: {order['order_id']}")
        
        # Mettre à jour inventaire
        update_inventory(order)
        
        # Publier notification
        notification = {
            'type': 'inventory_updated',
            'order_id': order['order_id'],
            'message': f"Inventory updated for order {order['order_id']}"
        }
        
        ch.basic_publish(
            exchange=NOTIFICATIONS_EXCHANGE,
            routing_key='',
            body=json.dumps(notification),
            properties=pika.BasicProperties(delivery_mode=2)
        )
        
        # ACK
        ch.basic_ack(delivery_tag=method.delivery_tag)
        print(f"[OK] Inventory updated for order {order['order_id']}")
        
    except Exception as e:
        print(f"[X] Error: {e}")
        ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)

def main():
    connection = get_connection()
    channel = connection.channel()
    channel.basic_qos(prefetch_count=1)
    
    channel.basic_consume(
        queue=INVENTORY_QUEUE,
        on_message_callback=callback,
        auto_ack=False
    )
    
    print('[PACKAGE] Inventory Service démarré. En attente...')
    channel.start_consuming()

if __name__ == '__main__':
    main()


ÉTAPE 7 : NOTIFICATION SERVICE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# notification_service.py
import pika
import json
from config import get_connection, NOTIFICATION_QUEUE

def send_notification(notification):
    """Envoyer notification (simulé)"""
    print(f"[EMAIL] Sending notification: {notification['message']}")
    # Simuler envoi email/SMS/push

def callback(ch, method, properties, body):
    try:
        notification = json.loads(body)
        print(f"\n[EMAIL] Received notification")
        
        send_notification(notification)
        
        ch.basic_ack(delivery_tag=method.delivery_tag)
        print(f"[OK] Notification sent")
        
    except Exception as e:
        print(f"[X] Error: {e}")
        ch.basic_ack(delivery_tag=method.delivery_tag)  # ACK quand même

def main():
    connection = get_connection()
    channel = connection.channel()
    
    channel.basic_consume(
        queue=NOTIFICATION_QUEUE,
        on_message_callback=callback,
        auto_ack=False
    )
    
    print('[EMAIL] Notification Service démarré. En attente...')
    channel.start_consuming()

if __name__ == '__main__':
    main()


ÉTAPE 8 : ANALYTICS SERVICE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# analytics_service.py
import pika
import json
from collections import defaultdict
from config import get_connection, ANALYTICS_QUEUE

# Stats
stats = {
    'total_orders': 0,
    'total_revenue': 0.0,
    'orders_by_status': defaultdict(int)
}

def update_analytics(order):
    """Mettre à jour les analytics"""
    stats['total_orders'] += 1
    stats['total_revenue'] += order['total']
    stats['orders_by_status'][order['status']] += 1
    
    print(f"\n[GRAPHIQUE] Analytics Updated:")
    print(f"  Total Orders: {stats['total_orders']}")
    print(f"  Total Revenue: ${stats['total_revenue']:.2f}")
    print(f"  By Status: {dict(stats['orders_by_status'])}")

def callback(ch, method, properties, body):
    try:
        order = json.loads(body)
        update_analytics(order)
        ch.basic_ack(delivery_tag=method.delivery_tag)
        
    except Exception as e:
        print(f"[X] Error: {e}")
        ch.basic_ack(delivery_tag=method.delivery_tag)

def main():
    connection = get_connection()
    channel = connection.channel()
    
    channel.basic_consume(
        queue=ANALYTICS_QUEUE,
        on_message_callback=callback,
        auto_ack=False
    )
    
    print('[GRAPHIQUE] Analytics Service démarré. En attente...')
    channel.start_consuming()

if __name__ == '__main__':
    main()


ÉTAPE 9 : LANCER LE SYSTÈME
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Terminal 1 : Setup
python config.py

# Terminal 2 : Payment Service
python payment_service.py

# Terminal 3 : Inventory Service
python inventory_service.py

# Terminal 4 : Notification Service
python notification_service.py

# Terminal 5 : Analytics Service
python analytics_service.py

# Terminal 6 : Web API
python web_api.py

# Terminal 7 : Tester
curl -X POST http://localhost:5000/orders \
  -H "Content-Type: application/json" \
  -d '{
    "customer_id": 12345,
    "items": [
      {"product_id": 1, "quantity": 2, "price": 50.00},
      {"product_id": 2, "quantity": 1, "price": 30.00}
    ]
  }'


RÉSULTAT ATTENDU :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. Web API publie "order.created"
2. Payment Service reçoit, traite paiement
3. Payment Service publie "order.paid"
4. Inventory Service reçoit, met à jour stock
5. Inventory Service publie notification
6. Notification Service reçoit, envoie notification
7. Analytics Service reçoit tous les events, met à jour stats


AVANTAGES DE CETTE ARCHITECTURE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Découplage total des services
[OK] Scaling indépendant (lancer N instances de chaque service)
[OK] Résilience (si un service crash, messages en queue)
[OK] Retry automatique (via NACK)
[OK] Async processing (API répond immédiatement)
[OK] Extensibilité (ajouter nouveau service = créer queue + bind)
"""


# [OK] PARTIE 13 : CAS D'USAGE RÉELS

"""
┌────────────────────────────────────────────────────────────────────────┐
│              CAS D'USAGE RABBITMQ DANS LE MONDE RÉEL                   │
└────────────────────────────────────────────────────────────────────────┘

1. E-COMMERCE : ORDER PROCESSING PIPELINE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

FLUX :
Customer -> Web -> Queue -> Payment -> Queue -> Fulfillment -> Queue -> Shipping

TOPOLOGIE :
orders.created -> payment_queue
orders.paid -> inventory_queue -> shipping_queue
orders.* -> analytics_queue

BÉNÉFICES :
[OK] Peak handling (Black Friday)
[OK] Retry automatique si paiement échoue
[OK] Async (client n'attend pas)
[OK] Audit trail complet


2. MICROSERVICES COMMUNICATION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ARCHITECTURE :
┌─────────┐      ┌─────────┐      ┌─────────┐
│ Service │─────[BLACK_RIGHT-POINTING_TRIANGLE]│RabbitMQ │─────[BLACK_RIGHT-POINTING_TRIANGLE]│ Service │
│    A    │      │         │      │    B    │
└─────────┘      └─────────┘      └─────────┘

PATTERNS :
- Event-driven architecture
- CQRS (Command Query Responsibility Segregation)
- Event sourcing

EXEMPLE :
user.registered -> [email_service, crm_service, analytics_service]


3. IMAGE/VIDEO PROCESSING PIPELINE
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

FLUX :
Upload -> Queue -> Resize -> Queue -> Watermark -> Queue -> CDN

WORK QUEUES :
-> N workers par étape
-> Scaling dynamique selon la charge
-> Priority queue pour VIP users

TOPOLOGIE :
image.uploaded -> resize_queue
image.resized -> watermark_queue
image.processed -> cdn_queue


4. LOG AGGREGATION & MONITORING
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

ARCHITECTURE :
Apps -> RabbitMQ -> [Elasticsearch, S3, Alert Service]

ROUTING :
logs.*.error -> alert_queue
logs.*.* -> elasticsearch_queue
logs.* -> s3_archive_queue

BÉNÉFICES :
[OK] Centralized logging
[OK] Alerting temps réel
[OK] Archive long-terme


5. IoT DATA INGESTION
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

SCÉNARIO :
Millions de capteurs -> RabbitMQ (MQTT) -> Processing

TOPOLOGIE :
sensor.temperature.* -> temp_processor_queue
sensor.humidity.* -> humidity_processor_queue
sensor.# -> analytics_queue

PROTOCOLE :
-> MQTT plugin (rabbitmq_mqtt)
-> Léger pour IoT devices


6. SCHEDULED TASKS & CRON JOBS
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PROBLÈME :
-> Cron jobs sur un seul serveur = SPOF
-> Pas de retry si échec

SOLUTION :
Scheduler -> RabbitMQ -> Workers

EXEMPLE :
# Publisher (cron)
*/5 * * * * python publish_task.py --task daily_report

# Workers (N instances)
python worker.py

BÉNÉFICES :
[OK] Distributed cron
[OK] Retry automatique
[OK] Scaling workers


7. WEBHOOK RELAY
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PROBLÈME :
-> Webhook externe (Stripe, GitHub) attend réponse rapide
-> Traitement peut être long

SOLUTION :
Webhook -> Queue -> Async Processing

CODE :
@app.route('/webhooks/stripe', methods=['POST'])
def stripe_webhook():
    # Publier immédiatement
    channel.basic_publish(
        exchange='webhooks',
        routing_key='stripe.payment',
        body=request.data
    )
    return '', 200  # Réponse immédiate

# Worker traite async
def process_webhook(data):
    # Traitement long...
    pass


BONNES PRATIQUES PAR CAS D'USAGE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

E-COMMERCE :
[OK] Quorum queues pour commandes
[OK] DLX pour erreurs
[OK] TTL pour éviter queues infinies
[OK] Priority pour VIP

MICROSERVICES :
[OK] Topic exchange
[OK] Event versioning
[OK] Schema validation
[OK] Idempotence

IMAGE PROCESSING :
[OK] Work queues
[OK] Prefetch adapté à la durée de traitement
[OK] S3 pour images (pas dans messages)
[OK] Progress tracking via metadata

LOGS :
[OK] Batch publishing
[OK] Compression
[OK] Sampling (pas tous les logs)
[OK] Retention policies

IoT :
[OK] MQTT plugin
[OK] Batching
[OK] Aggregation avant queue
[OK] Time-series DB (InfluxDB)
"""


# [OK] PARTIE 14 : COMPARAISON AVEC AUTRES TECHNOLOGIES

"""
┌────────────────────────────────────────────────────────────────────────┐
│              RABBITMQ vs KAFKA vs REDIS vs ACTIVEMQ                    │
└────────────────────────────────────────────────────────────────────────┘

TABLEAU COMPARATIF DÉTAILLÉ :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

┌──────────────┬─────────────┬────────────┬──────────────┬─────────────┐
│Caractéristique│  RabbitMQ  │   Kafka    │ Redis Pub/Sub│  ActiveMQ   │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Type          │Message      │Distributed │Cache +       │Message      │
│              │Broker       │Log         │Pub/Sub       │Broker       │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Protocole     │AMQP 0-9-1   │Proprietary │Redis Protocol│JMS, AMQP    │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Langage       │Erlang       │Scala/Java  │C             │Java         │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Performance   │10-50K msg/s │1M+ msg/s   │100K+ msg/s   │10K msg/s    │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Latence       │1-10 ms      │10-50 ms    │<1 ms         │5-20 ms      │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Persistence   │Oui (opt-in) │Oui (natif) │Non (défaut)  │Oui          │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Routing       │Très flexible│Basique     │Basique       │Flexible     │
│              │(4 types     │(topics)    │(patterns)    │(selectors)  │
│              │exchanges)   │            │              │             │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Ordre garanti │Oui (queue)  │Oui (part.) │Non           │Oui          │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Replay msg    │Non (limité) │Oui         │Non           │Non          │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Retention     │Jusqu'à ACK  │Time-based  │Fire-forget   │Jusqu'à ACK  │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Scaling       │Clustering   │Partitions  │Cluster/      │Clustering   │
│              │             │(natif)     │Sentinel      │             │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Consumer      │Pull/Push    │Pull        │Push          │Pull/Push    │
│model         │             │            │              │             │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Use case      │Microservices│Streaming   │Cache,        │Enterprise   │
│principal     │Task queues  │Event log   │Real-time     │JEE apps     │
│              │RPC          │Analytics   │notifications │             │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Complexité    │Moyenne      │Haute       │Faible        │Moyenne      │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Maturité      │Très mature  │Mature      │Très mature   │Très mature  │
│              │(2007)       │(2011)      │(2009)        │(2004)       │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Écosystème    │Excellent    │Excellent   │Excellent     │Bon          │
├──────────────┼─────────────┼────────────┼──────────────┼─────────────┤
│Monitoring    │Excellent    │Bon         │Bon           │Bon          │
│              │(Management  │(JMX,       │(Redis CLI)   │(JMX)        │
│              │UI)          │Prometheus) │              │             │
└──────────────┴─────────────┴────────────┴──────────────┴─────────────┘


RABBITMQ vs KAFKA : DIFFÉRENCES FONDAMENTALES
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PHILOSOPHIE :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

RABBITMQ : Traditional Message Broker
-> Smart broker, dumb consumer
-> Push model
-> Messages supprimés après ACK
-> Queues comme abstraction

KAFKA : Distributed Log
-> Dumb broker, smart consumer
-> Pull model
-> Messages conservés (retention)
-> Partitions comme abstraction


QUAND UTILISER RABBITMQ :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Microservices communication
[OK] Task queues (background jobs)
[OK] RPC (Request/Reply)
[OK] Routing complexe (topic, headers)
[OK] Priority queues
[OK] Messages < 1MB
[OK] Garanties de livraison strictes
[OK] Simplicité (setup, ops)
[OK] Latence faible (<10ms)

EXEMPLES :
-> Order processing pipeline
-> Email sending service
-> Image processing workers
-> Webhook relay
-> Notification service


QUAND UTILISER KAFKA :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Event streaming (millions msg/s)
[OK] Log aggregation
[OK] Real-time analytics
[OK] Event sourcing
[OK] Replay messages nécessaire
[OK] Big Data pipelines
[OK] Time-series data
[OK] CDC (Change Data Capture)

EXEMPLES :
-> User activity tracking
-> Log aggregation (apps -> Kafka -> Elasticsearch)
-> Metrics collection
-> Stream processing (Kafka Streams, Flink)
-> Event store


QUAND UTILISER REDIS PUB/SUB :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Ultra faible latence (<1ms)
[OK] Fire-and-forget (pas de persistence)
[OK] Real-time notifications
[OK] WebSocket broadcasts
[OK] Chat applications
[OK] Live updates

EXEMPLES :
-> Chat rooms
-> Live dashboard updates
-> Real-time notifications
-> Cache invalidation


QUAND UTILISER ACTIVEMQ :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Enterprise Java (JMS standard)
[OK] Legacy systems
[OK] STOMP, MQTT support
[OK] JMS features (selectors, etc.)

[ATTENTION] Moins performant que RabbitMQ
[ATTENTION] Moins actif (développement)


MIGRATION RABBITMQ -> KAFKA (ou inverse) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

PAS NÉCESSAIREMENT OPPOSÉS :
-> Peuvent coexister !
-> RabbitMQ pour microservices sync
-> Kafka pour event streaming

ARCHITECTURE HYBRIDE :
Microservices <--> RabbitMQ <--> Kafka <--> Analytics

EXEMPLE :
1. Order Service publie vers RabbitMQ (queue payment)
2. Payment Service traite, publie event vers Kafka
3. Analytics consomme de Kafka pour dashboards
"""


# [OK] RÉCAPITULATIF FINAL

"""
┌────────────────────────────────────────────────────────────────────────┐
│                  RÉCAPITULATIF COMPLET RABBITMQ                        │
└────────────────────────────────────────────────────────────────────────┘

CONCEPTS CLÉS :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. MESSAGE BROKER :
   -> Intermédiaire entre services
   -> Découplage, async, résilience

2. COMPOSANTS :
   -> Producer : envoie messages
   -> Exchange : route messages
   -> Queue : stocke messages
   -> Consumer : reçoit messages
   -> Binding : règle de routing

3. TYPES D'EXCHANGES :
   -> Direct : routing exact
   -> Fanout : broadcast
   -> Topic : pattern matching
   -> Headers : routing par headers

4. DURABILITÉ :
   -> Exchange durable
   -> Queue durable
   -> Message persistent (delivery_mode=2)
   -> Manual ACK

5. FIABILITÉ :
   -> Publisher Confirms
   -> Consumer ACKs
   -> Dead Letter Exchange
   -> Retry patterns

6. PERFORMANCE :
   -> Prefetch count
   -> Batching
   -> Connection pooling
   -> Lazy queues pour gros volumes

7. HAUTE DISPONIBILITÉ :
   -> Clustering
   -> Quorum queues
   -> Load balancing

8. SÉCURITÉ :
   -> Authentification (users)
   -> Autorisation (permissions)
   -> TLS/SSL
   -> Virtual hosts


QUAND UTILISER RABBITMQ :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] OUI SI :
-> Architecture microservices
-> Task queues (background jobs)
-> Routing complexe
-> Garanties de livraison strictes
-> RPC patterns
-> Messages < 1MB
-> Setup simple

[X] NON SI :
-> Streaming massif (utiliser Kafka)
-> Replay messages essentiel (utiliser Kafka)
-> Latence ultra-faible (<1ms) (utiliser Redis)
-> Simple cache (utiliser Redis)


CHECKLIST PRODUCTION :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

CONFIGURATION :
[ ] Exchange/Queue durable=True
[ ] delivery_mode=2 (persistent)
[ ] Publisher Confirms activé
[ ] Manual ACK (auto_ack=False)
[ ] prefetch_count configuré
[ ] DLX configuré
[ ] TTL configuré (si applicable)

SÉCURITÉ :
[ ] Guest désactivé
[ ] TLS/SSL activé
[ ] Users avec permissions minimales
[ ] Firewall configuré
[ ] Vhosts séparés

HAUTE DISPONIBILITÉ :
[ ] Cluster 3+ nœuds
[ ] Quorum queues pour données critiques
[ ] Load balancer (HAProxy)
[ ] pause_minority configuré

MONITORING :
[ ] Management UI accessible
[ ] Prometheus metrics activé
[ ] Alertes configurées (queue depth, memory, etc.)
[ ] Logs centralisés
[ ] Health checks

PERFORMANCE :
[ ] SSD pour persistent messages
[ ] RAM suffisante (4GB+ par nœud)
[ ] File descriptors augmentés
[ ] vm_memory_high_watermark configuré


COMMANDES ESSENTIELLES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

# Status
sudo rabbitmqctl status

# Lister queues
sudo rabbitmqctl list_queues name messages consumers

# Purger queue
sudo rabbitmqctl purge_queue queue_name

# Créer user
sudo rabbitmqctl add_user myuser password
sudo rabbitmqctl set_permissions -p / myuser ".*" ".*" ".*"

# Plugins
sudo rabbitmq-plugins enable rabbitmq_management

# Cluster
sudo rabbitmqctl join_cluster rabbit@node1

# Logs
tail -f /var/log/rabbitmq/rabbit@hostname.log


RESSOURCES POUR ALLER PLUS LOIN :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[DOCS] DOCUMENTATION OFFICIELLE :
-> https://www.rabbitmq.com/documentation.html
-> Très complète, avec tutoriels

[COURS] TUTORIELS OFFICIELS :
-> https://www.rabbitmq.com/getstarted.html
-> 6 tutoriels progressifs

[GUIDE] LIVRES :
-> "RabbitMQ in Action" - Manning
-> "RabbitMQ Essentials" - Packt

[OUTILS] OUTILS :
-> Management UI (http://localhost:15672)
-> rabbitmqctl (CLI)
-> Prometheus + Grafana (monitoring)

[MOVIE_CAMERA] CHAÎNES YOUTUBE :
-> RabbitMQ Official Channel
-> CloudAMQP Tutorials

[SPEECH_BALLOON] COMMUNAUTÉ :
-> RabbitMQ mailing list
-> Stack Overflow (#rabbitmq)
-> RabbitMQ Slack


ERREURS COMMUNES À ÉVITER :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[X] Auto-ACK en production
[X] Pas de DLX
[X] Pas de TTL (queues infinies)
[X] Trop de connexions (pas de pooling)
[X] Messages > 1MB
[X] Pas de monitoring
[X] Guest en production
[X] Pas de TLS
[X] Classic queues en cluster (utiliser quorum)
[X] Ignorer les alarms (memory, disk)


CE QUE TU SAIS MAINTENANT :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

[OK] Pourquoi utiliser RabbitMQ
[OK] Architecture complète (Erlang, WiredTiger, etc.)
[OK] Tous les types d'exchanges
[OK] Patterns de messaging
[OK] Persistence et fiabilité
[OK] Clustering et HA
[OK] Performance et optimisations
[OK] Sécurité
[OK] Monitoring
[OK] Cas d'usage réels
[OK] Comparaison avec Kafka/Redis

Tu es maintenant prêt à utiliser RabbitMQ en production ! [RAPIDE][RABBIT_FACE]


PROCHAINES ÉTAPES :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

1. Installer RabbitMQ (Docker recommandé)
2. Suivre les tutoriels officiels (6 tutorials)
3. Créer un projet personnel (ex: task queue)
4. Expérimenter avec les exchanges
5. Tester en cluster
6. Mettre en production (checklist ci-dessus)

Bon développement avec RabbitMQ ! [RABBIT_FACE][FORCE]
"""