Files
Jakub Zych 037dc53030 feat(lagoon): per-query collation for OrderBy, drop the database locale check
- lagoon.OrderBy takes variadic lagoon.OrderOption values; lagoon.Collate(name)
  emits a validated, double-quoted COLLATE clause (e.g. "pl-x-icu")
- remove the exported CheckLocale and the ICU pl-PL check from Open and Use
- framework test containers and per-test databases are plain PostgreSQL
- lagoon README, root README and docs pages drop the locale requirement;
  queries-and-pagination gains a "Sorting with a collation" section
  backed by ExampleCollate
2026-10-01 09:51:03 +02:00
..

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.

Channel authorization is transport-neutral too. Plugins register a lighthouse.Authorizer per channel namespace on the service's lighthouse.Registry. The driver's subscribe endpoint asks the authorizer of the channel's namespace on every subscribe, so a user who loses access is denied the next time the client subscribes. Nothing is cached.

Model broadcasts follow the WinterCMS BroadcastableModel trait with one change. A create, update or delete of a broadcastable model enqueues a River job (through conga) inside the write's own transaction, so nothing is published for a write that rolls back. The job publishes after commit, with one attempt and best effort: a failed publish is logged and never touches the write. Suppression is per model type and scoped to a context, and lighthouse.Service.Emit publishes one explicit summary event instead.

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, the token route handler, and the subscribe proxy 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.

  • Channel rules. A channel is namespace:entity:id, optionally prefixed once with presence:.

    • lighthouse.ParseChannel returns the namespace. It returns "" for a presence:presence: prefix or for more than three segments. The lookup is byte-exact and case-sensitive.
    • lighthouse.ChannelID returns segment 1 converted with PHP's (int) cast (lighthouse.PHPInt): 5abc is 5, abc is 0, and out-of-range values saturate.
    • lighthouse.FormatChannels lowercases channel names and applies the broadcast namespace prefix.
  • Authorizer registry: lighthouse.Registry (from lighthouse.Service.Registry) maps namespaces to a lighthouse.Authorizer or lighthouse.AuthorizerFunc. Registering an empty namespace, a namespace that contains :, a nil authorizer or a namespace twice is an error. lighthouse.Registry.Namespaces is sorted. An authorizer returns lighthouse.Allowed (optionally with info, capabilities and overrides) or lighthouse.Denied with an internal reason that only reaches the logs. It reads the realtime client id with lighthouse.ClientID.

  • Model broadcasts. A model broadcasts when its pointer type implements lighthouse.Broadcastable (BroadcastChannels(ctx, tx)), or when a lighthouse.Binding is registered for it with lighthouse.Bind. A binding keeps payload code out of the model package. Optional overrides:

    • lighthouse.BroadcastPayloader or Binding.Payload replaces the default payload {"model":…,"actor":…,"timestamp":"…+00:00","ttl":60}.
    • lighthouse.BroadcastAliaser or Binding.Alias replaces the alias.
    • lighthouse.BroadcastFilter or Binding.ShouldBroadcast can veto an action.
    • lighthouse.BroadcastTTLer or Binding.TTL replaces the ttl.

    The event name is {action}.{alias} lowercased: lighthouse.ActionCreated, lighthouse.ActionUpdated or lighthouse.ActionDeleted, then an alias that defaults to <plugin>.<model> (the Go package name, or the parent directory of a models package, and the type name). The payload builder receives a lighthouse.Event with the action, the lighthouse.Actor, the timestamp and the ttl. A soft delete counts as a delete. A delete's channels and payload are computed from a fresh read of the row before it is deleted, so deleting a model that holds only its id still broadcasts. An empty channel list means no broadcast.

  • Transactional delivery. GORM callbacks (lighthouse.CallbackAfterCreate, lighthouse.CallbackAfterUpdate, lighthouse.CallbackSnapshot and lighthouse.CallbackAfterDelete) are installed through lagoon.OnDatabase. The after-write callbacks run after the model's own after hook and before GORM commits the transaction it opens for a single-statement write, so they enqueue a lighthouse.BroadcastArgs job on the write's *sql.Tx in every case (an explicit transaction or a single Create, Save or Delete), on the realtime.broadcast_queue queue with MaxAttempts 1 and the realtime.broadcast_timeout timeout. Channel and payload queries and the enqueue run inside a savepoint, so a failure (also one a channel or payload function swallows) is rolled back to it, logged at Warn with channels and event (never the payload), and the write goes on. A write with a zero primary key, such as Model(&T{}).Where(…).Updates(…), is not broadcast; bulk paths suppress and emit instead. The null driver, or a driver whose Enabled reports false (Centrifugo without an API key), gets no jobs.

  • The broadcast job lowercases the channels and adds the realtime.broadcast_namespace prefix unless a channel already has it. It then publishes to one channel or broadcasts to several. A failure is logged as realtime: broadcast failed and is not retried. Delivery order across separate jobs is not guaranteed. The payload travels inside the job as a JSON string, so its key order survives Postgres JSONB.

  • Suppression: lighthouse.WithoutBroadcasting silences one model type for writes made with the context it hands to its function. Other types still broadcast, and a write through an outer context is not suppressed. lighthouse.Service.Emit enqueues one lighthouse.Broadcast on the caller's transaction and returns its error. Together they turn N row events into one summary event.

  • Centrifugo driver (centrifugo.Driver, driver name centrifugo):

    • centrifugo.Client POSTs publish, broadcast, presence and unsubscribe calls with Authorization: apikey <key> 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.

    • centrifugo.ProxyHandler is the subscribe proxy endpoint. Centrifugo reads a non-200 status as an internal error, so every answer is HTTP 200. The checks run in this order:

      1. X-Centrifugo-Secret must equal realtime.centrifugo.proxy_secret, compared in constant time. An empty configured secret denies every subscribe.
      2. An empty or "0" user denies. The user may arrive as a JSON string or number; any other type counts as empty.
      3. A missing channel denies. Centrifugo always sends one.
      4. The channel's namespace must have a registered authorizer.
      5. The authorizer receives the user id and the full original channel.

      An allow answers {"result":{"info":…}}, with an empty info encoded as []. A presence: channel also gets allow (the authorizer's capabilities, or ["prs"]) and override. The override starts from the defaults presence and join_leave true and force_push_join_leave false, then applies the authorizer's overrides. Every deny answers {"error":{"code":403,"message":"Access denied"}} and logs Subscription denied at Warn with the reason, and never either secret. The request body is capped at 64 KiB.

Usage

An application selects the driver in config/realtime.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, installs a user lookup and registers its channel authorizers:

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
	})
	// Allow room:{id} to members only; re-checked on every subscribe.
	return svc.Registry().Register("room", lighthouse.AuthorizerFunc(
		func(ctx context.Context, userID uint, channel string) lighthouse.Result {
			if isRoomMember(ctx, userID, lighthouse.ChannelID(channel)) {
				return lighthouse.Allowed(nil)
			}
			return lighthouse.Denied("not a room member")
		}))
}

