DevOps & Scaling

Планирование аварийного восстановления для конвейеров решения CAPTCHA

Сбой воркера, сетевая партиция или скомпрометированный API-ключ — рано или поздно один из этих сценариев случится с конвейером решения CAPTCHA. Разница между пятиминутным инцидентом и потерянным днём в том, спроектирован ли конвейер с планом аварийного восстановления (DR): постоянным хранением задач, чёткими целями RPO/RTO и runbook, по которому команда действует без импровизации, а не разбирается в проде методом проб и ошибок.

Цели DR: сколько данных и времени вы готовы потерять

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

Метрика Определение Цель для конвейера CAPTCHA
RPO (Recovery Point Objective) Максимально допустимая потеря данных < 5 минут задач в очереди
RTO (Recovery Time Objective) Максимальное время восстановления сервиса < 15 минут
MTTR (Mean Time to Recovery) Среднее время восстановления < 10 минут

Не рассматривайте конвейер как единый чёрный ящик — разбейте его на компоненты и присвойте цели каждому отдельно:

  • какие данные можно спокойно пересчитать позже, а какие пользовательские флоу требуют немедленного восстановления;
  • какой компонент — приём задач, очередь, воркеры, доставка результата — критичнее для конкретного RPO/RTO из таблицы выше;
  • кто в команде вправе объявить failover, кто проверяет, что восстановление прошло корректно, и кто принимает решение вернуться на основной маршрут.

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

Сценарии сбоев, к которым готовится план

Scenario 1: Worker crash         → Restart workers, replay queue
Scenario 2: Queue data loss      → Restore from persistent backup
Scenario 3: Network partition    → Failover to secondary region
Scenario 4: API key compromised  → Rotate key, update workers
Scenario 5: Config corruption    → Rollback to last known good

Для команд, которые разворачивают воркеры в европейских регионах или используют хостинг в Казахстане и Центральной Азии, сценарий сетевой партиции ощущается острее: более высокий RTT между регионами напрямую увеличивает время failover. Выбирайте вторичный регион по реальной задержке, а не только по цене — протестируйте RTT заранее, а не во время инцидента.

Почему очередь в памяти не переживает сбой

Никогда не решайте CAPTCHA из очереди, которая существует только в оперативной памяти процесса. Сохраняйте задачи на диск или во внешнее хранилище — тогда перезапуск воркера не превращается в потерю данных.

Python: постоянная очередь задач на SQLite

import os
import json
import time
import sqlite3
import threading
import requests
from datetime import datetime

API_KEY = os.environ["CAPTCHAAI_API_KEY"]


class PersistentTaskQueue:
    """SQLite-backed task queue that survives crashes."""

    def __init__(self, db_path="captcha_tasks.db"):
        self.db_path = db_path
        self.conn = sqlite3.connect(db_path, check_same_thread=False)
        self.lock = threading.Lock()
        self._init_db()

    def _init_db(self):
        self.conn.execute("""
            CREATE TABLE IF NOT EXISTS tasks (
                id TEXT PRIMARY KEY,
                payload TEXT NOT NULL,
                status TEXT DEFAULT 'pending',
                created_at TEXT DEFAULT CURRENT_TIMESTAMP,
                started_at TEXT,
                completed_at TEXT,
                result TEXT,
                attempts INTEGER DEFAULT 0
            )
        """)
        self.conn.commit()

    def enqueue(self, task_id, payload):
        with self.lock:
            self.conn.execute(
                "INSERT INTO tasks (id, payload) VALUES (?, ?)",
                (task_id, json.dumps(payload))
            )
            self.conn.commit()

    def dequeue(self):
        with self.lock:
            cursor = self.conn.execute(
                "SELECT id, payload FROM tasks "
                "WHERE status = 'pending' ORDER BY created_at LIMIT 1"
            )
            row = cursor.fetchone()
            if not row:
                return None

            task_id, payload = row
            self.conn.execute(
                "UPDATE tasks SET status = 'processing', "
                "started_at = ?, attempts = attempts + 1 WHERE id = ?",
                (datetime.utcnow().isoformat(), task_id)
            )
            self.conn.commit()
            return {"id": task_id, "payload": json.loads(payload)}

    def complete(self, task_id, result):
        with self.lock:
            self.conn.execute(
                "UPDATE tasks SET status = 'completed', "
                "completed_at = ?, result = ? WHERE id = ?",
                (datetime.utcnow().isoformat(), json.dumps(result), task_id)
            )
            self.conn.commit()

    def fail(self, task_id, error):
        with self.lock:
            # Requeue if under retry limit
            cursor = self.conn.execute(
                "SELECT attempts FROM tasks WHERE id = ?", (task_id,)
            )
            row = cursor.fetchone()
            if row and row[0] < 3:
                self.conn.execute(
                    "UPDATE tasks SET status = 'pending' WHERE id = ?",
                    (task_id,)
                )
            else:
                self.conn.execute(
                    "UPDATE tasks SET status = 'failed', "
                    "result = ? WHERE id = ?",
                    (json.dumps({"error": error}), task_id)
                )
            self.conn.commit()

    def recover_stale(self, timeout_seconds=600):
        """Reset tasks stuck in 'processing' after a crash."""
        with self.lock:
            cutoff = datetime.utcnow().timestamp() - timeout_seconds
            self.conn.execute(
                "UPDATE tasks SET status = 'pending' "
                "WHERE status = 'processing' "
                "AND started_at < datetime(?, 'unixepoch')",
                (cutoff,)
            )
            count = self.conn.total_changes
            self.conn.commit()
            return count

    @property
    def stats(self):
        cursor = self.conn.execute(
            "SELECT status, COUNT(*) FROM tasks GROUP BY status"
        )
        return dict(cursor.fetchall())


