Services

Realtime

Publish model changes and events to realtime channels with lighthouse, authorize subscriptions per channel namespace, and run the Centrifugo driver.

On this page

lighthouse is the SummerCMS counterpart of the WinterCMS websockets plugin. The application publishes events to named channels; the frontend holds a connection to a realtime server, subscribes to channels and receives the events. SummerCMS does not run the connection server itself. The Centrifugo driver publishes to a Centrifugo server through its HTTP API, issues the connection tokens the frontend needs, and answers Centrifugo's subscribe checks.

Drivers#

realtime.driver selects the driver:

Driver Publishes
null (default) Nothing.
log To the log: channel names and the event, never the payload.
memory Into memory, readable with lighthouse.MemoryDriver.Publications, for tests.
centrifugo To Centrifugo, from the lighthouse/centrifugo package.

A driver registers itself from its package's init function, as database/sql drivers do, so the application imports the driver package for its side effect: _ ".../modules/lighthouse/centrifugo". An unknown driver name stops the start-up with the list of registered drivers. lighthouse.From returns the application's lighthouse.Service, which holds the driver, the authorizer registry and the broadcast settings.

Channels and authorization#

A channel name is namespace:entity:id, optionally prefixed once with presence:. A plugin registers a lighthouse.Authorizer per namespace on the service's lighthouse.Registry. The driver asks the namespace's authorizer on every subscribe, so a user who loses access is refused the next time the client subscribes; nothing is cached:

modules/lighthouse/example_test.go#ExampleRegistry_Register
app, err := newApp(map[string]any{"realtime.driver": "memory"})
if err != nil {
	fmt.Println(err)
	return
}
svc, err := lighthouse.From(app)
if err != nil {
	fmt.Println(err)
	return
}

// blog:{entity}:{id} channels are open to members of the blog only. The
// authorizer runs on every subscribe; nothing is cached.
err = svc.Registry().Register("blog", lighthouse.AuthorizerFunc(
	func(ctx context.Context, userID uint, channel string) lighthouse.Result {
		if isMember(ctx, userID, lighthouse.ChannelID(channel)) {
			return lighthouse.Allowed(nil)
		}
		return lighthouse.Denied("not a member of the blog")
	}))
if err != nil {
	fmt.Println(err)
	return
}

for _, sub := range []struct {
	user    uint
	channel string
}{{42, "blog:7"}, {42, "blog:8"}, {42, "presence:blog:7"}, {42, "shop:7"}} {
	ns, presence := lighthouse.ParseChannel(sub.channel)
	auth, ok := svc.Registry().Get(ns)
	if !ok {
		fmt.Println(sub.channel, "no authorizer")
		continue
	}
	res := auth.Authorize(context.Background(), sub.user, sub.channel)
	fmt.Printf("%d %s namespace=%s presence=%v allowed=%v reason=%q\n", sub.user, sub.channel, ns, presence, res.Allowed, res.Reason())
}
fmt.Println(lighthouse.ChannelID("blog:12abc"), lighthouse.FormatChannels("acme", []string{"Blog:7"}))
// Output:
// 42 blog:7 namespace=blog presence=false allowed=true reason=""
// 42 blog:8 namespace=blog presence=false allowed=false reason="not a member of the blog"
// 42 presence:blog:7 namespace=blog presence=true allowed=false reason="not a member of the blog"
// shop:7 no authorizer
// 12 [acme:blog:7]

The channel rules follow the WinterCMS plugin byte for byte:

  • lighthouse.ParseChannel returns the namespace and whether the channel is a presence channel. A doubled presence: prefix or more than three segments give an empty namespace, which no authorizer matches.
  • lighthouse.ChannelID reads segment 1 with PHP's (int) cast: 12abc is 12. For a presence: channel, segment 1 is the namespace, so the ID is 0, as the presence line of the example shows. An authorizer for presence channels must parse the ID itself.
  • lighthouse.FormatChannels lowercases channel names and applies the realtime.broadcast_namespace prefix.

A denial's reason goes to the log only; the client always sees the same refusal.

Mounting the driver's routes#

A driver may need HTTP routes. The Centrifugo driver has two: the token route, which signed-in users call, and the subscribe proxy, which Centrifugo calls. The application mounts them once, from a plugin's Routes, with lighthouse.Mount, choosing the middleware per surface:

modules/lighthouse/centrifugo/example_test.go#ExampleDriver_Routes
app, err := newApp(map[string]any{"realtime.driver": "centrifugo"})
if err != nil {
	fmt.Println(err)
	return
}
svc, err := lighthouse.From(app)
if err != nil {
	fmt.Println(err)
	return
}
// In a plugin's Routes method, r is the router the plugin receives.
r := surf.New(nil)
err = lighthouse.Mount(r, svc.Driver(), lighthouse.Surfaces{
	UserAuth:   surf.Use("acme.auth"),
	Middleware: surf.Use("throttle:60,1"),
})
if err != nil {
	fmt.Println(err)
	return
}
for _, rt := range r.Routes() {
	fmt.Println(rt.Method, rt.Pattern, rt.Middleware, "raw:", rt.Raw)
}

// A user route without a guard is refused.
err = lighthouse.Mount(surf.New(nil), svc.Driver(), lighthouse.Surfaces{})
fmt.Println(err != nil)
// Output:
// GET /api/realtime/token [acme.auth throttle:60,1] raw: false
// POST /api/realtime/subscribe [throttle:60,1] raw: true
// true

lighthouse.UserAuth routes get the UserAuth middleware, lighthouse.ServerToServer routes are mounted in a raw group, and Middleware is added to every route after the surface's own. A user route without a guard is refused, so the token route can never be exposed to anonymous callers. Switching drivers never changes the application's route declarations.

