4.6 KiB
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.
API
public interface MessageBus {
void publish(String channel, Object payload);
<T> Subscription subscribe(String channel, Class<T> type, Consumer<T> handler);
interface Subscription extends AutoCloseable {
@Override void close();
}
}
Publishing
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
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:
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:
{
"v": 1,
"from": "lobby-3f2a91bd",
"auth": "<token from config>",
"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 eventschat.global- broadcast channelsrpc:req:<service>,rpc:resp:<replyTo>- reserved by RPC, don't touchnetwork:events- reserved by NetworkCore for hello/goodbye, don't touch
Auth filtering
If network.auth_token is set (see 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.
core.getMessageBus().publish("chat.global", chatMessage);
Target a specific server
Encode the target in the payload, let receivers filter:
class TargetedAction {
String targetServerId;
String action;
Map<String, Object> 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.<server-id>.
Avoid receiving your own messages
Filter on serverId:
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.