Skip to content

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:

PieceContractWhat it does
StoreWebhookStoreWhere subscriptions live. Answers "who wants invoice.paid for tenant acme?"
DelivererWebhookDelivererSigns the body, validates the URL against SSRF, POSTs it, retries transient failures
ManagerWebhookManager (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:

ts
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:

ts
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:

ts
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 purpose

Event patterns match like this:

  • 'invoice.paid' — that exact event only
  • 'invoice.*' — any event starting with invoice.
  • '*' or '**' — every event
ts
import { matchesEvent } from '@basaltkit/webhooks'
matchesEvent(['invoice.*'], 'invoice.paid') // true

Only 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:

ts
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 days

It 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:

ts
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):

ts
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).

ts
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:

ts
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.

ts
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 createdAt order.
  • Per-tenant cap. A flush delivers up to concurrency entries in parallel (default 8), but one tenant never has more than tenantConcurrency deliveries in flight (default ceil(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:

ts
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 attempts and records lastError, then backs off (exponential from 1 s, capped at 60 s, tracked per relay process). After maxAttempts (default 10) the entry is dead: it stays in the store with its lastError and is never flushed again. Outbox's onDead callback fires once — by default it writes to console.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's onFlushError exists 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:

ts
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:

ts
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 ​

  • timeoutMs is 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 to maxRetries (default 3, so four attempts in total).
  • Client errors (4xx) are not retried — a wrong URL or auth won't fix itself on retry. The result carries error: 'HTTP 404'. 408 and 429 are still marked retryable: true, so the outbox tries them again later.
  • Redirects are refused, not followed: a 3xx ends the delivery with error: '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):
ts
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.

ts
webhooksPlugin({ secret, ssrf: { allowedPorts: [443] } })        // exactly these ports
webhooksPlugin({ secret, ssrf: { allowedPorts: [443, 6379] } })  // your own policy
webhooksPlugin({ secret, ssrf: { allowedPorts: 'any' } })        // policy off

The 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:

ts
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.

ts
// 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:

ts
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).

ts
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.

bash
pnpm add @basaltkit/webhooks @basaltkit/webhooks-prisma
pnpm basalt prisma:sync --push   # add the WebhookEndpoint model + create the table
ts
import { 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? ​

StorePackageUse when
Memory@basaltkit/webhooksDev and tests (lost on restart)
SQLite@basaltkit/webhooks-sqliteA single node, zero dependencies, local file
Prisma@basaltkit/webhooks-prismaPostgres/MySQL, multiple instances share subscriptions

Writing your own store ​

Implement the WebhookStore contract — four methods — over any backend:

ts
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).

OptionTypeDefaultPurpose
storeWebhookStoreMemoryWebhookStoreWhere subscriptions live — swap for webhooks-sqlite/webhooks-prisma, or endpoints vanish on restart
delivererWebhookDelivererbuilt from these optionsBring your own (shared with an outbox relay, or a test double)
eventsstring[][] (off)Domain-event patterns to auto-dispatch. Non-empty makes the plugin depend on basalt:events
maxEndpointsPerDispatchnumber | false100Most active endpoints of one scope (tenant, or tenant-agnostic) per event; over it, that scope is refused whole (fan-out cap)
dispatchConcurrencynumber16Deliveries one dispatch runs at once
onFanOutExceeded(info) => voidconsole.warnCalled once per refused scope with { event, tenantId, endpoints, limit }. Must not throw
secretstring—Default HMAC signing secret (min 16 chars) for tenant-agnostic endpoints. Tenant endpoints always use their own
allowSharedSecretbooleanfalseOpt-out: sign a tenant endpoint that has no own secret with the default secret (otherwise refused)
allowUnsignedbooleanfalseOpt-out: send unsigned when there is no secret at all (otherwise refused)
maxRetriesnumber3Retries after the first attempt; only 5xx/network/timeout are retried
backoffMsnumber500Base wait, doubled per attempt (500 ms, 1 s, 2 s, …)
timeoutMsnumber10_000Per-attempt deadline covering DNS resolution and the request; running out counts as a transient failure
ssrfSsrfGuardOptions | falseonThe delivery-URL guard (below). false disables it entirely
fetchImpltypeof fetchbuilt-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)
fetchImplPinsAddressbooleanfalseDeclares 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>setTimeoutInjectable backoff sleep (tests)
now() => numberDate.now()/1000Injectable clock in seconds, used for the signature timestamp

SsrfGuardOptions (the ssrf option) ​

OptionTypeDefaultPurpose
allowPrivateHostsbooleanfalseTrusted self-hosted delivery to internal hosts. Skips validation and address pinning
allowedSchemesstring[]['https:', 'http:']Permitted URL schemes — narrow to ['https:'] to refuse plaintext endpoints
allowedPortsnumber[] | 'any'80, 443, >= 1024 minus DEFAULT_BLOCKED_PORTSThe 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) ​

