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
This commit is contained in:
Jakub Zych
2026-09-30 12:36:07 +02:00
parent 79fd705680
commit 211c413273
6 changed files with 736 additions and 6 deletions

View File

@@ -16,6 +16,8 @@ A driver may need HTTP endpoints, such as a token route for signed-in users or a
Channel authorization is transport-neutral too. Plugins register a `lighthouse.Authorizer` per channel namespace on the service's `lighthouse.Registry`. The driver's subscribe endpoint asks the authorizer of the channel's namespace on every subscribe, so a user who loses access is denied the next time the client subscribes. Nothing is cached. Channel authorization is transport-neutral too. Plugins register a `lighthouse.Authorizer` per channel namespace on the service's `lighthouse.Registry`. The driver's subscribe endpoint asks the authorizer of the channel's namespace on every subscribe, so a user who loses access is denied the next time the client subscribes. Nothing is cached.
Model broadcasts follow the WinterCMS `BroadcastableModel` trait with one change. A create, update or delete of a broadcastable model enqueues a River job (through `conga`) inside the write's own transaction, so nothing is published for a write that rolls back. The job publishes after commit, with one attempt and best effort: a failed publish is logged and never touches the write. Suppression is per model type and scoped to a context, and `lighthouse.Service.Emit` publishes one explicit summary event instead.
The `centrifugo` sub-package is the Centrifugo driver. It has a hand-rolled `net/http` client for the Centrifugo HTTP API, a token issuer with the claims of the WinterCMS `JwtTokenGenerator`, the token route handler, and the subscribe proxy handler. The `centrifugo` sub-package is the Centrifugo driver. It has a hand-rolled `net/http` client for the Centrifugo HTTP API, a token issuer with the claims of the WinterCMS `JwtTokenGenerator`, the token route handler, and the subscribe proxy handler.
## Features ## Features
@@ -30,6 +32,16 @@ The `centrifugo` sub-package is the Centrifugo driver. It has a hand-rolled `net
- `lighthouse.ChannelID` returns segment 1 converted with PHP's `(int)` cast (`lighthouse.PHPInt`): `5abc` is 5, `abc` is 0, and out-of-range values saturate. - `lighthouse.ChannelID` returns segment 1 converted with PHP's `(int)` cast (`lighthouse.PHPInt`): `5abc` is 5, `abc` is 0, and out-of-range values saturate.
- `lighthouse.FormatChannels` lowercases channel names and applies the broadcast namespace prefix. - `lighthouse.FormatChannels` lowercases channel names and applies the broadcast namespace prefix.
- Authorizer registry: `lighthouse.Registry` (from `lighthouse.Service.Registry`) maps namespaces to a `lighthouse.Authorizer` or `lighthouse.AuthorizerFunc`. Registering an empty namespace, a namespace that contains `:`, a nil authorizer or a namespace twice is an error. `lighthouse.Registry.Namespaces` is sorted. An authorizer returns `lighthouse.Allowed` (optionally with info, capabilities and overrides) or `lighthouse.Denied` with an internal reason that only reaches the logs. It reads the realtime client id with `lighthouse.ClientID`. - Authorizer registry: `lighthouse.Registry` (from `lighthouse.Service.Registry`) maps namespaces to a `lighthouse.Authorizer` or `lighthouse.AuthorizerFunc`. Registering an empty namespace, a namespace that contains `:`, a nil authorizer or a namespace twice is an error. `lighthouse.Registry.Namespaces` is sorted. An authorizer returns `lighthouse.Allowed` (optionally with info, capabilities and overrides) or `lighthouse.Denied` with an internal reason that only reaches the logs. It reads the realtime client id with `lighthouse.ClientID`.
- Model broadcasts. A model broadcasts when its pointer type implements `lighthouse.Broadcastable` (`BroadcastChannels(ctx, tx)`), or when a `lighthouse.Binding` is registered for it with `lighthouse.Bind`. A binding keeps payload code out of the model package. Optional overrides:
- `lighthouse.BroadcastPayloader` or `Binding.Payload` replaces the default payload `{"model":…,"actor":…,"timestamp":"…+00:00","ttl":60}`.
- `lighthouse.BroadcastAliaser` or `Binding.Alias` replaces the alias.
- `lighthouse.BroadcastFilter` or `Binding.ShouldBroadcast` can veto an action.
- `lighthouse.BroadcastTTLer` or `Binding.TTL` replaces the ttl.
The event name is `{action}.{alias}` lowercased: `lighthouse.ActionCreated`, `lighthouse.ActionUpdated` or `lighthouse.ActionDeleted`, then an alias that defaults to `<plugin>.<model>` (the Go package name, or the parent directory of a `models` package, and the type name). The payload builder receives a `lighthouse.Event` with the action, the `lighthouse.Actor`, the timestamp and the ttl. A soft delete counts as a delete. A delete's channels and payload are computed from a fresh read of the row before it is deleted, so deleting a model that holds only its id still broadcasts. An empty channel list means no broadcast.
- Transactional delivery. GORM callbacks (`lighthouse.CallbackAfterCreate`, `lighthouse.CallbackAfterUpdate`, `lighthouse.CallbackSnapshot` and `lighthouse.CallbackAfterDelete`) are installed through `lagoon.OnDatabase`. They enqueue a `lighthouse.BroadcastArgs` job on the write's `*sql.Tx`, on the `realtime.broadcast_queue` queue with MaxAttempts 1 and the `realtime.broadcast_timeout` timeout. Channel and payload queries and the enqueue run inside a savepoint, so a failure is rolled back to it, logged at Warn with channels and event (never the payload), and the write goes on. A write with a zero primary key, such as `Model(&T{}).Where(…).Updates(…)`, is not broadcast; bulk paths suppress and emit instead. The null driver, or a driver whose `Enabled` reports false (Centrifugo without an API key), gets no jobs.
- The broadcast job lowercases the channels and adds the `realtime.broadcast_namespace` prefix unless a channel already has it. It then publishes to one channel or broadcasts to several. A failure is logged as `realtime: broadcast failed` and is not retried. Delivery order across separate jobs is not guaranteed. The payload travels inside the job as a JSON string, so its key order survives Postgres JSONB.
- Suppression: `lighthouse.WithoutBroadcasting` silences one model type for writes made with the context it hands to its function. Other types still broadcast, and a write through an outer context is not suppressed. `lighthouse.Service.Emit` enqueues one `lighthouse.Broadcast` on the caller's transaction and returns its error. Together they turn N row events into one summary event.
- Centrifugo driver (`centrifugo.Driver`, driver name `centrifugo`): - Centrifugo driver (`centrifugo.Driver`, driver name `centrifugo`):
- `centrifugo.Client` POSTs `publish`, `broadcast`, `presence` and `unsubscribe` calls with `Authorization: apikey <key>` and a 5 s timeout. The publish body is `{"channel":…,"data":{"event":…,"payload":…,"timestamp":"…+00:00"}}`, with an empty payload sent as `[]`. Any 2xx status counts as success. With an empty API key nothing is sent and the call returns `centrifugo.ErrNotConfigured`. The key never appears in logs or errors. - `centrifugo.Client` POSTs `publish`, `broadcast`, `presence` and `unsubscribe` calls with `Authorization: apikey <key>` and a 5 s timeout. The publish body is `{"channel":…,"data":{"event":…,"payload":…,"timestamp":"…+00:00"}}`, with an empty payload sent as `[]`. Any 2xx status counts as success. With an empty API key nothing is sent and the call returns `centrifugo.ErrNotConfigured`. The key never appears in logs or errors.
- `centrifugo.TokenIssuer` signs HS256 tokens with five generators, the same as the WinterCMS generator: `ForUser` (claims `sub`, `exp`, `info` with only `name`), `Subscription`, `Anonymous` (`sub` "" and a 5-minute lifetime), `ForIdentifier` (an empty `info` is encoded as `[]`) and `SubscriptionForIdentifier`. It refuses to sign with an empty secret. - `centrifugo.TokenIssuer` signs HS256 tokens with five generators, the same as the WinterCMS generator: `ForUser` (claims `sub`, `exp`, `info` with only `name`), `Subscription`, `Anonymous` (`sub` "" and a 5-minute lifetime), `ForIdentifier` (an empty `info` is encoded as `[]`) and `SubscriptionForIdentifier`. It refuses to sign with an empty secret.
@@ -98,6 +110,39 @@ func (p *Plugin) Routes(r pact.Router) error {
} }
``` ```
A model package stays free of realtime code; the plugin binds the model at Boot:
```go
err := lighthouse.Bind[models.Post](svc, lighthouse.Binding[models.Post]{
Alias: "blog.post",
Channels: func(ctx context.Context, tx *gorm.DB, p *models.Post) ([]string, error) {
return []string{"blog:" + strconv.FormatUint(uint64(p.BlogID), 10)}, nil
},
})
```
A bulk import suppresses the per-row events and publishes one summary after commit:
```go
err := lighthouse.WithoutBroadcasting[models.Post](ctx, func(ctx context.Context) error {
return lagoon.Transaction(ctx, gdb, func(ctx context.Context, tx *gorm.DB) error {
for _, p := range posts {
if err := tx.Create(&p).Error; err != nil {
return err
}
}
return svc.Emit(ctx, tx, lighthouse.Broadcast{
Channels: []string{"blog:7"},
Event: "blog.bulk_updated",
Payload: struct {
Reason string `json:"reason"`
Count int `json:"count"`
}{"import", len(posts)},
})
})
})
```
Tests select the memory driver and read what was published: Tests select the memory driver and read what was published:
```go ```go
@@ -129,6 +174,15 @@ for _, pub := range mem.Publications() {
| `lighthouse.Registry`, `lighthouse.NewRegistry` | Namespace to authorizer map: `Register`, `Get`, `Namespaces`. | | `lighthouse.Registry`, `lighthouse.NewRegistry` | Namespace to authorizer map: `Register`, `Get`, `Namespaces`. |
| `lighthouse.ParseChannel`, `lighthouse.ChannelID`, `lighthouse.PHPInt`, `lighthouse.FormatChannels` | Channel rules. | | `lighthouse.ParseChannel`, `lighthouse.ChannelID`, `lighthouse.PHPInt`, `lighthouse.FormatChannels` | Channel rules. |
| `lighthouse.WithClientID`, `lighthouse.ClientID` | The realtime client id of a subscribe request, carried in the context. | | `lighthouse.WithClientID`, `lighthouse.ClientID` | The realtime client id of a subscribe request, carried in the context. |
| `lighthouse.Action`, `lighthouse.ActionCreated`, `lighthouse.ActionUpdated`, `lighthouse.ActionDeleted` | Broadcast actions. |
| `lighthouse.Event` | `Action`, `Actor`, `Timestamp`, `TTL` of a change. |
| `lighthouse.Broadcastable`, `lighthouse.BroadcastPayloader`, `lighthouse.BroadcastAliaser`, `lighthouse.BroadcastFilter`, `lighthouse.BroadcastTTLer` | The model-method broadcast contract. |
| `lighthouse.Binding`, `lighthouse.Bind` | Broadcast contract registered from outside the model package. |
| `lighthouse.WithoutBroadcasting` | Suppresses one model type for writes made with the given context. |
| `lighthouse.Broadcast`, `lighthouse.Service.Emit` | One explicit event enqueued on the caller's transaction. |
| `lighthouse.BroadcastArgs` | The River job (kind `summer.broadcast`). |
| `lighthouse.DefaultTTL` | The default payload ttl, 60 seconds. |
| `lighthouse.CallbackSnapshot`, `lighthouse.CallbackAfterCreate`, `lighthouse.CallbackAfterUpdate`, `lighthouse.CallbackAfterDelete` | Names of the GORM callbacks. |
| `lighthouse.User`, `lighthouse.UserLookup` | A user id with a display name, and the application's lookup. | | `lighthouse.User`, `lighthouse.UserLookup` | A user id with a display name, and the application's lookup. |
| `lighthouse.Actor`, `lighthouse.SystemActor` | Who caused a broadcast: `{"user_id":…,"name":…}`. | | `lighthouse.Actor`, `lighthouse.SystemActor` | Who caused a broadcast: `{"user_id":…,"name":…}`. |
| `lighthouse.DurationSetting(cfg, path)` | Reads a duration string or an integer number of seconds. | | `lighthouse.DurationSetting(cfg, path)` | Reads a duration string or an integer number of seconds. |
@@ -144,7 +198,7 @@ for _, pub := range mem.Publications() {
| `centrifugo.TokenIssuer`, `centrifugo.NewTokenIssuer` | HS256 token generators: `ForUser`, `Subscription`, `Anonymous`, `ForIdentifier`, `SubscriptionForIdentifier`, `Configured`. | | `centrifugo.TokenIssuer`, `centrifugo.NewTokenIssuer` | HS256 token generators: `ForUser`, `Subscription`, `Anonymous`, `ForIdentifier`, `SubscriptionForIdentifier`, `Configured`. |
| `centrifugo.TokenHandler(svc, issuer)` | The token route handler. | | `centrifugo.TokenHandler(svc, issuer)` | The token route handler. |
| `centrifugo.ProxyHandler(svc, cfg)` | The subscribe proxy handler. | | `centrifugo.ProxyHandler(svc, cfg)` | The subscribe proxy handler. |
| `centrifugo.Driver`, `centrifugo.NewDriver`, `centrifugo.DriverName` | The `lighthouse.Driver`, with `Client`, `Issuer` and `Config` accessors. | | `centrifugo.Driver`, `centrifugo.NewDriver`, `centrifugo.DriverName` | The `lighthouse.Driver`, with `Client`, `Issuer`, `Config` and `Enabled` (an API key is set). |
| `centrifugo.ErrNotConfigured` | Returned when the API key or token secret an operation needs is empty. | | `centrifugo.ErrNotConfigured` | Returned when the API key or token secret an operation needs is empty. |
## Configuration ## Configuration
@@ -166,7 +220,8 @@ for _, pub := range mem.Publications() {
## Dependencies ## Dependencies
- `backpack`, `bouncer`, `compass`, `pact` and `wire` from this repository; the centrifugo driver also uses `surf` for the client IP. - `backpack`, `bouncer`, `compass`, `conga` (the broadcast job), `lagoon` (callback installation), `pact` and `wire` from this repository; the centrifugo driver also uses `surf` for the client IP.
- `gorm.io/gorm` (broadcast callbacks).
- `github.com/golang-jwt/jwt/v5` (centrifugo token signing). - `github.com/golang-jwt/jwt/v5` (centrifugo token signing).
- The Centrifugo client is plain `net/http`; no Centrifugo SDK is used. - The Centrifugo client is plain `net/http`; no Centrifugo SDK is used.
@@ -176,4 +231,4 @@ for _, pub := range mem.Publications() {
go test ./modules/lighthouse/... go test ./modules/lighthouse/...
``` ```
A test selects `realtime.driver: memory` and reads `lighthouse.MemoryDriver.Publications`, or points `realtime.centrifugo.api_url` at an `httptest` server to see the exact Centrifugo requests. A test selects `realtime.driver: memory` and reads `lighthouse.MemoryDriver.Publications`, or points `realtime.centrifugo.api_url` at an `httptest` server to see the exact Centrifugo requests. Broadcast tests need a running `conga` worker (`conga.StartWorker`) to deliver the jobs.

View File

@@ -0,0 +1,479 @@
package lighthouse
import (
"bytes"
"context"
"database/sql"
"encoding/json"
"fmt"
"log/slog"
"reflect"
"strings"
"time"
"git.golem15.com/golem15/summercms/modules/wire"
"gorm.io/gorm"
"gorm.io/gorm/schema"
)
// Action is a model change that is broadcast.
type Action string
const (
ActionCreated Action = "created"
ActionUpdated Action = "updated"
ActionDeleted Action = "deleted"
)
// DefaultTTL is the ttl of the default broadcast payload, in seconds.
const DefaultTTL = 60
// Names of the GORM callbacks that broadcast model changes.
const (
CallbackSnapshot = "lighthouse:snapshot"
CallbackAfterCreate = "lighthouse:after_create"
CallbackAfterUpdate = "lighthouse:after_update"
CallbackAfterDelete = "lighthouse:after_delete"
)
const (
snapshotKey = "lighthouse:snapshot"
savepoint = "lighthouse_broadcast"
)
// Event describes the change being broadcast. Payload builders receive it.
type Event struct {
Action Action
Actor Actor
Timestamp wire.Time
TTL int
}
// Broadcastable is implemented (on the pointer receiver) by models that
// broadcast their creates, updates and deletes. An empty channel list means
// no broadcast.
type Broadcastable interface {
BroadcastChannels(ctx context.Context, tx *gorm.DB) ([]string, error)
}
// BroadcastPayloader replaces the default payload {model, actor, timestamp,
// ttl}.
type BroadcastPayloader interface {
BroadcastPayload(ctx context.Context, tx *gorm.DB, ev Event) (any, error)
}
// BroadcastAliaser replaces the default alias of the event name
// (<plugin>.<model>).
type BroadcastAliaser interface {
BroadcastAlias() string
}
// BroadcastFilter can veto the broadcast of an action.
type BroadcastFilter interface {
ShouldBroadcast(Action) bool
}
// BroadcastTTLer replaces the default ttl of 60 seconds.
type BroadcastTTLer interface {
BroadcastTTL() int
}
// Binding makes T broadcastable without methods on T, for payloads that
// need packages a model package must not import. Channels is required;
// every other field falls back to the defaults of Broadcastable.
type Binding[T any] struct {
// Alias replaces the default <plugin>.<model> of the event name.
Alias string
// Channels returns the channels of m; empty means no broadcast.
Channels func(ctx context.Context, tx *gorm.DB, m *T) ([]string, error)
// Payload replaces the default {model, actor, timestamp, ttl} payload.
Payload func(ctx context.Context, tx *gorm.DB, m *T, ev Event) (any, error)
// ShouldBroadcast can veto an action.
ShouldBroadcast func(Action) bool
// TTL replaces the default ttl of 60 seconds.
TTL int
}
// Bind registers the broadcast binding of model type T (a struct type).
// A nil Channels function or a second binding for T is an error. A bound
// type is broadcast through its binding even when it also implements
// Broadcastable.
func Bind[T any](svc *Service, b Binding[T]) error {
if svc == nil {
return fmt.Errorf("lighthouse: Bind on a nil service")
}
typ := reflect.TypeFor[T]()
if typ.Kind() != reflect.Struct {
return fmt.Errorf("lighthouse: Bind needs a struct model type, got %s", typ)
}
if b.Channels == nil {
return fmt.Errorf("lighthouse: Bind %s: Channels is nil", typ)
}
h := &handler{
alias: b.Alias,
ttl: b.TTL,
channels: func(ctx context.Context, tx *gorm.DB, m any) ([]string, error) {
return b.Channels(ctx, tx, m.(*T))
},
}
if h.alias == "" {
h.alias = defaultAlias(typ)
}
if h.ttl <= 0 {
h.ttl = DefaultTTL
}
if b.Payload != nil {
h.payload = func(ctx context.Context, tx *gorm.DB, m any, ev Event) (any, error) {
return b.Payload(ctx, tx, m.(*T), ev)
}
}
if b.ShouldBroadcast != nil {
h.should = func(_ any, a Action) bool { return b.ShouldBroadcast(a) }
}
svc.mu.Lock()
defer svc.mu.Unlock()
if svc.bindings == nil {
svc.bindings = map[reflect.Type]*handler{}
}
if _, dup := svc.bindings[typ]; dup {
return fmt.Errorf("lighthouse: %s is already bound", typ)
}
svc.bindings[typ] = h
return nil
}
// handler is the type-erased broadcast contract of one model type. The
// model argument is always a pointer to the model struct.
type handler struct {
alias string
ttl int
channels func(ctx context.Context, tx *gorm.DB, m any) ([]string, error)
payload func(ctx context.Context, tx *gorm.DB, m any, ev Event) (any, error)
should func(m any, a Action) bool
}
var broadcastableType = reflect.TypeFor[Broadcastable]()
// handlerFor returns the binding of typ, else a handler over the model's
// Broadcastable methods, else nil.
func (s *Service) handlerFor(typ reflect.Type) *handler {
s.mu.RLock()
h := s.bindings[typ]
s.mu.RUnlock()
if h != nil {
return h
}
ptr := reflect.PointerTo(typ)
if !ptr.Implements(broadcastableType) {
return nil
}
h = &handler{
alias: defaultAlias(typ),
ttl: DefaultTTL,
channels: func(ctx context.Context, tx *gorm.DB, m any) ([]string, error) {
return m.(Broadcastable).BroadcastChannels(ctx, tx)
},
should: func(m any, a Action) bool {
if f, ok := m.(BroadcastFilter); ok {
return f.ShouldBroadcast(a)
}
return true
},
}
zero := reflect.New(typ).Interface()
if a, ok := zero.(BroadcastAliaser); ok {
if alias := a.BroadcastAlias(); alias != "" {
h.alias = alias
}
}
if t, ok := zero.(BroadcastTTLer); ok {
if ttl := t.BroadcastTTL(); ttl > 0 {
h.ttl = ttl
}
}
if _, ok := zero.(BroadcastPayloader); ok {
h.payload = func(ctx context.Context, tx *gorm.DB, m any, ev Event) (any, error) {
return m.(BroadcastPayloader).BroadcastPayload(ctx, tx, ev)
}
}
return h
}
// defaultAlias is <plugin>.<model>: the last Go package path segment (the
// one before it when the last is "models") and the lowercased type name.
func defaultAlias(typ reflect.Type) string {
parts := strings.Split(typ.PkgPath(), "/")
plugin := "unknown"
if n := len(parts); n > 0 && parts[n-1] != "" {
plugin = parts[n-1]
if plugin == "models" && n > 1 {
plugin = parts[n-2]
}
}
return strings.ToLower(plugin + "." + typ.Name())
}
func (h *handler) eventName(a Action) string {
return strings.ToLower(string(a) + "." + h.alias)
}
func (h *handler) allows(m any, a Action) bool {
return h.should == nil || h.should(m, a)
}
// defaultPayload is the WinterCMS BroadcastableModel payload.
type defaultPayload struct {
Model any `json:"model"`
Actor Actor `json:"actor"`
Timestamp wire.Time `json:"timestamp"`
TTL int `json:"ttl"`
}
// installCallbacks registers the broadcast callbacks on gdb, replacing
// earlier ones so a handle shared by several apps broadcasts through the
// most recent service.
func (s *Service) installCallbacks(gdb *gorm.DB) error {
cb := gdb.Callback()
if cb.Create().Get(CallbackAfterCreate) == nil {
if err := cb.Create().After("gorm:after_create").Register(CallbackAfterCreate, s.afterCreate); err != nil {
return err
}
} else if err := cb.Create().Replace(CallbackAfterCreate, s.afterCreate); err != nil {
return err
}
if cb.Update().Get(CallbackAfterUpdate) == nil {
if err := cb.Update().After("gorm:after_update").Register(CallbackAfterUpdate, s.afterUpdate); err != nil {
return err
}
} else if err := cb.Update().Replace(CallbackAfterUpdate, s.afterUpdate); err != nil {
return err
}
if cb.Delete().Get(CallbackSnapshot) == nil {
if err := cb.Delete().Before("gorm:before_delete").Register(CallbackSnapshot, s.snapshot); err != nil {
return err
}
} else if err := cb.Delete().Replace(CallbackSnapshot, s.snapshot); err != nil {
return err
}
if cb.Delete().Get(CallbackAfterDelete) == nil {
if err := cb.Delete().After("gorm:after_delete").Register(CallbackAfterDelete, s.afterDelete); err != nil {
return err
}
} else if err := cb.Delete().Replace(CallbackAfterDelete, s.afterDelete); err != nil {
return err
}
return nil
}
// broadcasting reports whether the driver publishes anything. The null
// driver, or a driver whose Enabled reports false (Centrifugo without an
// API key), gets no broadcast jobs at all.
func (s *Service) broadcasting() bool {
switch d := s.driver.(type) {
case nil:
return false
case nullDriver:
return false
case interface{ Enabled() bool }:
return d.Enabled()
default:
return true
}
}
// target is one written model: its handler and a pointer to it.
type target struct {
h *handler
typ reflect.Type
model any
value reflect.Value
}
// targets returns the broadcastable, non-suppressed models of the
// statement that have a primary key. A batch update through an empty model
// has a zero key and is skipped: bulk paths suppress and Emit instead.
func (s *Service) targets(db *gorm.DB) []target {
if db.Error != nil || db.Statement == nil || db.Statement.Schema == nil || !s.broadcasting() {
return nil
}
sch := db.Statement.Schema
h := s.handlerFor(sch.ModelType)
if h == nil || suppressed(db.Statement.Context, sch.ModelType) {
return nil
}
pk := sch.PrioritizedPrimaryField
if pk == nil {
return nil
}
ctx := db.Statement.Context
var out []target
add := func(v reflect.Value) {
for v.Kind() == reflect.Pointer {
if v.IsNil() {
return
}
v = v.Elem()
}
if v.Kind() != reflect.Struct || v.Type() != sch.ModelType {
return
}
if _, zero := pk.ValueOf(ctx, v); zero {
return
}
if !v.CanAddr() {
c := reflect.New(v.Type()).Elem()
c.Set(v)
v = c
}
out = append(out, target{h: h, typ: sch.ModelType, model: v.Addr().Interface(), value: v})
}
rv := db.Statement.ReflectValue
switch rv.Kind() {
case reflect.Slice, reflect.Array:
for i := 0; i < rv.Len(); i++ {
add(rv.Index(i))
}
default:
add(rv)
}
return out
}
func (s *Service) afterCreate(db *gorm.DB) { s.afterWrite(db, ActionCreated) }
func (s *Service) afterUpdate(db *gorm.DB) { s.afterWrite(db, ActionUpdated) }
func (s *Service) afterWrite(db *gorm.DB, action Action) {
for _, t := range s.targets(db) {
if !t.h.allows(t.model, action) {
continue
}
s.inSavepoint(db, func(tx *gorm.DB) error {
args, ok, err := s.prepare(tx, t, action)
if err != nil || !ok {
return err
}
return s.enqueue(tx, args)
})
}
}
// snapshot runs before a delete: it reloads each row, computes its channels
// and deleted payload while the row still exists, and keeps them for
// afterDelete.
func (s *Service) snapshot(db *gorm.DB) {
var pending []BroadcastArgs
for _, t := range s.targets(db) {
if !t.h.allows(t.model, ActionDeleted) {
continue
}
s.inSavepoint(db, func(tx *gorm.DB) error {
full := reflect.New(t.typ)
pk, _ := db.Statement.Schema.PrioritizedPrimaryField.ValueOf(tx.Statement.Context, t.value)
if err := tx.Unscoped().Where(fmt.Sprintf("%s = ?", quoteColumn(db, db.Statement.Schema.PrioritizedPrimaryField)), pk).Take(full.Interface()).Error; err == nil {
t.model = full.Interface()
}
args, ok, err := s.prepare(tx, t, ActionDeleted)
if err == nil && ok {
pending = append(pending, args)
}
return err
})
}
if len(pending) > 0 {
db.InstanceSet(snapshotKey, pending)
}
}
// afterDelete enqueues the snapshots of a delete that succeeded.
func (s *Service) afterDelete(db *gorm.DB) {
if db.Error != nil {
return
}
v, ok := db.InstanceGet(snapshotKey)
if !ok {
return
}
pending, _ := v.([]BroadcastArgs)
for _, args := range pending {
s.inSavepoint(db, func(tx *gorm.DB) error { return s.enqueue(tx, args) })
}
}
func quoteColumn(db *gorm.DB, f *schema.Field) string {
return db.Statement.Quote(f.DBName)
}
// prepare computes the channels, event name and payload of one change.
// ok is false when the model has no channels.
func (s *Service) prepare(tx *gorm.DB, t target, action Action) (BroadcastArgs, bool, error) {
ctx := tx.Statement.Context
channels, err := t.h.channels(ctx, tx, t.model)
if err != nil {
return BroadcastArgs{}, false, fmt.Errorf("channels: %w", err)
}
if len(channels) == 0 {
return BroadcastArgs{}, false, nil
}
ev := Event{Action: action, Actor: s.Actor(ctx), Timestamp: wire.Time{Time: time.Now()}, TTL: t.h.ttl}
var payload any
if t.h.payload != nil {
payload, err = t.h.payload(ctx, tx, t.model, ev)
if err != nil {
return BroadcastArgs{}, false, fmt.Errorf("payload: %w", err)
}
} else {
payload = defaultPayload{Model: t.model, Actor: ev.Actor, Timestamp: ev.Timestamp, TTL: ev.TTL}
}
raw, err := marshalPayload(payload)
if err != nil {
return BroadcastArgs{}, false, err
}
return BroadcastArgs{Channels: channels, Event: t.h.eventName(action), Payload: raw}, true, nil
}
// inSavepoint runs fn on a fresh session of db's connection. Inside a
// transaction fn runs in a savepoint, so a failed query or enqueue is rolled
// back to it and never aborts the write. Failures are logged, not returned.
func (s *Service) inSavepoint(db *gorm.DB, fn func(tx *gorm.DB) error) {
ctx := db.Statement.Context
if ctx == nil {
ctx = context.Background()
}
tx := db.Session(&gorm.Session{NewDB: true, Context: ctx})
_, inTx := tx.Statement.ConnPool.(*sql.Tx)
if inTx {
if err := tx.SavePoint(savepoint).Error; err != nil {
s.Logger().Warn("realtime: broadcast skipped", slog.String("error", err.Error()))
return
}
}
var err error
func() {
defer func() {
if r := recover(); r != nil {
err = fmt.Errorf("panic: %v", r)
}
}()
err = fn(tx)
}()
if err != nil {
s.Logger().Warn("realtime: broadcast skipped", slog.String("error", err.Error()))
if inTx {
tx.RollbackTo(savepoint)
}
return
}
if inTx {
tx.Exec("RELEASE SAVEPOINT " + savepoint)
}
}
func marshalPayload(v any) (json.RawMessage, error) {
var buf bytes.Buffer
enc := json.NewEncoder(&buf)
enc.SetEscapeHTML(false)
if err := enc.Encode(v); err != nil {
return nil, fmt.Errorf("lighthouse: encode payload: %w", err)
}
return json.RawMessage(bytes.TrimSuffix(buf.Bytes(), []byte("\n"))), nil
}

View File

@@ -50,6 +50,10 @@ func NewDriver(svc *lighthouse.Service, cfg Config, hc *http.Client) *Driver {
// Name returns "centrifugo". // Name returns "centrifugo".
func (d *Driver) Name() string { return DriverName } 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. // Config returns the driver's configuration.
func (d *Driver) Config() Config { return d.cfg } func (d *Driver) Config() Config { return d.cfg }

93
modules/lighthouse/job.go Normal file
View File

@@ -0,0 +1,93 @@
package lighthouse
import (
"context"
"encoding/json"
"log/slog"
"git.golem15.com/golem15/summercms/modules/conga"
"gorm.io/gorm"
)
// BroadcastArgs is the River job of one broadcast. Channels are not yet
// namespaced; the worker formats them.
type BroadcastArgs struct {
Channels []string `json:"channels"`
Event string `json:"event"`
Payload json.RawMessage `json:"payload"`
}
// Kind is the River job kind, "summer.broadcast".
func (BroadcastArgs) Kind() string { return "summer.broadcast" }
// broadcastArgsJSON is the stored form of BroadcastArgs. River keeps job
// args in a JSONB column, and JSONB reorders object keys, so the payload
// travels as a JSON string to keep its key order byte for byte.
type broadcastArgsJSON struct {
Channels []string `json:"channels"`
Event string `json:"event"`
Payload string `json:"payload"`
}
// MarshalJSON stores Payload as a JSON string.
func (a BroadcastArgs) MarshalJSON() ([]byte, error) {
return json.Marshal(broadcastArgsJSON{Channels: a.Channels, Event: a.Event, Payload: string(a.Payload)})
}
// UnmarshalJSON reads the stored form written by MarshalJSON.
func (a *BroadcastArgs) UnmarshalJSON(b []byte) error {
var w broadcastArgsJSON
if err := json.Unmarshal(b, &w); err != nil {
return err
}
a.Channels, a.Event = w.Channels, w.Event
a.Payload = nil
if w.Payload != "" {
a.Payload = json.RawMessage(w.Payload)
}
return nil
}
// registerJob registers the broadcast job on the app's job manager: one
// attempt on the broadcast queue with the broadcast timeout.
func (s *Service) registerJob() error {
m, err := conga.From(s.app)
if err != nil {
return err
}
return m.Register(conga.Job(s.deliver, conga.OnQueue(s.queue), conga.MaxAttempts(1), conga.Timeout(s.timeout)))
}
// deliver publishes one broadcast: the channels are lowercased and
// namespaced, one channel is a publish and several a broadcast. A failure
// is logged and swallowed, so River never retries it.
func (s *Service) deliver(ctx context.Context, args BroadcastArgs) error {
channels := FormatChannels(s.namespace, args.Channels)
if len(channels) == 0 || s.driver == nil {
return nil
}
var err error
if len(channels) == 1 {
err = s.driver.Publish(ctx, channels[0], args.Event, args.Payload)
} else {
err = s.driver.Broadcast(ctx, channels, args.Event, args.Payload)
}
if err != nil {
s.Logger().Warn("realtime: broadcast failed",
slog.Any("channels", channels), slog.String("event", args.Event), slog.String("error", err.Error()))
}
return nil
}
// enqueue inserts the broadcast job on db's transaction when there is one.
func (s *Service) enqueue(db *gorm.DB, args BroadcastArgs) error {
m, err := conga.From(s.app)
if err != nil {
return err
}
ctx := context.Background()
if db != nil && db.Statement != nil && db.Statement.Context != nil {
ctx = db.Statement.Context
}
return m.Enqueue(ctx, db, args, conga.EnqueueOpts{Queue: s.queue, MaxAttempts: 1})
}

View File

@@ -8,14 +8,18 @@
package lighthouse package lighthouse
import ( import (
"database/sql"
"fmt" "fmt"
"log/slog" "log/slog"
"reflect"
"strings" "strings"
"sync" "sync"
"time" "time"
"git.golem15.com/golem15/summercms/modules/backpack" "git.golem15.com/golem15/summercms/modules/backpack"
"git.golem15.com/golem15/summercms/modules/compass" "git.golem15.com/golem15/summercms/modules/compass"
"git.golem15.com/golem15/summercms/modules/lagoon"
"gorm.io/gorm"
) )
const ( const (
@@ -39,8 +43,9 @@ type Service struct {
queue string queue string
timeout time.Duration timeout time.Duration
mu sync.RWMutex mu sync.RWMutex
lookup UserLookup lookup UserLookup
bindings map[reflect.Type]*handler
} }
// From returns the app's Service, building and publishing it on first use. // From returns the app's Service, building and publishing it on first use.
@@ -48,7 +53,10 @@ type Service struct {
// realtime.broadcast_namespace (default ""), realtime.broadcast_queue // realtime.broadcast_namespace (default ""), realtime.broadcast_queue
// (default "broadcasts") and realtime.broadcast_timeout (default 5s; an // (default "broadcasts") and realtime.broadcast_timeout (default 5s; an
// integer number of seconds or a duration string) and builds the driver // integer number of seconds or a duration string) and builds the driver
// through its registered factory. An unknown driver name is an error. // through its registered factory. An unknown driver name is an error. It
// also registers the broadcast job with the app's conga manager (so call it
// from Boot, before a worker starts) and installs the broadcast GORM
// callbacks through lagoon.OnDatabase.
func From(app *backpack.App) (*Service, error) { func From(app *backpack.App) (*Service, error) {
if app == nil { if app == nil {
return nil, fmt.Errorf("lighthouse: app is nil") return nil, fmt.Errorf("lighthouse: app is nil")
@@ -75,6 +83,16 @@ func From(app *backpack.App) (*Service, error) {
return nil, fmt.Errorf("lighthouse: driver %s returned nil", name) return nil, fmt.Errorf("lighthouse: driver %s returned nil", name)
} }
svc.driver = d svc.driver = d
if err := svc.registerJob(); err != nil {
return nil, fmt.Errorf("lighthouse: register broadcast job: %w", err)
}
// The broadcast callbacks need the GORM handle, which serve publishes
// after plugins boot.
if err := lagoon.OnDatabase(app, func(_ *sql.DB, gdb *gorm.DB) error {
return svc.installCallbacks(gdb)
}); err != nil {
return nil, fmt.Errorf("lighthouse: install broadcast callbacks: %w", err)
}
if err := app.Publish(svc); err != nil { if err := app.Publish(svc); err != nil {
if existing, ok := app.Lookup[*Service](); ok && existing != nil { if existing, ok := app.Lookup[*Service](); ok && existing != nil {
return existing, nil return existing, nil

View File

@@ -0,0 +1,81 @@
package lighthouse
import (
"context"
"fmt"
"reflect"
"gorm.io/gorm"
)
type suppressKey struct{}
// WithoutBroadcasting runs fn with a ctx in which writes of model type T
// are not broadcast. Other types still broadcast. Only writes made with the
// ctx handed to fn are suppressed, so fn must write through
// gdb.WithContext(ctx) (or lagoon.Transaction(ctx, ...)); a write through an
// outer ctx still broadcasts. Pair it with Service.Emit to publish one
// summary event for a bulk write.
func WithoutBroadcasting[T any](ctx context.Context, fn func(ctx context.Context) error) error {
if ctx == nil {
ctx = context.Background()
}
typ := reflect.TypeFor[T]()
for typ.Kind() == reflect.Pointer {
typ = typ.Elem()
}
prev, _ := ctx.Value(suppressKey{}).(map[reflect.Type]struct{})
next := make(map[reflect.Type]struct{}, len(prev)+1)
for t := range prev {
next[t] = struct{}{}
}
next[typ] = struct{}{}
return fn(context.WithValue(ctx, suppressKey{}, next))
}
func suppressed(ctx context.Context, typ reflect.Type) bool {
if ctx == nil {
return false
}
set, _ := ctx.Value(suppressKey{}).(map[reflect.Type]struct{})
_, ok := set[typ]
return ok
}
// Broadcast is an explicit event published by Service.Emit.
type Broadcast struct {
// Channels are the channel names before namespacing.
Channels []string
// Event is the event name, for example "collection.bulk_updated".
Event string
// Payload is encoded as JSON; use an ordered struct for a stable key
// order.
Payload any
}
// Emit enqueues one broadcast job for b on db's transaction, so it is
// published only when that transaction commits. No channels, or a driver
// that publishes nothing, enqueues nothing. Unlike model broadcasts, a
// failure is returned: the caller asked for this event.
func (s *Service) Emit(ctx context.Context, db *gorm.DB, b Broadcast) error {
if s == nil {
return fmt.Errorf("lighthouse: Emit on a nil service")
}
if len(b.Channels) == 0 || !s.broadcasting() {
return nil
}
if b.Event == "" {
return fmt.Errorf("lighthouse: Emit without an event name")
}
raw, err := marshalPayload(b.Payload)
if err != nil {
return err
}
if ctx == nil {
ctx = context.Background()
}
if db != nil {
db = db.WithContext(ctx)
}
return s.enqueue(db, BroadcastArgs{Channels: append([]string(nil), b.Channels...), Event: b.Event, Payload: raw})
}