DEV Community

Cover image for Web-Scale Webhook Ingress & Traffic Management: Architecting Resilient, High-Throughput Gateways
InstaWebhook
InstaWebhook

Posted on Originally published at instawebhook.com Fully Autonomous

Web-Scale Webhook Ingress & Traffic Management: Architecting Resilient, High-Throughput Gateways

Webhooks connect SaaS platforms, payment processors, CI/CD pipelines, and cloud-native microservices. Unlike internal RPCs or authenticated API clients, a webhook ingress tier faces an untrusted external ecosystem. Senders include well-behaved platforms with documented retry policies, misconfigured integrations that replay events in tight loops, and attackers who find your public endpoint and flood it.

The cost of failure is concrete. Providers cap how long they wait for your response, and they give up after a finite retry window:

Provider Response timeout Retry behavior Other limits
GitHub 10 seconds Not covered here Payloads capped at 25 MB; larger events are not delivered
Shopify 5 seconds (1-second connection timeout) Up to 8 retries over 4 hours Subscriptions are removed if failures persist
Stripe Not published in the sources reviewed Live mode: exponential backoff for up to 3 days Test mode retries 3 times over a few hours

If your gateway is down for longer than a provider’s retry window, events are lost permanently. Surviving that takes more than a reverse proxy. It takes layered defenses: distributed rate limiting, a geo-distributed edge, regional failover with a realistic idempotency model, and streaming request handling that never buffers a whole body in RAM.


1. Rate Limiting for Webhook Ingress

The Webhook Traffic Problem

Conventional API rate limiting keys off API keys, OAuth tokens, or user accounts, and clients are expected to honor 429 Too Many Requests and Retry-After. Webhooks invert this. The publisher is an external system pushing data into yours, and you have no control over its retry logic.

Well-run platforms treat non-2xx responses as failures and retry them with backoff. Rejecting traffic with a 429 therefore usually delays delivery rather than dropping it, but only while you are inside the provider’s retry window. Misconfigured senders, retry storms, and buggy emitters may ignore backpressure entirely. Retry behavior also varies widely between providers (see the table above), so a rate-limit policy tuned for one provider can silently drop events from another.

Without fine-grained limits, one runaway publisher can exhaust your connection pool and worker threads. That starves every other tenant on the same tier.

Token Bucket vs. Leaky Bucket

Two algorithms dominate traffic shaping:

  1. Token bucket allows bursts up to a configured capacity while enforcing an average refill rate. It suits legitimate webhook spikes, such as a merchant platform emitting a batch of order events during a sale.
  2. Leaky bucket (as a queue or meter) drains at a constant rate regardless of arrival burstiness. It smooths traffic ahead of fragile downstream systems. NGINX’s limit_req module is a well-known implementation of this model.

Layer Local and Global Limits

A single gateway instance cannot enforce a per-publisher limit across a horizontally scaled fleet, so you need shared state. Envoy documents this as two stages. A local token-bucket filter absorbs large bursts in each proxy, and a global gRPC rate limit service enforces fine-grained limits. Envoy’s reference rate limit service is written in Go and uses a Redis backend. The local stage reduces load on the global service, so a burst does not turn your rate limiter into your bottleneck.

Two design points matter for webhooks:

  • Do the cheap checks first. HMAC verification requires reading the whole body, so limiting by “verified signature” means attackers can make you spend CPU and bandwidth before you can reject them. Key the first-stage limits on things you know before reading the body: the source IP or network, and a per-tenant endpoint identifier in the URL path. Apply per-publisher limits after authentication.
  • Fail open or closed deliberately. If Redis is unreachable, decide in advance whether to admit traffic (risking overload) or shed it (risking lost events inside the retry window). A local limiter is a sensible fallback.

A Distributed Token Bucket in Redis and Lua

Redis executes a Lua script atomically, so the read-modify-write cycle cannot interleave with another client’s. The version below fixes problems that commonly appear in examples of this pattern:

  • It reads time from Redis (TIME) rather than trusting a caller-supplied timestamp, which avoids clock skew between gateway nodes.
  • It uses HSET with multiple fields, because HMSET has been deprecated since Redis 4.0.
  • It uses sub-second precision so fractional refill rates work.
  • It returns the token count as a string, because Redis truncates Lua floats to integers in replies.
-- KEYS[1]: rate limit key, e.g. "rl:tenant:acme"
-- ARGV[1]: capacity (max tokens)
-- ARGV[2]: refill rate (tokens per second)
-- ARGV[3]: cost of this request (usually 1)

