Script icon
Process icon
If icon
SlackIncomingWebhook icon
Fail icon
Get icon
TransformValue icon
Loop icon
Subflow icon
Set icon
Log icon
Webhook icon

Multi-Tenant Webhook Gateway with Dedupe and Subflow Fanout

Harden webhooks in Kestra: HMAC gate, Redis dedupe, JSONata reshape, and Loop fanout to matched subflows with Slack alerts for forgeries and dead events.

Categories
Infrastructure

Every integration team eventually builds the same brittle thing: a webhook URL per provider, signature checks copy-pasted and half-remembered, retries that reprocess the same order twice, and fanout logic buried in one service that nobody wants to touch. This blueprint replaces that with one gateway flow. Signatures are verified in constant time before anything else runs, Redis remembers which deliveries already happened so provider retries are free, JSONata normalizes the event into a stable envelope, and a Loop fans each verified event out to its matched subflows - individually, so one broken consumer never stalls the rest.

How it works

  1. inbound_webhook (io.kestra.plugin.core.trigger.Webhook) is the single shared endpoint; headers arrive as {{ trigger.headers }} and the raw body as {{ trigger.body }}.
  2. verify_and_route (io.kestra.plugin.scripts.python.Script on the Process runner) recomputes HMAC-SHA256(secret, body) and compares with hmac.compare_digest - constant time, never == - then extracts tenant, event type, event id (falling back to a payload hash so anonymous events still dedupe) and matches the event against the routes table, including the "*" catch-all. It emits everything through Kestra's outputs protocol.
  3. enforce_signature (io.kestra.plugin.core.flow.If) sends forgeries to a Slack alert and a failed execution before any routing happens.
  4. route_gate (second If) separates "verified but nobody consumes this" - logged and failed loudly, because a silent dead event is a wiring bug waiting for an outage.
  5. seen_check (io.kestra.plugin.redis.string.Get, failedOnMissing: false) reads the dedupe key; dedupe_gate (third If) skips replays with a one-line log. First deliveries pass through reshape (io.kestra.plugin.transform.jsonata.TransformValue) into the standard envelope and into fanout (io.kestra.plugin.core.flow.Loop), where dispatch starts each matched route as a io.kestra.plugin.core.flow.Subflow and mark_seen (io.kestra.plugin.redis.string.Set) records the delivery for future retries.
  6. announce_fanout posts the fanout summary to Slack; the flow-level errors handler covers infrastructure failures with distinct wording.

What you get

  • One endpoint instead of a webhook per provider - signature policy in one place.
  • Exactly-once dispatch semantics via Redis dedupe keys, safe against provider retries.
  • A routing table as a JSON input: add consumers without touching the flow.
  • Fanout isolation: each consumer is its own subflow execution with its own failure path.
  • tenant, event_type, and routes_matched outputs for observability and billing meters.

Who it's for

  • Platform teams standardizing inbound events (Stripe, GitHub, internal services) behind one gate.
  • Multi-tenant SaaS products where each customer's events must be isolated, deduped, and routed.
  • Integration engineers tired of re-implementing signature checks in every consumer service.

Why orchestrate this with Kestra

Gateway logic embedded in a service dies with that service's deploy cycle, and a cron poller can't react in real time. Kestra gives you a webhook trigger, native Redis and JSONata plugin tasks, subflow fanout with per-consumer execution history, and an execution record per delivery - so "did the tenant get it, and which consumers ran" is a query, not a log archaeology session.

Prerequisites

  • Redis reachable from Kestra (any version; used for dedupe keys only).
  • Target subflows must accept tenant, event_type, and event_id inputs (or delete the inputs block on dispatch).
  • Senders that sign the raw body (GitHub, Stripe style) and share the secret with WEBHOOK_SIGNING_SECRET.

Secrets

  • WEBHOOK_SIGNING_SECRET: shared HMAC-SHA256 signing secret.
  • REDIS_URL: Redis connection URL, e.g. redis://host:6379/0.
  • SLACK_WEBHOOK_URL: Slack incoming webhook for forgery, fanout, and failure messages.

Quick start

  1. Add the three secrets above and replace the webhook key.
  2. Edit routes to your real consumers.
  3. Send a signed test event (compute sha256= HMAC over the raw body) and watch verify → dedupe → fanout.
  4. Replay the same event - expect log_duplicate and no dispatches.
  5. Send a bad signature - expect the Slack intrusion alert and a failed execution.

How to extend

  • Per-tenant secrets: extend verify_and_route to load secret('WEBHOOK_' + tenant.upper()) for customer-specific signing.
  • Replay buffer: push rejected or unrouted events to a Redis list (redis.list.ListPush) and add a schedule that re-dispatches them once routes exist.
  • Backpressure: set wait: true on dispatch and add allowFailure semantics per consumer when ordering matters.
  • Meter usage: append tenant and routes_matched to a jdbc table for per-tenant event billing.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.