DevOps & Scaling

Rolling-обновления парка воркеров для решения CAPTCHA

Задача, уже отправленная в in.php, живёт на стороне сервиса ещё десятки секунд: reCAPTCHA v2 решается менее чем за 60 с, Cloudflare Turnstile — менее чем за 10 с. Если в этот момент погасить воркер, поток всё равно занят, а забрать готовый токен из res.php уже некому. Отсюда единственное правило поочерёдного обновления: сначала перестать принимать новые задачи, дождаться завершения активных и только потом подменять версию. Ниже — рабочая схема: дренаж, проверка работоспособности после каждого шага и автоматический откат, если новая версия не прошла проверку.

Что ломается при наивном деплое

Три сценария, превращающих обновление парка в инцидент:

  • Перезапуск всего парка разом. Пропускная способность падает до нуля, очередь копится, а после подъёма парк получает пиковую нагрузку и сам себе создаёт тайм-ауты.
  • Остановка воркера с активными задачами. ID задачи теряется вместе с процессом. Поток на стороне CaptchaAI остаётся занятым до завершения решения, но результат никто не заберёт — задачу придётся отправлять заново.
  • Деплой без проверки после старта. Процесс поднялся, порт слушает, а API-ключ подтянулся из старого конфига. Формально «зелёный» воркер отдаёт ошибки на каждой задаче, и вы узнаете об этом из графика успешных решений, а не из отчёта о деплое.

Поочерёдное обновление закрывает все три случая: воркеры меняются по одному, каждый проходит дренаж, а переход к следующему возможен только после успешной проверки предыдущего.

Схема поочерёдного обновления

Workers: [W1-old] [W2-old] [W3-old] [W4-old]

Step 1:  [W1-drain] [W2-old]  [W3-old]  [W4-old]
Step 2:  [W1-NEW✓]  [W2-old]  [W3-old]  [W4-old]
Step 3:  [W1-NEW✓]  [W2-drain] [W3-old]  [W4-old]
Step 4:  [W1-NEW✓]  [W2-NEW✓]  [W3-old]  [W4-old]
  ...until all updated

Каждый воркер проходит одну и ту же последовательность состояний: runningdrainingupdatingrunning. Пока идёт дренаж, маршрутизатор задач обязан исключить воркер из выдачи — иначе он получит новую задачу ровно в тот момент, когда собирается остановиться.

Python: оркестратор с дренажем и health-гейтом

Оркестратор ниже хранит состояние каждого воркера, направляет задачи только в состояние RUNNING, а при неуспешной проверке откатывает уже обновлённые экземпляры.

import os
import time
import signal
import threading
import requests
from dataclasses import dataclass, field
from enum import Enum

API_KEY = os.environ["CAPTCHAAI_API_KEY"]


class WorkerState(Enum):
    RUNNING = "running"
    DRAINING = "draining"
    STOPPED = "stopped"
    UPDATING = "updating"


@dataclass
class Worker:
    worker_id: str
    version: str
    state: WorkerState = WorkerState.RUNNING
    active_tasks: int = 0
    tasks_completed: int = 0
    session: requests.Session = field(default_factory=requests.Session)

    def solve(self, task):
        if self.state != WorkerState.RUNNING:
            return {"error": "WORKER_NOT_ACCEPTING"}

        self.active_tasks += 1
        try:
            result = self._do_solve(task)
            self.tasks_completed += 1
            return result
        finally:
            self.active_tasks -= 1

    def _do_solve(self, task):
        resp = self.session.post("https://ocr.captchaai.com/in.php", data={
            "key": API_KEY,
            "method": task.get("method", "userrecaptcha"),
            "googlekey": task["sitekey"],
            "pageurl": task["pageurl"],
            "json": 1
        })
        data = resp.json()
        if data.get("status") != 1:
            return {"error": data.get("request")}

        captcha_id = data["request"]
        for _ in range(60):
            time.sleep(5)
            result = self.session.get(
                "https://ocr.captchaai.com/res.php",
                params={
                    "key": API_KEY,
                    "action": "get",
                    "id": captcha_id,
                    "json": 1
                }
            ).json()
            if result.get("status") == 1:
                return {"solution": result["request"]}
            if result.get("request") != "CAPCHA_NOT_READY":
                return {"error": result.get("request")}
        return {"error": "TIMEOUT"}

    def drain(self, timeout=120):
        """Stop accepting tasks and wait for active tasks to complete."""
        self.state = WorkerState.DRAINING
        start = time.time()
        while self.active_tasks > 0:
            if time.time() - start > timeout:
                print(f"Worker {self.worker_id}: drain timeout with "
                      f"{self.active_tasks} tasks remaining")
                break
            time.sleep(1)
        self.state = WorkerState.STOPPED

    @property
    def is_healthy(self):
        return self.state == WorkerState.RUNNING


