diff --git a/README.md b/README.md index 6a8a9d2..597158b 100644 --- a/README.md +++ b/README.md @@ -92,6 +92,7 @@ The `migrate`, `migrate:status`, `migrate:rollback`, `serve` and admin commands | [festival](modules/festival/README.md) | Typed, synchronous event bus with listener priorities, payload collection and stop-when-handled dispatch. | | [fetchguard](modules/fetchguard/README.md) | Guarded outbound HTTPS fetcher that blocks private and reserved addresses and enforces host, size and timeout limits. | | [lagoon](modules/lagoon/README.md) | Postgres data layer: the shared GORM connection, per-plugin migrations, model helpers and file attachments. | +| [lighthouse](modules/lighthouse/README.md) | Transport-neutral realtime: a publisher interface with pluggable drivers, subscribe-time channel authorization, and model broadcasts enqueued in the write transaction. | | [pact](modules/pact/README.md) | Capability interfaces that compiled plugins implement to contribute routes, config, migrations, middleware, commands, admin screens, translations, mail templates and jobs. | | [party](modules/party/README.md) | Compiled plugin registry that orders plugins by their dependencies and runs their Register and Boot lifecycle. | | [phrasebook](modules/phrasebook/README.md) | Namespaced translation catalogs loaded from plugin YAML, with locale fallback, placeholder interpolation and CLDR pluralization. | diff --git a/modules/lighthouse/README.md b/modules/lighthouse/README.md new file mode 100644 index 0000000..b69ebc1 --- /dev/null +++ b/modules/lighthouse/README.md @@ -0,0 +1,151 @@ +# lighthouse + +Transport-neutral realtime: a publisher interface with pluggable drivers, subscribe-time channel authorization, and model broadcasts enqueued in the write transaction. + +`import "git.golem15.com/golem15/summercms/modules/lighthouse"` + +`import _ "git.golem15.com/golem15/summercms/modules/lighthouse/centrifugo"` + +## Overview + +lighthouse is the SummerCMS counterpart of the WinterCMS websockets plugin. The package itself knows no transport. It owns the interfaces application code writes against, and a driver package supplies the transport. A driver registers itself from its `init` function, the way `database/sql` drivers do, and the application picks one with `realtime.driver`. + +`lighthouse.From` builds the app-scoped `lighthouse.Service` on first use and publishes it on the app. The service holds the selected driver, the application's user lookup, and the broadcast settings. + +A driver may need HTTP endpoints, such as a token route for signed-in users or a callback the realtime server calls. It declares them as `lighthouse.Route` values, each tagged with a `lighthouse.Surface`. The application mounts them once with `lighthouse.Mount` and decides the guard, group and rate-limit bucket per surface. Switching drivers never edits the application's route file. + +The `centrifugo` sub-package is the Centrifugo driver. It has a hand-rolled `net/http` client for the Centrifugo HTTP API, a token issuer with the claims of the WinterCMS `JwtTokenGenerator`, and the token route handler. + +## Features + +- Driver selection by `realtime.driver`. The built-in drivers are `null` (the default; it discards everything), `log` (logs channel names and the event, never the payload) and `memory` (`lighthouse.MemoryDriver` records every `lighthouse.Publication` for tests). An unknown name is a boot error that lists the registered drivers. +- Third-party drivers: `lighthouse.RegisterDriver` with a `lighthouse.DriverFactory`. A duplicate name panics at init. +- The `lighthouse.Publisher` interface (`Publish` for one channel, `Broadcast` for several) and the `lighthouse.Driver` interface, which adds `Name` and `Routes`. +- Route mounting by surface: `lighthouse.UserAuth`, `lighthouse.ServerToServer` and `lighthouse.Public`. `lighthouse.Mount` puts user and public routes in `Group` and server-to-server routes in `GroupRaw`, each with the surface middleware followed by `lighthouse.Surfaces.Middleware`. A `lighthouse.UserAuth` route with no user middleware is refused, so a token route can never be mounted without a guard. Every route is validated before any is registered. +- Users and actors: the application installs a `lighthouse.UserLookup` with `lighthouse.Service.SetUserLookup`. `lighthouse.Service.User` loads a `lighthouse.User` (id and display name). `lighthouse.Service.Actor` returns the `lighthouse.Actor` of a request, and `lighthouse.SystemActor` when there is no signed-in user or the principal is a backend admin. +- Centrifugo driver (`centrifugo.Driver`, driver name `centrifugo`): + - `centrifugo.Client` POSTs `publish`, `broadcast`, `presence` and `unsubscribe` calls with `Authorization: apikey ` and a 5 s timeout. The publish body is `{"channel":…,"data":{"event":…,"payload":…,"timestamp":"…+00:00"}}`, with an empty payload sent as `[]`. Any 2xx status counts as success. With an empty API key nothing is sent and the call returns `centrifugo.ErrNotConfigured`. The key never appears in logs or errors. + - `centrifugo.TokenIssuer` signs HS256 tokens with five generators, the same as the WinterCMS generator: `ForUser` (claims `sub`, `exp`, `info` with only `name`), `Subscription`, `Anonymous` (`sub` "" and a 5-minute lifetime), `ForIdentifier` (an empty `info` is encoded as `[]`) and `SubscriptionForIdentifier`. It refuses to sign with an empty secret. + - `centrifugo.TokenHandler` serves the token route. It answers 401 `{"error":"Unauthorized"}` when no user is signed in, 503 `{"error":"WebSocket not configured"}` when the token secret is empty, and otherwise 200 `{"token":"…"}`. It sends `Cache-Control: no-cache, private` and no trailing newline. + +## Usage + +An application selects the driver in `config/realtime.yaml`: + +```yaml +driver: centrifugo +centrifugo: + token_secret: "" # set with SUMMER_REALTIME__CENTRIFUGO__TOKEN_SECRET +``` + +A plugin imports the driver package for its side effect, builds the service at Boot and installs a user lookup: + +```go +package acme + +import ( + "context" + + "git.golem15.com/golem15/summercms/modules/backpack" + "git.golem15.com/golem15/summercms/modules/lighthouse" + _ "git.golem15.com/golem15/summercms/modules/lighthouse/centrifugo" +) + +func (p *Plugin) Boot(app *backpack.App) error { + svc, err := lighthouse.From(app) + if err != nil { + return err + } + p.realtime = svc + svc.SetUserLookup(func(ctx context.Context, id uint) (lighthouse.User, bool, error) { + return lookupAcmeUser(ctx, id) // the application's own user model + }) + return nil +} +``` + +and mounts the driver's routes once: + +```go +func (p *Plugin) Routes(r pact.Router) error { + return lighthouse.Mount(r, p.realtime.Driver(), lighthouse.Surfaces{ + UserAuth: surf.Use("jwt.auth"), + ServerToServer: surf.Use(), + Middleware: surf.Use("throttle:acme-realtime"), + }) +} +``` + +Tests select the memory driver and read what was published: + +```go +mem := svc.Driver().(*lighthouse.MemoryDriver) +for _, pub := range mem.Publications() { + fmt.Println(pub.Method, pub.Channels, pub.Event) +} +``` + +## API reference + +### lighthouse + +| Identifier | Description | +|------------|-------------| +| `lighthouse.From(app)` | The app's `*lighthouse.Service`, built and published on first use. | +| `lighthouse.Service` | The realtime service: `Driver`, `Logger`, `Namespace`, `Queue`, `Timeout`, `SetUserLookup`, `User`, `Actor`. | +| `lighthouse.Publisher` | `Publish(ctx, channel, event, payload)` and `Broadcast(ctx, channels, event, payload)`. | +| `lighthouse.Driver` | `lighthouse.Publisher` plus `Name()` and `Routes()`. | +| `lighthouse.DriverFactory` | `func(app, svc) (lighthouse.Driver, error)`. | +| `lighthouse.RegisterDriver(name, factory)` | Registers a driver from an `init` function. | +| `lighthouse.MemoryDriver`, `lighthouse.NewMemoryDriver`, `lighthouse.Publication` | The recording driver and its records (`Method`, `Channels`, `Event`, `Payload`, `Timestamp`). | +| `lighthouse.Route` | `Name`, `Method`, `Path`, `Surface`, `Handler`. | +| `lighthouse.Surface`, `lighthouse.UserAuth`, `lighthouse.ServerToServer`, `lighthouse.Public` | Who calls a route. | +| `lighthouse.Surfaces` | Application middleware per surface plus `Middleware` for every route. | +| `lighthouse.Mount(r, driver, surfaces)` | Registers a driver's routes. | +| `lighthouse.User`, `lighthouse.UserLookup` | A user id with a display name, and the application's lookup. | +| `lighthouse.Actor`, `lighthouse.SystemActor` | Who caused a broadcast: `{"user_id":…,"name":…}`. | +| `lighthouse.DurationSetting(cfg, path)` | Reads a duration string or an integer number of seconds. | +| `lighthouse.DefaultDriver`, `lighthouse.DefaultQueue`, `lighthouse.DefaultTimeout` | Defaults of `realtime.driver`, `realtime.broadcast_queue` and `realtime.broadcast_timeout`. | + +### lighthouse/centrifugo + +| Identifier | Description | +|------------|-------------| +| `centrifugo.Config`, `centrifugo.LoadConfig` | The `realtime.centrifugo.*` settings with their defaults. | +| `centrifugo.Client`, `centrifugo.NewClient` | HTTP API client: `Publish`, `Broadcast`, `Presence`, `Unsubscribe`, `Enabled`, `DebugInfo`. | +| `centrifugo.DebugInfo` | `api_url`, `enabled`, `api_key_set`. | +| `centrifugo.TokenIssuer`, `centrifugo.NewTokenIssuer` | HS256 token generators: `ForUser`, `Subscription`, `Anonymous`, `ForIdentifier`, `SubscriptionForIdentifier`, `Configured`. | +| `centrifugo.TokenHandler(svc, issuer)` | The token route handler. | +| `centrifugo.Driver`, `centrifugo.NewDriver`, `centrifugo.DriverName` | The `lighthouse.Driver`, with `Client`, `Issuer` and `Config` accessors. | +| `centrifugo.ErrNotConfigured` | Returned when the API key or token secret an operation needs is empty. | + +## Configuration + +| Key | Default | Description | +|-----|---------|-------------| +| `realtime.driver` | `null` | `null`, `log`, `memory`, or a registered driver such as `centrifugo`. | +| `realtime.broadcast_namespace` | `""` | Prefix applied to broadcast channel names. | +| `realtime.broadcast_queue` | `broadcasts` | River queue of broadcast jobs. | +| `realtime.broadcast_timeout` | `5` | Broadcast job timeout, in seconds or as a duration string. | +| `realtime.centrifugo.api_url` | `http://127.0.0.1:8001/api` | Centrifugo HTTP API base. | +| `realtime.centrifugo.api_key` | `""` | HTTP API key; empty disables publishing. | +| `realtime.centrifugo.token_secret` | `""` | HS256 token secret; empty makes the token route answer 503. | +| `realtime.centrifugo.token_ttl` | `3600` | Token lifetime, in seconds or as a duration string. | +| `realtime.centrifugo.ws_url` | `/ws` | WebSocket URL of the Centrifugo server. | +| `realtime.centrifugo.proxy_secret` | `""` | Expected `X-Centrifugo-Secret` of subscribe proxy calls. | +| `realtime.centrifugo.token_path` | `/api/realtime/token` | Path of the token route. | +| `realtime.centrifugo.subscribe_path` | `/api/realtime/subscribe` | Path of the subscribe proxy route. | + +## Dependencies + +- `backpack`, `bouncer`, `compass`, `pact` and `wire` from this repository. +- `github.com/golang-jwt/jwt/v5` (centrifugo token signing). +- The Centrifugo client is plain `net/http`; no Centrifugo SDK is used. + +## Testing + +```bash +go test ./modules/lighthouse/... +``` + +A test selects `realtime.driver: memory` and reads `lighthouse.MemoryDriver.Publications`, or points `realtime.centrifugo.api_url` at an `httptest` server to see the exact Centrifugo requests. diff --git a/modules/lighthouse/centrifugo/client.go b/modules/lighthouse/centrifugo/client.go new file mode 100644 index 0000000..6416427 --- /dev/null +++ b/modules/lighthouse/centrifugo/client.go @@ -0,0 +1,203 @@ +package centrifugo + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "io" + "log/slog" + "net/http" + "strconv" + "time" + + "git.golem15.com/golem15/summercms/modules/wire" +) + +// requestTimeout bounds every Centrifugo API call, as the WinterCMS client +// does. +const requestTimeout = 5 * time.Second + +// maxResponseBytes caps how much of a response body is read. +const maxResponseBytes = 1 << 20 + +// ErrNotConfigured is returned when the secret or key an operation needs +// is empty: publishing without an API key, or signing without a token +// secret. +var ErrNotConfigured = errors.New("centrifugo: not configured") + +// Client calls the Centrifugo HTTP API. It is safe for concurrent use. +type Client struct { + apiURL string + apiKey string + hc *http.Client + log *slog.Logger + now func() time.Time +} + +// DebugInfo is the connection summary Client.DebugInfo returns. It never +// carries the API key. +type DebugInfo struct { + APIURL string `json:"api_url"` + Enabled bool `json:"enabled"` + APIKeySet bool `json:"api_key_set"` +} + +// NewClient returns a client for cfg.APIURL and cfg.APIKey. hc may be nil +// for a client with a 5s timeout; every request is also bounded by 5s. +func NewClient(cfg Config, hc *http.Client) *Client { + if hc == nil { + hc = &http.Client{Timeout: requestTimeout} + } + return &Client{ + apiURL: cfg.APIURL, + apiKey: cfg.APIKey, + hc: hc, + log: slog.Default(), + now: time.Now, + } +} + +// Enabled reports whether an API key is configured. +func (c *Client) Enabled() bool { return c != nil && c.apiKey != "" } + +// DebugInfo returns the API URL and whether the client is enabled. +func (c *Client) DebugInfo() DebugInfo { + if c == nil { + return DebugInfo{} + } + return DebugInfo{APIURL: c.apiURL, Enabled: c.Enabled(), APIKeySet: c.apiKey != ""} +} + +type eventData struct { + Event string `json:"event"` + Payload json.RawMessage `json:"payload"` + Timestamp wire.Time `json:"timestamp"` +} + +type publishRequest struct { + Channel string `json:"channel"` + Data eventData `json:"data"` +} + +type broadcastRequest struct { + Channels []string `json:"channels"` + Data eventData `json:"data"` +} + +type presenceRequest struct { + Channel string `json:"channel"` +} + +type unsubscribeRequest struct { + User string `json:"user"` + Channel string `json:"channel"` +} + +// Publish POSTs {"channel":…,"data":{"event":…,"payload":…,"timestamp":…}} +// to {api_url}/publish. An empty payload is sent as []. Any 2xx status is +// success, including Centrifugo's 200 responses that carry an error body. +func (c *Client) Publish(ctx context.Context, channel, event string, payload json.RawMessage) error { + if !c.Enabled() { + return ErrNotConfigured + } + body := publishRequest{Channel: channel, Data: c.data(event, payload)} + _, err := c.post(ctx, "/publish", body) + return err +} + +// Broadcast POSTs the same data with "channels" to {api_url}/broadcast. No +// request is sent for an empty channel list. +func (c *Client) Broadcast(ctx context.Context, channels []string, event string, payload json.RawMessage) error { + if !c.Enabled() { + return ErrNotConfigured + } + if len(channels) == 0 { + return nil + } + body := broadcastRequest{Channels: channels, Data: c.data(event, payload)} + _, err := c.post(ctx, "/broadcast", body) + return err +} + +// Presence returns result.presence of {api_url}/presence for channel. It +// returns an empty map, never nil, when the client is disabled or the call +// fails; the error says why. +func (c *Client) Presence(ctx context.Context, channel string) (map[string]any, error) { + out := map[string]any{} + if !c.Enabled() { + return out, ErrNotConfigured + } + raw, err := c.post(ctx, "/presence", presenceRequest{Channel: channel}) + if err != nil { + return out, err + } + var resp struct { + Result struct { + Presence map[string]any `json:"presence"` + } `json:"result"` + } + if err := json.Unmarshal(raw, &resp); err != nil { + c.log.Warn("centrifugo: presence response is not JSON", slog.String("channel", channel)) + return out, fmt.Errorf("centrifugo: presence: %w", err) + } + if resp.Result.Presence != nil { + out = resp.Result.Presence + } + return out, nil +} + +// Unsubscribe POSTs {"user":"","channel":…} to {api_url}/unsubscribe. +func (c *Client) Unsubscribe(ctx context.Context, userID uint, channel string) error { + if !c.Enabled() { + return ErrNotConfigured + } + body := unsubscribeRequest{User: strconv.FormatUint(uint64(userID), 10), Channel: channel} + _, err := c.post(ctx, "/unsubscribe", body) + return err +} + +func (c *Client) data(event string, payload json.RawMessage) eventData { + if len(bytes.TrimSpace(payload)) == 0 { + payload = json.RawMessage("[]") + } + return eventData{Event: event, Payload: payload, Timestamp: wire.Time{Time: c.now()}} +} + +// post sends body and returns the response body of a 2xx answer. The API +// key appears only in the Authorization header, never in logs or errors. +func (c *Client) post(ctx context.Context, method string, body any) ([]byte, error) { + var buf bytes.Buffer + enc := json.NewEncoder(&buf) + enc.SetEscapeHTML(false) + if err := enc.Encode(body); err != nil { + return nil, fmt.Errorf("centrifugo: %s: encode: %w", method, err) + } + if ctx == nil { + ctx = context.Background() + } + ctx, cancel := context.WithTimeout(ctx, requestTimeout) + defer cancel() + req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.apiURL+method, bytes.NewReader(bytes.TrimSuffix(buf.Bytes(), []byte("\n")))) + if err != nil { + return nil, fmt.Errorf("centrifugo: %s: %w", method, err) + } + req.Header.Set("Authorization", "apikey "+c.apiKey) + req.Header.Set("Content-Type", "application/json") + resp, err := c.hc.Do(req) + if err != nil { + return nil, fmt.Errorf("centrifugo: %s: %w", method, err) + } + defer resp.Body.Close() + raw, _ := io.ReadAll(io.LimitReader(resp.Body, maxResponseBytes)) + if resp.StatusCode < 200 || resp.StatusCode > 299 { + return nil, fmt.Errorf("centrifugo: %s: HTTP %d", method, resp.StatusCode) + } + if bytes.Contains(raw, []byte(`"error"`)) { + // Centrifugo reports API errors as 200 with an error body. The + // WinterCMS client counts those as success; so does this one. + c.log.Debug("centrifugo: API answered with an error body", slog.String("method", method), slog.Int("status", resp.StatusCode)) + } + return raw, nil +} diff --git a/modules/lighthouse/centrifugo/config.go b/modules/lighthouse/centrifugo/config.go new file mode 100644 index 0000000..5dc3c06 --- /dev/null +++ b/modules/lighthouse/centrifugo/config.go @@ -0,0 +1,78 @@ +// Package centrifugo is the Centrifugo driver of lighthouse: an HTTP API +// client, a connection and subscription token issuer, and the subscribe +// proxy handler. Import it for its side effect to register the +// "centrifugo" realtime.driver. +package centrifugo + +import ( + "strings" + "time" + + "git.golem15.com/golem15/summercms/modules/compass" + "git.golem15.com/golem15/summercms/modules/lighthouse" +) + +// Default values of the realtime.centrifugo.* keys. +const ( + DefaultAPIURL = "http://127.0.0.1:8001/api" + DefaultTokenTTL = 3600 * time.Second + DefaultWSURL = "/ws" + DefaultTokenPath = "/api/realtime/token" + DefaultSubscribePath = "/api/realtime/subscribe" +) + +// Config is the realtime.centrifugo.* configuration. +type Config struct { + // APIURL is the Centrifugo HTTP API base, without a trailing method. + APIURL string + // APIKey authenticates publishes; empty disables them. + APIKey string + // TokenSecret signs connection and subscription tokens (HS256); empty + // makes the token route answer 503. + TokenSecret string + // TokenTTL is the lifetime of issued tokens. + TokenTTL time.Duration + // WSURL is the WebSocket URL the client connects to. + WSURL string + // ProxySecret is the X-Centrifugo-Secret header value the subscribe + // proxy expects; empty denies every subscribe. + ProxySecret string + // TokenPath and SubscribePath are the mounted route paths. + TokenPath string + SubscribePath string +} + +// LoadConfig reads realtime.centrifugo.* from c, filling the defaults. +// token_ttl is an integer number of seconds or a duration string. +func LoadConfig(c *compass.Config) Config { + cfg := Config{ + APIURL: DefaultAPIURL, + TokenTTL: DefaultTokenTTL, + WSURL: DefaultWSURL, + TokenPath: DefaultTokenPath, + SubscribePath: DefaultSubscribePath, + } + if c == nil { + return cfg + } + str := func(key string) string { return strings.TrimSpace(c.String("realtime.centrifugo." + key)) } + if v := str("api_url"); v != "" { + cfg.APIURL = strings.TrimSuffix(v, "/") + } + cfg.APIKey = str("api_key") + cfg.TokenSecret = str("token_secret") + cfg.ProxySecret = str("proxy_secret") + if d := lighthouse.DurationSetting(c, "realtime.centrifugo.token_ttl"); d > 0 { + cfg.TokenTTL = d + } + if v := str("ws_url"); v != "" { + cfg.WSURL = v + } + if v := str("token_path"); v != "" { + cfg.TokenPath = v + } + if v := str("subscribe_path"); v != "" { + cfg.SubscribePath = v + } + return cfg +} diff --git a/modules/lighthouse/centrifugo/driver.go b/modules/lighthouse/centrifugo/driver.go new file mode 100644 index 0000000..8662bb2 --- /dev/null +++ b/modules/lighthouse/centrifugo/driver.go @@ -0,0 +1,77 @@ +package centrifugo + +import ( + "context" + "encoding/json" + "net/http" + + "git.golem15.com/golem15/summercms/modules/backpack" + "git.golem15.com/golem15/summercms/modules/lighthouse" +) + +// DriverName is the realtime.driver value of this driver. +const DriverName = "centrifugo" + +func init() { + lighthouse.RegisterDriver(DriverName, func(app *backpack.App, svc *lighthouse.Service) (lighthouse.Driver, error) { + var cfg Config + if app != nil { + cfg = LoadConfig(app.Config) + } else { + cfg = LoadConfig(nil) + } + return NewDriver(svc, cfg, nil), nil + }) +} + +// Driver is the Centrifugo lighthouse.Driver: it publishes through Client +// and declares the token route (UserAuth) and the subscribe proxy route +// (ServerToServer). +type Driver struct { + svc *lighthouse.Service + cfg Config + client *Client + issuer *TokenIssuer +} + +// NewDriver builds the driver for svc from cfg. hc may be nil (see +// NewClient). +func NewDriver(svc *lighthouse.Service, cfg Config, hc *http.Client) *Driver { + c := NewClient(cfg, hc) + c.log = svc.Logger() + return &Driver{ + svc: svc, + cfg: cfg, + client: c, + issuer: NewTokenIssuer(cfg.TokenSecret, cfg.TokenTTL), + } +} + +// Name returns "centrifugo". +func (d *Driver) Name() string { return DriverName } + +// Config returns the driver's configuration. +func (d *Driver) Config() Config { return d.cfg } + +// Client returns the HTTP API client. +func (d *Driver) Client() *Client { return d.client } + +// Issuer returns the token issuer. +func (d *Driver) Issuer() *TokenIssuer { return d.issuer } + +// Publish sends event to one channel through Client.Publish. +func (d *Driver) Publish(ctx context.Context, channel, event string, payload json.RawMessage) error { + return d.client.Publish(ctx, channel, event, payload) +} + +// Broadcast sends event to several channels through Client.Broadcast. +func (d *Driver) Broadcast(ctx context.Context, channels []string, event string, payload json.RawMessage) error { + return d.client.Broadcast(ctx, channels, event, payload) +} + +// Routes returns GET token_path (UserAuth, TokenHandler). +func (d *Driver) Routes() []lighthouse.Route { + return []lighthouse.Route{ + {Name: "token", Method: http.MethodGet, Path: d.cfg.TokenPath, Surface: lighthouse.UserAuth, Handler: TokenHandler(d.svc, d.issuer)}, + } +} diff --git a/modules/lighthouse/centrifugo/handlers.go b/modules/lighthouse/centrifugo/handlers.go new file mode 100644 index 0000000..f84e4be --- /dev/null +++ b/modules/lighthouse/centrifugo/handlers.go @@ -0,0 +1,67 @@ +package centrifugo + +import ( + "bytes" + "encoding/json" + "net/http" + + "git.golem15.com/golem15/summercms/modules/bouncer" + "git.golem15.com/golem15/summercms/modules/lighthouse" +) + +type errorBody struct { + Error string `json:"error"` +} + +type tokenBody struct { + Token string `json:"token"` +} + +// TokenHandler issues the connection token of the signed-in user. It must +// be mounted behind a user guard (the UserAuth surface). +// +// - no principal, or no user for it: 401 {"error":"Unauthorized"} +// - no token secret: 503 with the WinterCMS not-configured error body +// - otherwise 200 {"token":"…"} (see TokenIssuer.ForUser) +func TokenHandler(svc *lighthouse.Service, issuer *TokenIssuer) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + p, ok := bouncer.User(r.Context()) + if !ok || p == nil || p.ID == 0 { + writeJSON(w, http.StatusUnauthorized, errorBody{Error: "Unauthorized"}) + return + } + u, found, err := svc.User(r.Context(), p.ID) + if err != nil || !found { + writeJSON(w, http.StatusUnauthorized, errorBody{Error: "Unauthorized"}) + return + } + if !issuer.Configured() { + writeJSON(w, http.StatusServiceUnavailable, errorBody{Error: "WebSocket not configured"}) + return + } + tok, err := issuer.ForUser(u) + if err != nil { + svc.Logger().Error("realtime: token signing failed", "error", err) + writeJSON(w, http.StatusInternalServerError, errorBody{Error: "Internal server error"}) + return + } + writeJSON(w, http.StatusOK, tokenBody{Token: tok}) + } +} + +// writeJSON writes v with no trailing newline and no HTML escaping, plus the +// Content-Type and Cache-Control headers of a Laravel JSON response. +func writeJSON(w http.ResponseWriter, status int, v any) { + var buf bytes.Buffer + enc := json.NewEncoder(&buf) + enc.SetEscapeHTML(false) + if err := enc.Encode(v); err != nil { + w.WriteHeader(http.StatusInternalServerError) + return + } + h := w.Header() + h.Set("Content-Type", "application/json") + h.Set("Cache-Control", "no-cache, private") + w.WriteHeader(status) + _, _ = w.Write(bytes.TrimSuffix(buf.Bytes(), []byte("\n"))) +} diff --git a/modules/lighthouse/centrifugo/token.go b/modules/lighthouse/centrifugo/token.go new file mode 100644 index 0000000..5252111 --- /dev/null +++ b/modules/lighthouse/centrifugo/token.go @@ -0,0 +1,145 @@ +package centrifugo + +import ( + "bytes" + "encoding/json" + "strconv" + "time" + + "git.golem15.com/golem15/summercms/modules/lighthouse" + "github.com/golang-jwt/jwt/v5" +) + +// anonymousTTL is the lifetime of Anonymous tokens. +const anonymousTTL = 300 * time.Second + +// TokenIssuer signs Centrifugo connection and subscription tokens with +// HS256. Its claims match the WinterCMS JwtTokenGenerator. It is safe for +// concurrent use. +type TokenIssuer struct { + secret []byte + ttl time.Duration + // Now is the clock used for exp; nil means time.Now. + Now func() time.Time +} + +// NewTokenIssuer returns an issuer for secret with tokens valid for ttl +// (DefaultTokenTTL when ttl is not positive). +func NewTokenIssuer(secret string, ttl time.Duration) *TokenIssuer { + if ttl <= 0 { + ttl = DefaultTokenTTL + } + return &TokenIssuer{secret: []byte(secret), ttl: ttl} +} + +// Configured reports whether a signing secret is set. +func (i *TokenIssuer) Configured() bool { return i != nil && len(i.secret) > 0 } + +type userInfo struct { + Name *string `json:"name"` +} + +type userClaims struct { + Sub string `json:"sub"` + Exp int64 `json:"exp"` + Info userInfo `json:"info"` +} + +type channelClaims struct { + Sub string `json:"sub"` + Channel string `json:"channel"` + Exp int64 `json:"exp"` +} + +type anonymousClaims struct { + Sub string `json:"sub"` + Exp int64 `json:"exp"` +} + +type identifierClaims struct { + Sub string `json:"sub"` + Exp int64 `json:"exp"` + Info json.RawMessage `json:"info"` +} + +// ForUser returns a connection token with exactly the claims sub (the user +// id as a string), exp (now + TTL) and info {"name": u.Name}. info carries +// nothing else: no email, no other ids. +func (i *TokenIssuer) ForUser(u lighthouse.User) (string, error) { + return i.sign(userClaims{Sub: userSub(u.ID), Exp: i.exp(i.ttl), Info: userInfo{Name: u.Name}}) +} + +// Subscription returns a subscription token with the claims sub, channel +// and exp. +func (i *TokenIssuer) Subscription(u lighthouse.User, channel string) (string, error) { + return i.sign(channelClaims{Sub: userSub(u.ID), Channel: channel, Exp: i.exp(i.ttl)}) +} + +// Anonymous returns a connection token with sub "" and exp now + 5 minutes. +func (i *TokenIssuer) Anonymous() (string, error) { + return i.sign(anonymousClaims{Sub: "", Exp: i.exp(anonymousTTL)}) +} + +// ForIdentifier returns a connection token for a non-user identifier with +// the claims sub, exp and info. An empty info is encoded as [], as the +// WinterCMS generator's empty PHP array is. +func (i *TokenIssuer) ForIdentifier(identifier string, info map[string]any) (string, error) { + raw := json.RawMessage("[]") + if len(info) > 0 { + b, err := marshal(info) + if err != nil { + return "", err + } + raw = b + } + return i.sign(identifierClaims{Sub: identifier, Exp: i.exp(i.ttl), Info: raw}) +} + +// SubscriptionForIdentifier returns a subscription token for a non-user +// identifier with the claims sub, channel and exp. +func (i *TokenIssuer) SubscriptionForIdentifier(identifier, channel string) (string, error) { + return i.sign(channelClaims{Sub: identifier, Channel: channel, Exp: i.exp(i.ttl)}) +} + +func (i *TokenIssuer) exp(ttl time.Duration) int64 { + now := time.Now + if i != nil && i.Now != nil { + now = i.Now + } + return now().Add(ttl).Unix() +} + +func (i *TokenIssuer) sign(claims any) (string, error) { + if !i.Configured() { + return "", ErrNotConfigured + } + raw, err := marshal(claims) + if err != nil { + return "", err + } + return jwt.NewWithClaims(jwt.SigningMethodHS256, orderedClaims(raw)).SignedString(i.secret) +} + +func userSub(id uint) string { return strconv.FormatUint(uint64(id), 10) } + +func marshal(v any) ([]byte, error) { + var buf bytes.Buffer + enc := json.NewEncoder(&buf) + enc.SetEscapeHTML(false) + if err := enc.Encode(v); err != nil { + return nil, err + } + return bytes.TrimSuffix(buf.Bytes(), []byte("\n")), nil +} + +// orderedClaims keeps the claim order of the struct it was marshalled +// from. The jwt.Claims methods are never used for signing. +type orderedClaims json.RawMessage + +func (c orderedClaims) MarshalJSON() ([]byte, error) { return c, nil } +func (orderedClaims) GetExpirationTime() (*jwt.NumericDate, error) { return nil, nil } +func (orderedClaims) GetIssuedAt() (*jwt.NumericDate, error) { return nil, nil } +func (orderedClaims) GetNotBefore() (*jwt.NumericDate, error) { return nil, nil } +func (orderedClaims) GetIssuer() (string, error) { return "", nil } +func (orderedClaims) GetSubject() (string, error) { return "", nil } +func (orderedClaims) GetAudience() (jwt.ClaimStrings, error) { return nil, nil } diff --git a/modules/lighthouse/drivers.go b/modules/lighthouse/drivers.go new file mode 100644 index 0000000..300ab85 --- /dev/null +++ b/modules/lighthouse/drivers.go @@ -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 +} diff --git a/modules/lighthouse/lighthouse.go b/modules/lighthouse/lighthouse.go new file mode 100644 index 0000000..cf1904e --- /dev/null +++ b/modules/lighthouse/lighthouse.go @@ -0,0 +1,173 @@ +// 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 + 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, + 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 +} + +// 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() +} diff --git a/modules/lighthouse/route.go b/modules/lighthouse/route.go new file mode 100644 index 0000000..17e7b14 --- /dev/null +++ b/modules/lighthouse/route.go @@ -0,0 +1,142 @@ +package lighthouse + +import ( + "fmt" + "net/http" + "strings" + + "git.golem15.com/golem15/summercms/modules/pact" +) + +// Surface says who calls a driver route, so the application can put it +// behind the right guard and group. +type Surface int + +const ( + // UserAuth routes are called by a signed-in browser user; they mount + // behind the application's user guard. + UserAuth Surface = iota + 1 + // ServerToServer routes are called by the realtime server itself; they + // mount in a raw group without the house envelope. + ServerToServer + // Public routes need no authentication. + Public +) + +// String returns the surface name. +func (s Surface) String() string { + switch s { + case UserAuth: + return "UserAuth" + case ServerToServer: + return "ServerToServer" + case Public: + return "Public" + default: + return fmt.Sprintf("Surface(%d)", int(s)) + } +} + +// Route is an HTTP endpoint declared by a driver. +type Route struct { + // Name identifies the route inside the driver (for example "token"). + Name string + // Method is GET, POST, PUT, PATCH or DELETE. + Method string + // Path is the absolute request path. + Path string + Surface Surface + Handler http.HandlerFunc +} + +// Surfaces holds the application's middleware for each surface. Middleware +// is appended after the surface middleware on every route, so a guard in +// UserAuth runs before a throttle in Middleware. +type Surfaces struct { + UserAuth []string + ServerToServer []string + Public []string + Middleware []string +} + +// Mount registers the driver's routes on r. UserAuth and Public routes go +// through r.Group and ServerToServer routes through r.GroupRaw, each with +// the surface middleware followed by s.Middleware. A nil driver, or one +// without routes, mounts nothing. A UserAuth route with an empty +// s.UserAuth is an error: the application must never expose a user route +// without a guard. Every route is validated before any is registered. +func Mount(r pact.Router, d Driver, s Surfaces) error { + if r == nil { + return fmt.Errorf("lighthouse: router is nil") + } + if d == nil { + return nil + } + routes := d.Routes() + for _, rt := range routes { + if err := validateRoute(d, rt, s); err != nil { + return err + } + } + for _, rt := range routes { + mw := append(append([]string{}, surfaceMiddleware(rt.Surface, s)...), s.Middleware...) + add := func(g pact.Router) { addRoute(g, rt) } + if rt.Surface == ServerToServer { + r.GroupRaw("/", mw, add) + } else { + r.Group("/", mw, add) + } + } + return nil +} + +func validateRoute(d Driver, rt Route, s Surfaces) error { + where := fmt.Sprintf("lighthouse: driver %s route %q", d.Name(), rt.Name) + if rt.Handler == nil { + return fmt.Errorf("%s has no handler", where) + } + if !strings.HasPrefix(rt.Path, "/") { + return fmt.Errorf("%s path %q must be absolute", where, rt.Path) + } + switch strings.ToUpper(rt.Method) { + case http.MethodGet, http.MethodPost, http.MethodPut, http.MethodPatch, http.MethodDelete: + default: + return fmt.Errorf("%s has unsupported method %q", where, rt.Method) + } + switch rt.Surface { + case UserAuth: + if len(s.UserAuth) == 0 { + return fmt.Errorf("%s is a UserAuth route but Surfaces.UserAuth is empty; mount it behind a user guard", where) + } + case ServerToServer, Public: + default: + return fmt.Errorf("%s has unknown surface %v", where, rt.Surface) + } + return nil +} + +func surfaceMiddleware(surface Surface, s Surfaces) []string { + switch surface { + case UserAuth: + return s.UserAuth + case ServerToServer: + return s.ServerToServer + default: + return s.Public + } +} + +func addRoute(g pact.Router, rt Route) { + switch strings.ToUpper(rt.Method) { + case http.MethodGet: + g.Get(rt.Path, rt.Handler) + case http.MethodPost: + g.Post(rt.Path, rt.Handler) + case http.MethodPut: + g.Put(rt.Path, rt.Handler) + case http.MethodPatch: + g.Patch(rt.Path, rt.Handler) + case http.MethodDelete: + g.Delete(rt.Path, rt.Handler) + } +} diff --git a/modules/lighthouse/users.go b/modules/lighthouse/users.go new file mode 100644 index 0000000..8027a5b --- /dev/null +++ b/modules/lighthouse/users.go @@ -0,0 +1,71 @@ +package lighthouse + +import ( + "context" + + "git.golem15.com/golem15/summercms/modules/bouncer" +) + +// User is the identity a realtime driver needs: the id and display name. +type User struct { + ID uint + Name *string +} + +// UserLookup loads a user by id. found is false when no such user exists. +type UserLookup func(ctx context.Context, id uint) (user User, found bool, err error) + +// SetUserLookup installs the application's user lookup. +func (s *Service) SetUserLookup(fn UserLookup) { + if s == nil { + return + } + s.mu.Lock() + s.lookup = fn + s.mu.Unlock() +} + +// User loads user id through the installed lookup. Without a lookup it +// returns User{ID: id} with a nil Name and found true. +func (s *Service) User(ctx context.Context, id uint) (User, bool, error) { + if s == nil { + return User{ID: id}, true, nil + } + s.mu.RLock() + fn := s.lookup + s.mu.RUnlock() + if fn == nil { + return User{ID: id}, true, nil + } + return fn(ctx, id) +} + +// Actor is who caused a broadcast. It marshals as {"user_id":…,"name":…}. +type Actor struct { + UserID *uint `json:"user_id"` + Name *string `json:"name"` +} + +// SystemActor is the actor of writes made without a signed-in user: +// {"user_id":null,"name":"System"}. +func SystemActor() Actor { + name := "System" + return Actor{Name: &name} +} + +// Actor returns the actor of the request in ctx. No principal, or a +// backend-admin principal, gives SystemActor: the WinterCMS frontend user is +// empty in backend requests. Otherwise it is the principal's id and the +// looked-up name (nil when the lookup fails or finds nothing). +func (s *Service) Actor(ctx context.Context) Actor { + p, ok := bouncer.User(ctx) + if !ok || p == nil || p.Backend { + return SystemActor() + } + id := p.ID + a := Actor{UserID: &id} + if u, found, err := s.User(ctx, id); err == nil && found { + a.Name = u.Name + } + return a +}