Composable webhook delivery engine for Go.
Relay is a library — not a service. Import it into your Go application to get tenant-scoped webhook endpoints, dynamic event type definitions, guaranteed delivery with signature verification, and replay capabilities.
- Dynamic webhook definitions — Register event types at runtime with optional JSON Schema validation
- Composable store pattern — Plug in PostgreSQL, SQLite, Redis, MongoDB, or in-memory backends. Implement the
store.Storeinterface for anything else. - HMAC-SHA256 signatures — Every delivery is signed. Receivers verify authenticity using the
signaturepackage. - Exponential backoff retries — Configurable schedule (default: 5s → 30s → 2m → 15m → 2h). Failed deliveries land in the dead letter queue.
- Per-endpoint rate limiting — Token bucket limiter prevents overloading downstream services
- Admin HTTP API — Full CRUD for event types, endpoints, events, deliveries, and DLQ replay
- OpenTelemetry + Prometheus — Traces per delivery span, counters, latency histograms, and gauges out of the box
- Multi-tenant by default — Every endpoint and event is scoped to a tenant ID
go get github.com/xraph/relayRequires Go 1.22 or later.
package main
import (
"context"
"encoding/json"
"log"
"github.com/xraph/relay"
"github.com/xraph/relay/catalog"
"github.com/xraph/relay/endpoint"
"github.com/xraph/relay/event"
"github.com/xraph/relay/store/memory"
)
func main() {
ctx := context.Background()
// 1. Create a Relay instance with a store backend.
r, err := relay.New(
relay.WithStore(memory.New()),
)
if err != nil {
log.Fatal(err)
}
// 2. Register an event type in the catalog.
r.RegisterEventType(ctx, catalog.WebhookDefinition{
Name: "order.created",
Description: "Fired when a new order is placed",
Version: "2025-01-01",
})
// 3. Create a webhook endpoint for a tenant.
r.Endpoints().Create(ctx, endpoint.Input{
TenantID: "tenant-acme",
URL: "https://acme.example.com/webhook",
EventTypes: []string{"order.*"}, // glob pattern
})
// 4. Send an event — Relay fans out to all matching endpoints.
r.Send(ctx, &event.Event{
Type: "order.created",
TenantID: "tenant-acme",
Data: json.RawMessage(`{"order_id":"ORD-001","amount":99.99}`),
})
// 5. Start the delivery engine and stop gracefully.
r.Start(ctx)
defer r.Stop(ctx)
}All options are set via functional options on relay.New():
| Option | Default | Description |
|---|---|---|
WithStore(s) |
required | Persistence backend (memory.New(), postgres.New(db), sqlite.New(db), redis.New(kv), mongo.New(db)) |
WithLogger(l) |
slog.Default() |
Structured logger |
WithConcurrency(n) |
10 |
Delivery worker goroutines |
WithPollInterval(d) |
1s |
How often the engine checks for pending deliveries |
WithBatchSize(n) |
50 |
Max deliveries dequeued per poll cycle |
WithRequestTimeout(d) |
30s |
HTTP timeout per delivery attempt |
WithMaxRetries(n) |
5 |
Maximum delivery attempts before moving to DLQ |
WithRetrySchedule(s) |
5s, 30s, 2m, 15m, 2h |
Backoff intervals between retries |
WithShutdownTimeout(d) |
30s |
Grace period for in-flight deliveries on shutdown |
WithCacheTTL(d) |
30s |
Catalog in-memory cache TTL |
Receivers verify incoming webhooks using the signature package:
import "github.com/xraph/relay/signature"
func handleWebhook(w http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
sig := r.Header.Get("X-Relay-Signature") // "v1=<hex>"
ts, _ := strconv.ParseInt(r.Header.Get("X-Relay-Timestamp"), 10, 64)
if !signature.Verify(body, endpointSecret, ts, sig) {
http.Error(w, "invalid signature", http.StatusUnauthorized)
return
}
// Process the verified webhook...
}Every delivery includes these headers:
| Header | Description |
|---|---|
X-Relay-Signature |
v1=<hmac-sha256-hex> computed over timestamp.body |
X-Relay-Timestamp |
Unix timestamp (seconds) of the delivery attempt |
X-Relay-Event-ID |
The event's TypeID (e.g. evt_01h6rz...) |
Content-Type |
application/json |
Relay refuses to sign with an empty secret, and refuses to deliver to an endpoint that has one.
Before this, an endpoint with no secret still received deliveries carrying a
well-formed X-Relay-Signature, computed with the empty string as the key.
That signature verified against the same empty key, so a receiver calling
signature.Verify correctly still accepted it, and anybody who knew the
payload and the timestamp could reproduce it. Nothing in the header
distinguished it from a real one.
Endpoints().Create has always generated a secret when you do not supply one,
so an endpoint reaches this state only through the store interface directly.
To find out whether you have any, run this before upgrading:
unsigned, err := r.Endpoints().ListUnsigned(ctx, "")An empty tenant checks every tenant; pass one to check just that one. It reads
every endpoint you have, however many, so an empty result really does mean
none. Rotate a secret onto whatever it returns with Endpoints().RotateSecret.
Deliveries to an endpoint with no secret now fail with endpoint has no signing secret and land in the dead letter queue. They burn the whole retry
schedule getting there, because a missing secret reads as a network-class
failure to the retrier and a network-class failure is worth retrying. A
missing secret is not, so those retries are wasted. The DLQ entry names the
cause, which is what matters when you are working out why an endpoint stopped.
Mount the admin HTTP handler to manage webhooks at runtime:
import "github.com/xraph/relay/api"
handler := api.NewHandler(r.Store(), r.Catalog(), r.Endpoints(), r.DLQ(), logger)
mux.Handle("/webhooks/", http.StripPrefix("/webhooks", handler))| Method | Path | Description |
|---|---|---|
| POST | /event-types |
Register an event type |
| GET | /event-types |
List event types |
| GET | /event-types/{name} |
Get event type by name |
| DELETE | /event-types/{name} |
Deprecate an event type |
| POST | /endpoints |
Create an endpoint |
| GET | /endpoints |
List endpoints |
| GET | /endpoints/{id} |
Get endpoint |
| PUT | /endpoints/{id} |
Update endpoint |
| DELETE | /endpoints/{id} |
Delete endpoint |
| PATCH | /endpoints/{id}/enable |
Enable endpoint |
| PATCH | /endpoints/{id}/disable |
Disable endpoint |
| POST | /endpoints/{id}/rotate-secret |
Rotate signing secret |
| GET | /endpoints/{id}/deliveries |
List deliveries for endpoint |
| POST | /events |
Create an event |
| GET | /events |
List events |
| GET | /events/{id} |
Get event |
| GET | /dlq |
List DLQ entries |
| POST | /dlq/{id}/replay |
Replay a single DLQ entry |
| POST | /dlq/replay |
Bulk replay DLQ entries |
| GET | /stats |
Get delivery statistics |
import "github.com/xraph/relay/store/memory"
r, _ := relay.New(relay.WithStore(memory.New()))import (
"github.com/xraph/grove"
"github.com/xraph/grove/drivers/pgdriver"
"github.com/xraph/relay/store/postgres"
)
pgdb := pgdriver.New()
pgdb.Open(ctx, "postgres://localhost:5432/mydb?sslmode=disable")
db, _ := grove.Open(pgdb)
store := postgres.New(db)
store.Migrate(ctx) // creates relay_* tables
r, _ := relay.New(relay.WithStore(store))import (
"github.com/xraph/grove"
"github.com/xraph/grove/drivers/sqlitedriver"
"github.com/xraph/relay/store/sqlite"
)
sdb := sqlitedriver.New()
// busy_timeout makes a writer wait for the lock instead of failing at once
// with "database is locked". Relay writes from several goroutines (ten
// delivery workers by default), so you want it. It has to go in the DSN:
// it is per connection, and the DSN is the only place that reaches them all.
sdb.Open(ctx, "file:relay.db?_pragma=busy_timeout(5000)")
db, _ := grove.Open(sdb)
store := sqlite.New(db)
store.Migrate(ctx) // creates relay_* tables
r, _ := relay.New(relay.WithStore(store))import (
"github.com/xraph/grove/kv"
"github.com/xraph/grove/kv/drivers/redisdriver"
redisstore "github.com/xraph/relay/store/redis"
)
rdb := redisdriver.New()
rdb.Open(ctx, "redis://localhost:6379")
kvStore, _ := kv.Open(rdb)
store := redisstore.New(kvStore)
// Builds the every-tenant endpoint index, and checks that the kv store maps keys
// in a way relay can follow. Don't ignore this error: it's the only place that
// check shows up.
if err := store.Migrate(ctx); err != nil {
log.Fatal(err)
}
r, _ := relay.New(relay.WithStore(store))You can give it a kv store with a namespace hook, so two relays share one redis without sharing data:
kvStore, _ := kv.Open(rdb, kv.WithHook(middleware.NewNamespace("app1")))Every key relay touches lands under the prefix. That covers the records, the
endpoint, event, delivery and DLQ indexes, the idempotency keys, the replay
claim, the delivery queue, the migration markers and the wake channel. Two
stores on one redis with different namespaces can use the same tenant ids and
never see each other's rows. Each namespace builds its own indexes, so every
new namespace needs its own Migrate, and an every-tenant list returns
ErrEndpointIndexNotBuilt until it has run.
You need grove v1.7.0 or later, and relay's redis store won't build against anything older. Earlier versions ignored the hook on deletes and on most commands besides get and set, so a deleted endpoint stayed in redis with its secret.
The namespace has to cover every command. A hook scoped to reads only, or to
writes only, would write a record under one name and look for it under another,
and relay can't tell which name is the real one. Migrate returns
ErrKVKeysInconsistent for that before it writes anything. It does the same
for a rewrite that isn't a prefix, such as a suffix, because a scan can't undo
it. ErrKVRewritesKeys is the old name for the same error. It still matches,
and it's deprecated.
import (
"github.com/xraph/grove"
"github.com/xraph/grove/drivers/mongodriver"
"github.com/xraph/relay/store/mongo"
)
mdb := mongodriver.New()
mdb.Open(ctx, "mongodb://localhost:27017/relay")
db, _ := grove.Open(mdb)
store := mongo.New(db)
store.Migrate(ctx) // creates indexes
r, _ := relay.New(relay.WithStore(store))| Package | Description |
|---|---|
relay |
Root package — Relay engine, Send(), Start()/Stop(), functional options |
catalog |
Event type registry with in-memory cache and JSON Schema validation |
endpoint |
Webhook endpoint CRUD service with secret rotation |
event |
Event entity and store interface |
delivery |
Delivery engine, HTTP sender, retry logic with exponential backoff |
dlq |
Dead letter queue with replay and bulk operations |
id |
TypeID-based identity — single ID struct with prefix constants |
signature |
HMAC-SHA256 signing and verification |
ratelimit |
Token bucket rate limiter per endpoint |
observability |
Prometheus metrics and OpenTelemetry tracing |
api |
HTTP admin API handlers (Go 1.22+ ServeMux) |
store |
Composite Store interface (catalog + endpoint + event + delivery + dlq) |
store/memory |
In-memory store for testing |
store/postgres |
PostgreSQL backend using Grove ORM |
store/sqlite |
SQLite backend for embedded/edge deployments |
store/redis |
Redis backend using Grove KV |
store/mongo |
MongoDB backend |
extension |
Forge framework extension integration |
scope |
Multi-tenant context helpers |
┌─────────────────────────────────────────────┐
│ relay.Relay │
│ Send() → validate → persist → fan-out │
│ Start() / Stop() │
├────────────┬────────────┬───────────────────┤
│ Catalog │ Endpoint │ Delivery Engine │
│ (cache + │ Service │ (workers + poll │
│ validate) │ (CRUD) │ + retry + DLQ) │
├────────────┴────────────┴───────────────────┤
│ store.Store │
│ (catalog + endpoint + event + delivery + │
│ dlq interfaces composed) │
├────────┬────────┬───────┬───────┬─────────────┤
│Postgres│ SQLite │ Redis │ Mongo │ Memory │
│ (Grove)│(Grove) │(KV) │(Grove)│(testing) │
└────────┴────────┴───────┴───────┴─────────────┘
See the _examples/ directory:
- basic — Memory store, register type, create endpoint, send event, start engine
- dynamic-catalog — Mount admin API, register event types at runtime
- stripe-style — Webhook receiver with HMAC-SHA256 signature verification
See LICENSE for details.