Integrate, 6 minute read
Shared plugin messaging
Contents
Rift has one subject broker shared by routing hooks, HTTP hooks, Lua message
handlers, and external server plugins over QUIC. Native Rust components can also
use rift::messaging::Broker directly. Payloads are arbitrary bytes, so plugins
can choose JSON, protobuf, or their existing binary format.
Existing Bukkit/Paper plugins that use the BungeeCord plugin channel can enable BungeeCord compatibility for player queries and transfers over Minecraft connections, independently of this broker.
The API follows NATS-style subjects, queue groups and request/reply, with a separate retained-stream API for replay and acknowledgements. It uses Rift's own documented wire protocol; existing NATS clients do not speak this protocol. This is a single-process broker, without clustering, replication, transactions, or exactly-once processing. See the official NATS concepts and JetStream documentation for those systems' full semantics.
Local Rust plugins
Pass clones of the same broker to components. Creating a new broker creates an
independent bus. publish does not await network, disk or subscriber processing;
recv and request are asynchronous. Cloned messages share payload allocations.
use bytes::Bytes;
use rift::messaging::Broker;
let broker = Broker::default();
let mut events = broker.subscribe("game.player.*", None)?;
let result = broker.publish("game.player.join", Bytes::from_static(b"hello"))?;
assert_eq!(result.delivered, 1);
assert_eq!(events.recv().await?.payload, "hello");* matches one token; terminal > matches one or more tokens. Publishers use
literal subjects. Without a queue group every matching subscription receives a
copy. With a queue group only one matching member receives a given publication.
Drop a subscription to unregister it. Subscribe before publishing or declaring a
responder ready; core pub/sub has no history.
The broker bounds payload size, subject size, subscription count and subscriber
queue capacity. A full subscriber queue disconnects that subscriber and reports
SlowConsumer; a publish report counts slow consumers separately from successful
deliveries. Publishers do not wait for slow subscribers. Consumers that cannot
lose messages should use retained streams with explicit acknowledgement.
QUIC server plugins
Enable the listener explicitly in Lua configuration. It uses UDP, independent of Minecraft's TCP listener. Configure a certificate with the DNS/IP name clients will verify, and load token secrets through environment variables. Keep one connection and publisher open instead of reconnecting for each message.
messaging = {
listen = "127.0.0.1:4222",
certificate = "messaging-cert.pem",
private_key = "messaging-key.pem",
principals = {
{
name = "paper",
token_env = "RIFT_PAPER_TOKEN",
publish = { "game.>", "_INBOX.>" },
subscribe = { "game.>", "rift.events.>" },
},
{
name = "operator",
token_env = "RIFT_OPERATOR_TOKEN",
publish = { "rift.control.*" },
subscribe = { "rift.events.>" },
control = true,
},
},
}Set this field as rift.config.messaging = { ... } in a script, or inside a
legacy returned configuration table. Tokens must contain
32–1024 bytes without control characters. An external principal needs explicit
publish/subscribe grants. Control subjects additionally require control = true;
external clients cannot forge rift.events publications. Certificate and storage
paths are relative to Rift's working directory. Messaging configuration changes,
including credentials and stream limits, require a restart. Lua callback source
can reload while connections, subscriptions and queued messages remain intact.
use bytes::Bytes;
use rift::messaging::quic::{Client, tls_client_config};
let client = Client::connect(
"127.0.0.1:4222".parse()?, "localhost",
tls_client_config("messaging-cert.pem")?,
&std::env::var("RIFT_PAPER_TOKEN")?,
).await?;
let mut messages = client.subscribe("game.player.*", None).await?;
let mut publisher = client.publisher().await?;
publisher.publish("game.player.join", None, Bytes::from_static(b"hello")).await?;
publisher.flush().await?; // Barrier: preceding publications reached the broker.
let message = messages.recv().await?;
assert_eq!(message.payload, "hello");Client::publish waits for broker acceptance. Publisher::publish pipelines
unconfirmed writes for throughput; flush checks acceptance of earlier writes.
Neither means a live subscriber has finished processing. Each subscription and
request has its own QUIC stream. Other languages can implement the small binary
protocol using a QUIC library; there is no Java/JavaScript SDK in this change.
Lua plugins
Use rift.publish in on_route, on_http, and on_message. Async subscriptions
live outside Lua VMs; callbacks run on bounded blocking workers using the existing
Lua sandbox, instruction budget and execution deadline. See the
Lua configuration and callback reference for a complete
example, or start from examples/messaging.lua.
VM globals do not persist between callbacks. Rust/QUIC clients provide
the lower-overhead path when Lua VM creation is too expensive.
Rift control subjects
Send a JSON request using Client::request(subject, payload, timeout). Responses
are {"ok":true,"data":...} or {"ok":false,"error":"..."}. They use the same
handlers as Rift's existing administration service, even if that service's TCP
listener is disabled. The control bridge requires an _INBOX. reply subject;
request helpers create one automatically. Control requests are serialized.
| Subject | JSON payload | Result |
|---|---|---|
rift.control.status |
{} |
Connections, players, backends and maintenance state |
rift.control.maintenance |
{"enabled":true} |
Enable/disable new login maintenance |
rift.control.drain |
{"backend":"lobby","enabled":true} |
Enable/disable draining for a backend |
rift.control.transfer |
{"connection_id":123,"backend":"lobby"} |
Transfer an online player with existing protocol/access checks |
rift.control.reload |
{} |
Reload the selected configuration file |
rift.events.control reports command name and success, excluding status/player
details. rift.events.reload reports the active revision after successful reload.
rift.events.lifecycle reports ready and draining. These are ephemeral events;
query status to get current state. A timed-out command may already have executed;
check status before retrying a mutation. No automatic mutation retry is performed.
Retained streams
Streams are separate from live publication to keep storage out of the hot path.
Only explicit stream_publish / Stream::publish calls append a retained record;
ordinary broker publications do not persist, and stream publications do not
automatically fan out to live subscriptions. Publish once to each API when both
behaviors are desired; those two operations are not an atomic transaction.
-- Inside messaging:
streams = {
jobs = {
subjects = { "game.jobs.>" },
storage_path = "data/jobs", -- Omit for memory-only storage.
max_messages = 100000,
max_bytes = 67108864,
max_payload_bytes = 65536,
},
},The QUIC client exposes stream_publish, consumer, fetch, and ack.
ConsumerConfig selects a filter, starting sequence, acknowledgement deadline
and maximum in-flight records. Fetch returns sequence numbers and delivery counts;
acknowledge only after processing. Unacknowledged records become available again
after the deadline, and after reopening a file store. Remote durable names are
scoped to authenticated principal names. Reopen the same name with the same
configuration to resume. Use a distinct consumer per independent subscription or
share a name among workers in the same principal.
Retention bounds apply regardless of unacknowledged messages: records evicted by
count/byte limits cannot be redelivered. File storage persists records and consumer
state outside Tokio traffic threads; with the default sync_on_write = true,
acknowledgements wait for file synchronization. Disabling synchronization permits
loss of acknowledged writes on power failure. Unix also synchronizes directory
metadata; Windows has weaker power-loss guarantees for compaction renames. Disk
latency is consequently part of durable operations. Memory streams survive
client reconnects but disappear when Rift exits. File streams support one writer
process per storage directory, without replication or host-failure redundancy.
Measure latency and throughput
cargo test --locked --all-targets -- --test-threads=1
cargo run --locked --release --example messaging_bench -- 20000The benchmark reports local publish/receive and warm QUIC loopback round-trip p50/p95/p99 latency for 64-byte payloads, plus pipelined throughput. These figures exclude connection/TLS setup and persistence and check delivered payloads. Separate measurements cover memory stream publish/fetch/ack and synchronized file publication in the OS temporary directory (whose storage may be memory-backed). Run repeatedly on the target hardware under representative fanout and payloads. The benchmark is informational, with no timing thresholds in CI. Network distance, scheduling, TLS, fanout and disk synchronization cannot have zero latency.
To measure file persistence on a specific filesystem instead of a potentially
memory-backed OS temporary directory, point TMPDIR at a directory on that
filesystem. For example, from a source checkout on Unix:
mkdir -p target
TMPDIR="$PWD/target" cargo run --locked --release --example messaging_bench -- 20000File publication includes disk synchronization and blocking-worker dispatch; the core/QUIC numbers exclude persistence.