class RollingUpdateOrchestrator:
    def __init__(self, workers):
        self.workers = {w.worker_id: w for w in workers}
        self.lock = threading.Lock()

    def get_available_worker(self):
        """Route tasks only to RUNNING workers."""
        with self.lock:
            for worker in self.workers.values():
                if worker.state == WorkerState.RUNNING:
                    return worker
        return None

    def rolling_update(self, new_version, health_check_fn=None,
                       max_unavailable=1, drain_timeout=120):
        """Update workers one at a time with health gates."""
        worker_ids = list(self.workers.keys())
        updated = []
        failed = []

        for i in range(0, len(worker_ids), max_unavailable):
            batch = worker_ids[i:i + max_unavailable]

            for wid in batch:
                worker = self.workers[wid]
                print(f"[{wid}] Draining (v{worker.version})...")

                # Step 1: Drain active tasks
                worker.drain(timeout=drain_timeout)

                # Step 2: "Deploy" new version
                print(f"[{wid}] Deploying v{new_version}...")
                worker.state = WorkerState.UPDATING
                worker.version = new_version
                time.sleep(2)  # Simulate deployment

                # Step 3: Start and health check
                worker.state = WorkerState.RUNNING
                if health_check_fn:
                    healthy = health_check_fn(worker)
                    if not healthy:
                        print(f"[{wid}] Health check FAILED — rolling back")
                        failed.append(wid)
                        self._rollback(updated)
                        return {
                            "status": "rolled_back",
                            "failed_at": wid,
                            "updated": updated,
                        }

                updated.append(wid)
                print(f"[{wid}] Updated to v{new_version} ✓")

        return {"status": "complete", "updated": updated, "failed": failed}

    def _rollback(self, updated_ids):
        """Roll back already-updated workers."""
        for wid in updated_ids:
            worker = self.workers[wid]
            print(f"[{wid}] Rolling back...")
            worker.state = WorkerState.STOPPED
            time.sleep(1)
            worker.version = "rollback"
            worker.state = WorkerState.RUNNING

    @property
    def status(self):
        return {
            wid: {
                "version": w.version,
                "state": w.state.value,
                "active_tasks": w.active_tasks,
            }
            for wid, w in self.workers.items()
        }


# Create fleet
workers = [Worker(f"w{i}", "1.2.0") for i in range(6)]
orchestrator = RollingUpdateOrchestrator(workers)


def health_check(worker):
    """Verify worker can solve a test CAPTCHA."""
    # In production, send a real test task
    return worker.state == WorkerState.RUNNING


# Execute rolling update
result = orchestrator.rolling_update(
    new_version="1.3.0",
    health_check_fn=health_check,
    max_unavailable=1,
    drain_timeout=60
)
print(f"Rolling update result: {result}")

Три детали реализации. drain() не прерывает активные решения принудительно: он ждёт обнуления active_tasks, а тайм-аут срабатывает лишь как страховка от зависшей задачи. _do_solve опрашивает res.php каждые 5 с и отличает CAPCHA_NOT_READY от настоящей ошибки, поэтому «ещё решается» не путается со «сломалось». _rollback проходит по всем обновлённым воркерам, а не только по последнему, — иначе в парке останутся смешанные версии.

JavaScript: выкатка с прогрессом и порогом отказов

Тот же сценарий на Node.js, но с двумя дополнениями: счётчик прогресса и порог отказов, который останавливает выкатку без ручного вмешательства.

const axios = require("axios");

