DevOps & Scaling

RabbitMQ + CaptchaAI : intégration de la file d'attente de messages

Un worker qui s'arrête au milieu d'un lot de 500 résolutions ne doit faire perdre aucune tâche : c'est ce que RabbitMQ apporte devant l'API CaptchaAI. Le courtier écrit les tâches sur disque, ne les retire qu'après acquittement explicite du worker, et pousse vers un dead letter exchange (DLX) tout ce qui a échoué. Ce guide monte le pipeline complet en Python.


Quand une file de messages devient nécessaire

Tant qu'un script traite une page à la fois, un appel synchrone suffit. La file devient utile dès qu'un de ces symptômes apparaît :

  • le temps de résolution bloque le thread applicatif — un reCAPTCHA v2 se résout en moins de 60 s, une éternité dans une requête HTTP ;
  • plusieurs producteurs (crawler, tests QA, tâches planifiées) se disputent la même clé API ;
  • un redéploiement fait disparaître les tâches en cours, sans trace.

Exemple : une équipe QA à Lille rejoue chaque nuit un parcours de checkout sur trois environnements de recette. Les scénarios Playwright publient dans captcha.tasks et rendent la main aussitôt ; quatre workers hébergés chez OVHcloud consomment la file à leur rythme, et les messages d'un worker redéployé repartent vers un autre.


Ce que RabbitMQ apporte au pipeline de résolution

Mécanisme Ce que vous y gagnez
Files durables Les tâches survivent au redémarrage du courtier
Acquittement Aucune tâche perdue si un worker s'arrête brutalement
Dead letter exchange Les échecs isolés dans une file dédiée, pas noyés dans les logs
Files prioritaires Une résolution interactive passe devant un lot
Clés de routage Chaque type de CAPTCHA vers ses workers

Dimensionner les threads CaptchaAI avant les workers

Le nombre de workers n'a de sens qu'au regard de votre plan : CaptchaAI facture des threads simultanés, avec des résolutions illimitées par thread. Un consommateur en prefetch_count=1 ne détenant qu'une résolution à la fois, le parallélisme utile est plafonné par le plan — cinq workers avec BASIC ($15/mois, 5 threads), quinze avec STANDARD ($30/mois, 15 threads), cinquante avec ADVANCE ($90/mois, 50 threads). Au-delà, les messages attendent dans captcha.tasks : la file absorbe les pics au lieu de les transformer en erreurs.


Installer le courtier et le client Python

Une image Docker et deux paquets suffisent. Le port 15672 expose le tableau de bord : gardez-le hors d'atteinte depuis internet.

# Docker
docker run -d --hostname rabbitmq \
  -p 5672:5672 -p 15672:15672 \
  rabbitmq:3-management

# Python client
pip install pika requests

Le producteur : publier les tâches dans captcha.tasks

Le producteur ne résout rien. Il déclare la topologie — DLX, file d'échecs, file principale avec TTL de 5 minutes, file de résultats — puis publie un message persistant (delivery_mode=2) par CAPTCHA à traiter. Le champ priority fait passer une demande interactive devant un lot.

import json
import uuid
import pika


class CaptchaProducer:
    """Submit CAPTCHA tasks to RabbitMQ."""

    def __init__(self, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url),
        )
        self.channel = self.connection.channel()
        self._setup_queues()

    def _setup_queues(self):
        """Declare durable queues and exchanges."""
        # Dead letter exchange for failed tasks
        self.channel.exchange_declare(
            exchange="captcha.dlx",
            exchange_type="direct",
            durable=True,
        )
        self.channel.queue_declare(
            queue="captcha.failed",
            durable=True,
        )
        self.channel.queue_bind(
            queue="captcha.failed",
            exchange="captcha.dlx",
            routing_key="failed",
        )

        # Main task queue with dead letter routing
        self.channel.queue_declare(
            queue="captcha.tasks",
            durable=True,
            arguments={
                "x-dead-letter-exchange": "captcha.dlx",
                "x-dead-letter-routing-key": "failed",
                "x-message-ttl": 300000,  # 5 min TTL
            },
        )

        # Results queue
        self.channel.queue_declare(
            queue="captcha.results",
            durable=True,
        )

    def submit(self, method, params, priority=0):
        """Submit a CAPTCHA task."""
        task_id = str(uuid.uuid4())[:8]
        task = {
            "id": task_id,
            "method": method,
            "params": params,
        }

        self.channel.basic_publish(
            exchange="",
            routing_key="captcha.tasks",
            body=json.dumps(task),
            properties=pika.BasicProperties(
                delivery_mode=2,  # Persistent
                priority=priority,
                message_id=task_id,
            ),
        )
        return task_id

    def close(self):
        self.connection.close()


