Repository navigation
Custom StackExchange
- Redis and NATS (see Scale out) don't fit: you already run a different broker, or you want zero extra moving parts for a single-machine, multi-process deployment.
- You want to understand what a
StackExchangeactually has to do, before trusting a third-party one. - You're testing scale-out behavior without spinning up a real broker.
neffos.StackExchange is the interface Scale out is built on. Implement it, pass it to Server.UseStackExchange, and every Server.Broadcast and Server.Ask call routes through it instead of staying local.
type StackExchange interface {
OnConnect(c *Conn) error
OnDisconnect(c *Conn)
Publish(msgs []Message) bool
Subscribe(c *Conn, namespace string)
Unsubscribe(c *Conn, namespace string)
Ask(ctx context.Context, msg Message, token string) (Message, error)
NotifyAsk(msg Message, token string) error
}| Method | The server calls it when |
|---|---|
OnConnect(c) |
A connection finishes its handshake, before the server's own Server.OnConnect, so a rejection here also rejects the connection. |
OnDisconnect(c) |
The connection goes away, from any cause. |
Subscribe(c, namespace) |
This connection's namespace handshake completes (OnNamespaceConnected). |
Unsubscribe(c, namespace) |
This connection's namespace disconnects. |
Publish(msgs) |
Server.Broadcast is called; publishing replaces, not supplements, the local fan-out. |
Ask(ctx, msg, token) |
Server.Ask is called; the reply may come from a connection on a different server. |
NotifyAsk(msg, token) |
The core finds an incoming reply whose wait token it cannot match locally and this server has a StackExchange configured; it hands the reply to this method instead of dropping it. |
All seven are required; there is no default or embeddable base to leave some out.
type StackExchangeInitializer interface {
Init(Namespaces) error
}
type StackExchangeCloser interface {
Close() error
}Init runs once, from Server.UseStackExchange itself, with the server's registered namespaces, before any connection exists; a broker-backed exchange can open its connection here and fail the call at startup instead of on the first client. Init only runs through UseStackExchange; setting Server.StackExchange directly skips it.
Close runs once, from Server.Close or Server.Shutdown, whichever runs first, so a broker connection or a background goroutine gets released when the server shuts down. Both the redis and nats exchanges implement it; implement it on anything that holds a connection or a goroutine, and make it safe to call more than once; Server.Close is itself idempotent, and stacking more than one exchange (see "Registering" below) can run Close through the same exchange twice.
Server.Ask builds a wait token, then does something your StackExchange never has to know about: it marks a second copy of that token before it goes on the wire, by inserting one ! character right after the token's first character. Ask(ctx, msg, token) itself always receives the clean, unmarked token; msg.wait (not an exported field) carries the marked one.
Why this matters for your implementation:
- You publish
msgas given; the marked token travels inside it to whatever connection receives it. - When a server-side
Conndeserializes an incoming frame and finds the marker, it strips it, setsMessage.FromStackExchange = trueon the result, and keeps the clean token. So the client that answers never sees the marker: it just echoes the token it was given, now clean. - When that answer arrives back at a server, the clean token matches no local
Askcall (the real waiter is sitting in yourStackExchange.Ask, not in the core's own wait map). The core recognizes this, either because the message reportsFromStackExchangeor because an unmarked wait token that nothing local is waiting for falls through to the same place, and callsStackExchange.NotifyAsk(msg, token)with the clean token. -
NotifyAskis your job to route: delivermsgto whatever is waiting ontoken, by whatever means your exchange uses to park that wait (a local channel keyed by the token, as the walkthrough below does; a broker subscription named after the token, as redis and nats do).
A reply nobody is waiting for anymore (the Ask call already timed out, or two connections answered a namespace-wide ask) should be dropped by NotifyAsk, not treated as an error: Server.Ask only ever waits for the first one.
Server.Broadcast(exceptSender, msgs...) can be called with a sender to exclude, but with a StackExchange configured the message does not stay local, so a plain in-memory "skip this *Conn" check cannot work once the message reaches another server (a pointer from one process means nothing in another). Instead, neffos fills Message.FromExplicit with a string that identifies the excluded connection well enough to survive the trip (for every neffos.Sender: a *Conn, a *NSConn, a *Room, or neffos.Exclude(id) resolved against the local server's connections), and every connection's own Send refuses to deliver a message whose FromExplicit names it, wherever that connection happens to run. Your StackExchange.Publish does not need to interpret FromExplicit at all: publish msgs to every subscriber exactly as given, and let each connection's own write path decide whether it is the one being excluded.
_examples/07-scale-out/custom-stackexchange builds the simplest possible StackExchange: a bus, a broker that lives in process memory, shared by two or more *neffos.Server values running in the same Go program, with no network broker at all. Each server wraps the shared bus in its own exchange value (so log lines and any per-server state stay separate), and exchange is what actually implements StackExchange:
// bus is the broker: named channels with subscribed connections, and the
// waiting asks by token. Every exchange on the same bus sees every message.
type bus struct {
mu sync.Mutex
subs map[string]map[*neffos.Conn]struct{} // channel name to subscribers
asks map[string]chan []byte // ask token to its reply
}
// exchange is one server's StackExchange on the bus.
type exchange struct {
name string
bus *bus
}
var (
_ neffos.StackExchange = (*exchange)(nil)
_ neffos.StackExchangeInitializer = (*exchange)(nil)
_ neffos.StackExchangeCloser = (*exchange)(nil)
)
func connChannel(id string) string { return "conn." + id }
func namespaceChannel(ns string) string { return "namespace." + ns }
func channelOf(msg neffos.Message) string { // rooms travel on their namespace
if msg.To != "" {
return connChannel(msg.To)
}
return namespaceChannel(msg.Namespace)
}The seven required methods, plus the two optional ones, map onto the bus almost one for one:
// Init runs once, from Server.UseStackExchange.
func (e *exchange) Init(namespaces neffos.Namespaces) error {
log.Printf("exchange %s ready for namespaces %v", e.name, slices.Sorted(maps.Keys(namespaces)))
return nil
}
func (e *exchange) OnConnect(c *neffos.Conn) error {
e.bus.subscribe(connChannel(c.ID()), c)
return nil
}
func (e *exchange) OnDisconnect(c *neffos.Conn) {
e.bus.unsubscribeAll(c)
}
func (e *exchange) Subscribe(c *neffos.Conn, namespace string) {
e.bus.subscribe(namespaceChannel(namespace), c)
}
func (e *exchange) Unsubscribe(c *neffos.Conn, namespace string) {
e.bus.unsubscribe(namespaceChannel(namespace), c)
}
func (e *exchange) Publish(msgs []neffos.Message) bool {
for _, msg := range msgs {
e.bus.publish(channelOf(msg), msg.Serialize())
}
return true
}
// Ask publishes msg and waits for the reply that NotifyAsk sends with token.
func (e *exchange) Ask(ctx context.Context, msg neffos.Message, token string) (neffos.Message, error) {
replies := e.bus.wait(token)
defer e.bus.forget(token)
e.bus.publish(channelOf(msg), msg.Serialize())
select {
case <-ctx.Done():
return neffos.Message{}, ctx.Err()
case payload := <-replies:
reply := neffos.DeserializeMessage(neffos.TextMessage, payload, false, false)
return reply, reply.Err
}
}
// NotifyAsk runs on the server that received the client's reply.
func (e *exchange) NotifyAsk(msg neffos.Message, token string) error {
msg.ClearWait()
e.bus.reply(token, msg.Serialize())
return nil
}
// Close runs from Server.Close and Server.Shutdown. The bus belongs to the
// program, not to one server, so there is nothing to release here.
func (e *exchange) Close() error {
log.Printf("exchange %s closed", e.name)
return nil
}Notice:
-
OnConnectsubscribes the connection's own channel, keyed byconnChannel(c.ID()), so a later message withMessage.Toset can reach it throughchannelOf;OnDisconnectunsubscribes everything that connection held. -
Messages travel serialized.
bus.publishhandsmsg.Serialize()([]byte) to every subscriber, and each one decodes it itself withc.DeserializeMessage(seebus.publish's implementation) beforec.Write-ing it, settingFromStackExchange = trueon the way, exactly as the wire format expects. A real broker would do the same thing across a network instead of a map. -
Subscriptions are namespace-wide, not room-aware:
StackExchangehas no room-level hook, so a room broadcast still goes out onnamespaceChannel, and each connection's ownWritedrops it for anyone who never joined that room. The redis and nats exchanges make the same tradeoff (see Scale out). -
Ask/NotifyAsknever touch the wire token's marker.Askgets the cleantokenas described above, parks a channel for it, publishes, and waits;NotifyAskis handed the clean token too and just needs to find that channel, which is exactly whatbus.wait/bus.replydo.
The pattern above is also the easiest way to test any StackExchange without a real broker: build two servers in the same test, share one exchange between them, and assert that traffic crosses from one to the other. _examples/07-scale-out/custom-stackexchange/main_test.go does exactly this:
// startServers starts servers a and b on one bus and returns their URLs.
func startServers(t *testing.T) (a, b *neffos.Server, urlA, urlB string) {
t.Helper()
shared := newBus()
a, b = newServer("a", shared), newServer("b", shared)
for _, srv := range []*neffos.Server{a, b} {
ts := httptest.NewServer(srv)
t.Cleanup(ts.Close)
t.Cleanup(srv.Close)
url := "ws" + strings.TrimPrefix(ts.URL, "http")
if srv == a {
urlA = url
} else {
urlB = url
}
}
return a, b, urlA, urlB
}A test then dials a client into each server's URL, connects both to "chat", and asserts that a message emitted on one side reaches the client on the other, or that a.Ask finds a reply from a connection that only exists on b. This is the same shape the redis and nats exchanges are tested with in this repository, minus the broker: two httptest.NewServer values sharing one exchange, two sets of Go clients, one assertion that traffic crosses from one server to the other.
server.UseStackExchange(exc)returns a non-nil error only if exc implements StackExchangeInitializer and Init failed; otherwise the server starts using it right away. Calling UseStackExchange more than once does not replace the previous exchange, it wraps it, in registration order. OnConnect stops at the first error; Subscribe, Unsubscribe, Publish and Close always run on both, regardless of what the first one returned (Publish and Close keep the first failure); Ask and NotifyAsk are the odd ones out, falling back to the next exchange only when the first one returns an error. That is how a migration between two brokers, or a bus plus a real broker, can run side by side without a third interface.
- Scale out, the built-in redis and nats backends this interface also powers
- Scale out using Redis and Scale out using Nats, two real implementations to compare against
-
The ask method, the single-server version of
Ask/NotifyAsk -
Broadcast,
Message.Toand theexceptSenderparameterFromExplicitis derived from
Home | About | Project | Getting Started | Technical Docs | Copyright © 2019-2026 Gerasimos Maropoulos. Documentation terms of use.
Getting started
Concepts
Messaging
Production
Scale out