id: multi-tenant-webhook-gateway-fanout
namespace: company.team
description: |
One hardened webhook endpoint for every tenant: HMAC-SHA256 signature gate,
JSONata reshape, Redis-backed dedupe, and a Loop that fans each verified event
out to its matched subflows - duplicates and dead events are logged, never reprocessed.
inputs:
- id: routes
type: JSON
displayName: Event routing table
description: Map event_type to target subflow. Use "*" as the catch-all event.
defaults:
- event: order.created
namespace: company.team
flow: handle-order-event
- event: "*"
namespace: company.team
flow: archive-webhook-event
- id: signature_header
type: STRING
displayName: Signature header name
defaults: X-Hub-Signature-256
- id: sample_event
type: JSON
displayName: Sample event (manual runs)
description: Used when the flow runs without a webhook body.
defaults:
tenant: acme
event_type: order.created
event_id: evt_0001
data:
order_id: 1001
total: 49.99
tasks:
- id: verify_and_route
type: io.kestra.plugin.scripts.python.Script
description: Constant-time HMAC verification plus tenant, event, and route
matching in one pass.
taskRunner:
type: io.kestra.plugin.core.runner.Process
env:
SIGNING_SECRET: "{{ secret('WEBHOOK_SIGNING_SECRET') }}"
SIGNATURE_HEADER: "{{ inputs.signature_header }}"
HEADERS_JSON: "{{ trigger.headers | toJson }}"
BODY: "{{ trigger.body ?? inputs.sample_event | json }}"
ROUTES_JSON: "{{ inputs.routes | json }}"
script: |
import hashlib
import hmac
import json
import os
import uuid
body = os.environ.get("BODY", "")
secret = os.environ["SIGNING_SECRET"]
header_name = os.environ.get("SIGNATURE_HEADER", "X-Hub-Signature-256")
headers = {k.lower(): v for k, v in json.loads(os.environ.get("HEADERS_JSON", "{}")).items()}
routes = json.loads(os.environ.get("ROUTES_JSON", "[]"))
provided = (headers.get(header_name.lower()) or "").strip()
if not provided:
valid, reason = False, "missing header " + header_name
else:
candidate = provided[7:] if provided.lower().startswith("sha256=") else provided
expected = hmac.new(secret.encode("utf-8"), body.encode("utf-8"), hashlib.sha256).hexdigest()
valid = hmac.compare_digest(candidate.lower(), expected)
reason = "signature verified" if valid else "signature mismatch"
try:
payload = json.loads(body)
except Exception:
payload = {}
tenant = payload.get("tenant") or headers.get("x-tenant") or "default"
event_type = payload.get("event_type") or "unknown"
event_id = payload.get("event_id") or ("sha-" + hashlib.sha256(body.encode("utf-8")).hexdigest()[:16])
matched = []
if valid:
for route in routes:
if route.get("event") == event_type or route.get("event") == "*":
matched.append({
"namespace": route.get("namespace", "company.team"),
"flow": route.get("flow"),
"event_type": event_type,
"tenant": tenant,
"event_id": event_id,
})
if len(matched) >= 5:
break
print("webhook " + reason + " tenant=" + tenant + " event=" + event_type)
print("::" + json.dumps({"outputs": {
"valid": valid,
"reason": reason,
"tenant": tenant,
"event_type": event_type,
"event_id": event_id,
"dedupe_key": "webhook:gw:seen:" + tenant + ":" + str(event_id),
"matched": matched,
}}) + "::")
- id: enforce_signature
type: io.kestra.plugin.core.flow.If
description: Forgeries alert Slack and fail before any routing, dedupe, or
fanout happens.
condition: "{{ not (outputs.verify_and_route.vars.valid) }}"
then:
- id: alert_intrusion_attempt
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Report which check failed plus the execution to inspect.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: |
:no_entry: Gateway rejected a webhook: {{ outputs.verify_and_route.vars.reason }} (tenant {{ outputs.verify_and_route.vars.tenant }}, event {{ outputs.verify_and_route.vars.event_type }}). Execution {{ execution.id }}.
- id: reject_request
type: io.kestra.plugin.core.execution.Fail
description: Turn the execution red so monitoring sees rejected traffic, not
silent drops.
errorMessage: "Gateway rejected webhook: {{ outputs.verify_and_route.vars.reason }}"
- id: route_gate
type: io.kestra.plugin.core.flow.If
description: Verified events with at least one matching route proceed to dedupe
and fanout; unrouted events are logged and failed loudly.
condition: "{{ (outputs.verify_and_route.vars.matched | length) > 0 }}"
then:
- id: seen_check
type: io.kestra.plugin.redis.string.Get
description: Look up the dedupe key so replayed deliveries are skipped, not
re-dispatched.
url: "{{ secret('REDIS_URL') }}"
key: "{{ outputs.verify_and_route.vars.dedupe_key }}"
failedOnMissing: false
- id: dedupe_gate
type: io.kestra.plugin.core.flow.If
description: First delivery fans out; a non-empty seen value means this event
already ran.
condition: "{{ outputs.seen_check.value == null or outputs.seen_check.value ==
'' }}"
then:
- id: reshape
type: io.kestra.plugin.transform.jsonata.TransformValue
description: Normalize the verified event into the envelope every target subflow
receives.
from: "{{ trigger.body ?? inputs.sample_event }}"
expression: |
{
"tenant": tenant,
"event_type": event_type,
"event_id": event_id,
"forwarded_at": $now()
}
- id: fanout
type: io.kestra.plugin.core.flow.Loop
description: Fire every matched route as its own subflow - one bad tenant flow
never blocks the others.
values: "{{ outputs.verify_and_route.vars.matched }}"
tasks:
- id: dispatch
type: io.kestra.plugin.core.flow.Subflow
description: Start the target flow, passing the event envelope as inputs.
namespace: "{{ item.value.namespace }}"
flowId: "{{ item.value.flow }}"
wait: false
inputs:
tenant: "{{ item.value.tenant }}"
event_type: "{{ item.value.event_type }}"
event_id: "{{ item.value.event_id }}"
- id: mark_seen
type: io.kestra.plugin.redis.string.Set
description: Record the delivery so webhook retries of the same event are
deduplicated.
url: "{{ secret('REDIS_URL') }}"
key: "{{ outputs.verify_and_route.vars.dedupe_key }}"
value: "{{ execution.id }}"
- id: announce_fanout
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: |
:satellite: Gateway fanned out {{ outputs.verify_and_route.vars.event_type }} (tenant {{ outputs.verify_and_route.vars.tenant }}, id {{ outputs.verify_and_route.vars.event_id }}) to {{ outputs.verify_and_route.vars.matched | length }} flow(s) - execution {{ execution.id }}.
else:
- id: log_duplicate
type: io.kestra.plugin.core.log.Log
message: "Duplicate delivery skipped: {{
outputs.verify_and_route.vars.dedupe_key }} was already processed
(first seen as {{ outputs.seen_check.value }})."
else:
- id: log_unrouted
type: io.kestra.plugin.core.log.Log
message: "Verified event {{ outputs.verify_and_route.vars.event_type }} matched
no route in inputs.routes."
- id: fail_unrouted
type: io.kestra.plugin.core.execution.Fail
description: A verified event nobody consumes is a wiring bug - fail loudly so
it gets a route.
errorMessage: "Event {{ outputs.verify_and_route.vars.event_type }} (tenant {{
outputs.verify_and_route.vars.tenant }}) matched no route"
outputs:
- id: tenant
type: STRING
value: "{{ outputs.verify_and_route.vars.tenant }}"
- id: event_type
type: STRING
value: "{{ outputs.verify_and_route.vars.event_type }}"
- id: routes_matched
type: INT
value: "{{ outputs.verify_and_route.vars.matched | length }}"
errors:
- id: alert_on_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: ":rotating_light: Webhook gateway {{ flow.id }} FAILED (execution
{{ execution.id }}): {{ error.message }}"
triggers:
- id: inbound_webhook
type: io.kestra.plugin.core.trigger.Webhook
description: The single endpoint all tenants and providers post to - headers
carry the signature, body carries the event. Replace the key with your own
value.
key: replace-with-a-random-key