# MessageBus A thin wrapper over Redis pub/sub. Use it for fire-and-forget cross-server events (party invites, chat broadcasts, "player joined queue X"). For request/response, see [07_RPC.md](07_RPC.md). ## API ```java public interface MessageBus { void publish(String channel, Object payload); Subscription subscribe(String channel, Class type, Consumer handler); interface Subscription extends AutoCloseable { @Override void close(); } } ``` ## Publishing ```java PartyInvite invite = new PartyInvite(fromUuid, toUuid); core.getMessageBus().publish("party.invite", invite); ``` `payload` can be any class Gson can serialize. The bus wraps it in an envelope before sending. The call is non-blocking - it hands off to Lettuce's async API and returns immediately. There is no acknowledgement that another server received it. Pub/sub is at-most-once delivery: if no server is subscribed when you publish, the message disappears. ## Subscribing ```java MessageBus.Subscription sub = core.getMessageBus().subscribe( "party.invite", PartyInvite.class, invite -> handleInvite(invite) ); ``` Save the subscription if you want to unsubscribe later (`sub.close()`). NetworkCore's own subscriptions are cleaned up on shutdown, but yours aren't tied to anything by default. ### Thread safety - read this Callbacks run on an internal worker pool (`networkcore-bus-*` threads), **not** on a Hytale world thread. If your handler touches world or entity state, hop to the right thread first: ```java core.getMessageBus().subscribe("party.invite", PartyInvite.class, invite -> { World world = playerRef.getWorld(); world.execute(() -> { // Now safe to touch entities, components, etc. playerRef.sendMessage(Message.raw("Party invite from " + invite.from)); }); }); ``` Failing to hop will mostly work most of the time and then crash unpredictably under load. Always hop. The worker pool is sized 2 core / 8 max, queue 1024. A slow handler will not block other subscribers - it'll just back up its own queue. If the queue fills, the publisher thread runs the task itself (CallerRunsPolicy), which slows down message intake until the backlog clears. ## The envelope What actually goes over Redis: ```json { "v": 1, "from": "lobby-3f2a91bd", "auth": "", "ts": 1716610000000, "type": "com.example.PartyInvite", "body": { ... your payload ... } } ``` You don't construct this. The bus does it. `type` is informational only - dispatch is by channel, not by class. ## Channel naming A loose convention to keep things grep-able: - `party.invite`, `party.disband` - simple cross-server events - `chat.global` - broadcast channels - `rpc:req:`, `rpc:resp:` - reserved by RPC, don't touch - `network:events` - reserved by NetworkCore for hello/goodbye, don't touch ## Auth filtering If `network.auth_token` is set (see [03_Configuration.md](03_Configuration.md)), inbound messages whose envelope `auth` doesn't match are dropped silently (with a rate-limited warning log and an increment of `messagebus_rejected_total`). Outbound messages get stamped automatically. If a misconfigured server keeps spamming the bus with the wrong token, you'll see one log line every 10s rather than one per message. ## Performance and limits - Lettuce uses one connection for publish, one for subscribe. They share threading via Netty. - Pub/sub is in-memory in Redis. No persistence. If Redis restarts, in-flight messages are lost. - Redis pub/sub bandwidth is real but small - the practical bottleneck is usually subscriber-side processing, not network. - Payloads are JSON. Keep them small. Don't put 10MB worlds through the bus - put them in MySQL or object storage and bus a reference. ## Common patterns ### Broadcast to everyone Just publish - every server subscribed to the channel gets it. ```java core.getMessageBus().publish("chat.global", chatMessage); ``` ### Target a specific server Encode the target in the payload, let receivers filter: ```java class TargetedAction { String targetServerId; String action; Map data; } bus.subscribe("admin.action", TargetedAction.class, action -> { if (!action.targetServerId.equals(core.getServerId())) return; apply(action); }); ``` Or use a per-server channel: `admin.action.`. ### Avoid receiving your own messages Filter on `serverId`: ```java String myId = core.getServerId(); bus.subscribe("chat.global", ChatMessage.class, msg -> { if (myId.equals(msg.fromServerId)) return; relay(msg); }); ``` The bus does NOT filter your own publishes - all subscribers (including the publisher) see the message.