Model broadcasts#

A model broadcasts its creates, updates and deletes when a lighthouse.Binding is registered for it, or when its pointer type implements lighthouse.Broadcastable. A binding keeps realtime code out of the models package:

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

The event name is {action}.{alias}, here created.blog.post, and the default payload is {"model":...,"actor":...,"timestamp":"...+00:00","ttl":60}. A binding's Payload, ShouldBroadcast and TTL fields, or the matching model methods, replace the defaults.

Delivery is transactional. The write enqueues a broadcast job, through conga, inside its own transaction, so nothing is published for a write that rolls back, and the job publishes after the commit:

modules/lighthouse/example_test.go#write
return lagoon.Transaction(ctx, db, func(ctx context.Context, tx *gorm.DB) error {
	if err := tx.Create(&Post{BlogID: 7, Title: "Hello"}).Error; err != nil {
		return err
	}
	// The broadcast job is now queued in this transaction. It is
	// published only if the transaction commits.
	if fail {
		return fmt.Errorf("rolled back")
	}
	return nil
})

The job runs once, best effort: a failed publish is logged as realtime: broadcast failed and never affects the write. Delivery order across separate jobs is not guaranteed. A job worker must be running, in serve or in queue:work; see Queued jobs. A write without a primary key value, such as Model(&Post{}).Where(...).Updates(...), is not broadcast.

Bulk writes#

lighthouse.WithoutBroadcasting silences one model type for writes made with the context it hands to its function; other types still broadcast. lighthouse.Service.Emit enqueues one explicit event on the caller's transaction. Together they turn a thousand row events into one summary:

modules/lighthouse/example_test.go#bulk
return lighthouse.WithoutBroadcasting[Post](ctx, func(ctx context.Context) error {
	return lagoon.Transaction(ctx, db, func(ctx context.Context, tx *gorm.DB) error {
		for _, title := range titles {
			if err := tx.Create(&Post{BlogID: 7, Title: title}).Error; err != nil {
				return err
			}
		}
		return svc.Emit(ctx, tx, lighthouse.Broadcast{
			Channels: []string{"blog:7"},
			Event:    "blog.posts_imported",
			Payload: struct {
				Count int `json:"count"`
			}{len(titles)},
		})
	})
})

Only writes that use the context passed to the function are silenced, so write through it, as lagoon.Transaction does above.

The Centrifugo driver#

The driver reads realtime.centrifugo.* from config/realtime.yaml. Secrets go in the environment:

driver: centrifugo
centrifugo:
  api_url: http://127.0.0.1:8001/api
  ws_url: /ws

with SUMMER_REALTIME__CENTRIFUGO__API_KEY, SUMMER_REALTIME__CENTRIFUGO__TOKEN_SECRET and SUMMER_REALTIME__CENTRIFUGO__PROXY_SECRET set.

  • The token route (realtime.centrifugo.token_path, /api/realtime/token by default) answers a signed-in user with {"token":"..."}, an HS256 connection token signed with the token secret. It answers 401 without a user and 503 when the token secret is empty.
  • The subscribe proxy (realtime.centrifugo.subscribe_path) accepts a call only when its X-Centrifugo-Secret header equals the proxy secret, compared in constant time; an empty proxy secret refuses every subscribe. It then asks the channel's authorizer. Every answer is HTTP 200, as Centrifugo requires, with the decision in the body.
  • Publishing uses the HTTP API with the API key. With an empty API key nothing is sent and no broadcast jobs are queued.

The subscribe proxy runs the same authorizers as above:

modules/lighthouse/centrifugo/example_test.go#ExampleProxyHandler
svc, err := lighthouse.From(backpack.New(nil))
if err != nil {
	fmt.Println(err)
	return
}
// Members of blog 7 may subscribe to its channels.
err = svc.Registry().Register("blog", lighthouse.AuthorizerFunc(
	func(ctx context.Context, userID uint, channel string) lighthouse.Result {
		if userID == 42 && lighthouse.ChannelID(channel) == 7 {
			return lighthouse.Allowed(nil)
		}
		return lighthouse.Denied("not a member of the blog")
	}))
if err != nil {
	fmt.Println(err)
	return
}
// realtime.centrifugo.proxy_secret; set it through the environment.
proxy := centrifugo.ProxyHandler(svc, centrifugo.Config{ProxySecret: "test-only-proxy-secret"})

// What Centrifugo posts to the subscribe proxy.
subscribe := func(secret, user, channel string) {
	body := fmt.Sprintf(`{"client":"c1","user":%q,"channel":%q}`, user, channel)
	req := httptest.NewRequest(http.MethodPost, "/api/realtime/subscribe", strings.NewReader(body))
	req.Header.Set("X-Centrifugo-Secret", secret)
	rec := httptest.NewRecorder()
	proxy.ServeHTTP(rec, req)
	fmt.Println(rec.Code, strings.TrimSpace(rec.Body.String()))
}
subscribe("test-only-proxy-secret", "42", "blog:7")
subscribe("test-only-proxy-secret", "5", "blog:7")
subscribe("wrong-secret", "42", "blog:7")
// Output:
// 200 {"result":{"info":[]}}
// 200 {"error":{"code":403,"message":"Access denied"}}
// 200 {"error":{"code":403,"message":"Access denied"}}

Configure Centrifugo to call the subscribe proxy with the same secret, and keep its HTTP API on a private address.

websockets:health checks the connection to Centrifugo and prints the settings, never the API key:

./bin/acme websockets:health