# Usage
producer = CaptchaProducer()

task_id = producer.submit("userrecaptcha", {
    "googlekey": "SITE_KEY",
    "pageurl": "https://example.com",
}, priority=5)

print(f"Submitted: {task_id}")
producer.close()

Déclarer la topologie côté producteur évite les mauvaises surprises : le worker retrouve les mêmes files, avec les mêmes arguments. Si vous changez le TTL, supprimez la file avant de la redéclarer — RabbitMQ refuse une redéclaration dont les arguments diffèrent.


Le consommateur : le worker qui appelle l'API

Le worker enchaîne trois opérations : lire un message, envoyer la tâche à in.php, interroger res.php jusqu'au token. L'acquittement (basic_ack) n'intervient qu'après publication du résultat ; en cas d'exception, basic_nack(requeue=False) bascule le message vers le DLX.

import json
import os
import time
import pika
import requests


class CaptchaConsumer:
    """RabbitMQ consumer that solves CAPTCHAs."""

    def __init__(self, api_key, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
        self.api_key = api_key
        self.base = "https://ocr.captchaai.com"
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url),
        )
        self.channel = self.connection.channel()
        # Process one task at a time
        self.channel.basic_qos(prefetch_count=1)

    def start(self):
        """Start consuming tasks."""
        self.channel.basic_consume(
            queue="captcha.tasks",
            on_message_callback=self._handle_task,
        )
        print("Worker started. Waiting for tasks...")
        self.channel.start_consuming()

    def _handle_task(self, ch, method, properties, body):
        """Process a single CAPTCHA task."""
        task = json.loads(body)
        task_id = task["id"]
        print(f"Processing {task_id}...")

        try:
            token = self._solve(task["method"], task["params"])

            # Publish result
            result = {
                "task_id": task_id,
                "status": "success",
                "token": token,
            }
            ch.basic_publish(
                exchange="",
                routing_key="captcha.results",
                body=json.dumps(result),
                properties=pika.BasicProperties(delivery_mode=2),
            )

            # Acknowledge message (remove from queue)
            ch.basic_ack(delivery_tag=method.delivery_tag)
            print(f"{task_id} solved successfully")

        except Exception as e:
            print(f"{task_id} failed: {e}")

            # Reject and send to dead letter queue
            ch.basic_nack(
                delivery_tag=method.delivery_tag,
                requeue=False,  # Goes to DLX
            )

    def _solve(self, captcha_method, params, timeout=120):
        resp = requests.post(f"{self.base}/in.php", data={
            "key": self.api_key,
            "method": captcha_method,
            "json": 1,
            **params,
        }, timeout=30)
        result = resp.json()

        if result.get("status") != 1:
            raise RuntimeError(result.get("request"))

        captcha_id = result["request"]
        start = time.time()

        while time.time() - start < timeout:
            time.sleep(5)
            resp = requests.get(f"{self.base}/res.php", params={
                "key": self.api_key,
                "action": "get",
                "id": captcha_id,
                "json": 1,
            }, timeout=15)
            data = resp.json()
            if data["request"] != "CAPCHA_NOT_READY":
                if data.get("status") == 1:
                    return data["request"]
                raise RuntimeError(data["request"])

        raise TimeoutError("Solve timeout")


# Run worker
if __name__ == "__main__":
    consumer = CaptchaConsumer(os.environ["CAPTCHAAI_KEY"])
    consumer.start()

Deux réglages font la différence. prefetch_count=1 empêche un worker lent de réserver dix messages qu'il ne traitera pas ; le timeout de polling (120 s) reste au-dessus des délais annoncés — Turnstile se résout en moins de 10 s, un reCAPTCHA v2 en moins de 60 s. La clé API ne vit que dans la variable CAPTCHAAI_KEY, jamais dans le dépôt.


Récupérer les résultats côté application

La file captcha.results découple aussi la lecture : l'application vient chercher les tokens quand elle en a besoin, sans garder de connexion ouverte. Le collecteur enchaîne des basic_get jusqu'à réunir le nombre de résultats attendu ou atteindre son échéance.

import json
import pika


