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:
@@ -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.
|
||||
|
||||
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.
|
||||
|
||||
## 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.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`.
|
||||
- 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.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.
|
||||
@@ -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:
|
||||
|
||||
```go
|
||||
@@ -129,6 +174,15 @@ for _, pub := range mem.Publications() {
|
||||
| `lighthouse.Registry`, `lighthouse.NewRegistry` | Namespace to authorizer map: `Register`, `Get`, `Namespaces`. |
|
||||
| `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.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.Actor`, `lighthouse.SystemActor` | Who caused a broadcast: `{"user_id":…,"name":…}`. |
|
||||
| `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.TokenHandler(svc, issuer)` | The token route 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. |
|
||||
|
||||
## Configuration
|
||||
@@ -166,7 +220,8 @@ for _, pub := range mem.Publications() {
|
||||
|
||||
## 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).
|
||||
- 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/...
|
||||
```
|
||||
|
||||
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.
|
||||
|
||||
479
modules/lighthouse/broadcast.go
Normal file
479
modules/lighthouse/broadcast.go
Normal 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
|
||||
}
|
||||
@@ -50,6 +50,10 @@ func NewDriver(svc *lighthouse.Service, cfg Config, hc *http.Client) *Driver {
|
||||
// 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 }
|
||||
|
||||
|
||||
93
modules/lighthouse/job.go
Normal file
93
modules/lighthouse/job.go
Normal 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})
|
||||
}
|
||||
@@ -8,14 +8,18 @@
|
||||
package lighthouse
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"reflect"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"git.golem15.com/golem15/summercms/modules/backpack"
|
||||
"git.golem15.com/golem15/summercms/modules/compass"
|
||||
"git.golem15.com/golem15/summercms/modules/lagoon"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -39,8 +43,9 @@ type Service struct {
|
||||
queue string
|
||||
timeout time.Duration
|
||||
|
||||
mu sync.RWMutex
|
||||
lookup UserLookup
|
||||
mu sync.RWMutex
|
||||
lookup UserLookup
|
||||
bindings map[reflect.Type]*handler
|
||||
}
|
||||
|
||||
// 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
|
||||
// (default "broadcasts") and realtime.broadcast_timeout (default 5s; an
|
||||
// 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) {
|
||||
if app == 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)
|
||||
}
|
||||
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 existing, ok := app.Lookup[*Service](); ok && existing != nil {
|
||||
return existing, nil
|
||||
|
||||
81
modules/lighthouse/suppress.go
Normal file
81
modules/lighthouse/suppress.go
Normal 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})
|
||||
}
|
||||
Reference in New Issue
Block a user