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