Tutorials

Создание шины событий решения CAPTCHA с помощью Node.js и CaptchaAI

Если за отправку CAPTCHA в CaptchaAI, логирование, сбор метрик и повторные попытки в вашем Node.js-приложении отвечает один и тот же кусок кода, любое изменение логики решения превращается в правку сразу нескольких мест. Шина событий на базе стандартного EventEmitter разводит эти обязанности: submit-код только генерирует события жизненного цикла задачи (submitted, pending, solved, failed, timeout), а логирование, метрики и retry-логика подписываются на них независимо друг от друга — без единой строки жёсткой связки между модулями.

Ниже — рабочая реализация на JavaScript, её Python-аналог и типичные грабли, с которыми сталкиваются при переходе с прямого polling на событийную модель.

Как устроена шина событий

[CaptchaBus]
   ├── emit("submitted", { taskId, type, pageurl })
   ├── emit("pending", { taskId, elapsed })
   ├── emit("solved", { taskId, solution, duration })
   ├── emit("failed", { taskId, error, duration })
   └── emit("timeout", { taskId, elapsed })
        ↓          ↓           ↓
   [Logger]    [Metrics]   [Retry Handler]

Шина — это EventEmitter, через который проходят пять событий одного жизненного цикла задачи: заявка отправлена, идёт ожидание, решение получено, ошибка, истёк тайм-аут. Каждый слушатель — логгер, счётчик метрик, обработчик повторов — регистрируется сам по себе. Добавление нового слушателя (например, для сбора метрик) не требует ни строчки правок в коде отправки или опроса.

Класс CaptchaBus на JavaScript

Класс наследует EventEmitter и инкапсулирует весь цикл: отправку в in.php, опрос res.php и генерацию событий на каждом шаге. Идентификатор taskId присваивается локально в момент отправки — это ваш собственный ключ для сопоставления событий, отдельный от captchaId, который возвращает CaptchaAI.

const EventEmitter = require("events");
const axios = require("axios");

class CaptchaBus extends EventEmitter {
  constructor(apiKey, options = {}) {
    super();
    this.apiKey = apiKey;
    this.pollInterval = options.pollInterval || 5000;
    this.maxWait = options.maxWait || 300000; // 5 minutes
    this.pending = new Map();
  }

  async submit(params) {
    const { method, sitekey, pageurl, ...extra } = params;
    const taskId = `task_${Date.now()}_${Math.random().toString(36).slice(2, 8)}`;

    const submitParams = {
      key: this.apiKey,
      method: method || "userrecaptcha",
      googlekey: sitekey,
      pageurl: pageurl,
      json: 1,
      ...extra,
    };

    try {
      const resp = await axios.post(
        "https://ocr.captchaai.com/in.php",
        null,
        { params: submitParams }
      );

      if (resp.data.status !== 1) {
        this.emit("failed", {
          taskId,
          error: resp.data.request,
          duration: 0,
        });
        return null;
      }

      const captchaId = resp.data.request;
      const startTime = Date.now();

      this.emit("submitted", {
        taskId,
        captchaId,
        method: method || "userrecaptcha",
        pageurl,
      });

      // Start polling
      this._poll(taskId, captchaId, startTime);
      return taskId;
    } catch (err) {
      this.emit("failed", { taskId, error: err.message, duration: 0 });
      return null;
    }
  }

  async _poll(taskId, captchaId, startTime) {
    const check = async () => {
      const elapsed = Date.now() - startTime;

      if (elapsed > this.maxWait) {
        this.emit("timeout", { taskId, elapsed });
        return;
      }

      this.emit("pending", { taskId, elapsed });

      try {
        const resp = await axios.get("https://ocr.captchaai.com/res.php", {
          params: {
            key: this.apiKey,
            action: "get",
            id: captchaId,
            json: 1,
          },
        });

        if (resp.data.status === 1) {
          this.emit("solved", {
            taskId,
            captchaId,
            solution: resp.data.request,
            duration: Date.now() - startTime,
          });
        } else if (resp.data.request === "CAPCHA_NOT_READY") {
          setTimeout(check, this.pollInterval);
        } else {
          this.emit("failed", {
            taskId,
            error: resp.data.request,
            duration: Date.now() - startTime,
          });
        }
      } catch (err) {
        this.emit("failed", {
          taskId,
          error: err.message,
          duration: Date.now() - startTime,
        });
      }
    };

    setTimeout(check, this.pollInterval);
  }
}

module.exports = CaptchaBus;

Опрос запускается автоматически внутри submit() — вызывающему коду не нужно вручную вести цикл ожидания.

Как подключить обработчики событий

Каждый слушатель решает одну задачу и ничего не знает о других. В примере ниже один обработчик пишет лог в консоль, второй — накапливает метрики (submitted, solved, failed, суммарная длительность). Оба подписаны на одни и те же события независимо друг от друга.

const CaptchaBus = require("./captcha-bus");

