Temporal Workflow Activity में कैप्चा कैसे हल करें

Temporal में कैप्चा हल करने का काम किसी Activity में जाता है। workflow method में नहीं, workflow जिस helper को कॉल करता है उसमें नहीं, बल्कि किसी Activity में। workflow कोड हर बार workflow के फिर से चलने पर history से replay होता है, इसलिए उसका deterministic होना ज़रूरी है: कोई network calls नहीं, कोई randomness नहीं, घड़ी नहीं पढ़नी। हल करना एक साथ ये तीनों काम करता है। एक बार वह किसी Activity में आ गया तो बाक़ी बस एक retry policy और timing का एक अनुशासन है, और पूरा काम लगभग चालीस लाइनों का है।
आपको क्या चाहिए
- Python 3.10 या उससे नया, Temporal Python SDK और CapSkip Python SDK।
- जुड़ने के लिए एक Temporal Service। यहाँ लोकल dev server और Temporal Cloud दोनों काम करते हैं।
- सुरक्षित form का पेज URL, और उसकी sitekey।
- CapSkip को Local mode में रखिए जब worker और सॉल्वर एक ही मशीन साझा करते हों, या Server mode में जब न करते हों। दोनों का विवरण यहाँ मिलता है: कनेक्शन सेटिंग्स.
# pip install temporalio pip install -U temporalio capskip httpx # A local service to develop against. temporal server start-dev
हल का काम workflow कोड में क्यों नहीं रह सकता
worker के दोबारा शुरू होने, किसी deploy या हफ़्ते भर की नींद के बाद अपनी state फिर से बनाने के लिए Temporal किसी workflow की history को replay करता है। इससे हर बार वही जवाब निकले, इसके लिए workflow कोड का deterministic होना ज़रूरी है। SDK साफ़ बताता है कि इससे क्या क्या बाहर हो जाता है: कोई network IO नहीं, कोई threading नहीं, कोई randomness नहीं, processes को कोई बाहरी call नहीं, कोई global state बदलाव नहीं। वह तो workflow कोड को ऐसे sandbox में चलाता है जो हर run पर modules दोबारा import करता है, और इसीलिए activity वाले imports को एक pass-through block में लपेटा जाता है।
कैप्चा हल करना इस नियम को तीन बार तोड़ता है। यह एक network call है, जो token लौटता है वह हर प्रयास में अलग होता है, और उसमें कितना समय लगेगा यह मशीन पर निर्भर करता है। इसे workflow method में रखिए और development में यह चलता हुआ दिखेगा, फिर जैसे ही कोई worker बीच run में दोबारा शुरू होगा, पहली ही बार non-determinism error दे देगा।
अच्छी बात यह है कि यह पाबंदी आपको कुछ देती भी है। किसी Activity को Temporal ख़ुद retry करता है, उस policy के साथ जो आप घोषित करते हैं, न कि किसी loop के साथ जो आप लिखते हैं, और उसका नतीजा history में दर्ज हो जाता है। इसलिए जो हल एक बार सफल हो गया वह replay पर दोबारा नहीं होता, और जिस चीज़ का जवाब सिर्फ़ एक बार काम आता है, उसके लिए आपको यही चाहिए।
चरण 1: solve वाली activity
Python का CapSkip SDK एक असली asyncio client देता है, कोई alias नहीं, इसलिए async activity यहाँ स्वाभाविक रूप से फिट बैठती है और किसी thread pool की ज़रूरत नहीं पड़ती। sitekey और page URL को arguments के रूप में लीजिए और 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"]एक ही method reCAPTCHA v2, Invisible, Enterprise और v3 को कवर करता है। ये variants अलग methods नहीं बल्कि उसी call के options हैं: invisible को 1, enterprise को 1, या version को v3 और साथ में एक action। Turnstile और GeeTest के अपने methods हैं, आकार वही है, और पूरी parameter सूची यहाँ है: CapSkip API डॉक्यूमेंटेशन.
अगर आप synchronous client इस्तेमाल करना चाहें, तो activity को एक सादा def होना पड़ेगा और worker को एक activity_executor चाहिए, क्योंकि Temporal sync activities को thread pool में चलाता है। ऊपर वाला async संस्करण इस सब से पूरी तरह बच जाता है, और यह उन गिनी चुनी जगहों में से एक है जहाँ Python SDK बाक़ियों से सचमुच बेहतर है।
चरण 2: ऐसे timeouts जो सॉल्वर के अपने timeout से लंबे हों
हर activity को एक start_to_close_timeout चाहिए, और यहीं लोग चुपचाप अपने ही solves तोड़ बैठते हैं। CapSkip reCAPTCHA, Turnstile और GeeTest पर 300 सेकंड तक poll करता है, और image कैप्चा पर 120 सेकंड तक। activity का timeout इससे कम रखिए और Temporal उस प्रयास को तब रद्द कर देगा जब सॉल्वर अभी काम कर ही रहा है, फिर retry करेगा, और एक ही form के लिए आपके दो solves चल रहे होंगे।
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"],
)non-retryable सूची exception class के नाम से मिलाई जाती है, और उसे भरना फ़ायदे का है। ValidationException का मतलब है कोई argument ग़ायब है या ग़लत बना है, और ApiException का मतलब है कि API ने request ठुकरा दी, आमतौर पर ऐसा sitekey जो उस page URL का है ही नहीं। दूसरे प्रयास में इनमें से कुछ भी नहीं सुधरता। NetworkException और TimeoutException ही वे दो हैं जो सचमुच retry के लायक हैं: पहले का मतलब है सॉल्वर चल नहीं रहा या host ग़लत है, दूसरे का मतलब है कि हल polling के timeout से ज़्यादा लंबा खिंच गया। चारों एक ही base से निकले हैं, इसलिए अगर आप सब कुछ एक ही जगह संभालना चाहें तो CapSkipError पकड़ लेना काम कर जाता है।
चरण 3: हल आख़िर में कीजिए, शुरू में नहीं
durable execution की वजह से token की समय समाप्ति यहाँ किसी भी दूसरे orchestrator के मुक़ाबले ज़्यादा आसानी से ग़लत हो जाती है। कोई Temporal workflow किसी signal का इंतज़ार कर सकता है, एक दिन सो सकता है और फिर चल पड़ सकता है, और उसकी दर्ज activity results history से ज्यों की त्यों वापस आती हैं। इसलिए जो workflow पहले हल कर लेता है, फिर किसी approval का इंतज़ार करता है, फिर submit करता है, वह ऐसा token replay करेगा जो कल बना था।
reCAPTCHA token एक ही बार स्वीकार होता है और लगभग दो मिनट में ख़त्म हो जाता है। workflow का क्रम ऐसा रखिए कि हल submit से ठीक पहले वाला step हो, और बीच में ऐसा कुछ न हो जो रुक सके। इस समस्या का आम रूप इस गाइड में बताया गया है: reCAPTCHA token की समय समाप्ति.
पूरा चलने वाला उदाहरण
तीन activities और एक workflow जो उन्हें क्रम से कॉल करता है। sitekey पढ़िए, हल कीजिए, submit कीजिए। हर activity को अपना timeout और अपनी retry policy मिलती है, और हर एक Temporal UI में अलग से दिखती है, इसलिए जब कुछ धीमा हो तो आप देख सकते हैं कि वह कौन सा step था।
# 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_codeworkflow ख़ुद क्रम तय करने से आगे कोई logic नहीं रखता। बात यही है: जो कुछ भी fail हो सकता है वह किसी activity में है, और workflow वह deterministic हिस्सा है जो replay में बचा रहता है।
# 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),
)दो arguments वाली activities पर args की सूची पर ध्यान दीजिए। एक अकेला positional argument सीधे भेजा जा सकता है, लेकिन एक से ज़्यादा को args से होकर जाना पड़ता है, और यही ग़लती इस SDK में सबसे आम पहली error है।
# 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 काम schedule करती है और history सहेजती है। वह आपका कोड कभी नहीं चलाती। आप अपने ही infrastructure के अंदर जो worker शुरू करते हैं, वह service से एक लंबा outbound connection बनाए रखता है और queue से tasks उठाता है। आपके network पर कभी कोई inbound रास्ता नहीं खुलता।
इसलिए CapSkip वाली उसी Windows मशीन पर चलने वाला worker 127.0.0.1:8080 को ठीक वैसे ही कॉल करता है जैसे आपकी मेज़ पर रखी कोई script करती, और Temporal Cloud पर भी यही सच रहता है। Temporal Cloud का cloud orchestration की परत है, compute की परत नहीं।
जैसे ही worker जगह बदलता है, host बदल जाता है और बाक़ी कुछ नहीं बदलता। किसी container में, किसी Linux VM पर या Kubernetes में चलने वाला worker loopback से किसी Windows सॉल्वर तक नहीं पहुँच सकता, इसलिए सॉल्वर Server mode पर चला जाता है। Local mode 127.0.0.1 से bind होता है और सिर्फ़ उसी डिवाइस को जवाब देता है। Server mode आपके network या public IP से bind होता है, और एक static public IP उस पते को स्थिर रखता है। दोनों modes में यह अब भी आपका हार्डवेयर है और अब भी बिना मीटर, इसलिए दिन में दस हज़ार बार चलने वाले 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)
जैसे ही सॉल्वर किसी network address पर सुनने लगे, key validation चालू कर दें, और हर worker fleet को उसकी अपनी key दें ताकि किसी एक को बाक़ियों को छुए बिना रद्द किया जा सके। दोनों modes को यहाँ विस्तार से समझाया गया है: CapSkip सेटअप गाइड.
आम errors और उनका मतलब
| आप जो देखते हैं | कारण | फिक्स |
|---|---|---|
| replay पर non-determinism error | हल को workflow कोड से कॉल किया गया | इसे किसी activity में ले जाइए और execute_activity से कॉल कीजिए |
| import पर RestrictedWorkflowAccessError | activity वाला module sandbox में import हो गया | उसे workflow.unsafe.imports_passed_through के अंदर import कीजिए |
| activity हल के बीच में ही रद्द हो जाती है, फिर retry होती है | start_to_close_timeout सॉल्वर के अपने timeout से छोटा है | reCAPTCHA, Turnstile और GeeTest के लिए इसे 300 सेकंड से ऊपर रखिए |
| form ऐसे token को ठुकरा देता है जो देखने में सही लगता है | workflow हल करने और submit करने के बीच में रुक गया | हल को submit से ठीक पहले वाला step बनाइए |
| एक ही विफलता पर चार प्रयास बर्बाद | किसी deterministic error को retry किया जा रहा है | ValidationException और ApiException को non-retryable के रूप में सूचीबद्ध कीजिए |
| हर प्रयास पर NetworkException | worker उस मशीन पर नहीं है जिस पर सॉल्वर चल रहा है | सॉल्वर को Server mode पर ले जाइए और host सेट कीजिए |
| किसी activity पर arguments को लेकर TypeError | दो positional arguments सीधे भेज दिए गए | उन्हें args parameter के ज़रिए एक list के रूप में भेजिए |
| SDK से TimeoutException | solve recaptchaTimeout से ज़्यादा लंबा चला | उसे 300 सेकंड के डिफ़ॉल्ट से ऊपर बढ़ाएँ |
FAQ
क्या कोई Temporal Cloud workflow सचमुच 127.0.0.1 को कॉल कर सकता है?
हाँ, क्योंकि Temporal Cloud आपका कोड नहीं चलाता। उसे आपका worker चलाता है, जहाँ भी आपने उसे शुरू किया हो, और वह बाहर की ओर service से जुड़ता है। उस worker पर loopback का मतलब है worker की अपनी मशीन, इसलिए उसी मशीन पर मौजूद सॉल्वर सामान्य रूप से जवाब देता है। dev server से Cloud पर जाने पर इसमें कुछ नहीं बदलता।
क्या solve वाली activity को heartbeat करना चाहिए?
इससे कोई फ़ायदा नहीं होगा, क्योंकि SDK की call token आने तक रुकी रहती है और उसके अंदर ऐसी कोई जगह नहीं है जहाँ से ख़बर दी जा सके। इसके बजाय activity को असली जगह वाला start_to_close_timeout दीजिए, और विफल प्रयास को policy से retry होने दीजिए। Heartbeating उन activities के लिए है जो आपके अपने नियंत्रण वाले काम पर loop करती हैं।
सॉल्वर पर बोझ डाले बिना पूरा batch कैसे हल करूँ?
worker पर max_concurrent_activities सेट कीजिए, या हल करने के काम को उसकी अपनी task queue और कम सीमा वाला अपना worker दीजिए। इससे throttling वहीं होती है जहाँ काम चलता है, जो workflow शुरू होने के समय को फैलाने की कोशिश से ज़्यादा भरोसेमंद है। एक ही process के अंदर async client asyncio से काम फैलाता है, और यह यहाँ बताया गया है: समानांतर में कैप्चा हल करने की गाइड.
यह Airflow में करने से कैसे अलग है?
Airflow एक DAG वाला scheduler है, और graph बनते समय network call करने से आपको कोई नहीं रोकता, जो footguns का एक अलग ही सेट है। Temporal durable execution है, इसलिए वहाँ पाबंदी determinism की है और जवाब हमेशा कोई activity ही होता है। सॉल्वर वाला हिस्सा दोनों में एक जैसा है, और Airflow वाला संस्करण यहाँ लिखा गया है: Airflow कैप्चा गाइड.
संक्षेप में
हल को किसी activity में रखिए, workflow method में कभी नहीं। उसे सॉल्वर की अपनी 300 सेकंड की सीमा से लंबा start_to_close_timeout दीजिए, ValidationException और ApiException को non-retryable बताइए, और workflow का क्रम ऐसा रखिए कि हल submit से ठीक पहले की आख़िरी चीज़ हो। worker को वहीं चलाइए जहाँ CapSkip है और host 127.0.0.1 ही बना रहता है। बाक़ी Python surface यहाँ है: Python कैप्चा सॉल्वर पेज, और वही तीन calls Node.js, PHP और C# में भी हैं, जैसा यहाँ सूचीबद्ध है: CAPTCHA solving SDK पेज। reCAPTCHA के options ख़ुद यहाँ हैं: reCAPTCHA v2 सॉल्वर पेज.
इस पर कोई schedule लगाने से पहले एक बात जान लेना ठीक रहेगा। CapSkip एक कैप्चा सॉल्वर है जो उस हार्डवेयर पर चलता है जो पहले से आपका है, इसलिए हर मिनट चलने वाले workflow की लागत उतनी ही है जितनी हफ़्ते में एक बार चलने वाले की।
