package lighthouse import ( "context" "encoding/json" "fmt" "log/slog" "sort" "sync" "time" "git.golem15.com/golem15/summercms/modules/backpack" ) // Publisher sends an event to realtime channels. Channels are final names: // the broadcast job lowercases them and applies the namespace before it // calls a Publisher. type Publisher interface { // Publish sends event to one channel. Publish(ctx context.Context, channel, event string, payload json.RawMessage) error // Broadcast sends event to several channels in one call. Broadcast(ctx context.Context, channels []string, event string, payload json.RawMessage) error } // Driver is a realtime transport. Routes lists the HTTP endpoints the // transport needs (token issuing, subscribe authorization); the application // mounts them with Mount. type Driver interface { Publisher // Name is the realtime.driver value that selects the driver. Name() string // Routes returns the driver's HTTP routes; nil when it has none. Routes() []Route } // DriverFactory builds a driver for an app. svc is the Service being built: // its Registry, User lookup and Logger are available to handlers, but its // Driver is not set yet. type DriverFactory func(app *backpack.App, svc *Service) (Driver, error) // driverTable is the init-time driver registry, like database/sql's: it is // written only from package init functions and read when a Service is built. var driverTable = struct { mu sync.RWMutex factories map[string]DriverFactory }{factories: map[string]DriverFactory{}} // RegisterDriver makes a driver available under name. Call it from the // driver package's init function. A duplicate name, an empty name or a nil // factory panics. func RegisterDriver(name string, f DriverFactory) { if name == "" { panic("lighthouse: RegisterDriver with an empty name") } if f == nil { panic("lighthouse: RegisterDriver " + name + " with a nil factory") } driverTable.mu.Lock() defer driverTable.mu.Unlock() if _, dup := driverTable.factories[name]; dup { panic("lighthouse: RegisterDriver called twice for driver " + name) } driverTable.factories[name] = f } func driverFactory(name string) (DriverFactory, bool) { driverTable.mu.RLock() defer driverTable.mu.RUnlock() f, ok := driverTable.factories[name] return f, ok } func driverNames() []string { driverTable.mu.RLock() defer driverTable.mu.RUnlock() names := make([]string, 0, len(driverTable.factories)) for n := range driverTable.factories { names = append(names, n) } sort.Strings(names) return names } func init() { RegisterDriver("null", func(*backpack.App, *Service) (Driver, error) { return nullDriver{}, nil }) RegisterDriver("log", func(app *backpack.App, svc *Service) (Driver, error) { return &logDriver{log: svc.Logger()}, nil }) RegisterDriver("memory", func(*backpack.App, *Service) (Driver, error) { return NewMemoryDriver(), nil }) } // nullDriver discards every publication. Broadcast jobs are not even // enqueued while it is selected. type nullDriver struct{} func (nullDriver) Name() string { return "null" } func (nullDriver) Routes() []Route { return nil } func (nullDriver) Publish(context.Context, string, string, json.RawMessage) error { return nil } func (nullDriver) Broadcast(context.Context, []string, string, json.RawMessage) error { return nil } // logDriver logs the channel names and event of each publication at Info. // It never logs the payload. type logDriver struct { log *slog.Logger } func (d *logDriver) Name() string { return "log" } func (d *logDriver) Routes() []Route { return nil } func (d *logDriver) Publish(_ context.Context, channel, event string, _ json.RawMessage) error { d.log.Info("realtime publish", slog.String("channel", channel), slog.String("event", event)) return nil } func (d *logDriver) Broadcast(_ context.Context, channels []string, event string, _ json.RawMessage) error { d.log.Info("realtime broadcast", slog.Any("channels", channels), slog.String("event", event)) return nil } // Publication is one call recorded by MemoryDriver. type Publication struct { // Method is "publish" or "broadcast". Method string Channels []string Event string Payload json.RawMessage Timestamp time.Time } // MemoryDriver records publications for tests and payload comparisons. It // is safe for concurrent use. type MemoryDriver struct { mu sync.Mutex pubs []Publication } // NewMemoryDriver returns an empty memory driver. func NewMemoryDriver() *MemoryDriver { return &MemoryDriver{} } // Name returns "memory". func (d *MemoryDriver) Name() string { return "memory" } // Routes returns nil: the memory driver has no HTTP surface. func (d *MemoryDriver) Routes() []Route { return nil } // Publish records a single-channel publication. func (d *MemoryDriver) Publish(_ context.Context, channel, event string, payload json.RawMessage) error { return d.record("publish", []string{channel}, event, payload) } // Broadcast records a multi-channel publication. func (d *MemoryDriver) Broadcast(_ context.Context, channels []string, event string, payload json.RawMessage) error { return d.record("broadcast", channels, event, payload) } func (d *MemoryDriver) record(method string, channels []string, event string, payload json.RawMessage) error { if d == nil { return fmt.Errorf("lighthouse: memory driver is nil") } d.mu.Lock() defer d.mu.Unlock() d.pubs = append(d.pubs, clonePublication(Publication{ Method: method, Channels: channels, Event: event, Payload: payload, Timestamp: time.Now(), })) return nil } // Publications returns a copy of the recorded publications in call order. func (d *MemoryDriver) Publications() []Publication { if d == nil { return nil } d.mu.Lock() defer d.mu.Unlock() out := make([]Publication, len(d.pubs)) for i, p := range d.pubs { out[i] = clonePublication(p) } return out } func clonePublication(p Publication) Publication { p.Channels = append([]string(nil), p.Channels...) p.Payload = append(json.RawMessage(nil), p.Payload...) return p }