Costruire un receiver di webhook affidabile per la WhatsApp Cloud API su larga scala
Webhook della WhatsApp Cloud API su larga scala: verifica della firma sul body grezzo, risposta 200 immediata, deduplica, ordine per contatto e media.
Un receiver di webhook affidabile per la WhatsApp Cloud API fa quattro cose durante la richiesta, e nient’altro: verifica l’header X-Hub-Signature-256 sul body grezzo della richiesta, salva il payload in modo durevole, risponde subito con un 200 e passa il lavoro all’elaborazione asincrona. Tutto ciò che rende difficili le integrazioni con WhatsApp (deduplicare consegne at-least-once, mantenere in ordine ogni conversazione, riconciliare aggiornamenti di stato fuori ordine, rispettare la finestra di 24 ore, scaricare i media) avviene dopo il 200, dove un passaggio lento non può spingere Meta a ritentare e moltiplicare il tuo carico. Ho passato sei anni in Turn.io a lavorare su una piattaforma WhatsApp che gestisce milioni di messaggi al giorno, e questo è il modo in cui strutturerei il receiver in Elixir e Phoenix.
In breve
- Conserva il body grezzo e verifica
X-Hub-Signature-256(HMAC-SHA256 con l’app secret) esattamente su quei byte. - Durante la richiesta: verifica, salva, rispondi 200. Restituisci un 5xx solo se non sei riuscito a salvare, così Meta ritenta.
- I webhook sono at-least-once: deduplica i messaggi in arrivo per id del messaggio con un indice unique.
- Elabora i messaggi di ogni contatto in ordine, uno alla volta: un processo per contatto, oppure una coda partizionata per contatto.
- Gli aggiornamenti di stato arrivano fuori ordine: fai avanzare lo stato di un messaggio solo in avanti.
- Tieni traccia della finestra di assistenza clienti di 24 ore per ogni contatto e passa ai template quando sei fuori.
- Scarica i media subito, in un job in background: l’URL di download ha vita breve.
Cosa ti manda Meta
I payload dei webhook arrivano raggruppati. Un singolo POST contiene un array entry, ogni entry ha delle changes, e ogni change ha un value che può contenere messages (messaggi in arrivo dagli utenti), statuses (aggiornamenti sui messaggi che hai inviato), contacts e metadata con il phone_number_id a cui appartiene l’evento. Non dare per scontato un messaggio per richiesta, e non dare per scontato che una richiesta contenga un solo tipo di evento.
Prima di tutto questo, Meta verifica il tuo endpoint con una richiesta GET che contiene hub.mode=subscribe, un hub.verify_token che hai scelto tu e un hub.challenge che devi restituire tale e quale. Il verify token serve solo per questo handshake. È l’app secret a firmare ogni POST. Sono valori diversi, e confonderli è un classico bug del primo giorno.
Verificare X-Hub-Signature-256 sul body grezzo
La firma è un HMAC-SHA256 del body della richiesta calcolato con il tuo app secret, inviato come sha256=<hex digest>. È calcolata esattamente sui byte inviati da Meta. Quando un controller Phoenix entra in azione, Plug.Parsers ha già consumato il body e decodificato il JSON, e ricodificarlo non riproduce gli stessi byte. Quindi conservane una copia durante il parsing, con un body reader personalizzato:
defmodule MyAppWeb.CacheRawBody do
@moduledoc "Keeps the raw request body for webhook signature verification."
@webhook_path "/webhooks/whatsapp"
def read_body(conn, opts) do
case Plug.Conn.read_body(conn, opts) do
{:ok, chunk, conn} -> {:ok, chunk, store(conn, chunk)}
{:more, chunk, conn} -> {:more, chunk, store(conn, chunk)}
{:error, _reason} = error -> error
end
end
# Conserva una copia solo per la route dei webhook, così le altre richieste non ne pagano il costo.
defp store(%Plug.Conn{request_path: @webhook_path} = conn, chunk) do
Plug.Conn.assign(conn, :raw_body, [conn.assigns[:raw_body] || "", chunk])
end
defp store(conn, _chunk), do: conn
end
Collegalo ai parser nel tuo endpoint:
plug Plug.Parsers,
parsers: [:urlencoded, :multipart, :json],
pass: ["*/*"],
body_reader: {MyAppWeb.CacheRawBody, :read_body, []},
json_decoder: Phoenix.json_library()
Poi verifica la firma in un plug nella pipeline dei webhook, con un confronto a tempo costante:
defmodule MyAppWeb.Plugs.VerifyWhatsAppSignature do
import Plug.Conn
def init(opts), do: opts
def call(conn, _opts) do
secret = Application.fetch_env!(:my_app, :whatsapp_app_secret)
raw_body = IO.iodata_to_binary(conn.assigns[:raw_body] || "")
digest = :crypto.mac(:hmac, :sha256, secret, raw_body) |> Base.encode16(case: :lower)
with [signature] <- get_req_header(conn, "x-hub-signature-256"),
true <- Plug.Crypto.secure_compare("sha256=" <> digest, String.downcase(signature)) do
conn
else
_ -> conn |> send_resp(401, "invalid signature") |> halt()
end
end
end
Se le firme falliscono in produzione ma passano in locale, cerca qualcosa tra Meta e la tua applicazione che riscrive il body: un proxy che decomprime o ricodifica il JSON, o un middleware che fa il parsing del body prima che lo veda il tuo reader.
Rispondere subito, elaborare dopo
Meta ritenta le consegne che falliscono o vanno in timeout, con frequenza decrescente, fino a diversi giorni secondo la sua documentazione. È un bene per la durabilità e un male per il carico: se il tuo handler è lento perché una dipendenza a valle è lenta, le richieste vanno in timeout, Meta ritenta, e i retry arrivano mentre sei ancora lento. La soluzione è rendere il percorso della richiesta banalmente economico.
Con Oban, la tabella dei job è un’ottima inbox. Inserisci un job per ogni POST (non uno per messaggio) e rispondi:
defmodule MyAppWeb.WhatsAppWebhookController do
use MyAppWeb, :controller
alias MyApp.Workers.ProcessWhatsAppWebhook
def verify(conn, %{"hub.mode" => "subscribe", "hub.verify_token" => token, "hub.challenge" => challenge}) do
if Plug.Crypto.secure_compare(token, Application.fetch_env!(:my_app, :whatsapp_verify_token)) do
send_resp(conn, 200, challenge)
else
send_resp(conn, 403, "")
end
end
def create(conn, _params) do
case Oban.insert(ProcessWhatsAppWebhook.new(%{"payload" => conn.body_params})) do
{:ok, _job} -> send_resp(conn, 200, "")
{:error, _reason} -> send_resp(conn, 500, "")
end
end
end
Il 500 è voluto: se non sei riuscito a salvare l’evento, vuoi che Meta lo mandi di nuovo. Qualsiasi cosa possa fallire dopo il salvataggio è un tuo problema da ritentare, non di Meta.
Il job poi distribuisce il payload:
defmodule MyApp.Workers.ProcessWhatsAppWebhook do
use Oban.Worker, queue: :whatsapp_webhooks, max_attempts: 10
alias MyApp.WhatsApp.{Inbound, Statuses}
@impl Oban.Worker
def perform(%Oban.Job{args: %{"payload" => payload}}) do
for %{"changes" => changes} <- Map.get(payload, "entry", []),
%{"value" => value} <- changes do
phone_number_id = get_in(value, ["metadata", "phone_number_id"])
Enum.each(Map.get(value, "messages", []), &Inbound.handle(phone_number_id, &1))
Statuses.apply_batch(Map.get(value, "statuses", []))
end
:ok
end
end
Deduplicare: i webhook sono at-least-once
Riceverai lo stesso messaggio in arrivo più di una volta: dopo un retry, dopo un timeout in cui l’avevi elaborato ma Meta non ha visto il 200, e ogni tanto senza un motivo apparente. Ogni messaggio in arrivo ha un id (la stringa wamid.…). Mettici sopra un indice unique e lascia che sia il database a fare da arbitro:
def handle(phone_number_id, %{"id" => wamid, "from" => wa_id, "timestamp" => ts} = message) do
attrs = %{
wamid: wamid,
wa_id: wa_id,
phone_number_id: phone_number_id,
payload: message,
sent_at: DateTime.from_unix!(String.to_integer(ts))
}
case Repo.insert(InboundMessage.changeset(attrs), on_conflict: :nothing, conflict_target: :wamid) do
# Con on_conflict: :nothing, un duplicato torna indietro senza chiave primaria.
{:ok, %InboundMessage{id: nil}} -> :duplicate
{:ok, inbound} -> ContactServer.dispatch(inbound)
end
end
Una sottigliezza: se l’insert va a buon fine e poi l’elaborazione va in crash, il job ritentato vede un duplicato e lo salta, e il messaggio non viene mai gestito. Aggiungi una colonna processed_at, impostala quando la logica della conversazione ha finito, e fai recuperare a un job periodico le righe salvate ma mai elaborate. “Visto” e “gestito” sono stati diversi.
Se la tabella dei messaggi è partizionata, un indice unique sul solo id del messaggio non è possibile (gli indici unique devono includere la chiave di partizione). Lo risolve una piccola tabella di deduplicazione con wamid come chiave; ne parlo nell’articolo sul partizionamento di Postgres.
Mantenere in ordine ogni conversazione
Un utente che invia in rapida successione “ciao”, “ho bisogno di aiuto” e una foto si aspetta che il bot li veda in quell’ordine. Due worker concorrenti che elaborano due messaggi dello stesso utente prima o poi sbaglieranno, e una macchina a stati di un chatbot che elabora due input insieme è una race condition con un’interfaccia utente.
Su un singolo nodo, l’approccio corretto più semplice è un processo per contatto:
defmodule MyApp.Conversations.ContactServer do
use GenServer, restart: :transient
@idle_timeout :timer.minutes(5)
def dispatch(%{wa_id: wa_id} = inbound) do
pid =
case DynamicSupervisor.start_child(MyApp.ContactSupervisor, {__MODULE__, wa_id}) do
{:ok, pid} -> pid
{:error, {:already_started, pid}} -> pid
end
GenServer.cast(pid, {:inbound, inbound})
end
def start_link(wa_id), do: GenServer.start_link(__MODULE__, wa_id, name: via(wa_id))
defp via(wa_id), do: {:via, Registry, {MyApp.ContactRegistry, wa_id}}
@impl true
def init(wa_id), do: {:ok, %{wa_id: wa_id}, @idle_timeout}
@impl true
def handle_cast({:inbound, inbound}, state) do
:ok = MyApp.Bot.handle_inbound(state.wa_id, inbound)
MyApp.WhatsApp.Inbound.mark_processed(inbound)
{:noreply, state, @idle_timeout}
end
@impl true
def handle_info(:timeout, state), do: {:stop, :normal, state}
end
I messaggi di un contatto vengono gestiti rigorosamente uno dopo l’altro; contatti diversi girano in parallelo; i processi inattivi spariscono dopo qualche minuto. In un cluster ti serve che ogni evento di un contatto arrivi nello stesso posto: fai l’hash dell’id del contatto verso una partizione posseduta da un solo nodo, usa un registry distribuito, oppure una coda di job partizionata per contatto con concorrenza uno per partizione (Oban Pro supporta limiti partizionati proprio per casi come questo). Scegline uno e rendilo esplicito, perché “di solito lo stesso nodo” non è un ordinamento.
Nemmeno l’ordine di arrivo è garantito. Quando i messaggi arrivano a uno o due secondi l’uno dall’altro, trattenerli per un attimo e ordinarli per il timestamp del payload aiuta, ma quel timestamp ha una risoluzione di un secondo, quindi a parità di valore si torna all’ordine di arrivo.
I webhook di stato arrivano fuori ordine
Per i messaggi che invii riceverai stati come sent, delivered, read e failed, e non necessariamente in quest’ordine: read può arrivare prima di delivered. Tratta lo stato come monotono. Assegna un rango agli stati e applica un aggiornamento solo se fa avanzare il messaggio, e applicali in batch, perché gli stati arrivano a un multiplo del tuo ritmo di invio. Descrivo l’SQL per farlo nell’articolo sugli invii massivi.
Se imposti biz_opaque_callback_data al momento dell’invio, lo ritrovi nei webhook di stato di quel messaggio, il che è comodo per collegare gli stati ai tuoi record senza una ricerca in più.
La finestra di assistenza clienti di 24 ore
Quando un utente ti scrive si apre una finestra di assistenza clienti di 24 ore, durante la quale puoi rispondere con messaggi liberi. Fuori dalla finestra puoi inviare solo messaggi template approvati, e un invio libero fallisce (errore 131047). Quando l’utente risponde, si apre una nuova finestra.
Non scoprirlo dagli errori. Salva il timestamp dell’ultimo messaggio in arrivo di ogni contatto (è già nella tua tabella dei messaggi in arrivo) e controllalo prima di inviare. Lascia un margine: una risposta che il bot genera a 23 ore e 59 minuti potrebbe arrivare dopo la chiusura della finestra. Progetta i flussi in modo che tutto ciò che potrebbe partire più tardi, come promemoria, follow-up o un operatore che risponde la mattina dopo, abbia un template pronto.
Media: scaricali subito
I messaggi multimediali in arrivo non contengono il file. Contengono un id del media (più tipo MIME e un hash). Per ottenere i byte servono due richieste autenticate: recuperi i metadati del media tramite l’id, che restituiscono un URL a vita breve (Meta lo documenta come valido per cinque minuti), poi scarichi da quell’URL con lo stesso bearer token.
def download_media(media_id) do
version = Application.fetch_env!(:my_app, :graph_api_version)
auth = [auth: {:bearer, Application.fetch_env!(:my_app, :whatsapp_token)}]
with {:ok, %{status: 200, body: %{"url" => url, "mime_type" => mime}}} <-
Req.get("https://graph.facebook.com/#{version}/#{media_id}", auth),
{:ok, %{status: 200, body: bytes}} <- Req.get(url, auth ++ [decode_body: false]) do
{:ok, bytes, mime}
end
end
Fallo in un job in background dedicato, poco dopo l’arrivo del webhook, e copia il file nel tuo storage. Non salvare l’URL: sarà scaduto quando qualcuno ci cliccherà sopra. E non farlo nemmeno durante la richiesta: il download di un video pesante è esattamente il passaggio lento che manda in timeout i webhook. Limiti di dimensione e scansione antivirus vanno nello stesso job.
Cosa si rompe su larga scala
- Picchi di risposte dopo le campagne. Invii un messaggio a un milione di persone e una parte di loro risponde nel giro di pochi minuti. Il traffico dei webhook dipende dal traffico in uscita. Fai load test del percorso in ingresso al picco che genererà la tua campagna più grande.
- Handler lenti che diventano tempeste di retry. Qualsiasi lavoro sincrono nel percorso della richiesta prima o poi rallenta, e la lentezza diventa timeout, retry e duplicati. Limita il percorso della richiesta a verificare, salvare, rispondere.
- Esaurimento del pool di connessioni al database. Le richieste dei webhook, la coda dei job e i worker delle conversazioni competono tutti per lo stesso pool. Dimensionalo e tieni d’occhio il tempo di attesa nel pool, non solo il tempo delle query.
- Contatti caldi. Un singolo numero bloccato in un loop (spesso un altro bot) può inondare il processo di un contatto. Aggiungi rate limit per contatto e rilevamento dei loop.
- Deploy. I riavvii progressivi interrompono l’elaborazione in corso. Va bene se tutto è salvato prima del 200 e l’elaborazione è idempotente; se non lo è, è perdita di dati.
- Loggare i payload interi. Contengono numeri di telefono e contenuto dei messaggi. Soprattutto nei servizi sanitari, trattali come dati sensibili.
Costruire sulla WhatsApp Business Platform
Il receiver dei webhook è la porta d’ingresso di qualsiasi integrazione con WhatsApp, e la regola più importante è anche la più semplice: verifica, salva, rispondi, e fai tutto il resto dopo. Se stai costruendo sulla WhatsApp Cloud API, o la tua integrazione attuale fatica sotto carico, aiuto i team a costruire sulla WhatsApp Business Platform.