const API_KEY = process.env.CAPTCHAAI_API_KEY;

class RollingUpdater {
  constructor(workerCount, currentVersion) {
    this.workers = Array.from({ length: workerCount }, (_, i) => ({
      id: `worker-${i}`,
      version: currentVersion,
      state: "running",
      activeTasks: 0,
    }));
    this.progress = { total: workerCount, completed: 0, failed: 0 };
  }

  async update(newVersion, options = {}) {
    const {
      maxUnavailable = 1,
      drainTimeout = 60000,
      healthCheckRetries = 3,
    } = options;

    console.log(
      `Starting rolling update: v${this.workers[0].version} → v${newVersion}`
    );

    for (let i = 0; i < this.workers.length; i += maxUnavailable) {
      const batch = this.workers.slice(i, i + maxUnavailable);

      for (const worker of batch) {
        try {
          // Drain
          console.log(`[${worker.id}] Draining...`);
          worker.state = "draining";
          await this.waitForDrain(worker, drainTimeout);

          // Deploy
          console.log(`[${worker.id}] Deploying v${newVersion}...`);
          worker.state = "updating";
          worker.version = newVersion;

          // Health check
          worker.state = "running";
          const healthy = await this.healthCheck(worker, healthCheckRetries);

          if (!healthy) {
            worker.state = "failed";
            this.progress.failed++;
            console.log(`[${worker.id}] FAILED health check`);

            if (this.progress.failed > Math.floor(this.workers.length * 0.25)) {
              console.log("Too many failures — aborting rolling update");
              return { status: "aborted", progress: this.progress };
            }
            continue;
          }

          this.progress.completed++;
          console.log(
            `[${worker.id}] Updated ✓ (${this.progress.completed}/${this.progress.total})`
          );
        } catch (err) {
          console.error(`[${worker.id}] Error: ${err.message}`);
          this.progress.failed++;
        }
      }
    }

    return { status: "complete", progress: this.progress };
  }

  async waitForDrain(worker, timeout) {
    const start = Date.now();
    while (worker.activeTasks > 0 && Date.now() - start < timeout) {
      await new Promise((r) => setTimeout(r, 1000));
    }
  }

  async healthCheck(worker, retries) {
    for (let attempt = 0; attempt < retries; attempt++) {
      try {
        const resp = await axios.get("https://ocr.captchaai.com/res.php", {
          params: { key: API_KEY, action: "getbalance", json: 1 },
          timeout: 10000,
        });
        if (resp.data.status === 1) return true;
      } catch {
        // Retry
      }
      await new Promise((r) => setTimeout(r, 5000));
    }
    return false;
  }
}

// Execute
const updater = new RollingUpdater(8, "1.2.0");
updater
  .update("1.3.0", { maxUnavailable: 2, drainTimeout: 30000 })
  .then((result) => console.log("Result:", JSON.stringify(result, null, 2)));

Проверка здесь дешёвая: запрос action=getbalance к res.php подтверждает, что новый экземпляр видит API-ключ и достучался до сервиса, и не занимает поток решением. Если доля неудачных воркеров превышает 25 % парка, выкатка прерывается — это дешевле, чем обнаружить сломанную версию на всех восьми экземплярах.

Как подобрать drain_timeout и maxUnavailable

Обе величины выводятся из фактов, а не из привычки:

  • drain_timeout = максимальное время решения используемого типа CAPTCHA плюс запас. Для reCAPTCHA v2 (SLA — менее 60 с) разумно взять 120 с; для парка, который работает только с Cloudflare Turnstile (менее 10 с), хватит 30–40 с. Если воркеры стоят в европейском регионе или в Казахстане, а часть трафика идёт через нестабильный канал, добавьте запас на сетевые повторы.
  • maxUnavailable = сколько воркеров вы готовы вывести из работы одновременно. Для парка меньше 10 экземпляров — один за раз. Для крупного парка — 10–25 % от общего числа. Половину парка выводить нельзя: оставшейся мощности не хватит на текущую нагрузку.

