A kit

Realtime

Dedikált, horizontálisan skálázható gateway, ami élőben pusholja a notificationöket a csatlakozott klienseknek: csak szerver→kliens, centrifuge-ra épülve.

A broadcast csatorna egy notificationt a recipiens élő kapcsolatára kézbesít, abban a pillanatban, amikor történik: egy "poke" a durable database csatorna tetején, nem annak helyettesítője. Ez az oldal a kapcsolatokat tartó gatewayről szól.

Miért külön gateway

A notificationök küldése bármely processből történhet, gyakran a workerből, egy listenerből. A kliens-kapcsolatok tartása más feladat, más skálázási alakkal: sok, hosszú élettartamú stream, egy csatlakozott userenként, aminek túl kell élnie a példányok szabad jövés-menését. A kit ezeket szétválasztja: a notify/broadcast egy pici publisher, amit bármely process importálhat (szerver nélkül, állapot nélkül); a realtime egy dedikált binary (cmd/realtime, az add realtime scaffoldolja), ami csak kapcsolatokat tart és relay-eli, amit publikálnak.

A driver-seam szándékos: a notify/broadcast-on és a realtime-on kívül semmi nem importálja közvetlenül a centrifuge-ot, így egy másik transport később lecserélheti anélkül, hogy az alkalmazáskódhoz nyúlna.

Publikálás

import "github.com/gp-system/notify/broadcast"

pub := broadcast.MustNewPublisher(ctx, cfg.Realtime.BroadcastValkey(), cfg.Realtime.Broadcast)
notifier := notify.NewHub(
    notify.NewMailChannel(mailer),
    notifydatabase.NewChannel(store),
    notify.NewBroadcastChannel(pub, cfg.Realtime.Broadcast.ChannelPrefix),
)

A Publisher egy publish-only centrifuge.Node, ami osztozik a projekt Valkey brokerén; sosem fogad kapcsolatot, így egy worker-process semmi mástól nem függ, csak egy hálózati klienstől. A publikálás fire-and-forget: egy sikertelen publish visszaadódik a hívónak, nincs újrapróbálva: a durability a database csatornából jön, egy kliens, aki lemarad egy élő push-ról, reconnectkor onnan backfillel. Az insert-majd-publish kontraktust maga a notify.Hub.Send már kikényszeríti (database csatorna a broadcast előtt, ugyanaz a delivery ID mindkettőben).

A gateway

