Webhooks
@basaltkit/webhooks delivers outbound webhooks: signed payloads, retries with backoff, per-tenant subscriptions, and automatic dispatch from your domain events. It is decoupled from your domain — nothing in your code needs to know an endpoint exists — and from the transport, because delivery is a plain signed POST any receiver can verify.
Mental model
Three pieces, each replaceable:
| Piece | Contract | What it does |
|---|---|---|
| Store | WebhookStore | Where subscriptions live. Answers "who wants invoice.paid for tenant acme?" |
| Deliverer | WebhookDeliverer | Signs the body, validates the URL against SSRF, POSTs it, retries transient failures |
| Manager | WebhookManager (token WEBHOOKS) | Register/list/unregister endpoints, and dispatch(event, data) to every match |
dispatch is the whole flow: the manager asks the store for the endpoints matching the event and the tenant, then hands each one to the deliverer and returns one DeliveryResult per endpoint. Nothing is persisted about the attempt — if you need an audit trail, store the results yourself; if you need delivery to survive a crash, use the outbox (below).
Tenant scoping is anti-widening throughout: an ambient ctx().tenant.id always wins over a caller-supplied tenantId, so client input can never widen or switch the scope.
Quickstart
webhooksPlugin registers a WebhookManager under the WEBHOOKS token. The only near-required option is a default signing secret:
import { createApp } from '@basaltkit/core'
import { webhooksPlugin, WEBHOOKS } from '@basaltkit/webhooks'
const app = await createApp({
plugins: [
webhooksPlugin({ secret: process.env.WEBHOOK_SECRET }),
],
}).boot()
const hooks = app.container.get(WEBHOOKS)
await hooks.register({ url: 'https://customer.example.com/hooks', events: ['invoice.*'] })
await hooks.dispatch('invoice.paid', { id: 'in_1', amount: 42 })The default secret must be at least 16 characters (MIN_WEBHOOK_SECRET_LENGTH; generateWebhookSecret() makes a strong one) — a shorter one refuses to boot. Deliveries are never sent unsigned by default: with no per-endpoint secret and no default secret, the delivery is refused (error: 'no signing secret; refusing unsigned delivery'). allowUnsigned: true is the explicit opt-out.
The default secret signs only tenant-agnostic endpoints. Every tenant-bound endpoint is signed with its own secret — register() generates one when you don't pass it — because a secret shared by every tenant would let one tenant forge webhooks that another tenant's receiver accepts. A tenant-bound endpoint stored without its own secret is refused (allowSharedSecret: true is the opt-out for legacy data).
Managing subscriptions
An endpoint is a subscription: a destination URL plus the event patterns it wants. Register, list and remove them through the manager:
const endpoint = await hooks.register({
url: 'https://customer.example.com/hooks',
events: ['invoice.*', 'user.created'], // patterns: exact, prefix `x.*`, or `*`
tenantId: 'acme', // bound from ctx() when a tenant is in context
secret: 'whsec_acme_...',// optional; generated for tenant endpoints when omitted
active: true, // set false to disable without deleting
})
endpoint.secret // returned ONCE — hand it to the customer now
await hooks.list() // every endpoint (no tenant in context) — secrets redacted
await hooks.list('acme') // only tenant "acme"'s endpoints
await hooks.unregister(endpoint.id)register() returns the stored endpoint including its signing secret — generated as whsec_… (32 random bytes) for a tenant-bound endpoint, or for any endpoint when there is no default secret. list() never returns secrets: each item is a WebhookEndpointView with hasSecret: boolean instead, so a management route can't leak them. To change a secret without breaking the receiver, use rotateSecret() (below).
register() validates the endpoint before storing it, so a subscription that could never be delivered to is refused up front instead of failing on every event: the url must parse as an absolute URL with a scheme the deliverer allows (http:/https:, or your ssrf.allowedSchemes) and a port its port policy allows, a secret you pass must be at least 16 characters, and events must be a non-empty list of non-empty patterns — otherwise it throws WebhookEndpointInvalidError (WEBHOOK_ENDPOINT_INVALID, 400). Whether the host is public is still decided at delivery, where DNS is resolved and the connection pinned. A caller-supplied id that another scope (another tenant, or global vs tenant) already holds throws WebhookEndpointIdInUseError (WEBHOOK_ENDPOINT_ID_IN_USE, 409); every bundled store enforces this in its own write, so two concurrent registrations of the same id can't overwrite each other.
Scoping is anti-widening: inside a request with a tenant in context, register, list, unregister and dispatch are forced to that tenant — a caller-supplied tenantId (which may carry client input) can never widen or switch the scope. The explicit argument and the system-wide behavior above apply only where there is no ambient tenant (jobs, CLI, single-tenant apps). unregister is a no-op — not an error — for an endpoint owned by another tenant. The manager re-verifies ownership itself (the endpoint must appear in that tenant's own list()) before calling store.remove, so even a store whose remove ignores the tenantId argument can't be used for a cross-tenant delete.
Fail-closed when tenancy is active. When tenancyPlugin is registered (its 'tenancy:active' marker), a management call with no tenant at all — no tenant in context and no explicit tenantId — throws WebhookTenantRequiredError (WEBHOOKS_TENANT_REQUIRED) instead of running unscoped. That stops a central route from creating a global endpoint that would receive every tenant's events, or from listing/deleting every tenant's endpoints. Deliberate system operations say so explicitly:
await hooks.register({ url, events }, { system: true }) // global endpoint, on purpose
await hooks.list(undefined, { system: true }) // every tenant's endpoints
await hooks.unregister(id, { tenantId: 'acme' }) // scoped delete off the request path
await hooks.unregister(id, { system: true }) // unscoped delete, on purposeEvent patterns match like this:
'invoice.paid'— that exact event only'invoice.*'— any event starting withinvoice.'*'or'**'— every event
import { matchesEvent } from '@basaltkit/webhooks'
matchesEvent(['invoice.*'], 'invoice.paid') // trueOnly active !== false endpoints receive deliveries; flipping active to false is the reversible way to stop a flapping customer endpoint.
Rotating a signing secret
rotateSecret() gives an endpoint a new secret without breaking its receiver. For a grace window, every delivery is signed with both the new and the old secret: t=…,v1=<new>,v1=<old>. A receiver that still verifies with the old secret keeps working, and it can switch whenever it is ready:
const { secret } = await hooks.rotateSecret(endpoint.id, {
graceSeconds: 7 * 86_400, // default 24 h; at most 30 days; 0 = cut over now
// secret: 'whsec_…', // optional; generated when omitted (min 16 chars)
})
// hand `secret` to the customer — the old one stops signing after 7 daysIt is scoped like unregister: an endpoint outside the ambient (or given) tenant throws WebhookEndpointNotFoundError (WEBHOOK_ENDPOINT_NOT_FOUND, 404), and with tenancy active and no tenant it needs { system: true }. The previous secret is never returned (and list() redacts both). Use graceSeconds: 0 after a leak — the old secret stops signing at once. Re-registering the endpoint under the same id also ends a rotation in progress. An endpoint that signs with the plugin-wide default secret has no own secret to rotate (WebhookEndpointInvalidError): rotate the default in your configuration, or register the endpoint with its own secret.
The rotation lives on the endpoint as previousSecret + previousSecretExpiresAt; the deliverer ignores a previous secret with no expiry, or once the expiry has passed. webhooks-sqlite adds the columns itself. webhooks-prisma needs the two fields in your schema (basalt prisma:sync) before the first rotation.
Dispatching events
dispatch(event, data, tenantId?) finds every subscribed endpoint (the tenant's own plus global ones) and delivers a signed POST to each, returning one DeliveryResult per endpoint:
const results = await hooks.dispatch('invoice.paid', { id: 'in_1', amount: 42 }, 'acme')
// [{ endpointId: '...', ok: true, status: 200, attempts: 1 }]Scoping is fail-closed. A dispatch with no tenant — no tenant in ctx() and no explicit tenantId, as from a scheduler job, a billing webhook or a central route — reaches only tenant-agnostic endpoints, never a tenant-bound one, so one tenant's event data can't fan out to another tenant's endpoints. A deliberate system-wide broadcast opts in explicitly (it is ignored inside a tenant context):
await hooks.dispatch('maintenance.scheduled', { at }, { allTenants: true })The manager re-applies the tenant filter to whatever the store returns, so a custom store that ignores the tenantId argument can't widen delivery.
Each result is { endpointId, ok, status?, attempts, error?, retryable? } — persist it for an audit trail. retryable (on failures) says whether trying again later can help. A delivery that throws unexpectedly (say, a malformed store row) becomes a failed result with error: 'internal delivery error' instead of rejecting the whole dispatch. Deliveries run in parallel — at most dispatchConcurrency (default 16) at a time — and dispatch resolves only when all of them have settled, so an endpoint that eats the full retry budget delays the whole call: dispatch from a job or the outbox rather than inline in a request handler.
Fan-out cap
A single dispatch refuses a scope — one tenant, or the tenant-agnostic endpoints — whose active endpoints subscribed to the event exceed maxEndpointsPerDispatch (default 100). Picking some endpoints would be arbitrary, so the whole scope is refused and none of its endpoints is sent to. Each gets a failed result (attempts: 0, retryable: false, error: 'fan-out cap exceeded: …'), and onFanOutExceeded is called once for that scope (default console.warn). Other scopes in the same dispatch — the global endpoints next to a tenant's, or other tenants in an allTenants broadcast — are unaffected. The outbox treats the refusal as a permanent failure (onPermanentFailure).
webhooksPlugin({
secret,
maxEndpointsPerDispatch: 25, // or false to disable
dispatchConcurrency: 8,
onFanOutExceeded: ({ event, tenantId, endpoints, limit }) =>
metrics.increment('webhooks.fanout_refused', { event, tenantId }),
})What the recipient receives
content-type: application/json
x-basalt-event: invoice.paid
x-basalt-delivery: 5f0c…-uuid
x-basalt-signature: t=1712345678,v1=<hmac-sha256(t.body)>
{"id":"5f0c…-uuid","event":"invoice.paid","endpointId":"ep_…","data":{"id":"in_1","amount":42},"sentAt":"2026-08-07T10:00:00.000Z"}id (also in x-basalt-delivery) is unique per delivery and stable across that delivery's retries — dedupe on it to make a replay within the tolerance window harmless. Through the outbox it is derived from the outbox entry id and the endpoint id, so it also stays the same across outbox retries and restarts; dispatch(event, data, { idempotencyKey }) gives your own jobs the same guarantee. Every attempt is signed with its own timestamp t, so a retry after a long backoff still falls inside the receiver's tolerance; the body (with its id and sentAt) is identical on every attempt. endpointId names the subscription it was signed for. Both are inside the signed body, so neither can be altered without breaking the signature.
Auto-dispatch from domain events
Wire the bus once and matching domain events fan out to subscribed endpoints automatically — tenant-scoped from the request context and fire-and-forget, so the emitter never blocks on HTTP. This requires @basaltkit/events:
import { createApp } from '@basaltkit/core'
import { defineEvent, EVENTS, eventsPlugin } from '@basaltkit/events'
import { webhooksPlugin } from '@basaltkit/webhooks'
const app = await createApp({
plugins: [
eventsPlugin(),
webhooksPlugin({
secret: process.env.WEBHOOK_SECRET,
events: ['invoice.*', 'user.created'], // domain events to forward
}),
],
}).boot()
const InvoicePaid = defineEvent<{ amount: number }>('invoice.paid')
await app.container.get(EVENTS).emit(InvoicePaid, { amount: 42 })
// → delivered to every endpoint subscribed to "invoice.*"WARNING
Auto-dispatch needs eventsPlugin() registered — the plugin declares that dependency when events is non-empty. The tenant comes from ctx().tenant.id; when you emit outside a request (e.g. in a job) there's no tenant in context, so the event reaches only global endpoints. Dispatch manually with an explicit tenantId when you need tenant scoping off the request path.
Fire-and-forget also means silent: the listener does void dispatch(...), so a failed delivery never reaches the emitter and nothing retries after the in-process budget is spent. That is the trade — see the outbox below when losing an event is not acceptable.
Durable integration events (outbox)
Auto-dispatch above is fire-and-forget — a failed delivery or a crash between "committed" and "delivered" loses the event. For guaranteed delivery use the outbox: domain events are first written to a transactional store, then a relay publishes them to subscribers with retries (at-least-once). Requires eventsPlugin.
import { webhooksPlugin, webhookOutboxPlugin, webhookOutboxDispatch, WEBHOOKS } from '@basaltkit/webhooks'
import { eventsPlugin, OUTBOX } from '@basaltkit/events'
createApp({
plugins: [
eventsPlugin(),
webhooksPlugin({ store }), // no `events:` here — the outbox captures them
webhookOutboxPlugin({
events: ['invoice.*', 'user.created'], // patterns to capture (default '**')
// store: new MyDurableOutboxStore(), // durable in production (default in-memory)
intervalMs: 5000, // relay poll; 0 = flush manually via OUTBOX
batchSize: 50, // entries per flush
maxAttempts: 10, // then the entry is left dead
concurrency: 8, // entries delivered in parallel per flush
tenantConcurrency: 4, // max in-flight deliveries per tenant
dispatchTimeoutMs: 10_000, // max wait per entry before moving on
}),
],
})webhookOutboxDispatch re-queues an entry only for a transient failure (network error, timeout, 5xx, 408/429). The retry skips endpoints that already accepted the entry and re-sends the same delivery id to the rest — derived from the entry id and the endpoint id, so it holds across restarts too. Permanent failures (SSRF-blocked URL, redirect, other 4xx, refused signing secret) are reported through onPermanentFailure (default console.warn) and never re-queue the entry: retrying can't fix them and would only re-deliver to healthy endpoints. Delivery stays at-least-once (the skip list lives in memory), so subscribers dedupe on id. An entry recorded with no tenant in context reaches only tenant-agnostic endpoints, like dispatch.
One tenant's failing or hanging endpoint can't stall everyone else, however many events it emits:
- Fair selection. Entries backing off don't occupy the batch (the relay over-fetches past them), and when one tenant's backlog fills a whole page the relay queries again excluding the tenants already seen, then interleaves the batch round-robin by tenant. Each tenant keeps its own
createdAtorder. - Per-tenant cap. A flush delivers up to
concurrencyentries in parallel (default 8), but one tenant never has more thantenantConcurrencydeliveries in flight (defaultceil(concurrency / 2)), across flushes. - Bounded flush. A flush waits at most
dispatchTimeoutMs(default 10 s) on one entry. A slower delivery keeps running detached — not cancelled, not failed, not re-sent meanwhile — and its outcome is recorded when it settles, so the relay keeps ticking for every other tenant.
Set concurrency: 1 for strictly sequential delivery. Resolve the OUTBOX token to relay yourself, e.g. from a queue worker instead of the timer:
await container.get(OUTBOX).flush(webhookOutboxDispatch(container.get(WEBHOOKS)))Back the outbox with a durable OutboxStore (your DB) so nothing is lost across restarts — the whole point of the pattern. See Persistence.
Dead entries and flush failures
Two different failures, handled in two different places:
- A per-entry dispatch failure increments the entry's
attemptsand recordslastError, then backs off (exponential from 1 s, capped at 60 s, tracked per relay process). AftermaxAttempts(default 10) the entry is dead: it stays in the store with itslastErrorand is never flushed again.Outbox'sonDeadcallback fires once — by default it writes toconsole.error, because a silently dropped integration event is the worst possible outcome. - A flush-level failure (the store's
pending()itself throws) is not an entry problem at all.outboxPlugin'sonFlushErrorexists for it.
webhookOutboxPlugin doesn't expose onDead
It forwards only maxAttempts, concurrency, tenantConcurrency and dispatchTimeoutMs to the Outbox it builds, so dead entries go to console.error. A store-level flush failure goes to its onFlushError (default console.error). When you need to page on a dead event, wire the outbox yourself with outboxPlugin from @basaltkit/events and webhookOutboxDispatch as its dispatch:
import { eventsPlugin, outboxPlugin } from '@basaltkit/events'
import {
WebhookDeliverer,
WebhookManager,
webhookOutboxDispatch,
webhooksPlugin,
} from '@basaltkit/webhooks'
const deliverer = new WebhookDeliverer({ secret: process.env.WEBHOOK_SECRET })
const webhooks = new WebhookManager(store, deliverer) // `store` is your WebhookStore
createApp({
plugins: [
eventsPlugin(),
webhooksPlugin({ store, deliverer }), // the WEBHOOKS token gets the same pieces
outboxPlugin({
store: outboxStore,
captureEvents: ['invoice.*', 'user.created'],
intervalMs: 5000,
dispatch: webhookOutboxDispatch(webhooks),
maxAttempts: 10,
onDead: (entry, error) => pager.page(`webhook outbox dead: ${entry.event}`, error),
onFlushError: (error) => logger.error({ error }, 'outbox flush failed'),
}),
],
})Register one of the two — both claim the OUTBOX token.
Signing & verification
Each delivery carries X-Basalt-Signature: t=<unix>,v1=<hmac-sha256(t.body)> — the same scheme Stripe uses. Receivers recompute the HMAC over timestamp.body and compare in constant time, rejecting stale timestamps to prevent replay:
import { verifySignature } from '@basaltkit/webhooks'
// in your receiver, over the RAW request body (not a re-serialized object):
const valid = verifySignature(
req.headers['x-basalt-signature'] as string,
rawBody,
process.env.WEBHOOK_SECRET!,
300, // tolerance in seconds (default) — rejects timestamps older than this
)
if (!valid) return res.status(400).end()Verify over the raw body
The HMAC is computed over the exact bytes sent. If your framework parses JSON and you re-JSON.stringify it, the bytes change and verification fails. Capture the raw body (e.g. Express express.raw()) before parsing.
signPayload(body, secret, timestampSeconds) produces the same header if you need to sign manually. verifySignature returns false for a malformed header, a missing v1, a timestamp outside the tolerance, a mismatched digest, or an empty/unset or shorter-than-16-character secret (so a receiver whose WEBHOOK_SECRET env var is missing fails closed instead of accepting an HMAC computed with an empty key). A receiver can treat it as a single boolean. The one exception is a configuration bug: a toleranceSeconds that is not a finite number ≥ 0 (e.g. Number(process.env.UNSET) → NaN) or a non-finite nowSeconds throws a RangeError, because it would otherwise make every timestamp "fresh" and silently disable replay protection.
Each tenant endpoint has its own secret: a receiver verifies with the secret register() returned for its endpoint, not with the app-wide default.
A header may carry several v1= entries — a sender rotating its secret signs with both the new and the old one (t=…,v1=<new>,v1=<old>), as Stripe does, and as this deliverer does during a rotateSecret() grace window. verifySignature returns true when any of them matches your secret, so a receiver keeps working whether it has already switched secrets or not. Unknown schemes (e.g. v0=) are ignored; a duplicate t is rejected.
Delivery semantics
timeoutMsis one deadline per attempt that covers resolving the host and the request together. A DNS server that never answers fails the attempt (error: 'host resolution timed out', retryable) instead of holding the delivery — or the outbox flush waiting on it — indefinitely. The next attempt resolves again.- Transient failures (
5xx, network errors, timeouts) retry with exponential backoff —500ms,1s,2s, … up tomaxRetries(default3, so four attempts in total). - Client errors (
4xx) are not retried — a wrong URL or auth won't fix itself on retry. The result carrieserror: 'HTTP 404'.408and429are still markedretryable: true, so the outbox tries them again later. - Redirects are refused, not followed: a
3xxends the delivery witherror: 'redirect refused'. Following one would let a compliant public URL bounce the request to an internal address. - Only the status line is read. The response body is discarded and the connection closed as soon as the status is known, so a receiver that trickles an endless body can't hold sockets open.
- Tune the deliverer through the plugin options (they pass straight through to the
WebhookDeliverer):
webhooksPlugin({
secret: process.env.WEBHOOK_SECRET,
maxRetries: 5, // retries after the first attempt (default 3)
backoffMs: 500, // base wait, doubled each attempt (default 500)
timeoutMs: 10_000, // per-attempt deadline, DNS + request (default 10s)
})For durable, distributed retries that survive a restart mid-delivery, drive dispatch() from @basaltkit/queue instead of relying on the in-process retry loop — see Queues & jobs.
The SSRF guard
Endpoint URLs are supplied by customers, so every delivery URL is treated as hostile input. Before the first attempt the deliverer resolves the hostname once and refuses the delivery if the scheme isn't http:/https:, or if any resolved address is loopback, private (10/8, 172.16/12, 192.168/16), link-local (including the 169.254.169.254 cloud-metadata address), CGNAT, IPv6 ULA, documentation (192.0.2/24, 198.51.100/24, 203.0.113/24, 2001:db8::/32, 3fff::/20), benchmarking (198.18/15), the deprecated 6to4 relay anycast (192.88.99/24), multicast, or otherwise reserved (240/4). IPv6 is judged over the parsed address, so every spelling counts: an IPv6 literal that embeds an IPv4 address — IPv4-mapped ([::ffff:127.0.0.1], which URL parsing rewrites to [::ffff:7f00:1]), IPv4-compatible, NAT64 (64:ff9b::/96) or 6to4 (2002::/16) — is judged by that IPv4 address, and Teredo, local-use NAT64, discard and documentation ranges are refused.
Port policy
The public-address check does not help when the target is someone else's exposed Redis. So the port is checked too, before any DNS lookup. A delivery goes only to 80, 443, or a port from 1024 up that is not in DEFAULT_BLOCKED_PORTS. That list covers the ports registered to databases, caches, message brokers, cluster control planes, proxies and remote-admin services: Postgres 5432, MySQL 3306, Redis 6379, memcached 11211, MongoDB 27017, Docker 2375, etcd 2379, kubelet 10250, Kafka 9092, Elasticsearch 9200, Squid 3128, and others. Other privileged ports (22, 25, 110, 445, …) are refused. None of these services receives webhooks, and several speak text protocols that a crafted POST body can drive cross-protocol. Common receiver ports (3000, 8080, 8443, …) are allowed.
webhooksPlugin({ secret, ssrf: { allowedPorts: [443] } }) // exactly these ports
webhooksPlugin({ secret, ssrf: { allowedPorts: [443, 6379] } }) // your own policy
webhooksPlugin({ secret, ssrf: { allowedPorts: 'any' } }) // policy offThe policy applies at register() (WebhookEndpointInvalidError) and on every delivery (port N is not allowed, attempts: 0, retryable: false). It also applies with allowPrivateHosts, where internal services are all the more exposed. Redirects are never followed, so a 3xx can't reach a blocked port either. A malformed allowedPorts throws when the deliverer is built.
The socket is then pinned to the address that was validated, so a hostile authoritative DNS can't return a public IP to the check and an internal IP at connect time (DNS rebinding). The Host header and TLS SNI still carry the original hostname, so vhosts and certificate validation are unaffected.
A blocked URL is a permanent configuration error, not a transient one: the result is { ok: false, attempts: 0, retryable: false, error: 'Refusing to deliver webhook…' } and nothing is retried. When the verdict comes from DNS (the host does not resolve, or resolves to a private address) the error is one generic message — host does not resolve to an allowed public address — with no address and no way to tell the two apart: whoever registers endpoints could otherwise map your internal DNS. The resolved address is kept on WebhookUrlBlockedError.resolvedAddress for server-side logs when you call resolveAndValidate yourself.
A custom fetchImpl does not pin by itself
The default transport is not global fetch: it is a built-in client that connects to the validated IP. A custom fetchImpl (proxy, instrumentation) receives that IP on its init object under PINNED_ADDRESS, but plain fetch ignores it and resolves the hostname again — reopening the rebind window. Keep pinning by delegating to the exported pinnedFetch, then declare it:
import { pinnedFetch } from '@basaltkit/webhooks'
webhooksPlugin({
secret,
fetchImpl: async (url, init) => {
const started = Date.now()
try { return await pinnedFetch(url, init) } finally { metrics.observe(Date.now() - started) }
},
fetchImplPinsAddress: true,
})Without fetchImplPinsAddress, a custom fetchImpl triggers a one-time process warning (BASALT_WEBHOOKS_UNPINNED_FETCH) and the deliverer re-resolves and re-validates the host before every retry. That narrows the window but can't close it — an unpinned client resolves on its own at connect time. Rewriting the URL to the IP is not done for you: plain fetch can't set TLS SNI separately, so certificate validation would break for https endpoints.
// Self-hosted setup that must deliver to an internal host:
webhooksPlugin({ secret, ssrf: { allowPrivateHosts: true } })
// HTTPS only (reject http:// endpoints at registration-time delivery):
webhooksPlugin({ secret, ssrf: { allowedSchemes: ['https:'] } })
// Turn the guard off entirely — don't, unless every URL is yours:
webhooksPlugin({ secret, ssrf: false })allowPrivateHosts: true skips validation and pinning, so the operator's own resolver is honoured at connect time. assertDeliverableUrl(url) is exported if you want to reject a bad URL at registration time — with a clear error to the customer — instead of at first delivery.
Exposing endpoint management over HTTP
The package ships no HTTP routes: who may manage a tenant's endpoints is an app decision, and it is a privileged one (an endpoint is an outbound data export). Build them on the neutral route() so they serve identically on Fastify, Express and Hono:
import { route } from '@basaltkit/http'
import { ctx, type Container } from '@basaltkit/core'
import { WEBHOOKS } from '@basaltkit/webhooks'
import { z } from 'zod'
const hooks = () => (ctx().container as Container).get(WEBHOOKS)
export const webhookRoutes = () => [
route({
method: 'GET',
url: '/webhooks/endpoints',
meta: { auth: true, teamRole: 'admin' },
// list() never includes signing secrets (`hasSecret` instead)
async handler() { return { data: await hooks().list() } },
}),
route({
method: 'POST',
url: '/webhooks/endpoints',
meta: { auth: true, teamRole: 'admin' },
body: z.object({ url: z.string().url(), events: z.array(z.string()).min(1) }),
async handler({ body, reply }) {
// tenantId is forced from ctx() — never read it from the body. The response
// carries the endpoint's generated signing secret — the only time it is shown.
return reply.code(201).send(await hooks().register(body))
},
}),
]meta.teamRole needs teamsPlugin; meta.auth needs authPlugin. Declaring either without its plugin refuses to boot with UnguardedRouteMetaError (HTTP_UNGUARDED_ROUTE_META) rather than serving the route unguarded — see the adapters guide.
Durable subscription stores
The default MemoryWebhookStore forgets every endpoint on restart — after a redeploy nobody is subscribed and events silently stop. In production, swap in a durable store. The WebhookStore contract is identical across backends, so it's a one-line change.
SQLite (single node, zero dependencies)
@basaltkit/webhooks-sqlite persists subscriptions in a local file over Node's built-in node:sqlite (Node 22.5+; flag-free on Node 24).
import { webhooksPlugin } from '@basaltkit/webhooks'
import { sqliteWebhookStore } from '@basaltkit/webhooks-sqlite'
const webhooks = sqliteWebhookStore('./data/webhooks.db') // ':memory:' by default
webhooksPlugin({ store: webhooks.store, secret: process.env.WEBHOOK_SECRET })sqliteWebhookStore() opens (or creates) the database, applies an idempotent schema (adding the secret-rotation columns to a table created by an earlier version), and returns { store, db } — the raw db handle is exposed if you need it.
Prisma (Postgres/MySQL, multi-instance)
@basaltkit/webhooks-prisma shares one set of subscriptions across instances on the database you already run. Bring your own PrismaClient; the package ships a reference model.
pnpm add @basaltkit/webhooks @basaltkit/webhooks-prisma
pnpm basalt prisma:sync --push # add the WebhookEndpoint model + create the tableimport { webhooksPlugin } from '@basaltkit/webhooks'
import { prismaWebhookStore } from '@basaltkit/webhooks-prisma'
import { PrismaClient } from '@prisma/client'
const prisma = new PrismaClient()
const webhooks = prismaWebhookStore(prisma)
webhooksPlugin({ store: webhooks.store, secret: process.env.WEBHOOK_SECRET })prisma:sync discovers every installed @basaltkit/*-prisma package and merges its models into your schema.prisma. Wire the store before its model exists and it fails fast, naming the missing model — see Persistence.
The model's previousSecret / previousSecretExpiresAt columns hold a secret rotation. The store writes them only when an endpoint is rotating, so a schema that predates them keeps working until you call rotateSecret() — sync and migrate before that.
Which backend?
| Store | Package | Use when |
|---|---|---|
| Memory | @basaltkit/webhooks | Dev and tests (lost on restart) |
| SQLite | @basaltkit/webhooks-sqlite | A single node, zero dependencies, local file |
| Prisma | @basaltkit/webhooks-prisma | Postgres/MySQL, multiple instances share subscriptions |
Writing your own store
Implement the WebhookStore contract — four methods — over any backend:
import { type WebhookStore, type WebhookEndpoint, matchesEvent } from '@basaltkit/webhooks'
class MyWebhookStore implements WebhookStore {
// active, tenant-scoped, event-pattern matched (use matchesEvent)
async forEvent(event: string, tenantId?: string): Promise<WebhookEndpoint[]> { /* … */ return [] }
async add(endpoint: Omit<WebhookEndpoint, 'id'> & { id?: string }): Promise<WebhookEndpoint> { /* … */ throw 0 }
// when tenantId is given, remove ONLY if that tenant owns the endpoint
async remove(id: string, tenantId?: string): Promise<void> { /* … */ }
async list(tenantId?: string): Promise<WebhookEndpoint[]> { /* … */ return [] }
}
webhooksPlugin({ store: new MyWebhookStore(), secret: process.env.WEBHOOK_SECRET })Two rules the built-in stores follow and yours must too: forEvent returns endpoints whose tenantId matches or is undefined (global endpoints get everything), and it skips active: false. Called with no tenantId (or null / '') it must fail closed and return tenant-agnostic endpoints only — never every tenant's. A deliberate allTenants dispatch reads the endpoints through list() instead, and the manager re-filters every result, so a store that gets this wrong still can't widen delivery. remove(id, tenantId) must be a silent no-op when the endpoint belongs to someone else (the manager also checks ownership before calling it). A SQL row's NULL secret / tenantId may come back as null: the deliverer and manager treat null exactly like an absent field (a tenant-agnostic endpoint, signed with the default secret).
To support secret rotation, persist previousSecret and previousSecretExpiresAt (a Date), and clear both when add() receives them explicitly set to undefined: that is how a re-register ends a rotation. A store that drops them still works; its rotations are immediate cut-overs.
Options reference
webhooksPlugin(options)
Everything except store, deliverer, events and the three fan-out options is forwarded to the WebhookDeliverer it constructs (and ignored if you pass your own deliverer).
| Option | Type | Default | Purpose |
|---|---|---|---|
store | WebhookStore | MemoryWebhookStore | Where subscriptions live — swap for webhooks-sqlite/webhooks-prisma, or endpoints vanish on restart |
deliverer | WebhookDeliverer | built from these options | Bring your own (shared with an outbox relay, or a test double) |
events | string[] | [] (off) | Domain-event patterns to auto-dispatch. Non-empty makes the plugin depend on basalt:events |
maxEndpointsPerDispatch | number | false | 100 | Most active endpoints of one scope (tenant, or tenant-agnostic) per event; over it, that scope is refused whole (fan-out cap) |
dispatchConcurrency | number | 16 | Deliveries one dispatch runs at once |
onFanOutExceeded | (info) => void | console.warn | Called once per refused scope with { event, tenantId, endpoints, limit }. Must not throw |
secret | string | — | Default HMAC signing secret (min 16 chars) for tenant-agnostic endpoints. Tenant endpoints always use their own |
allowSharedSecret | boolean | false | Opt-out: sign a tenant endpoint that has no own secret with the default secret (otherwise refused) |
allowUnsigned | boolean | false | Opt-out: send unsigned when there is no secret at all (otherwise refused) |
maxRetries | number | 3 | Retries after the first attempt; only 5xx/network/timeout are retried |
backoffMs | number | 500 | Base wait, doubled per attempt (500 ms, 1 s, 2 s, …) |
timeoutMs | number | 10_000 | Per-attempt deadline covering DNS resolution and the request; running out counts as a transient failure |
ssrf | SsrfGuardOptions | false | on | The delivery-URL guard (below). false disables it entirely |
fetchImpl | typeof fetch | built-in pinned transport (not global fetch) | Injected HTTP client; it receives the validated address on the init object under the exported PINNED_ADDRESS symbol. Plain fetch ignores it — delegate to pinnedFetch to keep pinning (see the SSRF guard) |
fetchImplPinsAddress | boolean | false | Declares that your fetchImpl honours PINNED_ADDRESS (e.g. wraps pinnedFetch): silences the unpinned warning and skips the per-retry re-validation |
sleep | (ms) => Promise<void> | setTimeout | Injectable backoff sleep (tests) |
now | () => number | Date.now()/1000 | Injectable clock in seconds, used for the signature timestamp |
SsrfGuardOptions (the ssrf option)
| Option | Type | Default | Purpose |
|---|---|---|---|
allowPrivateHosts | boolean | false | Trusted self-hosted delivery to internal hosts. Skips validation and address pinning |
allowedSchemes | string[] | ['https:', 'http:'] | Permitted URL schemes — narrow to ['https:'] to refuse plaintext endpoints |
allowedPorts | number[] | 'any' | 80, 443, >= 1024 minus DEFAULT_BLOCKED_PORTS | The port policy. An array allows exactly those ports; 'any' turns it off. Applies with allowPrivateHosts too |
lookup | (host) => Promise<{ address, family? }[]> | dns.lookup(host, { all: true }) | Injected resolver (tests) |
webhookOutboxPlugin(options)
| Option | Type | Default | Purpose |
|---|---|---|---|
store | OutboxStore | MemoryOutboxStore | Durable outbox — in-memory defeats the pattern's whole purpose |
events | string[] | ['**'] (all) | Domain-event patterns captured into the outbox |
intervalMs | number | 5000 | Relay poll interval. 0 disables the timer — relay manually through OUTBOX |
batchSize | number | 50 | Entries delivered per flush |
maxAttempts | number | 10 | Attempts before an entry is left dead (never flushed again) |
concurrency | number | 8 | Entries delivered in parallel per flush, so one slow endpoint can't block the batch |
tenantConcurrency | number | ceil(concurrency / 2) | Most deliveries one tenant may have in flight at once, across flushes |
dispatchTimeoutMs | number | false | 10_000 | Max wait per entry before the flush moves on; the delivery continues detached and its outcome is still recorded |
onFlushError | (error) => void | console.error | A timer/shutdown flush failed at the store level. Must never throw |
onPermanentFailure | (entry, failures) => void | console.warn | An entry's delivery failed permanently for some endpoints; they are not retried. Must never throw |
No onDead here — use outboxPlugin from @basaltkit/events when you need it, as shown above. The plugin depends on both basalt:webhooks and basalt:events, and drains the outbox once on shutdown (best-effort).
Signing & SSRF helpers
| Export | Signature | Purpose |
|---|---|---|
signPayload | (body, secret | secrets[], timestampSeconds) => string | Builds t=…,v1=… — sign a payload by hand; one v1 per secret when given several (current first) |
verifySignature | (header, body, secret, toleranceSeconds = 300, nowSeconds?) => boolean | Constant-time verify in a receiver; true if any v1 matches; false for a secret under 16 chars; throws RangeError only for an invalid tolerance/clock |
generateWebhookSecret | () => string | A fresh whsec_… secret (32 random bytes) |
MIN_WEBHOOK_SECRET_LENGTH | 16 | Minimum secret length enforced on both ends |
assertDeliverableUrl | (url, options?) => Promise<void> | Reject an SSRF-unsafe URL at registration time; throws WebhookUrlBlockedError |
resolveAndValidate | (url, options?) => Promise<ValidatedTarget> | The same check, returning the resolved addresses and the one to pin |
isPrivateIp | (ip) => boolean | The range predicate itself; anything that isn't a public IP literal is true |
isPortAllowed | (port, allowedPorts?) => boolean | The port policy predicate |
DEFAULT_BLOCKED_PORTS | readonly number[] | Ports the default policy refuses at or above 1024 |
matchesEvent | (patterns, event) => boolean | The pattern matcher, for a custom store's forEvent |
webhookOutboxDispatch | (webhooks, options?) => OutboxDispatch | Adapts a WebhookManager into an outbox dispatch; throws only on transient failures, with stable delivery ids |
pinnedFetch | (url, init) => Promise<Response> | fetch-compatible client over the pinned transport — the delegate for a custom fetchImpl |
deriveDeliveryId | (idempotencyKey, endpointId) => string | The deterministic delivery id used by the outbox / idempotencyKey |
Failure modes & troubleshooting
Most delivery problems are not exceptions — they come back on the DeliveryResult, because one bad endpoint must not fail the others:
| Outcome | error | attempts | When |
|---|---|---|---|
| SSRF refusal | Refusing to deliver webhook to <url>: <reason> (bad URL/scheme, private IP literal) or Refusing to deliver webhook: host does not resolve to an allowed public address. (DNS verdict) | 0 | Bad scheme, or the host is/resolves to a private, loopback, link-local, CGNAT, ULA or reserved address — or doesn't resolve |
| Blocked port | Refusing to deliver webhook to <url>: port N is not allowed. | 0 | The URL's port is outside the port policy |
| Fan-out cap | fan-out cap exceeded: … | 0 | The scope has more than maxEndpointsPerDispatch endpoints for this event; none of them was sent to |
| No usable secret | no signing secret; refusing unsigned delivery / tenant endpoint has no own secret; … / endpoint signing secret is too short … | 0 | Nothing to sign with, a tenant endpoint with only the shared secret, or a secret under 16 chars |
| Client error | HTTP 4xx | 1 | The receiver rejected it — never retried inline (408/429 are retryable for the outbox) |
| Redirect | redirect refused | 1 | The endpoint answered 3xx; following it would defeat the SSRF check |
| Transient | last network/timeout message (host resolution timed out when DNS ate the deadline) | maxRetries + 1 | 5xx, connection error or per-attempt timeout (DNS included), retried with backoff, still failing |
| Internal | internal delivery error | 0 | deliver() threw unexpectedly; logged server-side, the other endpoints are unaffected |
| Error | Code | HTTP | When |
|---|---|---|---|
WebhookUrlBlockedError | — (name only) | — | Thrown by assertDeliverableUrl / resolveAndValidate; inside deliver() it is caught and turned into the failed result above |
WebhookTenantRequiredError | WEBHOOKS_TENANT_REQUIRED | — | register / list / unregister with tenancy active and no tenant (context or explicit) without { system: true } |
WebhookEndpointInvalidError | WEBHOOK_ENDPOINT_INVALID | 400 | register() with a URL that doesn't parse or uses a scheme or port the deliverer refuses, a secret under 16 characters, or an empty events list — nothing is stored. Also rotateSecret() with a bad graceSeconds/secret, or on an endpoint without its own secret |
WebhookEndpointNotFoundError | WEBHOOK_ENDPOINT_NOT_FOUND | 404 | rotateSecret() on an endpoint that doesn't exist in the caller's scope |
WebhookEndpointIdInUseError | WEBHOOK_ENDPOINT_ID_IN_USE | 409 | register() / MemoryWebhookStore.add() with an id another scope already holds (the SQL stores throw their own error with the same code) |
UnknownTokenError | DI_UNKNOWN_TOKEN | — | container.get(WEBHOOKS) without webhooksPlugin registered |
UnguardedRouteMetaError | HTTP_UNGUARDED_ROUTE_META | boot | Your own endpoint-management routes declare meta.auth / meta.teamRole without the enforcing plugin |
- Every delivery fails with
attempts: 0in development — the SSRF guard is refusinglocalhost/127.0.0.1/ a.localname. Use a tunnel with a public hostname, orssrf: { allowPrivateHosts: true }in the dev config only. - Receivers report a bad signature although the secret matches — they are verifying over a re-serialized body. The HMAC covers the exact bytes; capture the raw body before parsing.
- Endpoints disappear after every deploy — still on
MemoryWebhookStore. Move towebhooks-sqliteorwebhooks-prisma. - Events emitted from a job only reach some endpoints — there is no tenant in
ctx()outside a request, so only global (tenant-less) endpoints match. Calldispatch(event, data, tenantId)explicitly (or run the job inside the tenant's context). - A request handler got slow after adding webhooks —
dispatchawaits every delivery, retries included (up to(maxRetries + 1) × timeoutMsper endpoint). Move it to the outbox or a queue job. - The outbox stops delivering an event and nothing is logged where you look — it hit
maxAttemptsand is dead; the defaultonDeadwrites toconsole.error. InspectlastErroron the entry, or wireoutboxPluginwith your ownonDead.
See also
- Queues & jobs — drive delivery from a queue for durable retries.
- Persistence — durable stores,
prisma:sync, the outbox store. - Teams — who may manage a tenant's endpoints.
- Multi-tenant SaaS cookbook — per-tenant endpoints in a real app.