如何在 Temporal 工作流的 Activity 中识别验证码

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

在 Temporal 里,验证码识别要放进 Activity。不是放在 workflow 方法里,也不是放在 workflow 调用的辅助函数里,就是放进 Activity。workflow 代码在每次工作流恢复时都会从历史记录中重放,所以它必须是确定性的:不能发网络请求,不能有随机性,不能读时钟。而一次识别同时占了这三条。放进 Activity 之后,剩下的就是一份重试策略和一点时序上的纪律,整件事大约四十行代码。

你需要什么

  • Python 3.10 或更新版本、Temporal Python SDK 和 CapSkip Python SDK。
  • 一个可连接的 Temporal Service。本地开发服务器或 Temporal Cloud 都可以。
  • 受保护表单的页面 URL,以及它的 sitekey。
  • worker 和识别工具在同一台机器上时,CapSkip 用 Local mode(本地模式),不在同一台时用 Server mode(服务器模式)。两种模式都说明在 连接设置.
# pip install temporalio
pip install -U temporalio capskip httpx

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

为什么识别不能放在 workflow 代码里

在 worker 重启、部署或睡了一周之后,Temporal 会重放工作流的历史记录来重建状态。要让两次重放得出同样的结果,workflow 代码必须是确定性的。SDK 明确列出了被排除的东西:不能有网络 IO,不能用线程,不能有随机性,不能调用外部进程,不能修改全局状态。它甚至把 workflow 代码放在一个沙箱里运行,每次运行都重新导入模块,这也是 activity 的导入语句要包在一个直通块里的原因。

一次验证码识别把这条规则违反了三遍。它是网络请求,每次返回的 token 都不一样,耗时还取决于机器。把它放进 workflow 方法,开发时看起来一切正常,但 worker 第一次在运行中途重启时就会报非确定性错误。

好消息是这条约束也给了你东西。Activity 由 Temporal 自己重试,用的是你声明的策略,而不是你手写的循环,而且它的结果会被记进历史。所以一次成功的识别在重放时不会再做一遍,这正是你对一次性答案所期望的行为。

第 1 步:识别 activity

CapSkip 的 Python SDK 提供的是真正的 asyncio 客户端,不是一层别名,所以 async activity 是最自然的选择,也不需要线程池。把 sitekey 和页面 URL 作为参数传入,返回 token。

# 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,worker 还需要配置 activity_executor,因为 Temporal 会在线程池里跑同步 activity。上面的异步版本完全绕开了这一点,这也是 Python SDK 确实比其他语言更好用的少数几处之一。

第 2 步:超时要比识别工具自身的更长

每个 activity 都需要 start_to_close_timeout,人们往往在这里悄无声息地毁掉自己的识别。CapSkip 对 reCAPTCHA、Turnstile 和极验(GeeTest)最多轮询 300 秒,对图片验证码最多 120 秒。把 activity 超时设得比这更短,Temporal 就会在识别工具还在干活时取消这次尝试,然后重试,于是一个表单就有两次识别同时在跑。

给 activity 留出比识别工具自身上限更多的余量。识别工具上限是五分钟时,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:前者说明识别工具没在运行或 host 填错了,后者说明识别耗时超过了轮询超时。这四个异常都继承自同一个基类,所以如果你更愿意在一个地方统一处理,捕获 CapSkipError 就够了。

第 3 步:最后再识别,不要一开始就识别

持久化执行让 token 过期这件事在这里比在任何其他编排工具里都更容易出错。Temporal 工作流可以等待一个信号、睡上一天再继续,而它记录下来的 activity 结果会原封不动地从历史里取回。所以一个先识别、再等待审批、最后才提交的工作流,重放出来的会是昨天生成的 token。

reCAPTCHA token 只能用一次,大约两分钟后过期。安排工作流的顺序,让识别紧挨在提交之前,中间不要有任何可能阻塞的步骤。这个问题的整体形态见指南 reCAPTCHA token 过期.

完整可运行示例

三个 activity,加上一个按顺序调用它们的 workflow。读取 sitekey,识别,提交。每个 activity 有自己的超时和自己的重试策略,并且在 Temporal UI 里各自单独显示,所以哪一步慢了一眼就能看出来。

# 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),
        )

注意那两个带两个参数的 activity 上的 args 列表。单个位置参数可以直接传,但超过一个就必须通过 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())

worker 在哪里运行,识别工具又在哪里运行

Temporal 把这两者切分得很干净,而且这对你有利。Temporal Service 负责调度任务并保存历史。它从不执行你的代码。你在自己基础设施里启动的 worker 会保持一条向外的长连接连到服务端,从队列里领取任务。你的网络上不需要开放任何入站连接。

