Integrations

Solving CAPTCHAs in Elixir Broadway Pipelines

This guide builds a Broadway pipeline that pulls CAPTCHA jobs from RabbitMQ or Amazon SQS, solves each one with CaptchaAI, submits the protected form from the same processor and records the outcome in Postgres. One constraint drives nearly every setting: a processor stays busy for the whole solve, 15 seconds before the first poll and up to 120 seconds in total. Processor concurrency is therefore your number of in-flight CAPTCHAs, and it must match your plan's thread count, for example BASIC ($15/month, 5 threads). The jobs are assumed to target forms in your own applications or portals you are authorized to automate.

What each Broadway option does here

Broadway runs producer, processors, then batchers. Here the processors do all the CaptchaAI work and spend the token, the :results batcher only writes rows, and failures leave through handle_failed/2. The Broadway documentation lists every option; these are the ones a CAPTCHA workload changes:

Option Value here Why
processors: [default: [concurrency: n]] plan threads, per node Each processor holds one CaptchaAI task from submit to token.
max_demand on processors 1 A processor never holds messages it cannot start; on SQS their visibility clock would keep running. BroadwayRabbitMQ's docs also suggest 1 for slow processing.
rate_limiting on the producer 5 messages per 1,000 ms Solve time bounds steady-state throughput anyway; the limiter caps fast-failing messages and the startup burst.
batchers: [results: ...] 50 per batch, 2,000 ms timeout One insert_all per batch; at five threads it usually flushes on the timeout.
shutdown 180_000 Longer than the slowest job, so a deploy drains in-flight solves.

Set processor concurrency explicitly. Broadway's default, System.schedulers_online() * 2, would tie your in-flight solves to the CPU count of whichever machine runs the release.

Dependencies

The pipeline needs Broadway, one producer package, Req for HTTP, Jason and Ecto. Plug is there for Req's test stubs. Leave it without only: :test: in a Phoenix app, Phoenix already needs Plug in every environment, and Mix refuses the narrower entry.

# mix.exs (deps only)
defp deps do
  [
    {:broadway, "~> 1.3"},
    {:broadway_rabbitmq, "~> 0.8"},
    # On Amazon SQS, use this instead of broadway_rabbitmq:
    # {:broadway_sqs, "~> 1.0"},
    {:req, "~> 0.7"},
    {:jason, "~> 1.4"},
    {:ecto_sql, "~> 3.13"},
    {:postgrex, ">= 0.0.0"},
    {:telemetry, "~> 1.3"},
    {:plug, "~> 1.16"}
  ]
end

Runtime configuration, including the API key, lives with each producer below.

A Req client for in.php and res.php

All CaptchaAI traffic goes through one module. A submit is a POST to https://ocr.captchaai.com/in.php with the form fields key, method=userrecaptcha, googlekey (the site key), pageurl and json=1. Success looks like {"status":1,"request":"<task id>"} and a rejection like {"status":0,"request":"ERROR_PAGEURL"}. The result comes from res.php with action=get, the task id and json=1; until the token is ready the answer is CAPCHA_NOT_READY, spelled without the T. CaptchaAI's reCAPTCHA v2 guide says to wait 15 to 20 seconds before the first poll, then poll every 5 seconds; this client waits 15, polls every 5 and gives up 120 seconds after submitting.

