Files
summercms/modules/lighthouse/centrifugo/driver.go
Jakub Zych 211c413273 feat(11-03): broadcast model writes through River jobs enqueued in the write transaction
- Broadcastable contract and Bind[T] bindings; event {action}.{alias},
  default {model, actor, timestamp, ttl} payload, delete snapshot taken
  before the row goes
- GORM callbacks installed via lagoon.OnDatabase enqueue a summer.broadcast
  job on the write's *sql.Tx inside a savepoint; failures are logged and
  never abort the write; zero-key batch writes are skipped
- WithoutBroadcasting[T] (ctx-scoped, per type) and Service.Emit for one
  summary event; the one-attempt job namespaces channels and publishes or
  broadcasts; the payload travels as a JSON string so JSONB keeps its order
- no jobs for the null driver or Centrifugo without an API key
2026-09-30 12:36:07 +02:00

84 lines
2.6 KiB
Go

package centrifugo
import (
"context"
"encoding/json"
"net/http"
"git.golem15.com/golem15/summercms/modules/backpack"
"git.golem15.com/golem15/summercms/modules/lighthouse"
)
// DriverName is the realtime.driver value of this driver.
const DriverName = "centrifugo"
func init() {
lighthouse.RegisterDriver(DriverName, func(app *backpack.App, svc *lighthouse.Service) (lighthouse.Driver, error) {
var cfg Config
if app != nil {
cfg = LoadConfig(app.Config)
} else {
cfg = LoadConfig(nil)
}
return NewDriver(svc, cfg, nil), nil
})
}
// Driver is the Centrifugo lighthouse.Driver: it publishes through Client
// and declares the token route (UserAuth) and the subscribe proxy route
// (ServerToServer).
type Driver struct {
svc *lighthouse.Service
cfg Config
client *Client
issuer *TokenIssuer
}
// NewDriver builds the driver for svc from cfg. hc may be nil (see
// NewClient).
func NewDriver(svc *lighthouse.Service, cfg Config, hc *http.Client) *Driver {
c := NewClient(cfg, hc)
c.log = svc.Logger()
return &Driver{
svc: svc,
cfg: cfg,
client: c,
issuer: NewTokenIssuer(cfg.TokenSecret, cfg.TokenTTL),
}
}
// Name returns "centrifugo".
func (d *Driver) Name() string { return DriverName }
// Enabled reports whether publishing is configured (an API key is set).
// Without it lighthouse enqueues no broadcast jobs.
func (d *Driver) Enabled() bool { return d.client.Enabled() }
// Config returns the driver's configuration.
func (d *Driver) Config() Config { return d.cfg }
// Client returns the HTTP API client.
func (d *Driver) Client() *Client { return d.client }
// Issuer returns the token issuer.
func (d *Driver) Issuer() *TokenIssuer { return d.issuer }
// Publish sends event to one channel through Client.Publish.
func (d *Driver) Publish(ctx context.Context, channel, event string, payload json.RawMessage) error {
return d.client.Publish(ctx, channel, event, payload)
}
// Broadcast sends event to several channels through Client.Broadcast.
func (d *Driver) Broadcast(ctx context.Context, channels []string, event string, payload json.RawMessage) error {
return d.client.Broadcast(ctx, channels, event, payload)
}
// Routes returns GET token_path (UserAuth, TokenHandler) and POST
// subscribe_path (ServerToServer, ProxyHandler).
func (d *Driver) Routes() []lighthouse.Route {
return []lighthouse.Route{
{Name: "token", Method: http.MethodGet, Path: d.cfg.TokenPath, Surface: lighthouse.UserAuth, Handler: TokenHandler(d.svc, d.issuer)},
{Name: "subscribe", Method: http.MethodPost, Path: d.cfg.SubscribePath, Surface: lighthouse.ServerToServer, Handler: ProxyHandler(d.svc, d.cfg)},
}
}