Практический пример. Агентство из Алматы держит парсинг на плане ADVANCE ($90/мес, 50 потоков) и обновляет 12 воркеров: при maxUnavailable=1 и дренаже до 120 с выкатка занимает около 25 минут, а мощность в худший момент проседает примерно на 8 %. Небольшая команда на плане BASIC ($15/мес, 5 потоков) обновляет строго по одному воркеру — потеря слота из пяти заметна сразу. Тарификация по потокам делает стоимость выкатки предсказуемой: вы просто дольше держите те же потоки.

Сравнение стратегий развёртывания

Стратегия Простой Скорость отката Сложность Когда подходит
Rolling (поочерёдная) Нет Умеренная Низкая Большинство парков воркеров
Blue-green Нет Мгновенная Средняя Критичные сервисы с жёстким SLA
Canary (канареечная) Нет Быстрая Высокая Крупные парки от 50 воркеров
Recreate (полная замена) Короткий Не применяется Минимальная Среды разработки и staging

Rolling — компромисс по умолчанию: нужен один парк вместо двух, а откат укладывается в те же шаги, что и выкатка.

Наблюдаемость во время выкатки

Минимальный набор сигналов, по которым видно, что обновление идёт нормально:

  1. Версия воркера в каждой строке лога — без неё невозможно понять, какая версия даёт ошибки.
  2. Доля успешных решений в разрезе версий — падение только на новой версии останавливает выкатку раньше, чем сработает порог отказов.
  3. Время решения (p95) по типам — рост на новой версии обычно указывает на изменившийся тайм-аут HTTP-клиента, а не на сервис.
  4. Число активных задач на воркер — если счётчик не обнуляется, дренаж уходит в тайм-аут и задачи теряются.

Отдельная гигиена: в логи выкатки не должны попадать персональные данные с целевых страниц. Пишите ID задачи, тип CAPTCHA и код ответа — и собирайте только то, что вы вправе обрабатывать (152-ФЗ «О персональных данных», GDPR-подход для трансграничных проектов).

Типичные сбои и что делать

Симптом Причина Что сделать
Задачи теряются во время обновления drain_timeout короче реального времени решения Поднимите тайм-аут до максимального времени решения типа CAPTCHA плюс запас
Проверка после старта стабильно не проходит Дефект в новой версии или подмена конфига Откатитесь и прогоните версию в staging с настоящей задачей
Выкатка тянется часами maxUnavailable=1 на большом парке Увеличьте до 2–3 воркеров для парков от 20 экземпляров
В парке остались разные версии Откат прошёл частично Ведите список обновлённых воркеров и откатывайте его целиком
ERROR_ZERO_BALANCE сразу после деплоя Новый экземпляр читает чужой или пустой API-ключ Проверяйте getbalance в health-гейте до возврата трафика

Частые вопросы

Что происходит с задачей, отправленной в in.php, если воркер остановлен?

Решение на стороне сервиса продолжается, но забирать результат из res.php некому: ID задачи хранился в памяти процесса. Поток освободится сам после завершения решения. Чтобы не терять такие задачи, сохраняйте ID во внешнее хранилище (Redis, очередь) до начала дренажа — тогда результат дочитает следующий воркер.

Как обновлять парк, если все потоки тарифного плана заняты?

Дренаж по определению временно снижает число одновременных задач, поэтому обновляйтесь в часы низкой нагрузки и держите maxUnavailable минимальным. Потоки тарифицируются как конкурентные слоты, а не за отдельное решение, так что растянутая выкатка не увеличивает счёт — она лишь удлиняет время, в течение которого часть слотов простаивает.

Достаточно ли проверить, что процесс запустился?

Нет. Запуск процесса не подтверждает ни доступность API, ни правильный ключ. Минимум — запрос getbalance; для критичных парков добавьте одну настоящую задачу (например, Cloudflare Turnstile, менее 10 с) и переходите к следующему воркеру только после успешно полученного токена.

Нужно ли откатывать весь парк, если не поднялся один воркер?

Зависит от причины. Единичный сетевой сбой — повторите проверку. Если новая версия не проходит health-гейт уже на первом воркере, откатывайте всех обновлённых: смешанный парк даёт плавающие ошибки, которые тяжело диагностировать в проде.

Следующие шаги

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