# lib/my_app/captcha_ai.ex
defmodule MyApp.CaptchaAI do
  @moduledoc "CaptchaAI client (in.php + res.php, json=1) shared by Broadway, Phoenix and Oban."

  @account ~w(ERROR_WRONG_USER_KEY ERROR_KEY_DOES_NOT_EXIST IP_BANNED)
  @transient ~w(ERROR_SERVER_ERROR ERROR_INTERNAL_SERVER_ERROR ERROR_CAPTCHA_UNSOLVABLE)

  @type error :: :not_ready | :timeout | {:api, String.t()} | {:transport, Exception.t()} | {:unexpected, term()}

  @spec solve(map()) :: {:ok, String.t()} | {:error, error()}
  def solve(params) do
    deadline = System.monotonic_time(:millisecond) + setting(:deadline_ms, 120_000)

    with {:ok, task_id} <- submit(params) do
      Process.sleep(setting(:first_wait_ms, 15_000))
      poll(task_id, deadline)
    end
  end

  def submit(params) do
    form = Map.merge(params, %{"key" => api_key(), "json" => 1})
    client() |> Req.post(url: "/in.php", form: form) |> decode()
  end

  def fetch_result(task_id) do
    params = [key: api_key(), action: "get", id: task_id, json: 1]
    client() |> Req.get(url: "/res.php", params: params) |> decode()
  end

  @doc "How a caller should react to an error: :transient, :permanent or :account."
  def classify({:api, "ERROR_ZERO_BALANCE"}) do
    # No free thread, or no active plan: threadsinfo tells the two apart. A thread can free
    # up between in.php and threadsinfo, so "no plan" needs free threads on two reads.
    if threads_free?() and threads_free?(5_000), do: :account, else: :transient
  end

  def classify({:api, code}) when code in @account, do: :account
  def classify({:api, code}) when code in @transient, do: :transient
  def classify({:api, _code}), do: :permanent
  def classify(_timeout_network_or_form_error), do: :transient

  defp poll(task_id, deadline) do
    result = fetch_result(task_id)
    delay = poll_again_in(result)

    cond do
      is_nil(delay) ->
        result

      # Give up before sleeping past the deadline, so no res.php call starts after it.
      System.monotonic_time(:millisecond) + delay >= deadline ->
        {:error, :timeout}

      true ->
        Process.sleep(delay)
        poll(task_id, deadline)
    end
  end

  # The task is still alive at CaptchaAI: ask about the same ID again, never resubmit.
  defp poll_again_in({:error, :not_ready}), do: setting(:poll_interval_ms, 5_000)
  defp poll_again_in({:error, {:transport, _}}), do: setting(:poll_interval_ms, 5_000)
  defp poll_again_in({:error, {:unexpected, text}}) when is_binary(text), do: setting(:poll_interval_ms, 5_000)
  defp poll_again_in({:error, {:api, "ERROR_INTERNAL_SERVER_ERROR"}}), do: 2 * setting(:poll_interval_ms, 5_000)
  defp poll_again_in(_result), do: nil

  defp decode({:ok, %Req.Response{body: body}}), do: parse(json(body))
  defp decode({:error, exception}), do: {:error, {:transport, exception}}

  defp parse(%{"status" => 1, "request" => answer}), do: {:ok, to_string(answer)}
  defp parse(%{"request" => "CAPCHA_NOT_READY"}), do: {:error, :not_ready}
  defp parse(%{"request" => code}) when is_binary(code), do: {:error, {:api, code}}
  # json=1 is not a promise: some errors arrive as bare text, and 5xx pages as HTML.
  defp parse(text) when is_binary(text) do
    if text =~ ~r/^(ERROR_[A-Z_]+|IP_BANNED)$/,
      do: {:error, {:api, text}},
      else: {:error, {:unexpected, String.slice(text, 0, 120)}}
  end

  defp parse(other), do: {:error, {:unexpected, other}}

  defp threads_free?(wait_ms \\ 0) do
    Process.sleep(wait_ms)

    case Req.get(client(), url: "/res.php", params: [key: api_key(), action: "threadsinfo"]) do
      {:ok, %Req.Response{body: body}} -> free_threads?(json(body))
      {:error, _exception} -> false
    end
  end

  # threadsinfo answers {"threads":"5","working_threads":3}: a string and an integer.
  defp free_threads?(%{"threads" => threads, "working_threads" => working}),
    do: to_int(working) < to_int(threads)

  defp free_threads?(_other), do: false

  # Req decodes JSON only when the response's content type says application/json.
  defp json(body) when is_binary(body) do
    case Jason.decode(body) do
      {:ok, %{} = map} -> map
      _ -> String.trim(body)
    end
  end

  defp json(body), do: body

  defp to_int(value) when is_integer(value), do: value
  defp to_int(value) when is_binary(value), do: String.to_integer(value)

  # :plug is set only in tests (Req.Test); production calls go to ocr.captchaai.com.
  defp client do
    Req.new([base_url: "https://ocr.captchaai.com", retry: false] ++ Keyword.take(config(), [:plug]))
  end

  defp api_key, do: Keyword.fetch!(config(), :api_key)
  defp setting(key, default), do: Keyword.get(config(), key, default)
  defp config, do: Application.fetch_env!(:my_app, __MODULE__)
end

Three decisions in this module matter under Broadway:

  • retry: false. Req's default retry: :safe_transient retries only GET and HEAD, up to three times, on 408, 429, 500, 502, 503 and 504 responses and on timed-out, refused or closed connections. It never retries POST, so a slow in.php call cannot create a second task either way. Turning retries off keeps a poll's worst-case time predictable; the poll loop retries the same task ID itself.
  • Defensive parsing. Even with json=1, CaptchaAI can answer a bare error code as text, a 5xx arrives as HTML, and Req decodes JSON only when the content type says so.
  • classify/1 is the retry policy every consumer shares. ERROR_ZERO_BALANCE means no free thread or no active plan, and res.php?action=threadsinfo tells the two apart: all threads working means busy, free threads mean the plan is the problem. The second read five seconds later matters under Broadway. With one read, a thread that finished between the in.php rejection and the threadsinfo call would look like a dead plan and stop the whole pipeline.