A cmd/realtime nem hordoz üzleti függőséget (nincs adatbázis-pool, nincs outbox), pont ezért skálázódik szabadon. WebSocketet exponál egy HTTP-streaming/SSE fallbackkal (mindkettő a centrifugal/centrifuge-on keresztül, ami a Centrifugo és a Grafana Live mögötti library) a /realtime/* alatt, minden kapcsolat JWT-jét a connect frame-ben authentikálja (nincs külön /broadcasting/auth endpoint; v1-ben az egyetlen privát csatorna egy kapcsolat saját subjectje, így elég egyszer, connectkor authentikálni), és minden kapcsolatot server-side feliratkoztat a saját user:<subject> csatornájára. A gyakori esetben nincs szükség kliens-oldali subscribe-logikára.

A megosztott topic-csatornák opt-in és explicit módon authorizáltak: minden kapcsolat JWT-je egy rbac.Identity-ra oldódik fel, és egy topic csak akkor iratkoztatható fel, ha valamelyik Allow/AllowPrefix callback jótáll érte.

cmd/realtime/main.go
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
})

Egy topic, aminek nincs regisztrált Allow callbackje, sosem iratkoztatható fel.

Ha egy csatorna-család teljes neve csak futásidőben ismert (egy dinamikus erőforrás soronként: egy tenant, egy projekt, egy repository), az Allow-val egyenként, induláskor való elnevezés nem lehetséges: helyette AllowPrefix használandó, aminek a callbackje megkapja a teljes topic stringet, hogy meg tudja állapítani, épp melyik konkrét instance-t kérik.

topics.AllowPrefix("project:", func(id *rbac.Identity, topic string) bool {
    projectID, _ := strings.CutPrefix(topic, "project:")
    return projectAccess.CanView(id.Subject, projectID)
})

Egy pontos Allow-egyezés mindig elsőbbséget élvez egy illeszkedő AllowPrefix-szel szemben; a prefixek közül az elsőként regisztrált, illeszkedő nyer.

Skálázás: a példányok szabadon jönnek-mennek

Minden gateway-példány stateless: csak azokat a csatornákat tartja, amikre a saját kapcsolatainak szükségük van, és a közös Valkey broker minden érintett példányhoz eljuttat egy publish-t, bármely processből, sticky session-igény nélkül.

  • Felskálázás: docker compose up -d --scale realtime=N; egy új példány nulla állapottal indul, és csak akkor gyűlnek a szubszkripciói, amikor kliensek landolnak rajta.
  • Leskálázás / rolling restart: a SIGTERM bekapcsol egy draining flaget (a /realtime/healthz és az új connect-kísérletek 503-at adnak, így a load balancer health checkje kiejti a példányt), a node.Shutdown minden tartott kapcsolatot reconnect-jogosult kóddal zár, a closerek lefutnak, a telemetria flush-öl. A kliens saját reconnect-logikája (backoff, beépítve a JS kliensbe) egy túlélő példányra viszi, ami frissen feliratkoztatja; ami a rés alatt publikálódott, azt a database csatornából állítja helyre, nem a gateway replay-eli.
  • Crash: a kliens oldala ugyanaz, mint egy graceful restartnál, csak a drain lépés nélkül.
  • Valkey restart: a reconnect és az újra-subscribe a centrifuge dolga; a publikálás a kiesés alatt logolva és eldobva (ugyanaz a fire-and-forget kontraktus), ugyanúgy helyreállítva, mint egy kliens saját reconnectje alatt lemaradt publish.

Config

VáltozóAlapértékJelentés
REALTIME_LISTEN_ADDR:3000konténer-belső listen address
REALTIME_DRAIN_TIMEOUT20sgraceful-shutdown budget; tartsd a deployment stop-grace-periódusa alatt
REALTIME_READ_HEADER_TIMEOUT5sSlowloris-védelem a connect handshake-en
REALTIME_CHANNEL_PREFIXa projekt neveValkey csatorna-névtér, közös minden notify/broadcast.Publisher-rel
REALTIME_MAX_CONNECTIONS0 (korlátlan)példányonkénti stream-plafon; felette a connect-kísérletek 503-at kapnak
VALKEY_*örököltújrahasznosítva a queue.Config-ból: a gateway brokere ugyanaz a Valkey, amit a queue használ
JWT_*örököltújrahasznosítva az auth.Config-ból: ugyanaz a token-issuer, amiben minden modul megbízik

Lásd az add realtime referenciát a scaffoldolt elemekért (a binary, a compose.yml Traefik-routolt service-e, a .env.example), és az add notification oldalt a broadcast csatorna egy notification Via-jába drótozásához.

Kliens

import { Centrifuge } from 'centrifuge'

const client = new Centrifuge('wss://your-app.example/realtime/connection/websocket', {
    getToken: async () => fetchAccessToken(), // ugyanaz a JWT, amit az API használ
})
client.on('publication', (ctx) => {
    // ctx.data: { id, name, payload, at }, az id visszaköti a database csatorna sorához
})
client.connect()

A centrifuge-js kezeli a reconnect-backoffot és a token-refresht; a server-side subscription a recipiens saját csatornájára azt jelenti, hogy a kliensnek nincs szüksége explicit subscribe hívásra a neki szóló notificationökhöz. Reconnectkor kérd le az olvasatlan notificationöket a saját modulod inbox-endpointjából (a database csatornára építve), hogy backfilleld, ami a kapcsolat megszakadása alatt lemaradt. A socket egy élő poke, a database az igazságforrás.

Kapcsolódó oldalak

  • notify/broadcast: a publish-only oldal, amit bármely process importálhat.
  • auth/rbac: az Identity, amire egy kapcsolat JWT-je feloldódik, és amit az Allow/AllowPrefix ellenőriz.
  • Alkalmazás-életciklus: a közös app.Run váz, amire a realtime.Run épül.
Copyright © 2026