Files
summercms/modules/lighthouse/job.go
Jakub Zych 211c413273 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
2026-09-30 12:36:07 +02:00

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})
}