- lighthouse: Registry of namespace authorizers (Result, Allowed, Denied), ParseChannel, ChannelID with PHP (int)-cast semantics (PHPInt, pinned by a php -r table test), FormatChannels, WithClientID/ClientID - centrifugo: ProxyHandler (constant-time X-Centrifugo-Secret, HTTP 200 generic deny, info [] on allow, presence allow/override merge, 64 KiB body cap) mounted as the ServerToServer subscribe route - README: proxy contract, registry and channel rules
185 lines
4.8 KiB
Go
185 lines
4.8 KiB
Go
// Package lighthouse is the transport-neutral realtime layer: a publisher
|
|
// interface with pluggable drivers, subscribe-time channel authorization,
|
|
// and model broadcasts enqueued in the write transaction.
|
|
//
|
|
// Models, authorizers and application code talk only to this package. A
|
|
// driver package (for example lighthouse/centrifugo) is imported for its
|
|
// side effect of registering itself, and is chosen with realtime.driver.
|
|
package lighthouse
|
|
|
|
import (
|
|
"fmt"
|
|
"log/slog"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"git.golem15.com/golem15/summercms/modules/backpack"
|
|
"git.golem15.com/golem15/summercms/modules/compass"
|
|
)
|
|
|
|
const (
|
|
// DefaultDriver is the driver used when realtime.driver is empty.
|
|
DefaultDriver = "null"
|
|
// DefaultQueue is the River queue of broadcast jobs when
|
|
// realtime.broadcast_queue is empty.
|
|
DefaultQueue = "broadcasts"
|
|
// DefaultTimeout is the broadcast job timeout when
|
|
// realtime.broadcast_timeout is empty.
|
|
DefaultTimeout = 5 * time.Second
|
|
)
|
|
|
|
// Service is the app-scoped realtime service. Get it with From.
|
|
type Service struct {
|
|
app *backpack.App
|
|
driver Driver
|
|
registry *Registry
|
|
log *slog.Logger
|
|
namespace string
|
|
queue string
|
|
timeout time.Duration
|
|
|
|
mu sync.RWMutex
|
|
lookup UserLookup
|
|
}
|
|
|
|
// From returns the app's Service, building and publishing it on first use.
|
|
// The first call reads realtime.driver (default "null"),
|
|
// realtime.broadcast_namespace (default ""), realtime.broadcast_queue
|
|
// (default "broadcasts") and realtime.broadcast_timeout (default 5s; an
|
|
// integer number of seconds or a duration string) and builds the driver
|
|
// through its registered factory. An unknown driver name is an error.
|
|
func From(app *backpack.App) (*Service, error) {
|
|
if app == nil {
|
|
return nil, fmt.Errorf("lighthouse: app is nil")
|
|
}
|
|
if s, ok := app.Lookup[*Service](); ok && s != nil {
|
|
return s, nil
|
|
}
|
|
svc := newService(app)
|
|
name := DefaultDriver
|
|
if app.Config != nil {
|
|
if v := strings.TrimSpace(app.Config.String("realtime.driver")); v != "" {
|
|
name = strings.ToLower(v)
|
|
}
|
|
}
|
|
factory, ok := driverFactory(name)
|
|
if !ok {
|
|
return nil, fmt.Errorf("lighthouse: unknown realtime.driver %q (registered: %s; a driver package must be imported to register itself)", name, strings.Join(driverNames(), ", "))
|
|
}
|
|
d, err := factory(app, svc)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("lighthouse: driver %s: %w", name, err)
|
|
}
|
|
if d == nil {
|
|
return nil, fmt.Errorf("lighthouse: driver %s returned nil", name)
|
|
}
|
|
svc.driver = d
|
|
if err := app.Publish(svc); err != nil {
|
|
if existing, ok := app.Lookup[*Service](); ok && existing != nil {
|
|
return existing, nil
|
|
}
|
|
return nil, fmt.Errorf("lighthouse: %w", err)
|
|
}
|
|
return svc, nil
|
|
}
|
|
|
|
func newService(app *backpack.App) *Service {
|
|
svc := &Service{
|
|
app: app,
|
|
registry: NewRegistry(),
|
|
log: loggerFromApp(app),
|
|
queue: DefaultQueue,
|
|
timeout: DefaultTimeout,
|
|
}
|
|
if app == nil || app.Config == nil {
|
|
return svc
|
|
}
|
|
c := app.Config
|
|
svc.namespace = strings.TrimSpace(c.String("realtime.broadcast_namespace"))
|
|
if v := strings.TrimSpace(c.String("realtime.broadcast_queue")); v != "" {
|
|
svc.queue = v
|
|
}
|
|
if d := DurationSetting(c, "realtime.broadcast_timeout"); d > 0 {
|
|
svc.timeout = d
|
|
}
|
|
return svc
|
|
}
|
|
|
|
// Driver returns the driver selected by realtime.driver.
|
|
func (s *Service) Driver() Driver {
|
|
if s == nil {
|
|
return nil
|
|
}
|
|
return s.driver
|
|
}
|
|
|
|
// Registry returns the channel-namespace authorizer registry that the
|
|
// driver's subscribe authorization consults.
|
|
func (s *Service) Registry() *Registry {
|
|
if s == nil {
|
|
return nil
|
|
}
|
|
return s.registry
|
|
}
|
|
|
|
// Logger returns the app logger the service and its driver log through.
|
|
func (s *Service) Logger() *slog.Logger {
|
|
if s == nil || s.log == nil {
|
|
return slog.Default()
|
|
}
|
|
return s.log
|
|
}
|
|
|
|
// Namespace returns realtime.broadcast_namespace.
|
|
func (s *Service) Namespace() string {
|
|
if s == nil {
|
|
return ""
|
|
}
|
|
return s.namespace
|
|
}
|
|
|
|
// Queue returns the River queue broadcast jobs are enqueued on.
|
|
func (s *Service) Queue() string {
|
|
if s == nil {
|
|
return DefaultQueue
|
|
}
|
|
return s.queue
|
|
}
|
|
|
|
// Timeout returns the per-job timeout of a broadcast job.
|
|
func (s *Service) Timeout() time.Duration {
|
|
if s == nil {
|
|
return DefaultTimeout
|
|
}
|
|
return s.timeout
|
|
}
|
|
|
|
// DurationSetting reads path as a duration string ("5s") or an integer
|
|
// number of seconds. It returns 0 when the key is missing or invalid.
|
|
func DurationSetting(c *compass.Config, path string) time.Duration {
|
|
if c == nil {
|
|
return 0
|
|
}
|
|
raw := strings.TrimSpace(c.String(path))
|
|
if raw == "" {
|
|
return 0
|
|
}
|
|
if d, err := time.ParseDuration(raw); err == nil && d > 0 {
|
|
return d
|
|
}
|
|
if n := c.Int(path); n > 0 {
|
|
return time.Duration(n) * time.Second
|
|
}
|
|
return 0
|
|
}
|
|
|
|
func loggerFromApp(app *backpack.App) *slog.Logger {
|
|
if app != nil {
|
|
if log, ok := app.Lookup[*slog.Logger](); ok && log != nil {
|
|
return log
|
|
}
|
|
}
|
|
return slog.Default()
|
|
}
|