Как решать капчу в Activity внутри Temporal Workflow

temporal captcha - How to Solve CAPTCHA in a Temporal Workflow Activity

В Temporal решение капчи живёт в Activity. Не в методе workflow, не во вспомогательной функции, которую workflow вызывает, а именно в Activity. Код workflow заново проигрывается из истории каждый раз, когда workflow возобновляется, поэтому он обязан быть детерминированным: никаких сетевых вызовов, никакой случайности, никакого чтения часов. Решение капчи нарушает все три пункта сразу. Как только оно оказалось в Activity, остаются политика повторов и немного дисциплины в тайминге, а весь код занимает около сорока строк.

Что понадобится

  • Python 3.10 или новее, Temporal Python SDK и CapSkip Python SDK.
  • Temporal Service, к которому можно подключиться. Подойдёт и локальный dev-сервер, и Temporal Cloud.
  • URL страницы с защищённой формой и её sitekey.
  • CapSkip в режиме Local, когда воркер и решатель стоят на одной машине, или в режиме Server, когда на разных. Оба описаны в разделе Настройки подключения.
# pip install temporalio
pip install -U temporalio capskip httpx

# A local service to develop against.
temporal server start-dev

Почему решение капчи не может жить в коде workflow

Temporal заново проигрывает историю workflow, чтобы восстановить его состояние после перезапуска воркера, деплоя или недельного сна. Чтобы это дважды давало один и тот же результат, код workflow должен быть детерминированным. SDK прямо перечисляет, что из-за этого запрещено: никакого сетевого ввода-вывода, никаких потоков, никакой случайности, никаких внешних вызовов процессов, никаких изменений глобального состояния. Более того, код workflow выполняется в песочнице, которая заново импортирует модули на каждый запуск, и поэтому импорты для activity оборачивают в блок, пропускающий их насквозь.

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

Хорошая новость в том, что это ограничение кое-что вам даёт. Activity повторяет сам Temporal по политике, которую вы объявляете, а не по циклу, который вы пишете, и результат записывается в историю. Поэтому успешное решение капчи никогда не повторяется при реплее, а это ровно то, что нужно для одноразового ответа.

Шаг 1: activity, которая решает капчу

CapSkip SDK для Python поставляется с настоящим asyncio-клиентом, а не с псевдонимом, поэтому асинхронная activity подходит естественным образом и не требует пула потоков. Принимайте sitekey и URL страницы аргументами и возвращайте токен.

# pip install capskip
import os
from temporalio import activity
from capskip import AsyncCapSkip

@activity.defn
async def solve_recaptcha(sitekey: str, page_url: str) -> str:
    solver = AsyncCapSkip(
        host=os.environ.get("CAPSKIP_HOST", "127.0.0.1"),
        port=8080,
        apiKey=os.environ.get("CAPSKIP_API_KEY", "capskip"),
    )
    result = await solver.recaptcha(sitekey=sitekey, url=page_url)
    return result["code"]

Один метод покрывает reCAPTCHA v2, Invisible, Enterprise и v3. Варианты задаются опциями того же вызова, а не отдельными методами: invisible равно 1, enterprise равно 1 или version со значением v3 вместе с action. У Turnstile и GeeTest свои методы такой же формы, а полный список параметров есть в документации CapSkip API.

Если вам удобнее синхронный клиент, activity придётся сделать обычной def, а воркеру понадобится activity_executor, потому что синхронные activity Temporal выполняет в пуле потоков. Асинхронный вариант выше избавляет от этого полностью, и это одно из немногих мест, где Python SDK действительно приятнее остальных.

Шаг 2: таймауты длиннее, чем у самого решателя

Каждой activity нужен start_to_close_timeout, и именно здесь люди тихо ломают себе решение капчи. CapSkip опрашивает до 300 секунд на reCAPTCHA, Turnstile и GeeTest и до 120 секунд на капчах с картинками. Поставьте таймаут activity ниже этих значений, и Temporal отменит попытку, пока решатель ещё работает, затем повторит её, и на одну форму у вас окажется два решения в работе.

Дайте activity запас над собственным лимитом решателя. Шесть минут против пятиминутного потолка решателя вполне комфортны.

# Longer than recaptchaTimeout, which defaults to 300s.
from datetime import timedelta
from temporalio.common import RetryPolicy

SOLVE_TIMEOUT = timedelta(minutes=6)

SOLVE_RETRIES = RetryPolicy(
    initial_interval=timedelta(seconds=5),
    backoff_coefficient=2.0,
    maximum_attempts=4,
    # These fail identically every time. Do not burn attempts.
    non_retryable_error_types=["ValidationException", "ApiException"],
)