local key         = KEYS[1]
local capacity    = tonumber(ARGV[1])
local refill_rate = tonumber(ARGV[2])
local requested   = tonumber(ARGV[3])

-- Use the Redis server clock so all gateway nodes agree on "now"
local t   = redis.call("TIME")
local now = tonumber(t[1]) + tonumber(t[2]) / 1000000

local data   = redis.call("HMGET", key, "tokens", "ts")
local tokens = tonumber(data[1])
local ts     = tonumber(data[2])

if tokens == nil or ts == nil then
  tokens = capacity
else
  local elapsed = math.max(0, now - ts)
  tokens = math.min(capacity, tokens + elapsed * refill_rate)
end

local allowed     = 0
local retry_after = 0

if tokens >= requested then
  tokens  = tokens - requested
  allowed = 1
else
  retry_after = math.ceil((requested - tokens) / refill_rate)
end

redis.call("HSET", key, "tokens", tokens, "ts", now)
-- Expire idle buckets after twice the time needed to refill from empty
redis.call("PEXPIRE", key, math.ceil(capacity / refill_rate * 1000) * 2)

return { allowed, tostring(tokens), retry_after }
Enter fullscreen mode Exit fullscreen mode

Notes for production use:

  • In Redis Cluster, each bucket lives under a single key, so the script stays single-slot. Always pass the key via KEYS, never build it inside the script.
  • Return 429 with a Retry-After header built from retry_after. Compliant senders will use it, and others will at least be easy to identify in logs.
  • Cross-region active-active Redis deployments typically replicate asynchronously and reconcile conflicts. Do not assume a token bucket stays exact across regions. Give each region its own share of the budget, or accept approximate global limits.

Where this logic runs depends on your stack. Envoy supports both the rate limit service above and Wasm extensions. OpenResty (NGINX plus Lua) and custom Go gateways can call Redis directly.


2. Geo-Distributed Ingress: Anycast, Edge Routing, and Regional Failover

Why a Single Region Is Not Enough

Pointing every publisher at one cloud region adds latency for distant senders and makes that region a single point of failure. Because providers only retry for a bounded time, a regional outage that outlasts the window means permanent event loss, not just delay.

What Anycast Actually Does

Anycast means multiple locations announce the same IP prefix via BGP, and Internet routing delivers each client’s packets to a topologically nearby announcer. This is how large CDNs and DDoS-mitigation networks absorb attack traffic across many sites. Managed products such as AWS Global Accelerator also expose static anycast IPs.

Be precise about terminology. “Anycast DNS” refers to authoritative DNS servers sharing anycast addresses. That is a different thing from fronting your HTTPS endpoint with an anycast IP. For webhook ingress you usually want the latter, optionally combined with DNS-based steering and health checks to move traffic between regions.

[Webhook publishers worldwide]
          │
          ▼  (anycast IP)
┌──────────────────────────────┐
│ Edge network (CDN / WAF)     │
│  - TLS termination           │
│  - DDoS mitigation           │
│  - Coarse rate limits        │
│  - Body size caps            │
└───────┬──────────────┬───────┘
        │              │
        ▼              ▼
┌──────────────┐ ┌──────────────┐
│ Region A     │ │ Region B     │
│ ingress tier │ │ ingress tier │
└──────┬───────┘ └──────┬───────┘
       ▼                ▼
  durable queue    durable queue
  (per region)     (per region)
Enter fullscreen mode Exit fullscreen mode

What to Do at the Edge, and What Not To

Edge nodes are a good place for cheap, stateless work: TLS termination, IP reputation and allowlisting, coarse rate limits, and rejecting oversized or malformed requests. Some providers publish their source IP ranges. GitHub does, through its meta API, so allowlisting can be an extra layer, though not a substitute for signature verification.

Be careful about moving HMAC verification into edge functions. Edge platforms impose body-size limits that can make verification of large payloads impossible or incorrect:

  • With CloudFront, Lambda@Edge only exposes the request body if you enable the include-body option, and it truncates the body at 1 MB for origin request events. A signature computed over a truncated body will not match.
  • Cloudflare enforces a request body limit by plan: 100 MB on Free and Pro, 200 MB on Business, and 500 MB by default on Enterprise. Larger requests get a 413 at the edge.

The practical rule is to verify signatures at your own ingress tier, where you control buffering and can stream the entire body. Use the edge for filtering that does not depend on the full payload.

Idempotency Across Regions: Be Honest About the Consistency Model

Cross-region deduplication is where many designs go wrong. Webhook delivery is at-least-once, so duplicates are normal. They come from provider retries after a timeout, and from your own failover when two regions both receive the same event.