# On startup: recover tasks that were processing during a crash
queue = PersistentTaskQueue()
recovered = queue.recover_stale(timeout_seconds=600)
print(f"Recovered {recovered} stale tasks after restart")

При старте процесс поднимает recover_stale() и возвращает в очередь задачи, застрявшие в статусе processing — это и есть автоматическое восстановление после падения воркера без ручного вмешательства.

JavaScript: менеджер восстановления с чекпоинтами

const axios = require("axios");
const fs = require("fs");

const API_KEY = process.env.CAPTCHAAI_API_KEY;

class DisasterRecoveryManager {
  constructor(checkpointDir = "./dr-checkpoints") {
    this.checkpointDir = checkpointDir;
    if (!fs.existsSync(checkpointDir)) {
      fs.mkdirSync(checkpointDir, { recursive: true });
    }
  }

  checkpoint(label, data) {
    const filename = `${this.checkpointDir}/${label}-${Date.now()}.json`;
    fs.writeFileSync(filename, JSON.stringify(data, null, 2));
    this.pruneOldCheckpoints(label, 10); // Keep last 10
    return filename;
  }

  restore(label) {
    const files = fs.readdirSync(this.checkpointDir)
      .filter((f) => f.startsWith(label) && f.endsWith(".json"))
      .sort()
      .reverse();

    if (files.length === 0) return null;
    const latest = fs.readFileSync(
      `${this.checkpointDir}/${files[0]}`, "utf8"
    );
    return JSON.parse(latest);
  }

  pruneOldCheckpoints(label, keep) {
    const files = fs.readdirSync(this.checkpointDir)
      .filter((f) => f.startsWith(label) && f.endsWith(".json"))
      .sort();

    while (files.length > keep) {
      const old = files.shift();
      fs.unlinkSync(`${this.checkpointDir}/${old}`);
    }
  }

  async healthCheck() {
    try {
      const resp = await axios.get("https://ocr.captchaai.com/res.php", {
        params: { key: API_KEY, action: "getbalance", json: 1 },
        timeout: 10000,
      });
      return {
        healthy: resp.data.status === 1,
        balance: parseFloat(resp.data.request || 0),
      };
    } catch (err) {
      return { healthy: false, error: err.message };
    }
  }
}

class ResilientSolver {
  constructor() {
    this.dr = new DisasterRecoveryManager();
    this.pendingTasks = [];
  }

  async solveBatch(tasks) {
    // Checkpoint before starting
    this.dr.checkpoint("batch-pending", {
      tasks,
      startedAt: new Date().toISOString(),
    });

    const results = [];
    for (const task of tasks) {
      try {
        const result = await this.solveSingle(task);
        results.push({ taskId: task.id, ...result });
      } catch (err) {
        results.push({ taskId: task.id, error: err.message });
      }

      // Checkpoint progress periodically
      if (results.length % 10 === 0) {
        this.dr.checkpoint("batch-progress", { results, remaining: tasks.length - results.length });
      }
    }

    // Final checkpoint
    this.dr.checkpoint("batch-complete", { results });
    return results;
  }

  async recover() {
    // Check for incomplete batch
    const progress = this.dr.restore("batch-progress");
    const pending = this.dr.restore("batch-pending");

    if (progress) {
      const completedIds = new Set(progress.results.map((r) => r.taskId));
      const remaining = pending?.tasks.filter((t) => !completedIds.has(t.id));
      console.log(
        `Recovering: ${progress.results.length} done, ${remaining?.length || 0} remaining`
      );
      return remaining || [];
    }

    if (pending) {
      console.log(`Recovering full batch: ${pending.tasks.length} tasks`);
      return pending.tasks;
    }

    return [];
  }