Список не подлежащих повтору ошибок сопоставляется по имени класса исключения, и заполнить его стоит. ValidationException означает пропущенный или неправильный аргумент, а ApiException означает, что API отклонил запрос, обычно из-за sitekey, который не относится к URL страницы. Ни то ни другое не исправится со второй попытки. NetworkException и TimeoutException это как раз те два случая, которые действительно заслуживают повтора: первый означает, что решатель не запущен или указан не тот хост, второй что решение не уложилось в таймаут опроса. Все четыре наследуются от общего базового класса, поэтому перехват CapSkipError работает, если вам удобнее обрабатывать всё в одном месте.

Шаг 3: решайте капчу в конце, а не в начале

Долговременное выполнение делает истечение срока токена здесь опаснее, чем в любом другом оркестраторе. Workflow в Temporal может ждать сигнала, спать сутки и продолжать дальше, а записанные результаты activity возвращаются из истории без изменений. Поэтому workflow, который решает капчу заранее, ждёт одобрения, а потом отправляет форму, повторно проиграет токен, выпущенный вчера.

Токен reCAPTCHA принимается один раз и истекает примерно через две минуты. Выстройте workflow так, чтобы решение капчи было шагом непосредственно перед отправкой и между ними не было ничего, что может заблокироваться. Общая картина этой проблемы разобрана в руководстве по истечению срока действия токена reCAPTCHA.

Полный рабочий пример

Три activity и workflow, который вызывает их по порядку. Прочитать sitekey, решить капчу, отправить форму. У каждой activity свой таймаут и своя политика повторов, и каждая отдельно видна в интерфейсе Temporal, поэтому, когда что-то тормозит, видно, на каком шаге это произошло.

# activities.py
import os, re, httpx
from temporalio import activity
from capskip import AsyncCapSkip

@activity.defn
async def read_sitekey(page_url: str) -> str:
    async with httpx.AsyncClient(timeout=30) as client:
        html = (await client.get(page_url)).text
    found = re.search(r'data-sitekey=["\']([^"\']+)', html)
    if not found:
        raise RuntimeError("No data-sitekey on the page.")
    return found.group(1)

@activity.defn
async def solve_recaptcha(sitekey: str, page_url: str) -> str:
    solver = AsyncCapSkip(host=os.environ.get("CAPSKIP_HOST", "127.0.0.1"))
    return (await solver.recaptcha(sitekey=sitekey, url=page_url))["code"]

@activity.defn
async def submit_form(page_url: str, token: str) -> int:
    async with httpx.AsyncClient(timeout=30) as client:
        reply = await client.post(
            page_url, data={"g-recaptcha-response": token}
        )
    return reply.status_code

Сам workflow не содержит логики, кроме порядка вызовов. В этом и смысл: всё, что может упасть, вынесено в activity, а workflow остаётся детерминированной частью, которая переживает реплей.

# workflow.py
from datetime import timedelta
from temporalio import workflow
from temporalio.common import RetryPolicy

with workflow.unsafe.imports_passed_through():
    from activities import read_sitekey, solve_recaptcha, submit_form

@workflow.defn
class SubmitProtectedForm:
    @workflow.run
    async def run(self, page_url: str) -> int:
        sitekey = await workflow.execute_activity(
            read_sitekey, page_url,
            start_to_close_timeout=timedelta(seconds=60),
        )
        # Solve directly before the submit. Tokens go stale.
        token = await workflow.execute_activity(
            solve_recaptcha, args=[sitekey, page_url],
            start_to_close_timeout=timedelta(minutes=6),
            retry_policy=RetryPolicy(
                maximum_attempts=4,
                non_retryable_error_types=["ValidationException"],
            ),
        )
        return await workflow.execute_activity(
            submit_form, args=[page_url, token],
            start_to_close_timeout=timedelta(seconds=60),
        )

Обратите внимание на список args у activity с двумя аргументами. Один позиционный аргумент можно передать напрямую, но больше одного придётся передавать через args, и именно на этом чаще всего спотыкаются в этом SDK в первый раз.

# worker.py - run this where CapSkip can be reached
import asyncio
from temporalio.client import Client
from temporalio.worker import Worker
from activities import read_sitekey, solve_recaptcha, submit_form
from workflow import SubmitProtectedForm

async def main():
    client = await Client.connect("localhost:7233")
    worker = Worker(
        client,
        task_queue="captcha-queue",
        workflows=[SubmitProtectedForm],
        activities=[read_sitekey, solve_recaptcha, submit_form],
    )
    await worker.run()

asyncio.run(main())

Где работает воркер и где работает решатель

Temporal разделяет эти две вещи чисто, и это вам на руку. Temporal Service планирует работу и хранит историю. Ваш код он не выполняет никогда. Воркер, который вы запускаете внутри своей инфраструктуры, держит длительное исходящее соединение с сервисом и забирает задачи из очереди. Никаких входящих подключений к вашей сети при этом не открывается.