and mounts the driver's routes once:

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"),
	})
}

A model package stays free of realtime code; the plugin binds the model at Boot:

err := lighthouse.Bind[models.Post](svc, lighthouse.Binding[models.Post]{
	Alias: "blog.post",
	Channels: func(ctx context.Context, tx *gorm.DB, p *models.Post) ([]string, error) {
		return []string{"blog:" + strconv.FormatUint(uint64(p.BlogID), 10)}, nil
	},
})

A bulk import suppresses the per-row events and publishes one summary after commit:

err := lighthouse.WithoutBroadcasting[models.Post](ctx, func(ctx context.Context) error {
	return lagoon.Transaction(ctx, gdb, func(ctx context.Context, tx *gorm.DB) error {
		for _, p := range posts {
			if err := tx.Create(&p).Error; err != nil {
				return err
			}
		}
		return svc.Emit(ctx, tx, lighthouse.Broadcast{
			Channels: []string{"blog:7"},
			Event:    "blog.bulk_updated",
			Payload: struct {
				Reason string `json:"reason"`
				Count  int    `json:"count"`
			}{"import", len(posts)},
		})
	})
})

Tests select the memory driver and read what was published:

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, Registry, 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.Authorizer, lighthouse.AuthorizerFunc Authorize(ctx, userID, channel) lighthouse.Result.
lighthouse.Result, lighthouse.Allowed, lighthouse.Denied A subscribe decision: Allowed, Info, Capabilities, Overrides and Reason().
lighthouse.Registry, lighthouse.NewRegistry Namespace to authorizer map: Register, Get, Namespaces.
lighthouse.ParseChannel, lighthouse.ChannelID, lighthouse.PHPInt, lighthouse.FormatChannels Channel rules.
lighthouse.WithClientID, lighthouse.ClientID The realtime client id of a subscribe request, carried in the context.
lighthouse.Action, lighthouse.ActionCreated, lighthouse.ActionUpdated, lighthouse.ActionDeleted Broadcast actions.
lighthouse.Event Action, Actor, Timestamp, TTL of a change.
lighthouse.Broadcastable, lighthouse.BroadcastPayloader, lighthouse.BroadcastAliaser, lighthouse.BroadcastFilter, lighthouse.BroadcastTTLer The model-method broadcast contract.
lighthouse.Binding, lighthouse.Bind Broadcast contract registered from outside the model package.
lighthouse.WithoutBroadcasting Suppresses one model type for writes made with the given context.
lighthouse.Broadcast, lighthouse.Service.Emit One explicit event enqueued on the caller's transaction.
lighthouse.BroadcastArgs The River job (kind summer.broadcast).
lighthouse.DefaultTTL The default payload ttl, 60 seconds.
lighthouse.CallbackSnapshot, lighthouse.CallbackAfterCreate, lighthouse.CallbackAfterUpdate, lighthouse.CallbackAfterDelete Names of the GORM callbacks.
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, plus TrustedProxies from http.trusted_proxies for logging client IPs.
centrifugo.Client, centrifugo.NewClient HTTP API client: Publish, Broadcast, Presence, Unsubscribe, Info (the connectivity probe; an error body is an error), 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.ProxyHandler(svc, cfg) The subscribe proxy handler.
centrifugo.Driver, centrifugo.NewDriver, centrifugo.DriverName The lighthouse.Driver, with Client, Issuer, Config and Enabled (an API key is set).
centrifugo.Commands(app), centrifugo.HealthCommandName The websockets:health command and its name.
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.

CLI commands

centrifugo.Commands(app) returns websockets:health for the application binary. An application adds it to the list its plugin returns from Commands.

Command Description
websockets:health With an empty realtime.centrifugo.api_key it prints Centrifugo not configured (API key missing) and exits 1 without sending a request. Otherwise it prints the API URL and calls the Centrifugo info API method. On success it prints Configuration OK and a Setting/Value table (API URL, Enabled, API Key Set); on any failure it prints Connection check failed: … and exits 1. The API key is never printed, only whether it is set.

Dependencies

  • backpack, bouncer, compass, conga (the broadcast job), lagoon (callback installation), pact and wire from this repository; the centrifugo driver also uses surf for the client IP and bonfire for its command.
  • gorm.io/gorm (broadcast callbacks).
  • github.com/golang-jwt/jwt/v5 (centrifugo token signing).
  • The Centrifugo client is plain net/http; no Centrifugo SDK is used.

Testing

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. Broadcast tests need a running conga worker (conga.StartWorker) to deliver the jobs.