- 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
82 lines
2.3 KiB
Go
82 lines
2.3 KiB
Go
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})
|
|
}
|