Поэтому воркер на той же машине с Windows, где стоит CapSkip, обращается к 127.0.0.1:8080 ровно так же, как это сделал бы скрипт на вашем столе, и на Temporal Cloud это остаётся верным. Облако в Temporal Cloud отвечает за оркестрацию, а не за вычисления.

Как только воркер переезжает, меняется хост и больше ничего. Воркер в контейнере, на виртуальной машине с Linux или в Kubernetes не достучится до решателя на Windows через loopback, поэтому решатель переключается в режим Server. Режим Local слушает 127.0.0.1 и отвечает только этому устройству. Режим Server слушает ваш сетевой или публичный IP, а статический публичный IP держит адрес постоянным. В обоих режимах это по-прежнему ваше железо и по-прежнему без счётчика, поэтому workflow, который выполняется десять тысяч раз в день, стоит столько же, сколько выполняющийся один раз.

# One environment variable, no code change.
# CAPSKIP_HOST=10.0.0.12 on the worker.
solver = AsyncCapSkip(host=os.environ["CAPSKIP_HOST"], port=8080)

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

Частые ошибки и что они означают

Что вы видитеПричинаИсправить
Ошибка недетерминированности при реплееРешение капчи вызвали из кода workflowВынесите его в activity и вызывайте через execute_activity
RestrictedWorkflowAccessError при импортеМодуль с activity был импортирован в песочницуИмпортируйте его внутри workflow.unsafe.imports_passed_through
Activity отменяется посреди решения, а затем повторяетсяstart_to_close_timeout короче собственного таймаута решателяПоставьте его больше 300 секунд для reCAPTCHA, Turnstile и GeeTest
Форма отклоняет токен, который выглядит правильнымWorkflow заблокировался между решением капчи и отправкойСделайте решение капчи шагом непосредственно перед отправкой
Четыре попытки сгорели на одной и той же ошибкеПовторяется детерминированная ошибкаУкажите ValidationException и ApiException как не подлежащие повтору
NetworkException на каждой попыткеВоркер стоит не на той машине, где работает решательПереключите решатель в режим Server и укажите хост
TypeError про аргументы activityДва позиционных аргумента передали напрямуюПередавайте их списком через параметр args
TimeoutException из SDKРешение заняло больше, чем recaptchaTimeoutПоднимите его выше значения по умолчанию в 300 секунд

FAQ

Может ли workflow в Temporal Cloud действительно обращаться к 127.0.0.1?

Да, потому что Temporal Cloud не выполняет ваш код. Его выполняет ваш воркер там, где вы его запустили, и он сам подключается наружу к сервису. Loopback на этом воркере означает машину самого воркера, поэтому решатель на этой машине отвечает как обычно. При переходе с dev-сервера на Cloud здесь ничего не меняется.

Нужен ли heartbeat у activity, которая решает капчу?

Полезного heartbeat не выйдет, потому что вызов SDK блокируется до прихода токена и внутри него нет точки, откуда отчитываться. Вместо этого дайте activity start_to_close_timeout с реальным запасом и позвольте политике повторить неудачную попытку. Heartbeat нужен для activity, которые крутят цикл по работе, которую вы контролируете.

Как решать пакет капч, не заваливая решатель?

Задайте max_concurrent_activities у воркера или выделите решению капчи собственную очередь задач и собственный воркер с низким лимитом. Это ограничивает нагрузку там, где работа выполняется, и это надёжнее, чем пытаться разносить запуски workflow по времени. Внутри одного процесса асинхронный клиент распараллеливает вызовы через asyncio, и это разобрано в руководстве по параллельному решению капчи.

Чем это отличается от того же самого в Airflow?

Airflow представляет собой планировщик с DAG, и ничто не мешает сделать сетевой вызов прямо во время построения графа, а это уже другой набор граблей. Temporal реализует долговременное выполнение, поэтому ограничением становится детерминированность, а ответом всегда является activity. Со стороны решателя всё одинаково, а вариант для Airflow разобран в руководстве по капче в Airflow.

Коротко

Помещайте решение капчи в activity, а не в метод workflow. Дайте ей start_to_close_timeout длиннее собственного потолка решателя в 300 секунд, отметьте ValidationException и ApiException как не подлежащие повтору и выстройте workflow так, чтобы решение капчи было последним шагом перед отправкой. Запустите воркер там, где стоит CapSkip, и хост останется 127.0.0.1. Остальная часть поверхности Python описана на странице сервиса распознавания капч для Python, а те же три вызова есть в Node.js, PHP и C#, как перечислено на на странице SDK для распознавания капчи. Сами опции reCAPTCHA описаны на странице сервиса распознавания reCAPTCHA v2.

Ещё одно, что стоит знать, прежде чем направлять на это расписание. CapSkip представляет собой распознавание капчи и работает на железе, которое у вас уже есть, поэтому workflow, который срабатывает каждую минуту, стоит столько же, сколько запускаемый раз в неделю.