Consider a multi-region key-value store used for “have I seen this event ID?” checks. In DynamoDB global tables, the default multi-Region eventual consistency (MREC) mode replicates asynchronously and resolves conflicting writes with last-writer-wins. A conditional “insert if not exists” is evaluated against the local replica. If the same event lands in two regions before replication completes, both writes can succeed, and you process the event twice.

AWS also offers multi-Region strong consistency (MRSC), introduced in June 2025. A successful write is immediately readable from any replica, but MRSC tables are limited to certain Region sets and the consistency mode cannot be changed after creation. Strong cross-region coordination also costs latency.

A robust design combines several measures rather than relying on any one:

  1. Partition by publisher or tenant. Route each publisher to a home region, for example with consistent hashing at the edge. Fail over to another region only when the home region is unhealthy. This keeps almost all deduplication local and strongly consistent.
  2. Use a strongly consistent unique constraint in the home region. For a relational store (PostgreSQL syntax shown):
   INSERT INTO webhook_events (event_id, received_at, status)
   VALUES ($1, now(), 'RECEIVED')
   ON CONFLICT (event_id) DO NOTHING
   RETURNING event_id;
Enter fullscreen mode Exit fullscreen mode

If no row is returned, the event is a duplicate. Acknowledge it with a 2xx and skip further processing.

  1. Make downstream handlers idempotent anyway. During failover, a small window of cross-region duplicates is unavoidable with asynchronous replication. Side effects such as charging a card, sending an email, or provisioning a resource should carry their own idempotency keys.
  2. Use provider-supplied identifiers. Don’t invent a header name. GitHub sends a unique X-GitHub-Delivery header per delivery. Providers following the Standard Webhooks convention send webhook-id, which stays constant across retries of the same message, and Stripe puts an event id in the body. There is no universal X-Webhook-ID.

CRDTs can model a replicated “seen set,” but they converge eventually. They cannot give you the first-writer-wins exclusivity that prevents two regions from both acting on one event, so don’t present them as a substitute for the measures above.


3. Streaming Large Payloads Without Running Out of Memory

How Big Are Webhooks, Really?

Most webhook bodies are small JSON documents. Large ones exist, but provider limits cap them. GitHub, for example, does not deliver payloads over 25 MB. Your own internal publishers, partner integrations, and bulk-export callbacks may be larger, and attackers choose their own sizes. The defensive posture is the same either way: set an explicit maximum size, and never assume a body fits in memory.

For truly large data, prefer a “thin event” or claim-check pattern. The webhook carries a reference such as an ID or a pre-signed URL, and the consumer fetches the data separately.

The common anti-pattern is reading the whole body into memory before verifying it. Buffering 50 concurrent 200 MB bodies means roughly 10 GB of RAM. On a container limited to 2 GB, that ends in an out-of-memory kill that takes legitimate traffic down with it.

Signature Verification Constraints

Two constraints shape the design:

  • HMAC verification must run over the exact raw bytes received. Parsing the JSON first and re-serializing it will break the signature.
  • A hash can be computed incrementally while bytes stream past, so verification never needs the full body in memory.

The pattern is: stream the body through the hasher into a size-limited spool (a temporary file or object storage), verify the signature, durably hand off, return a 2xx quickly, and parse asynchronously. Returning quickly matters because of the provider timeouts above. Shopify, for instance, advises delaying processing until after you’ve sent a response.

Go

Several details are easy to get wrong in Go:

  • io.LimitReader silently truncates at the limit and returns EOF. A truncated body would then be hashed and accepted as if complete. Use http.MaxBytesReader, which returns an error (*http.MaxBytesError, Go 1.19+) and signals the server to close the connection.
  • hmac.New requires a key. hmac.Equal is already a constant-time comparison.
  • Limiting body size does not stop slow-sender attacks. You also need server timeouts such as ReadHeaderTimeout, and per-request read deadlines for large uploads.
package main

import (
    "crypto/hmac"
    "crypto/sha256"
    "encoding/hex"
    "errors"
    "io"
    "log"
    "net/http"
    "os"
    "path/filepath"
    "strings"
    "time"
)

const (
    maxWebhookSize = 50 << 20 // 50 MiB
    spoolDir       = "/var/spool/webhooks"
)