const bus = new CaptchaBus(process.env.CAPTCHAAI_API_KEY, {
  pollInterval: 5000,
  maxWait: 120000,
});

// Logging listener
bus.on("submitted", (e) => {
  console.log(`[SUBMIT] ${e.taskId} → ${e.method} on ${e.pageurl}`);
});

bus.on("pending", (e) => {
  console.log(`[PENDING] ${e.taskId} — ${(e.elapsed / 1000).toFixed(1)}s`);
});

bus.on("solved", (e) => {
  console.log(
    `[SOLVED] ${e.taskId} in ${(e.duration / 1000).toFixed(1)}s — ${e.solution.substring(0, 30)}...`
  );
});

bus.on("failed", (e) => {
  console.error(`[FAILED] ${e.taskId} — ${e.error}`);
});

bus.on("timeout", (e) => {
  console.error(
    `[TIMEOUT] ${e.taskId} after ${(e.elapsed / 1000).toFixed(1)}s`
  );
});

// Metrics listener
const metrics = { submitted: 0, solved: 0, failed: 0, totalDuration: 0 };

bus.on("submitted", () => metrics.submitted++);
bus.on("solved", (e) => {
  metrics.solved++;
  metrics.totalDuration += e.duration;
});
bus.on("failed", () => metrics.failed++);

// Submit a CAPTCHA
bus.submit({
  sitekey: "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-",
  pageurl: "https://example.com",
});

Аналог на Python

Если основной сервис решения CAPTCHA написан на Python, а не на Node.js, тот же паттерн реализуется словарём слушателей и потоком (threading.Thread) для фонового опроса — принцип «emit → подписчики реагируют» переносится без изменений.

import os
import time
import threading
from collections import defaultdict
import requests


class CaptchaBus:
    def __init__(self, api_key, poll_interval=5, max_wait=300):
        self.api_key = api_key
        self.poll_interval = poll_interval
        self.max_wait = max_wait
        self._listeners = defaultdict(list)

    def on(self, event, callback):
        """Register a listener for an event."""
        self._listeners[event].append(callback)
        return self

    def emit(self, event, data):
        """Emit an event to all registered listeners."""
        for callback in self._listeners.get(event, []):
            try:
                callback(data)
            except Exception as e:
                print(f"Listener error on {event}: {e}")

    def submit(self, sitekey, pageurl, method="userrecaptcha", **extra):
        """Submit a CAPTCHA and begin tracking."""
        task_id = f"task_{int(time.time())}_{id(sitekey) % 10000}"

        resp = requests.post("https://ocr.captchaai.com/in.php", data={
            "key": self.api_key,
            "method": method,
            "googlekey": sitekey,
            "pageurl": pageurl,
            "json": 1,
            **extra
        })
        data = resp.json()

        if data.get("status") != 1:
            self.emit("failed", {
                "task_id": task_id,
                "error": data.get("request"),
                "duration": 0
            })
            return None

        captcha_id = data["request"]
        start_time = time.time()

        self.emit("submitted", {
            "task_id": task_id,
            "captcha_id": captcha_id,
            "method": method,
            "pageurl": pageurl
        })

        # Poll in a background thread
        thread = threading.Thread(
            target=self._poll,
            args=(task_id, captcha_id, start_time),
            daemon=True
        )
        thread.start()
        return task_id

    def _poll(self, task_id, captcha_id, start_time):
        while True:
            elapsed = time.time() - start_time

            if elapsed > self.max_wait:
                self.emit("timeout", {"task_id": task_id, "elapsed": elapsed})
                return

            time.sleep(self.poll_interval)
            self.emit("pending", {"task_id": task_id, "elapsed": elapsed})

            resp = requests.get("https://ocr.captchaai.com/res.php", params={
                "key": self.api_key,
                "action": "get",
                "id": captcha_id,
                "json": 1
            })
            data = resp.json()

            if data.get("status") == 1:
                self.emit("solved", {
                    "task_id": task_id,
                    "solution": data["request"],
                    "duration": time.time() - start_time
                })
                return
            elif data.get("request") != "CAPCHA_NOT_READY":
                self.emit("failed", {
                    "task_id": task_id,
                    "error": data.get("request"),
                    "duration": time.time() - start_time
                })
                return


# Usage
bus = CaptchaBus(os.environ["CAPTCHAAI_API_KEY"])

bus.on("submitted", lambda e: print(f"[SUBMIT] {e['task_id']}"))
bus.on("solved", lambda e: print(f"[SOLVED] {e['task_id']} in {e['duration']:.1f}s"))
bus.on("failed", lambda e: print(f"[FAILED] {e['task_id']} — {e['error']}"))
bus.on("timeout", lambda e: print(f"[TIMEOUT] {e['task_id']}"))

bus.submit("6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-", "https://example.com")

Автоматические повторы через отдельный обработчик

