- lighthouse: Service/From with realtime.driver selection, RegisterDriver registry, null/log/memory drivers, Route/Surface/Mount, users and actors - centrifugo: HTTP API client (apikey header, 2xx success, no request without a key), five-generator HS256 TokenIssuer, TokenHandler with the WinterCMS 401/503 bodies - module README and root modules table row
195 lines
5.9 KiB
Go
195 lines
5.9 KiB
Go
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
|
|
}
|