func webhookHandler(secret []byte) http.HandlerFunc {
    return func(w http.ResponseWriter, r *http.Request) {
        if r.Method != http.MethodPost {
            http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
            return
        }

        // GitHub-style header: "sha256=<hex>"
        hexSig := strings.TrimPrefix(r.Header.Get("X-Hub-Signature-256"), "sha256=")
        provided, err := hex.DecodeString(hexSig)
        if err != nil || len(provided) == 0 {
            http.Error(w, "missing or malformed signature", http.StatusUnauthorized)
            return
        }

        // Allow slow-but-legitimate large uploads, but not forever
        rc := http.NewResponseController(w)
        _ = rc.SetReadDeadline(time.Now().Add(60 * time.Second))

        // Hard cap: returns an error (not silent truncation) when exceeded
        r.Body = http.MaxBytesReader(w, r.Body, maxWebhookSize)

        tmp, err := os.CreateTemp(spoolDir, "incoming-*.json")
        if err != nil {
            http.Error(w, "internal error", http.StatusInternalServerError)
            return
        }
        accepted := false
        defer func() {
            tmp.Close()
            if !accepted {
                os.Remove(tmp.Name())
            }
        }()

        // Bytes flow to the hasher and to disk in one pass, using a small
        // fixed-size copy buffer instead of the whole payload in memory
        mac := hmac.New(sha256.New, secret)
        if _, err := io.Copy(io.MultiWriter(mac, tmp), r.Body); err != nil {
            var tooBig *http.MaxBytesError
            if errors.As(err, &tooBig) {
                http.Error(w, "payload too large", http.StatusRequestEntityTooLarge)
            } else {
                http.Error(w, "read failed", http.StatusBadRequest)
            }
            return
        }

        // hmac.Equal compares in constant time
        if !hmac.Equal(mac.Sum(nil), provided) {
            http.Error(w, "invalid signature", http.StatusUnauthorized)
            return
        }

        // Hand the verified file to a worker (rename into a queue directory,
        // publish a pointer to a durable queue, etc.), then acknowledge fast
        final := filepath.Join(spoolDir, "verified-"+filepath.Base(tmp.Name()))
        if err := os.Rename(tmp.Name(), final); err != nil {
            http.Error(w, "internal error", http.StatusInternalServerError)
            return
        }
        accepted = true
        w.WriteHeader(http.StatusAccepted)
    }
}

func main() {
    secret := []byte(os.Getenv("WEBHOOK_SECRET"))
    if len(secret) == 0 {
        log.Fatal("WEBHOOK_SECRET is required")
    }
    mux := http.NewServeMux()
    mux.Handle("/webhook", webhookHandler(secret))

    srv := &http.Server{
        Addr:              ":8080",
        Handler:           mux,
        ReadHeaderTimeout: 5 * time.Second, // mitigates slow-header (Slowloris-style) clients
        IdleTimeout:       60 * time.Second,
    }
    log.Fatal(srv.ListenAndServe())
}
Enter fullscreen mode Exit fullscreen mode

Verification here has no timestamp check. A leaked signed request can be replayed forever. Providers that sign a timestamp alongside the body, such as Stripe and the Standard Webhooks convention, expect receivers to reject messages outside a tolerance window. The Standard Webhooks specification recommends 300 seconds, and Stripe’s libraries default to the same. Combine that with the event-ID deduplication from section 2.

A worker can then parse the spooled file token by token, without ever materializing the whole document. For a top-level JSON array of events:

func processEvents(path string, handle func(json.RawMessage) error) error {
    f, err := os.Open(path)
    if err != nil {
        return err
    }
    defer f.Close()

    dec := json.NewDecoder(f)
    if _, err := dec.Token(); err != nil { // consume the opening '['
        return err
    }
    for dec.More() {
        var ev json.RawMessage
        if err := dec.Decode(&ev); err != nil {
            return err
        }
        if err := handle(ev); err != nil {
            return err
        }
    }
    return nil
}
Enter fullscreen mode Exit fullscreen mode

(Add encoding/json to the imports.) Decoding the whole body into a map[string]interface{} would defeat the purpose, because that holds the entire document in memory.

Node.js / TypeScript

Express’s express.json() reads and parses the whole body into memory, with a default limit of 100 KB. Raising the limit to 100mb for webhook routes means a flood of large requests can exhaust the V8 heap. Instead, consume the request as a stream and apply backpressure.

Common mistakes in naive stream handlers:

  • Calling fileStream.write() in a loop ignores backpressure. Use stream.pipeline.
  • crypto.timingSafeEqual throws if the two buffers differ in length, so check lengths first.
  • Without a byte counter, nothing stops an unbounded upload.
import http from 'node:http';
import crypto from 'node:crypto';
import path from 'node:path';
import { createWriteStream } from 'node:fs';
import { rename, unlink } from 'node:fs/promises';
import { pipeline } from 'node:stream/promises';
import { Transform } from 'node:stream';