  async solveSingle(task) {
    const resp = await axios.post("https://ocr.captchaai.com/in.php", null, {
      params: {
        key: API_KEY,
        method: "userrecaptcha",
        googlekey: task.sitekey,
        pageurl: task.pageurl,
        json: 1,
      },
    });

    if (resp.data.status !== 1) throw new Error(resp.data.request);

    const captchaId = resp.data.request;
    for (let i = 0; i < 60; i++) {
      await new Promise((r) => setTimeout(r, 5000));
      const poll = await axios.get("https://ocr.captchaai.com/res.php", {
        params: { key: API_KEY, action: "get", id: captchaId, json: 1 },
      });
      if (poll.data.status === 1) return { solution: poll.data.request };
      if (poll.data.request !== "CAPCHA_NOT_READY")
        throw new Error(poll.data.request);
    }
    throw new Error("TIMEOUT");
  }
}

// Start with recovery check
const solver = new ResilientSolver();
solver.recover().then((remaining) => {
  if (remaining.length > 0) {
    console.log(`Resuming ${remaining.length} tasks from checkpoint`);
    solver.solveBatch(remaining);
  }
});

Здесь recover() при старте процесса сначала ищет чекпоинт batch-progress, а если его нет — batch-pending, и возобновляет обработку только недостающих задач, не решая заново то, что уже успешно завершилось.

Runbook на случай сбоя конвейера

Runbook — это не документация «для галочки», а последовательность действий, которую дежурный инженер выполняет в 3 часа ночи, не думая.

RUNBOOK: CAPTCHA Pipeline Recovery
====================================

1. DETECT
   - Alert fires: [PagerDuty / Slack / Email]
   - Symptom: [Queue growing / Workers offline / Error spike]

2. ASSESS
   - Check worker health: curl http://workers/health
   - Check API status: GET /res.php?action=getbalance
   - Check queue depth: SELECT COUNT(*) FROM tasks WHERE status='pending'

3. RECOVER
   If: Workers crashed
     → Restart worker containers: docker-compose up -d workers
     → Run stale task recovery: recovery.py --recover-stale

   If: Network partition
     → Failover to secondary region
     → Update DNS or load balancer routing

   If: API key compromised
     → Generate new key at captchaai.com
     → Update secret store
     → Rolling restart workers

4. VERIFY
   - Confirm solve rate > 90%
   - Confirm queue draining
   - Confirm no duplicate solves

5. POST-MORTEM
   - Document root cause
   - Update runbook if needed

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

Как часто нужно делать чекпоинты очереди задач?

Каждые 5–10 завершённых задач или раз в 30 секунд — смотря что наступит раньше. Более частые чекпоинты снижают RPO, но добавляют нагрузку на диск; для конвейера на плане ADVANCE ($90/мес, 50 потоков) интервала в 30 секунд обычно достаточно.

Как выбрать RPO и RTO для конвейера решения CAPTCHA?

Отталкивайтесь от стоимости простоя для бизнеса, а не от технической возможности. Если конвейер обслуживает пользовательский флоу вроде регистрации, RTO должен укладываться в терпение пользователя ждать ответ — обычно до 15 минут. Для фоновой обработки и парсинга RTO в 30–60 минут допустим, а RPO определяется тем, сколько задач вы готовы решить повторно без потери денег.

Что делать, если сбой произошёл на стороне самого API CaptchaAI?

Складывайте задачи в локальную очередь и повторяйте запросы после восстановления API — circuit breaker с экспоненциальной задержкой (backoff) предотвращает шторм повторных запросов в момент, когда сервис только поднялся. Держите план на такой случай отдельно от плана на сбой собственной инфраструктуры: причины и действия разные.

Нужно ли тестировать план DR, если сбоев ещё не было?

Да — план, который ни разу не проверялся на практике, это предположение, а не план. Раз в квартал имитируйте падение воркера или временную недоступность API в staging и сравнивайте фактическое время восстановления с целевым RTO из таблицы выше.

Сколько потоков нужно держать в резерве на случай failover?

Заложите запас, а не тратьте все потоки тарифа на пиковую нагрузку. Если обычная работа занимает 40 из 50 потоков плана ADVANCE, оставшиеся 10 дают конвейеру пространство для догона очереди после восстановления без апгрейда тарифа посреди инцидента.

Типичные ошибки восстановления и как их избежать

Ниже — ошибки, которые чаще всего превращают плановый failover в незапланированный простой.

Проблема Причина Решение
Задачи теряются при сбое Очередь хранится только в оперативной памяти Используйте постоянную очередь (SQLite, Redis с AOF)
Дублирующиеся решения после восстановления Устаревшие задачи переобрабатываются без дедупликации Добавьте ключи идемпотентности, проверяйте, не решена ли задача уже
Восстановление занимает больше RTO Резервная копия базы данных устарела Увеличьте частоту чекпоинтов
Failover уходит не в тот регион Слишком большой TTL у DNS-записей Снизьте TTL до 60 секунд перед плановым failover

Что дальше

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