- 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
94 lines
2.9 KiB
Go
94 lines
2.9 KiB
Go
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})
|
|
}
|