feat(11-03): add the lighthouse realtime package and its Centrifugo driver
- 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
This commit is contained in:
194
modules/lighthouse/drivers.go
Normal file
194
modules/lighthouse/drivers.go
Normal file
@@ -0,0 +1,194 @@
|
||||
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
|
||||
}
|
||||
Reference in New Issue
Block a user