# 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`:

```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:

```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
	})
	// 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:

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

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

```go
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:

```go
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:

```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`, `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

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