所以和 CapSkip 在同一台 Windows 机器上的 worker,调用 127.0.0.1:8080 的方式和你桌上的脚本完全一样,在 Temporal Cloud 上也是如此。Temporal Cloud 里的云指的是编排层,不是计算层。

worker 一搬家,变的只有 host,别的都不变。跑在容器、Linux 虚拟机或 Kubernetes 里的 worker 无法通过回环地址访问 Windows 上的识别工具,所以识别工具要切到 Server mode(服务器模式)。Local mode(本地模式)绑定 127.0.0.1,只响应本机。Server mode 绑定你的内网或公网 IP,固定公网 IP 能让地址保持不变。两种模式下硬件都还是你自己的,也都不计次,所以一天跑一万次的工作流和只跑一次的成本一样。

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

一旦识别工具监听在网络地址上,就打开密钥校验,并给每一组 worker 集群分配各自的密钥,这样吊销其中一个不会影响其他。两种模式的完整步骤见 CapSkip 设置指南.

常见错误及其含义

你所看到的原因修复
重放时报非确定性错误识别是从 workflow 代码里调用的把它移进 activity,并用 execute_activity 调用
导入时报 RestrictedWorkflowAccessErroractivity 模块被导入到了沙箱里在 workflow.unsafe.imports_passed_through 里导入它
activity 在识别中途被取消,然后重试start_to_close_timeout 比识别工具自身的超时还短对 reCAPTCHA、Turnstile 和极验(GeeTest)把它设到 300 秒以上
表单拒绝了一个看起来没问题的 token工作流在识别和提交之间被阻塞了让识别成为紧挨着提交的那一步
同一个失败白白重试了四次一个确定性的错误被拿去重试了把 ValidationException 和 ApiException 列为不可重试
每次尝试都报 NetworkExceptionworker 不在运行识别工具的那台机器上把识别工具切到 Server mode,并设置 host
activity 报参数相关的 TypeError两个位置参数被直接传了进去把它们作为列表通过 args 参数传入
SDK 抛出 TimeoutException识别耗时超过了 recaptchaTimeout把它调到默认的 300 秒以上

常见问题

Temporal Cloud 上的工作流真的能调用 127.0.0.1 吗?

可以,因为 Temporal Cloud 并不运行你的代码。运行代码的是你的 worker,你在哪里启动它,它就在哪里,而且是向外连接到服务端。那台 worker 上的回环地址指的就是 worker 自己的机器,所以装在这台机器上的识别工具会正常响应。从开发服务器换到 Cloud,这一点不会有任何变化。

识别 activity 需要发心跳吗?

发不出有意义的心跳,因为 SDK 调用会一直阻塞到 token 返回,中间没有可以上报的位置。不如给 activity 设一个留足余量的 start_to_close_timeout,让失败的尝试交给重试策略处理。心跳适用于那些在你自己控制的任务上循环的 activity。

怎样批量识别又不把识别工具冲垮?

在 worker 上设置 max_concurrent_activities,或者给识别单独开一个任务队列和一个上限较低的 worker。这样限流发生在任务真正执行的地方,比试图把工作流的启动时间错开更可靠。在单个进程内,异步客户端会用 asyncio 并发展开,这部分内容见 并行识别验证码的指南.

这和在 Airflow 里做有什么不同?

Airflow 是一个带 DAG 的调度器,在构建图的过程中没什么能拦住你发网络请求,那是另一套坑。Temporal 是持久化执行,约束是确定性,答案永远是放进 activity。识别工具这一侧两边完全一样,Airflow 版本写在 Airflow 验证码指南.

简短版结论

把识别放进 activity,绝不要放在 workflow 方法里。给它一个比识别工具自身 300 秒上限更长的 start_to_close_timeout,把 ValidationException 和 ApiException 标为不可重试,并安排工作流的顺序,让识别成为提交之前的最后一步。把 worker 跑在 CapSkip 所在的机器上,host 就一直是 127.0.0.1。Python 的其余接口见 Python 验证码识别页面,同样的三个调用在 Node.js、PHP 和 C# 里都有,列在 验证码识别 SDK 页面。reCAPTCHA 本身的各项选项见 reCAPTCHA v2 识别页面.

在给它挂上定时任务之前,还有一点值得知道。CapSkip 是一款 验证码识别工具 ,跑在你已经拥有的硬件上,所以每分钟触发一次的工作流,和每周触发一次的成本相同。