The client covers types that return the token in request: reCAPTCHA v2 and v3 and Cloudflare Turnstile. reCAPTCHA Enterprise and Cloudflare Challenge answer in result with a user_agent you must reuse, so extend parse/1 before routing those jobs here.

Spending the token inside the processor

Google's verification docs say each reCAPTCHA response token is valid for two minutes and can be verified only once. That rules out solving in the processor and submitting the form from the batcher or a second queue, where a batch timeout or a backlog can age the token past its window. The processor that solved it spends it.

TargetForm loads the form page, keeps its cookies, and posts the job's fields plus g-recaptcha-response with them, so the token arrives in the session that loaded the challenge:

# lib/my_app/target_form.ex
defmodule MyApp.TargetForm do
  @moduledoc "Loads a form you own or may automate, then posts it with the token and the page's cookies."

  def open(page_url) do
    case Req.get(req(), url: page_url) do
      {:ok, %Req.Response{status: 200} = resp} ->
        cookies =
          resp
          |> Req.Response.get_header("set-cookie")
          |> Enum.map_join("; ", fn cookie -> cookie |> String.split(";") |> hd() end)

        {:ok, Req.merge(req(), headers: if(cookies == "", do: [], else: [cookie: cookies]))}

      {:ok, %Req.Response{status: status}} ->
        {:error, {:form, status}}

      {:error, exception} ->
        {:error, {:transport, exception}}
    end
  end

  def submit(session, action_url, fields) do
    case Req.post(session, url: action_url, form: fields) do
      {:ok, %Req.Response{status: status}} when status in 200..399 -> {:ok, status}
      {:ok, %Req.Response{status: status}} -> {:error, {:form, status}}
      {:error, exception} -> {:error, {:transport, exception}}
    end
  end

  defp req, do: Req.new([retry: false, redirect: false] ++ Application.get_env(:my_app, __MODULE__, []))
end

Opening the page first means an outage in your own app fails the job before it occupies a CaptchaAI thread. A 4xx from the form counts as transient: a redelivery gets a fresh solve, never a second post of a rejected token.

The pipeline module

handle_message/3 validates the job, solves, submits and hands a small result map to the :results batcher. Anything that goes wrong becomes Message.failed/2 with a class, and configuration maps the class to an acknowledgement, because RabbitMQ and SQS spell their ack options differently.

