Realtime
The broadcast channel delivers a notification to a recipient's live connection, the moment it happens: a "poke" on top of the durable database channel, not a replacement for it. This page covers the gateway that holds those connections.
Why a separate gateway
Sending notifications happens from any process, often the worker, in a listener. Holding client connections is a different job with a different scaling shape: many long-lived streams, one per connected user, that must survive instances coming and going freely. The kit keeps these apart: notify/broadcast is a tiny publisher any process imports (no server, no held state); realtime is a dedicated binary (cmd/realtime, scaffolded by add realtime) that only holds connections and relays what gets published.
The driver seam is deliberate: nothing outside notify/broadcast and realtime imports centrifuge directly, so a different transport can replace it later without touching application code.
Publishing
import "github.com/gp-system/gpsystem/notify/broadcast"
pub := broadcast.MustNewPublisher(ctx, cfg.Realtime.Redis, cfg.Realtime.Broadcast)
notifier := notify.NewHub(
notify.NewMailChannel(mailer),
notifydatabase.NewChannel(store),
notify.NewBroadcastChannel(pub, cfg.Realtime.Broadcast.ChannelPrefix),
)
Publisher is a publish-only centrifuge.Node sharing the project's Redis broker; it never accepts a connection, so a worker process depends on nothing but a network client. Publishing is fire-and-forget: a failed publish is returned to the caller, not retried. Durability lives in the database channel, and a client that misses a live push backfills from it on reconnect. Insert-then-publish is the contract notify.Hub.Send already enforces (database channel before broadcast, same delivery ID in both).
The gateway
cmd/realtime holds no business dependencies (no database pool, no outbox), which is exactly why it scales freely. It exposes WebSocket plus an HTTP-streaming/SSE fallback (all through centrifugal/centrifuge, the library behind Centrifugo and Grafana Live) under /realtime/*, authenticates each connection's JWT in the connect frame (no separate /broadcasting/auth endpoint; v1's only private channel is a connection's own subject, so authenticating once at connect is sufficient), and server-side subscribes every connection to its own user:<subject> channel. No client-side subscribe logic needed for the common case.
Shared topic channels are opt-in and explicitly authorized:
err := realtime.Run(ctx, cfg.Realtime, func(topics *realtime.TopicAuth) error {
topics.Allow("announcements", func(id *rbac.Identity) bool { return true })
topics.Allow("admin-alerts", func(id *rbac.Identity) bool { return id.HasRole("admin") })
return nil
})
A topic with no registered Allow callback is never subscribable.
Scaling: instances come and go freely
Every gateway instance is stateless: it only ever holds the channels its own connections need, and the shared Redis broker fans a publish out to every instance that needs it, from any process, without a sticky-session requirement.
- Scale up:
docker compose up -d --scale realtime=N; a new instance starts with zero state and accumulates subscriptions only as clients land on it. - Scale down / rolling restart:
SIGTERMflips a draining flag (/realtime/healthzand new connect attempts return 503, so the load balancer's health check evicts the instance),node.Shutdowndisconnects every held connection with a reconnect-eligible code, closers run, telemetry flushes. The client's own reconnect logic (backoff, built into the JS client) lands it on a surviving instance, which subscribes it fresh; anything published during the gap is recovered from the database channel, not replayed by the gateway. - Crash: the client side is identical to a graceful restart, just without the drain step.
- Redis restart: reconnect and resubscribe are
centrifuge's job; publishing during the outage is logged and dropped (the same fire-and-forget contract), recovered the same way as a missed publish during a client's own reconnect.
Config
| Variable | Default | Meaning |
|---|---|---|
REALTIME_LISTEN_ADDR | :3000 | container-internal listen address |
REALTIME_DRAIN_TIMEOUT | 20s | graceful-shutdown budget; keep under the deployment's stop-grace-period |
REALTIME_READ_HEADER_TIMEOUT | 5s | Slowloris defense on the connect handshake |
REALTIME_CHANNEL_PREFIX | the project name | Redis channel namespace shared with every notify/broadcast.Publisher |
REALTIME_MAX_CONNECTIONS | 0 (unlimited) | per-instance stream cap; over it, connect attempts get 503 |
REDIS_* | inherited | reused from queue.Config: the gateway's broker is the same Redis the queue uses |
JWT_* | inherited | reused from auth.Config: the same token issuer every module trusts |
See the add realtime reference for what gets scaffolded (the binary, compose.yml's Traefik-routed service, .env.example) and the add notification page for wiring the broadcast channel into a notification's Via.
Client
import { Centrifuge } from 'centrifuge'
const client = new Centrifuge('wss://your-app.example/realtime/connection/websocket', {
getToken: async () => fetchAccessToken(), // the same JWT the API uses
})
client.on('publication', (ctx) => {
// ctx.data: { id, name, payload, at }, id ties back to the database channel's row
})
client.connect()
centrifuge-js handles reconnect backoff and token refresh; the server-side subscription to the recipient's own channel means the client needs no explicit subscribe call for notifications addressed to it. On reconnect, fetch unread notifications from your module's own inbox endpoint (built on the database channel) to backfill anything missed while disconnected. The socket is a live poke, the database is the source of truth.