diff --git a/modules/lighthouse/README.md b/modules/lighthouse/README.md index 40d10b3..119497d 100644 --- a/modules/lighthouse/README.md +++ b/modules/lighthouse/README.md @@ -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 `.` (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 ` 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. diff --git a/modules/lighthouse/broadcast.go b/modules/lighthouse/broadcast.go new file mode 100644 index 0000000..0aeb756 --- /dev/null +++ b/modules/lighthouse/broadcast.go @@ -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 +// (.). +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 . 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 .: 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 +} diff --git a/modules/lighthouse/centrifugo/driver.go b/modules/lighthouse/centrifugo/driver.go index cea644d..5c5e7fb 100644 --- a/modules/lighthouse/centrifugo/driver.go +++ b/modules/lighthouse/centrifugo/driver.go @@ -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 } diff --git a/modules/lighthouse/job.go b/modules/lighthouse/job.go new file mode 100644 index 0000000..bab9983 --- /dev/null +++ b/modules/lighthouse/job.go @@ -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}) +} diff --git a/modules/lighthouse/lighthouse.go b/modules/lighthouse/lighthouse.go index 79e0ea9..4bcf7c6 100644 --- a/modules/lighthouse/lighthouse.go +++ b/modules/lighthouse/lighthouse.go @@ -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 diff --git a/modules/lighthouse/suppress.go b/modules/lighthouse/suppress.go new file mode 100644 index 0000000..70b2afe --- /dev/null +++ b/modules/lighthouse/suppress.go @@ -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}) +}