# lib/my_app/captcha_pipeline.ex
defmodule MyApp.CaptchaPipeline do
  use Broadway

  require Logger
  alias Broadway.Message
  alias MyApp.{CaptchaAI, TargetForm}

  def start_link(_opts) do
    config = Application.fetch_env!(:my_app, __MODULE__)

    Broadway.start_link(__MODULE__,
      name: __MODULE__,
      producer: [
        module: Keyword.fetch!(config, :producer),
        concurrency: 1,
        rate_limiting: [allowed_messages: 5, interval: 1_000]
      ],
      # One processor = one CaptchaAI task in flight = one plan thread.
      processors: [default: [concurrency: Keyword.fetch!(config, :threads), max_demand: 1]],
      batchers: [results: [batch_size: Keyword.get(config, :batch_size, 50), batch_timeout: 2_000]],
      shutdown: 180_000
    )
  end

  @impl true
  def handle_message(_processor, message, _context) do
    case Jason.decode(message.data) do
      {:ok, %{"job_id" => _, "pageurl" => _, "googlekey" => _, "submit_url" => _} = job} ->
        solve_and_submit(message, job)

      _ ->
        fail(message, :permanent, :malformed_job)
    end
  end

  defp solve_and_submit(message, job) do
    captcha = %{"method" => "userrecaptcha", "googlekey" => job["googlekey"], "pageurl" => job["pageurl"]}

    with {:ok, session} <- TargetForm.open(job["pageurl"]),
         {:ok, token} <- CaptchaAI.solve(captcha),
         fields = Map.put(job["fields"] || %{}, "g-recaptcha-response", token),
         {:ok, status} <- TargetForm.submit(session, job["submit_url"], fields) do
      message
      |> Message.update_data(fn _raw ->
        %{job_id: job["job_id"], http_status: status, completed_at: DateTime.utc_now()}
      end)
      |> Message.put_batcher(:results)
    else
      {:error, reason} -> fail(message, CaptchaAI.classify(reason), reason)
    end
  end

  defp fail(message, class, reason) do
    config = Application.fetch_env!(:my_app, __MODULE__)
    # RabbitMQ redelivers a requeued message at once; the pause gives transient errors room.
    if class == :transient, do: Process.sleep(Keyword.get(config, :transient_pause_ms, 10_000))

    message
    |> Message.configure_ack(on_failure: config |> Keyword.fetch!(:failure_acks) |> Keyword.fetch!(class))
    |> Message.failed({class, reason})
  end

  @impl true
  def handle_batch(:results, messages, _batch_info, _context) do
    rows = Enum.map(messages, & &1.data)
    # job_id is the primary key: a redelivered job that was already recorded is skipped.
    MyApp.Repo.insert_all("captcha_results", rows, on_conflict: :nothing, conflict_target: [:job_id])
    messages
  end

  @impl true
  def handle_failed(messages, _context) do
    for %Message{status: {:failed, {class, reason}}} <- messages do
      Logger.warning("captcha job failed: class=#{class} reason=#{inspect(reason)}")
      if class == :account, do: halt(reason)
    end

    messages
  end

  # A bad key or a dead plan fails every job the same way. Stop consuming and page someone.
  defp halt(reason) do
    Logger.error("CaptchaAI account error #{inspect(reason)}: stopping #{inspect(__MODULE__)}")

    Task.start(fn ->
      try do
        Broadway.stop(__MODULE__, :shutdown)
      catch
        :exit, _already_stopping -> :ok
      end
    end)
  end

  @doc false
  # BroadwayRabbitMQ :after_connect hook; it runs before the producer declares captcha_jobs.
  def declare_dead_letter(channel) do
    with :ok <- AMQP.Exchange.declare(channel, "captcha_jobs.dlx", :fanout, durable: true),
         {:ok, _} <- AMQP.Queue.declare(channel, "captcha_jobs.dlq", durable: true),
         :ok <- AMQP.Queue.bind(channel, "captcha_jobs.dlq", "captcha_jobs.dlx") do
      :ok
    end
  end
end

The results table uses the job ID as its primary key, which is what makes on_conflict: :nothing absorb redeliveries:

# priv/repo/migrations/20260929000000_create_captcha_results.exs
defmodule MyApp.Repo.Migrations.CreateCaptchaResults do
  use Ecto.Migration

  def change do
    create table(:captcha_results, primary_key: false) do
      add :job_id, :string, primary_key: true
      add :http_status, :integer, null: false
      add :completed_at, :utc_datetime_usec, null: false
    end
  end
end

The key protects the row, not the target. If insert_all raises, the whole batch fails with the producer's default on_failure, and those jobs run again, form post included. When a second post matters, make the target endpoint idempotent on job_id.

handle_failed/2 runs before Broadway acknowledges failed messages, so it logs the reason and, for account errors, stops the pipeline. Broadway.stop/3 runs in a separate task because a processor cannot wait for its own pipeline to shut down. Messages that raise inside handle_message/3 skip the {:failed, _} pattern; Broadway logs them and applies the producer's default on_failure.

How each failure is acknowledged

Outcome Class RabbitMQ on_failure SQS on_failure
Solved and the form accepted it none: goes to the :results batcher ack after insert_all delete after insert_all
ERROR_PAGEURL, ERROR_GOOGLEKEY, ERROR_WRONG_GOOGLEKEY, ERROR_BAD_TOKEN_OR_PAGEURL, ERROR_BAD_PARAMETERS, malformed job JSON permanent :reject (dead-lettered) {:nack, 0} (DLQ after maxReceiveCount)
ERROR_SERVER_ERROR, ERROR_INTERNAL_SERVER_ERROR, ERROR_CAPTCHA_UNSOLVABLE, the 120-second timeout, network errors, a form 4xx or 5xx transient :reject_and_requeue_once {:nack, 30}
ERROR_ZERO_BALANCE while every thread is busy transient :reject_and_requeue_once {:nack, 30}
ERROR_WRONG_USER_KEY, ERROR_KEY_DOES_NOT_EXIST, IP_BANNED, or ERROR_ZERO_BALANCE with threads free (no active plan) account: the pipeline stops :reject_and_requeue :noop

ERROR_CAPTCHA_UNSOLVABLE earns a bounded number of fresh submissions (one requeue on RabbitMQ, up to maxReceiveCount receives on SQS); the client never polls that task ID again. Account errors stop the pipeline instead of dead-lettering because a wrong key fails every job identically, and CaptchaAI bans the calling IP (IP_BANNED, lifted after 5 minutes) after repeated wrong-key requests. Rejecting jobs one by one would drain the queue into the DLQ while hammering the API. Replaying dead-lettered jobs after a fix is covered in the dead-letter queue guide for failed CAPTCHA tasks.