class ResultCollector:
    """Collect task results from the results queue."""

    def __init__(self, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url),
        )
        self.channel = self.connection.channel()
        self.results = {}

    def collect(self, expected_count, timeout=120):
        """Collect a specific number of results."""
        deadline = time.time() + timeout

        while len(self.results) < expected_count and time.time() < deadline:
            method, _, body = self.channel.basic_get(
                queue="captcha.results",
                auto_ack=True,
            )
            if body:
                result = json.loads(body)
                self.results[result["task_id"]] = result

            time.sleep(0.5)

        return self.results

Un token a une durée de vie courte côté site cible : consommez-le dans la foulée plutôt que de le stocker. Pour un traitement soumis au RGPD, limitez la journalisation à l'identifiant de tâche et à l'horodatage.


Router chaque type de CAPTCHA vers sa file

Un worker qui mélange CAPTCHA image (moins de 0,5 s) et reCAPTCHA v2 (moins de 60 s) produit des métriques illisibles. L'exchange direct captcha.types sépare les flux : une file et une clé de routage par type, des workers dimensionnés séparément.

# Setup exchanges and queues
channel.exchange_declare(
    exchange="captcha.types",
    exchange_type="direct",
    durable=True,
)

# Queue per type
for captcha_type in ["recaptcha", "turnstile", "image"]:
    channel.queue_declare(queue=f"captcha.{captcha_type}", durable=True)
    channel.queue_bind(
        queue=f"captcha.{captcha_type}",
        exchange="captcha.types",
        routing_key=captcha_type,
    )


# Submit with routing
def submit_routed(channel, captcha_type, task):
    channel.basic_publish(
        exchange="captcha.types",
        routing_key=captcha_type,
        body=json.dumps(task),
        properties=pika.BasicProperties(delivery_mode=2),
    )

Ces files correspondent aux types réellement pris en charge : reCAPTCHA v2 et v3, Cloudflare Turnstile, Cloudflare Challenge, GeeTest v3, CAPTCHA image/OCR et grilles d'images, plus CaptchaFox (bêta), Friendly Captcha (bêta) et Lemin (bêta). Inutile de prévoir une file hCaptcha ou FunCaptcha : ces types ne sont pas pris en charge, et GeeTest v4 reste annoncé comme à venir.


Exploitation : trois indicateurs suffisent

La profondeur de captcha.tasks dit si vos threads saturent, pas le courtier. Le débit de captcha.failed révèle un sitekey obsolète ou une pageurl incorrecte plus vite qu'un fichier de logs. L'âge du plus vieux message non acquitté détecte un worker bloqué mieux qu'une moyenne. Le tableau de bord de gestion affiche déjà ces trois valeurs : exportez-les vers votre supervision, et ajoutez un backoff exponentiel sur les erreurs réseau — sans lui, un incident transitoire vide la file vers le DLX en quelques secondes.


Dépannage

Problème Cause probable Correctif
Messages perdus après un crash File non durable, messages non persistants durable=True et delivery_mode=2
Worker bloqué sur une tâche Prefetch trop élevé, résolution longue prefetch_count=1 et timeout de polling explicite
captcha.failed qui grossit Sitekey ou pageurl erronés, solde épuisé Inspectez un message, corrigez, republiez
Connexion coupée pendant le polling Heartbeat AMQP expiré Heartbeat dans l'URL AMQP, reconnexion gérée
PRECONDITION_FAILED à la déclaration Arguments de file modifiés après coup Supprimez la file avant de la redéclarer

FAQ

Le plan CaptchaAI limite-t-il le nombre de messages traités par jour ?

Non. La facturation porte sur les threads simultanés, sans plafond journalier ni frais par résolution : seul le débit dépend du plan.

Que devient un message dont le TTL de 5 minutes expire ?

Il quitte captcha.tasks pour captcha.failed via le DLX, avec ses propriétés d'origine. C'est là que se décide un rejeu : une fois la cause corrigée, republier le message suffit.

Faut-il vraiment une file par type de CAPTCHA ?

Seulement quand deux types ont des temps de résolution très différents. Sinon, une file unique avec des priorités reste plus simple à exploiter.

Puis-je router des défis hCaptcha ou FunCaptcha vers ces workers ?

Non — ces deux types ne sont pas pris en charge par CaptchaAI, et GeeTest v4 reste annoncé comme à venir. Routez reCAPTCHA v2 et v3, Turnstile, Cloudflare Challenge, GeeTest v3 et les CAPTCHA image.


Guides connexes


Une file fiable de bout en bout — ouvrez un compte CaptchaAI et branchez-la sur RabbitMQ.

Les commentaires sont désactivés pour cet article.