Skip to content

Repository files navigation

Relay

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.

Features

  • 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.Store interface for anything else.
  • HMAC-SHA256 signatures — Every delivery is signed. Receivers verify authenticity using the signature package.
  • 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

Install

go get github.com/xraph/relay

Requires Go 1.22 or later.

Quick Start

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)
}

Configuration

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

Webhook Verification

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

Endpoints must have a signing secret

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.

Admin API

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

Routes

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

Store Backends

Memory (testing)

import "github.com/xraph/relay/store/memory"

r, _ := relay.New(relay.WithStore(memory.New()))

PostgreSQL

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

SQLite

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

Redis

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.

MongoDB

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 Index

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

Architecture

┌─────────────────────────────────────────────┐
│                  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)    │
└────────┴────────┴───────┴───────┴─────────────┘

Examples

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

License

See LICENSE for details.

About

Composable webhook delivery engine for Go

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Used by

Contributors

Languages