RabbitMQ: prefetch, requeue and dead-lettering

# config/runtime.exs
import Config

if config_env() != :test do
  config :my_app, MyApp.CaptchaAI, api_key: System.fetch_env!("CAPTCHAAI_API_KEY")

  # In-flight solves on THIS node. Across all nodes the total must not exceed your plan's threads.
  threads = String.to_integer(System.get_env("CAPTCHAAI_THREADS", "5"))

  config :my_app, MyApp.CaptchaPipeline,
    threads: threads,
    batch_size: 50,
    producer:
      {BroadwayRabbitMQ.Producer,
       queue: "captcha_jobs",
       connection: System.fetch_env!("AMQP_URL"),
       declare: [durable: true, arguments: [{"x-queue-type", :longstr, "quorum"}]],
       after_connect: &MyApp.CaptchaPipeline.declare_dead_letter/1,
       # Unacked messages = one per processor + up to a full results batch.
       qos: [prefetch_count: threads + 50],
       on_failure: :reject_and_requeue_once},
    failure_acks: [permanent: :reject, transient: :reject_and_requeue_once, account: :reject_and_requeue]
end

BroadwayRabbitMQ's default on_failure is :reject_and_requeue, and its producer documentation warns that it loops forever on a message that can never succeed. Here the producer default is :reject_and_requeue_once, and each failure overrides it per message through Message.configure_ack/2.

RabbitMQ stops delivering once a channel holds prefetch_count unacknowledged messages. BroadwayRabbitMQ defaults to 50 and asks for at least max_demand times the processor count, more with batchers. Messages waiting in the results batcher are still unacked, so threads + 50 keeps processors fed while a batch fills; more than that parks jobs on one node that another could be solving.

If a node dies mid-solve, RabbitMQ requeues its unacked deliveries when the channel closes and marks them redelivered. The job is solved again, which on a thread-based plan costs thread time rather than a per-solve fee, and the primary key keeps the result row single. :reject_and_requeue_once reads that redelivered flag, so a job that already came back after a crash is dead-lettered on its next transient failure.

The dead-letter exchange and queue are declared by declare_dead_letter/1 on every connect. A policy then attaches them to captcha_jobs; RabbitMQ recommends policies over hard-coded x- arguments because a policy can change without redeclaring the queue. The queue type is the one x- argument left in declare:, since it is fixed when the queue is created:

#!/usr/bin/env bash
# Dead-letter rejected jobs to captcha_jobs.dlx and cap redeliveries of any one job at 5.
set -euo pipefail

rabbitmqctl set_policy captcha-jobs '^captcha_jobs$' \
  '{"dead-letter-exchange": "captcha_jobs.dlx", "delivery-limit": 5}' \
  --priority 10 \
  --apply-to queues

On a quorum queue, delivery-limit is the backstop for any loop you did not anticipate: since RabbitMQ 4.0 it defaults to 20, and past it a message is dead-lettered or dropped. For broader exchange and routing patterns, see RabbitMQ and CaptchaAI message queue integration.

Amazon SQS: visibility timeout and redrive

Swap the dependency to broadway_sqs and delete declare_dead_letter/1, which calls the AMQP client. BroadwaySQS 1.0 talks to SQS through Req and finds AWS credentials through aws_credentials (environment, profile or instance role), so the config below needs only the region. Replace the pipeline block in config/runtime.exs:

# config/runtime.exs: SQS version of the MyApp.CaptchaPipeline block
config :my_app, MyApp.CaptchaPipeline,
  threads: threads,
  # BroadwaySQS recommends batches of 10 so acks go out as SQS batch requests.
  batch_size: 10,
  producer:
    {BroadwaySQS.Producer,
     queue_url: System.fetch_env!("CAPTCHA_QUEUE_URL"),
     config: [region: System.fetch_env!("AWS_REGION")],
     wait_time_seconds: 20,
     visibility_timeout: 180},
  # A nack already delays the retry, so no pause inside the processor.
  transient_pause_ms: 0,
  failure_acks: [permanent: {:nack, 0}, transient: {:nack, 30}, account: :noop]