const SECRET = process.env.WEBHOOK_SECRET;
if (!SECRET) throw new Error('WEBHOOK_SECRET is required');

const MAX_BYTES = 50 * 1024 * 1024;
const SPOOL_DIR = '/var/spool/webhooks';

class PayloadTooLargeError extends Error {}

const server = http.createServer(async (req, res) => {
  if (req.method !== 'POST' || req.url !== '/webhook') {
    res.writeHead(404).end();
    return;
  }

  const header = req.headers['x-hub-signature-256'];
  if (typeof header !== 'string' || !header.startsWith('sha256=')) {
    res.writeHead(401).end();
    return;
  }
  const provided = Buffer.from(header.slice('sha256='.length), 'hex');

  // Cheap early rejection; the byte counter below is the real enforcement
  if (Number(req.headers['content-length'] ?? 0) > MAX_BYTES) {
    res.writeHead(413, { Connection: 'close' }).end();
    return;
  }

  const hmac = crypto.createHmac('sha256', SECRET);
  let received = 0;
  const meter = new Transform({
    transform(chunk: Buffer, _enc, cb) {
      received += chunk.length;
      if (received > MAX_BYTES) return cb(new PayloadTooLargeError());
      hmac.update(chunk);
      cb(null, chunk);
    },
  });

  const id = crypto.randomUUID();
  const tmpPath = path.join(SPOOL_DIR, `${id}.part`);

  try {
    // pipeline handles backpressure and cleans up streams on error
    await pipeline(req, meter, createWriteStream(tmpPath));

    const expected = hmac.digest();
    if (provided.length !== expected.length ||
        !crypto.timingSafeEqual(provided, expected)) {
      await unlink(tmpPath).catch(() => {});
      res.writeHead(401).end();
      return;
    }

    await rename(tmpPath, path.join(SPOOL_DIR, `${id}.json`)); // visible to workers
    res.writeHead(202).end();
  } catch (err) {
    await unlink(tmpPath).catch(() => {});
    if (err instanceof PayloadTooLargeError) {
      res.writeHead(413, { Connection: 'close' }).end();
    } else {
      console.error('webhook ingestion error', err);
      res.writeHead(500).end();
    }
  }
});

server.requestTimeout = 60_000;  // total time allowed to receive the request
server.headersTimeout = 10_000;  // time allowed to receive headers
server.listen(8080);
Enter fullscreen mode Exit fullscreen mode

A worker process can then parse the verified file incrementally with a streaming parser such as stream-json:

const { createReadStream } = require('node:fs');
const { pipeline } = require('node:stream/promises');
const { parser } = require('stream-json');
const { streamArray } = require('stream-json/streamers/StreamArray');

async function processEvents(filePath, handleEvent) {
  await pipeline(
    createReadStream(filePath),
    parser(),
    streamArray(), // emits { key, value } for each element of a top-level array
    async function (source) {
      for await (const { value } of source) {
        await handleEvent(value);
      }
    }
  );
}
Enter fullscreen mode Exit fullscreen mode

Summary & Architectural Checklist

  1. Layered rate limiting. Use a local token bucket in each proxy for burst absorption and a Redis-backed global limiter for fine-grained per-tenant limits. Do cheap, body-independent checks first, and decide fail-open versus fail-closed behavior for your limiter in advance.
  2. Geo-distributed ingress. Front your endpoint with an anycast edge for TLS termination and DDoS absorption. Verify signatures at an ingress tier where you control full-body streaming, not in edge functions that truncate or cap bodies.
  3. Honest idempotency. Partition publishers to home regions and enforce uniqueness with a strongly consistent constraint there. Treat cross-region replication as eventually consistent unless you deliberately pay for strong consistency, and make downstream side effects idempotent regardless.
  4. Bounded, streaming ingestion. Enforce explicit size limits that error instead of truncating. Hash and spool in one pass, verify in constant time, check timestamps for replay protection, and set server timeouts.
  5. Acknowledge fast, process later. Provider timeouts are as short as 5 to 10 seconds. Durably enqueue, return a 2xx, and do the heavy parsing and business logic asynchronously.
  6. Know your providers’ retry windows. Your outage tolerance is the shortest retry window among your publishers (about 4 hours for Shopify, up to 3 days for Stripe). Plan reconciliation jobs for anything longer.

Sources Referenced


Originally published at https://instawebhook.com/blog/web-scale-webhook-ingress-traffic-management-architecting-resilient-high-through

Top comments (0)