OptionTypeDefaultPurpose
storeOutboxStoreMemoryOutboxStoreDurable outbox — in-memory defeats the pattern's whole purpose
eventsstring[]['**'] (all)Domain-event patterns captured into the outbox
intervalMsnumber5000Relay poll interval. 0 disables the timer — relay manually through OUTBOX
batchSizenumber50Entries delivered per flush
maxAttemptsnumber10Attempts before an entry is left dead (never flushed again)
concurrencynumber8Entries delivered in parallel per flush, so one slow endpoint can't block the batch
tenantConcurrencynumberceil(concurrency / 2)Most deliveries one tenant may have in flight at once, across flushes
dispatchTimeoutMsnumber | false10_000Max wait per entry before the flush moves on; the delivery continues detached and its outcome is still recorded
onFlushError(error) => voidconsole.errorA timer/shutdown flush failed at the store level. Must never throw
onPermanentFailure(entry, failures) => voidconsole.warnAn 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 ​

ExportSignaturePurpose
signPayload(body, secret | secrets[], timestampSeconds) => stringBuilds t=…,v1=… — sign a payload by hand; one v1 per secret when given several (current first)
verifySignature(header, body, secret, toleranceSeconds = 300, nowSeconds?) => booleanConstant-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() => stringA fresh whsec_… secret (32 random bytes)
MIN_WEBHOOK_SECRET_LENGTH16Minimum 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) => booleanThe range predicate itself; anything that isn't a public IP literal is true
isPortAllowed(port, allowedPorts?) => booleanThe port policy predicate
DEFAULT_BLOCKED_PORTSreadonly number[]Ports the default policy refuses at or above 1024
matchesEvent(patterns, event) => booleanThe pattern matcher, for a custom store's forEvent
webhookOutboxDispatch(webhooks, options?) => OutboxDispatchAdapts 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) => stringThe 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:

OutcomeerrorattemptsWhen
SSRF refusalRefusing 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)0Bad scheme, or the host is/resolves to a private, loopback, link-local, CGNAT, ULA or reserved address — or doesn't resolve
Blocked portRefusing to deliver webhook to <url>: port N is not allowed.0The URL's port is outside the port policy
Fan-out capfan-out cap exceeded: …0The scope has more than maxEndpointsPerDispatch endpoints for this event; none of them was sent to
No usable secretno signing secret; refusing unsigned delivery / tenant endpoint has no own secret; … / endpoint signing secret is too short …0Nothing to sign with, a tenant endpoint with only the shared secret, or a secret under 16 chars
Client errorHTTP 4xx1The receiver rejected it — never retried inline (408/429 are retryable for the outbox)
Redirectredirect refused1The endpoint answered 3xx; following it would defeat the SSRF check
Transientlast network/timeout message (host resolution timed out when DNS ate the deadline)maxRetries + 15xx, connection error or per-attempt timeout (DNS included), retried with backoff, still failing
Internalinternal delivery error0deliver() threw unexpectedly; logged server-side, the other endpoints are unaffected
ErrorCodeHTTPWhen
WebhookUrlBlockedError— (name only)—Thrown by assertDeliverableUrl / resolveAndValidate; inside deliver() it is caught and turned into the failed result above
WebhookTenantRequiredErrorWEBHOOKS_TENANT_REQUIRED—register / list / unregister with tenancy active and no tenant (context or explicit) without { system: true }
WebhookEndpointInvalidErrorWEBHOOK_ENDPOINT_INVALID400register() 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
WebhookEndpointNotFoundErrorWEBHOOK_ENDPOINT_NOT_FOUND404rotateSecret() on an endpoint that doesn't exist in the caller's scope
WebhookEndpointIdInUseErrorWEBHOOK_ENDPOINT_ID_IN_USE409register() / MemoryWebhookStore.add() with an id another scope already holds (the SQL stores throw their own error with the same code)
UnknownTokenErrorDI_UNKNOWN_TOKEN—container.get(WEBHOOKS) without webhooksPlugin registered
UnguardedRouteMetaErrorHTTP_UNGUARDED_ROUTE_METAbootYour own endpoint-management routes declare meta.auth / meta.teamRole without the enforcing plugin
  • Every delivery fails with attempts: 0 in development — the SSRF guard is refusing localhost / 127.0.0.1 / a .local name. Use a tunnel with a public hostname, or ssrf: { 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 to webhooks-sqlite or webhooks-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. Call dispatch(event, data, tenantId) explicitly (or run the job inside the tenant's context).
  • A request handler got slow after adding webhooks — dispatch awaits every delivery, retries included (up to (maxRetries + 1) × timeoutMs per 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 maxAttempts and is dead; the default onDead writes to console.error. Inspect lastError on the entry, or wire outboxPlugin with your own onDead.

See also ​

Released under the MIT License.