On SQS the visibility timeout is the whole story. Per the SQS visibility timeout guide, it starts when a message is delivered, not when a processor picks it up, and when it runs out other consumers can receive the message. If that happens mid-poll, a second node solves the same CAPTCHA. Size it from the worst case: 120 seconds of solving, the last res.php request that started before the deadline (Req's default receive timeout is 15 seconds), then your page load and form post at up to 15 seconds each, about 165 seconds. 180 leaves a margin.

BroadwaySQS leaves failed messages alone by default (:noop), so they reappear after the full visibility timeout. {:nack, 30} brings a transient failure back in 30 seconds, and {:nack, 0} returns a permanent one at once, so it reaches maxReceiveCount and the dead-letter queue after a few quick receives. SQS has no reject-to-DLQ call, so a broken job does run three times, three rejected in.php calls, before it lands there. Keep maxReceiveCount low for that reason. The queue and its DLQ:

#!/usr/bin/env bash
# captcha-jobs: 180 s visibility timeout, long polling, dead-letter after 3 receives.
set -euo pipefail

dlq_url=$(aws sqs create-queue --queue-name captcha-jobs-dlq --query QueueUrl --output text)
dlq_arn=$(aws sqs get-queue-attributes --queue-url "$dlq_url" \
  --attribute-names QueueArn --query Attributes.QueueArn --output text)

attributes=$(mktemp)
cat > "$attributes" <<EOF
{
  "VisibilityTimeout": "180",
  "ReceiveMessageWaitTimeSeconds": "20",
  "RedrivePolicy": "{\"deadLetterTargetArn\":\"${dlq_arn}\",\"maxReceiveCount\":\"3\"}"
}
EOF

aws sqs create-queue --queue-name captcha-jobs --attributes "file://${attributes}"
rm -f "$attributes"

SQS standard queues deliver at least once, even inside the visibility window, so the idempotent result write matters here as much as on RabbitMQ. Guarding the solve itself, not just the write, is covered in idempotent CAPTCHA solving.

Telemetry for solve latency and failures

Broadway wraps every handle_message/3 call in a [:broadway, :processor, :message] span. The :stop event carries the message as handle_message/3 returned it, so its status holds the failure class, and the duration covers the whole open, solve and submit sequence. This handler folds both into one event your reporter can aggregate:

# lib/my_app/captcha_telemetry.ex
defmodule MyApp.CaptchaTelemetry do
  @moduledoc "Turns Broadway's per-message events into one CAPTCHA job metric."

  @events [
    [:broadway, :processor, :message, :stop],
    [:broadway, :processor, :message, :exception]
  ]

  def attach do
    :telemetry.attach_many("captcha-pipeline", @events, &__MODULE__.handle_event/4, nil)
  end

  def handle_event(
        [:broadway, :processor, :message, kind],
        %{duration: duration},
        %{topology_name: MyApp.CaptchaPipeline} = meta,
        _config
      ) do
    outcome =
      case {kind, meta.message.status} do
        {:exception, _} -> :crashed
        {:stop, :ok} -> :solved
        {:stop, {:failed, {class, _reason}}} -> class
        {:stop, _other} -> :failed
      end

    :telemetry.execute(
      [:my_app, :captcha, :job],
      %{duration_ms: System.convert_time_unit(duration, :native, :millisecond)},
      %{outcome: outcome}
    )
  end

  def handle_event(_event, _measurements, _meta, _config), do: :ok
end

Point Telemetry.Metrics or your StatsD or Prometheus reporter at my_app.captcha.job.duration_ms, tagged by outcome. Alert on any :account outcome (the pipeline has stopped), on a rising share of :transient, and on a p95 creeping toward 180 seconds, which means the visibility timeout and shutdown window are about to be too small. Processors divided by the median seconds per job is a node's real capacity, the number to compare with the producer's 5-per-second limit.

Testing with Broadway.test_message and Req.Test

In the test environment the producer becomes Broadway.DummyProducer, and Broadway.test_message/3 pushes one message through, flushes its batch at once and sends {:ack, ref, successful, failed} back to the test process. Its acknowledger also reports every Message.configure_ack/2 call as {:configure, ref, options}, which lets a test assert the ack decision itself: a bad site key must be rejected, never requeued.

# config/test.exs
import Config

config :my_app, MyApp.CaptchaAI,
  api_key: "test-key",
  plug: {Req.Test, MyApp.CaptchaAI},
  first_wait_ms: 0,
  poll_interval_ms: 0

config :my_app, MyApp.TargetForm, plug: {Req.Test, MyApp.TargetForm}

config :my_app, MyApp.CaptchaPipeline,
  threads: 2,
  producer: {Broadway.DummyProducer, []},
  transient_pause_ms: 0,
  failure_acks: [permanent: :reject, transient: :reject_and_requeue_once, account: :reject_and_requeue]
# test/my_app/captcha_pipeline_test.exs
defmodule MyApp.CaptchaPipelineTest do
  # Processors are not children of the test process: stubs and the DB sandbox run shared.
  use ExUnit.Case, async: false

  alias Ecto.Adapters.SQL.Sandbox

  @job Jason.encode!(%{
         job_id: "job-1",
         pageurl: "https://app.example.test/signup",
         googlekey: "test-sitekey",
         submit_url: "https://app.example.test/signup"
       })

  setup do
    :ok = Sandbox.checkout(MyApp.Repo)
    Sandbox.mode(MyApp.Repo, {:shared, self()})
    Req.Test.set_req_test_to_shared()
    Req.Test.stub(MyApp.TargetForm, &Plug.Conn.send_resp(&1, 200, "ok"))
    :ok
  end

  test "a solved job is recorded through the results batcher" do
    Req.Test.stub(MyApp.CaptchaAI, fn
      %{request_path: "/in.php"} = conn -> Req.Test.json(conn, %{status: 1, request: "73185"})
      %{request_path: "/res.php"} = conn -> Req.Test.json(conn, %{status: 1, request: "token-abc"})
    end)

    ref = Broadway.test_message(MyApp.CaptchaPipeline, @job)
    assert_receive {:ack, ^ref, [%{data: %{job_id: "job-1", http_status: 200}}], []}, 1_000
  end

  test "a bad site key is dead-lettered, not requeued" do
    Req.Test.stub(MyApp.CaptchaAI, &Req.Test.json(&1, %{status: 0, request: "ERROR_WRONG_GOOGLEKEY"}))

    ref = Broadway.test_message(MyApp.CaptchaPipeline, @job)
    assert_receive {:configure, ^ref, [on_failure: :reject]}, 1_000
    assert_receive {:ack, ^ref, [], [%{status: {:failed, {:permanent, {:api, "ERROR_WRONG_GOOGLEKEY"}}}}]}, 1_000
  end
end

Run it with mix test; it assumes the usual pool: Ecto.Adapters.SQL.Sandbox setting for the test Repo. The zero waits in config/test.exs keep the 15-second first poll out of the suite.

Supervision and draining on deploy

# lib/my_app/application.ex
defmodule MyApp.Application do
  use Application

  @impl true
  def start(_type, _args) do
    MyApp.CaptchaTelemetry.attach()

    children = [
      MyApp.Repo,
      # :transient, so a deliberate stop after an account error stays stopped until a restart.
      Supervisor.child_spec(MyApp.CaptchaPipeline, restart: :transient)
    ]

    Supervisor.start_link(children, strategy: :one_for_one, name: MyApp.Supervisor)
  end
end

A supervisor stops its children in reverse start order, so when a release receives SIGTERM the pipeline drains before the Repo goes away and the last batch still gets written. Broadway's default shutdown is 30 seconds, shorter than one reCAPTCHA job can take; at 180 seconds, in-flight jobs finish and are acknowledged. Make sure your platform's grace period before SIGKILL is longer still. Otherwise the kill lands mid-drain and the unacked jobs come back: immediately on RabbitMQ, after the visibility timeout on SQS.

Calling the same client from Phoenix and Oban

The client contains nothing Broadway-specific, so request handlers and background jobs can share it along with its error classes. From a Phoenix controller:

# lib/my_app_web/controllers/captcha_check_controller.ex
defmodule MyAppWeb.CaptchaCheckController do
  use MyAppWeb, :controller

  # Internal smoke test after you rotate your own site key: does it still solve?
  def create(conn, %{"googlekey" => googlekey, "pageurl" => pageurl}) do
    params = %{"method" => "userrecaptcha", "googlekey" => googlekey, "pageurl" => pageurl}
    task = Task.async(fn -> MyApp.CaptchaAI.solve(params) end)

    case Task.yield(task, 140_000) || Task.shutdown(task, :brutal_kill) do
      {:ok, {:ok, _token}} -> json(conn, %{solvable: true})
      {:ok, {:error, reason}} -> conn |> put_status(502) |> json(%{solvable: false, reason: inspect(reason)})
      nil -> conn |> put_status(504) |> json(%{solvable: false, reason: "timeout"})
    end
  end
end

Keep this for internal callers: it holds the request open for up to about two minutes, and Task.yield/2 plus Task.shutdown/2 caps it even if a call hangs. In an Oban worker, split the work instead of blocking a job:

# lib/my_app/workers/captcha_job.ex
defmodule MyApp.Workers.CaptchaJob do
  # config :my_app, Oban, queues: [captcha: 5]  (5 = BASIC threads)
  use Oban.Worker, queue: :captcha, max_attempts: 3

  alias MyApp.CaptchaAI

  @impl Oban.Worker
  def perform(%Oban.Job{args: %{"ref" => ref, "task_id" => task_id, "deadline" => deadline}}) do
    case CaptchaAI.fetch_result(task_id) do
      {:ok, token} -> Phoenix.PubSub.broadcast(MyApp.PubSub, "captcha:" <> ref, {:captcha_solved, token})
      {:error, :not_ready} -> if System.os_time(:second) < deadline, do: {:snooze, 5}, else: {:cancel, :timeout}
      # Never poll an unsolvable task ID again; submitting afresh is the caller's decision.
      {:error, {:api, "ERROR_CAPTCHA_UNSOLVABLE"} = reason} -> {:cancel, reason}
      {:error, reason} -> retry_or_cancel(reason)
    end
  end

  def perform(%Oban.Job{args: %{"ref" => ref, "captcha" => params}}) do
    case CaptchaAI.submit(params) do
      {:ok, task_id} ->
        %{"ref" => ref, "task_id" => task_id, "deadline" => System.os_time(:second) + 120}
        |> new(schedule_in: 15)
        |> Oban.insert()

      {:error, reason} ->
        retry_or_cancel(reason)
    end
  end

  defp retry_or_cancel(reason), do: if(CaptchaAI.classify(reason) == :transient, do: {:error, reason}, else: {:cancel, reason})
end

The first run submits and schedules a poll job 15 seconds out; the poll job returns {:snooze, 5} until the token is ready. Oban's docs say a snooze rolls back the attempt, so max_attempts never ends a snoozing job; the deadline in its args does. max_attempts only bounds real errors: network failures and ERROR_INTERNAL_SERVER_ERROR retry the same poll, while ERROR_CAPTCHA_UNSOLVABLE cancels instead of asking about a dead task ID three times. Snoozed jobs don't occupy the queue, so captcha: 5 throttles concurrent calls rather than capping in-flight tasks. Nothing on the Oban side holds you to your plan's threads. Past the cap, CaptchaAI either rejects the submit with ERROR_ZERO_BALANCE, which classify/1 turns into an Oban retry, or the task just stays CAPCHA_NOT_READY longer, so keep the deadline. A fuller Phoenix design with LiveView is in running CAPTCHA solves in backend services, and the same ideas for Celery, Sidekiq, BullMQ and River are in CAPTCHA solving in background job queues.

Troubleshooting

Symptom Likely cause Fix
Queue depth grows while every processor sits in Process.sleep Working as designed: throughput is processors divided by solve time Add plan threads, not processors; see when to upgrade your thread limit
ERROR_ZERO_BALANCE although threads matches the plan Concurrency is per node (three nodes at 5 run 15 solves), or another service shares the key Set CAPTCHAAI_THREADS to plan threads divided by nodes, and compare working_threads from threadsinfo with your own in-flight count
The same job solved twice SQS visibility shorter than the slowest job, or a node killed before acking Visibility above the worst case, shutdown inside the platform's grace period, idempotent writes
Your form rejects a fresh token Token older than two minutes, posted without the page's cookies, or posted twice Open, solve and submit in one handle_message/3; never re-post a token
One failing job floods the logs :reject_and_requeue (the RabbitMQ default) or :noop on SQS for a permanent error Classify it permanent: :reject or {:nack, 0} plus a DLQ, with delivery-limit as a backstop
Pipeline gone while the app stays up An :account failure called Broadway.stop/3 Fix the logged key or plan problem, then restart the release

FAQ

Does this pipeline work for Turnstile or reCAPTCHA v3?

Yes, with a different task map. Turnstile uses method=turnstile with sitekey and pageurl; the token goes in cf-turnstile-response and is single-use. reCAPTCHA v3 keeps method=userrecaptcha, adds version=v3 and the page's action, and is sent without a proxy per CaptchaAI's proxy guide. Both answer in request, so the client works unchanged; add a type field to the job and build the map from it.

Why not solve CAPTCHAs in handle_batch to use batching?

CaptchaAI's API takes one task per in.php call and each task holds its own thread, so batching gains nothing on the API side. It would also age tokens by up to batch_timeout plus the slowest solve in the batch. Batch the database writes instead.

Size the pipeline to your plan

Pick the plan whose threads cover the in-flight solves you need across all nodes, set CAPTCHAAI_THREADS on each node to its share, and let Broadway's processors, acks and batchers do the rest. Plans and thread counts are on the CaptchaAI pricing page.

Comments are disabled for this article.