Retry-логика — хороший пример того, зачем вообще нужна шина событий: обработчик подписывается на failed и сам решает, стоит ли повторить отправку, без единой правки в основном классе CaptchaBus.

// Automatic retry on failure
bus.on("failed", async (e) => {
  if (e.retryCount >= 3) {
    console.error(`[GIVE UP] ${e.taskId} after 3 retries`);
    return;
  }

  console.log(`[RETRY] ${e.taskId} — attempt ${(e.retryCount || 0) + 1}`);
  await bus.submit({
    ...e.originalParams,
    _retryCount: (e.retryCount || 0) + 1,
  });
});

Promise-обёртка над шиной событий

Если остальной код проекта построен на async/await, а не на слушателях, оберните шину в промис: функция подписывается на solved/failed/timeout для конкретного taskId, снимает слушатели после первого срабатывания и резолвит или реджектит промис.

function solveCaptcha(bus, params) {
  return new Promise((resolve, reject) => {
    const taskId = bus.submit(params);

    function onSolved(e) {
      if (e.taskId === taskId) {
        cleanup();
        resolve(e.solution);
      }
    }

    function onFailed(e) {
      if (e.taskId === taskId) {
        cleanup();
        reject(new Error(e.error));
      }
    }

    function cleanup() {
      bus.removeListener("solved", onSolved);
      bus.removeListener("failed", onFailed);
      bus.removeListener("timeout", onFailed);
    }

    bus.on("solved", onSolved);
    bus.on("failed", onFailed);
    bus.on("timeout", onFailed);
  });
}

// Usage
const solution = await solveCaptcha(bus, {
  sitekey: "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-",
  pageurl: "https://example.com",
});

Что учесть при сохранении событий

Если вы пишете события в файл или БД для последующего аудита, сохраняйте только то, что реально нужно для отладки. pageurl из события submitted — это адрес страницы клиента, и в нём иногда встречаются query-параметры с идентифицирующими данными. Для аудитории из РФ это прямое следствие 152-ФЗ «О персональных данных» — собирайте только те данные, которые вы вправе обрабатывать; для международных команд разумно придерживаться той же дисциплины в духе GDPR.

Полезно также сверять число одновременно ожидающих задач с лимитом потоков вашего тарифа. План BASIC ($15/мес, 5 потоков) выдержит не более пяти задач в состоянии pending одновременно; при более высокой нагрузке имеет смысл перейти на ADVANCE ($90/мес, 50 потоков) — иначе новые заявки будут просто ждать своей очереди, а не теряться.

Типичные проблемы и их решение

Проблема Причина Решение
Обработчик не срабатывает Расхождение в имени события (например, solve вместо solved) Сверьте точные имена событий в emit/on
Предупреждение об утечке памяти Слишком много слушателей на одном событии Используйте setMaxListeners() или снимайте слушателей после использования
Событие pending заваливает консоль Слишком короткий pollInterval Увеличьте pollInterval минимум до 5000 мс
События теряются при повторной попытке При повторе генерируется новый taskId Передавайте исходные параметры дальше, чтобы восстановить состояние

Часто задаваемые вопросы

Нужна ли шина событий, если я решаю только одну CAPTCHA за раз?

Нет смысла. Для одиночного вызова прямой await вокруг submit/poll проще и читается лучше. Шина событий окупается, когда результат нужен нескольким независимым частям приложения одновременно — логированию, метрикам, UI, retry-логике.

Как подключить внешний брокер сообщений вроде Kafka или RabbitMQ?

EventEmitter живёт в рамках одного процесса. Если события должны попадать в другие сервисы, добавьте ещё один слушатель, который публикует их в брокер (bus.on("solved", (e) => producer.send(...))) — сама шина при этом не меняется, она просто получает ещё одного подписчика.

Как сохранить события для последующего аудита и отладки?

Подпишите отдельный слушатель, который пишет каждое событие в JSONL-файл или таблицу БД. Это не требует изменений в логике решения — см. раздел про 152-ФЗ и GDPR выше о том, какие поля стоит логировать, а какие лучше не сохранять.

Что произойдёт, если процесс Node.js перезапустится, пока задача в статусе pending?

Состояние живёт только в памяти процесса, поэтому запись о задаче будет потеряна вместе с процессом. Для устойчивости к перезапускам сохраняйте taskId/captchaId во внешнем хранилище (Redis, БД) при получении события submitted и восстанавливайте опрос из него при старте приложения.

Как не превысить лимит потоков тарифа при большом количестве задач?

Считайте события submitted без парного solved/failed/timeout — это и есть число задач, реально занимающих потоки прямо сейчас. Если счётчик регулярно упирается в лимит потоков вашего плана, это сигнал повысить тариф или добавить очередь с ограничением параллелизма перед submit().

Что дальше

Подключите шину событий к своему пайплайну и переходите от прямого опроса к событийной архитектуре — начните с быстрого старта выше.

Комментарии для этой статьи отключены.