High-availability controller, broker metadata coordinator, and OpenRaft-based election service for RocketMQ-Rust.
rocketmq-controller provides the controller runtime used by RocketMQ-Rust clusters. It owns controller bootstrap,
OpenRaft consensus, broker heartbeat tracking, master election, replica metadata, controller request processing,
persistent controller state, and optional metrics integration.
The crate exposes both a library API and the rocketmq-controller-rust binary.
| Area | What it provides |
|---|---|
| Controller bootstrap | rocketmq-controller-rust binary, CLI parsing, configuration loading, startup logging, remoting version setup, single-node bootstrap, and graceful shutdown. |
| Raft coordination | OpenRaft integration, gRPC raft transport, log store, state machine, cluster initialization, leader checks, and replicated controller events. |
| Broker coordination | Broker registration, heartbeats, inactive broker scanning, master election, sync-state-set management, broker ID allocation, broker cleanup, and role-change notifications. |
| Metadata management | Broker replica metadata, topic/config metadata, controller metadata queries, sync-state snapshots, and Java-compatible controller response models. |
| Request processing | Broker-facing remoting request processor for controller request codes such as elect-master, alter-sync-state-set, get-replica-info, broker heartbeat, and register broker. |
| Storage | RocksDB enabled by default, an opt-in file backend, and an in-memory backend for tests. |
| Observability | Controller request, election, heartbeat, and latency metrics, plus optional OpenTelemetry, OTLP, and Prometheus exporters. |
The controller starts from rocketmq-controller-rust, loads ControllerConfig, and builds ControllerManager.
ControllerManager owns broker-facing remoting, request processing, heartbeat tracking, role-change notifications,
OpenRaft-backed state replication, storage, snapshots, and optional metrics exporters.
The binary requires a non-empty rocketmqHome value after configuration loading. For normal runs, set
ROCKETMQ_HOME or provide rocketmqHome in the controller config file.
| Path | Purpose |
|---|---|
src/bin/controller_bootstrap.rs |
Binary entry point, logger setup, CLI/config loading, manager lifecycle, single-node bootstrap, and shutdown signal handling. |
src/cli.rs |
Clap-based CLI model and config-file loading through the config crate's config parser. |
src/config.rs |
Owns ControllerConfig, peer/address validation and storage settings, re-exported at the crate root. |
src/controller |
Controller trait implementations, ControllerManager, OpenRaft controller wrapper, heartbeat manager, and housekeeping service. |
src/openraft |
OpenRaft node manager, gRPC network, log store, state machine, storage bridge, and generated raft service glue. |
src/processor |
Controller request processor and domain processors for broker, topic, and metadata operations. |
src/manager |
Replica information manager, broker replica metadata, and sync-state models. |
src/openraft/state_machine.rs |
Broker, topic, config, and replica metadata stores. |
src/heartbeat |
Broker identity, live-info tracking, and default heartbeat manager. |
src/event |
Replicated controller event models and event serialization. |
src/storage |
Storage backend abstraction, default RocksDB backend, opt-in file backend, and in-memory backend for tests. |
src/qualification.rs |
T0-T5 failover timeline, PutOk recovery audit, confirm-offset boundary audit, and machine-readable qualification report. |
src/metrics |
Controller metric constants, request/election status enums, and metrics manager. |
proto |
gRPC protobuf definitions for controller and OpenRaft RPCs. |
examples |
Runnable examples for single-node raft, three-node raft, manager usage, metrics, and CLI parsing. |
tests |
Integration and contract tests for raft, snapshots, multi-node setup, request processor behavior, and metrics. |
- Stable Rust
1.95.0is the workspace minimum and the repository build toolchain. - Build this crate with the pinned repository toolchain from
../rust-toolchain.toml. ROCKETMQ_HOMEorrocketmqHomemust be set before starting the controller binary.- The default feature set enables RocksDB and the default configuration uses
storageBackend = "RocksDB". - Use
storageBackend = "File"only when building with the opt-instorage-filefeature.
Build the controller binary from the workspace root:
cargo build -p rocketmq-controller --bin rocketmq-controller-rust --releaseBuild with optional storage or metrics features:
cargo build -p rocketmq-controller --bin rocketmq-controller-rust --release --features storage-rocksdb
cargo build -p rocketmq-controller --bin rocketmq-controller-rust --release --features metrics
cargo build -p rocketmq-controller --bin rocketmq-controller-rust --release --features metrics-otlp
cargo build -p rocketmq-controller --bin rocketmq-controller-rust --release --features metrics-prometheusThe CLI loads TOML, JSON, YAML, and other formats supported by the config crate. Current file loading deserializes
directly into this crate's ControllerConfig. Its Serde defaults allow partial files; omitted fields retain their configured defaults.
Use camelCase field names:
rocketmqHome = "/opt/rocketmq"
configStorePath = "/opt/rocketmq/controller/controller.properties"
controllerType = "Raft"
scanNotActiveBrokerInterval = 5000
controllerThreadPoolNums = 16
controllerRequestThreadPoolQueueCapacity = 50000
mappedFileSize = 1073741824
controllerStorePath = ""
electMasterMaxRetryCount = 3
enableElectUncleanMaster = false
isProcessReadEvent = false
notifyBrokerRoleChanged = true
scanInactiveMasterInterval = 5000
raftScanWaitTimeoutMs = 1000
configBlackList = "configBlackList;configStorePath;maintenanceCheckpointRoot"
nodeId = 1
listenAddr = "127.0.0.1:60109"
controllerPeers = []
electionTimeoutMs = 1000
heartbeatIntervalMs = 300
storagePath = "/opt/rocketmq/controller/node-1"
storageBackend = "RocksDB"
enableElectUncleanMasterLocal = false
[[raftPeers]]
id = 1
addr = "127.0.0.1:60110"
[observability.metrics]
exporter = "disable"
[observability.traces]
exporter = "disable"
[observability.logs]
exporter = "disable"
[observability.otlp]
endpoint = "http://127.0.0.1:4317"
protocol = "grpc"
[observability.prometheus]
host = "127.0.0.1"
port = 5557
path = "/metrics"The legacy flat telemetry fields have been removed. Configure Controller observability only through the structured [observability] sections; present runtime environment variables override matching file values.
For a multi-node controller cluster, each node should use its own nodeId, listenAddr, and storagePath, while all
nodes share the same raftPeers list. listenAddr is the broker-facing remoting endpoint; each raftPeers.addr is the
advertised OpenRaft gRPC endpoint and therefore must use a separate port.
When the advertised Raft address is not bindable on the Pod or host, set
ROCKETMQ_CONTROLLER_RAFT_BIND_ADDR=<local-ip>:<raft-port> (for example, 0.0.0.0:60110). The advertised address remains
the matching raftPeers.addr; the override changes only the local listener. Invalid socket-address values fail startup.
A single-member configuration retains automatic bootstrap. Multi-member bootstrap is disabled unless
ROCKETMQ_CONTROLLER_AUTO_INITIALIZE_CLUSTER=true (or 1) is set. With that explicit opt-in, only the lowest configured
node ID initializes the full membership, and existing committed state is never reinitialized. Operators must still
verify that a leader and quorum formed; configuration and replica count alone are not quorum evidence.
Start a single-node controller with a config file:
# Linux/macOS
export ROCKETMQ_HOME=/opt/rocketmq
cargo run -p rocketmq-controller --bin rocketmq-controller-rust -- -c ./controller-node1.toml# Windows PowerShell
$env:ROCKETMQ_HOME = "C:\rocketmq"
cargo run -p rocketmq-controller --bin rocketmq-controller-rust -- -c .\controller-node1.tomlInspect CLI options:
cargo run -p rocketmq-controller --bin rocketmq-controller-rust -- --helpRun the single-node OpenRaft example:
cargo run -p rocketmq-controller --example single_nodeRun a three-node OpenRaft example in separate terminals:
cargo run -p rocketmq-controller --example three_node_cluster -- --node-id 1 --init
cargo run -p rocketmq-controller --example three_node_cluster -- --node-id 2
cargo run -p rocketmq-controller --example three_node_cluster -- --node-id 3Create and manage a controller directly from Rust:
use std::{future::Future, sync::Arc};
use rocketmq_controller::{ControllerConfig, ControllerManager, ControllerResult};
use rocketmq_observability::TelemetryHandle;
use rocketmq_runtime::ChildServiceContext;
async fn run_controller(
config: ControllerConfig,
context: ChildServiceContext,
shutdown: impl Future<Output = ()>,
) -> ControllerResult<()> {
let manager = Arc::new(
ControllerManager::new(config, context, TelemetryHandle::noop()).await?,
);
let result = async {
if !manager.initialize().await? {
return Err(rocketmq_error::Error::new(
&rocketmq_error::CONTROLLER_INTERNAL_FAILURE,
));
}
manager.start().await?;
shutdown.await;
Ok(())
}.await;
let shutdown_result = manager.shutdown().await;
result?;
shutdown_result
}The application supplies its owned ChildServiceContext and closes the RuntimeOwner after the controller stops. This function does not perform the binary's automatic cluster bootstrap; a new cluster needs an explicit controller().initialize_cluster(...) call. When authentication, authorization or maintenance is enabled, use a constructor accepting ControllerSecurity.
raftListenAddr, raftPeerEndpoints and controllerPeerEndpoints support separate bind and advertised endpoints. Do not mix the endpoint lists with legacy raftPeers/controllerPeers.
| Feature | Default | Purpose |
|---|---|---|
storage-file |
No | Enables the file-based controller storage backend for storageBackend = "File". |
storage-rocksdb |
Yes | Enables the default RocksDB backend for storageBackend = "RocksDB". |
metrics |
No | Enables controller metrics integration through rocketmq-observability. |
metrics-otlp |
No | Enables OTLP metrics export support. |
metrics-prometheus |
No | Enables Prometheus metrics export support. |
debug |
No | Reserved debug feature flag for controller builds. |
Additional flags include dev-single (enables storage-file), otel-traces, otel-logs,
otlp-traces and otlp-logs. Exporter features still require matching runtime configuration.
cargo run -p rocketmq-controller --example single_node
cargo run -p rocketmq-controller --example three_node_cluster -- --node-id 1 --init
cargo run -p rocketmq-controller --example controller_manager_basic
cargo run -p rocketmq-controller --example controller_manager_cluster
cargo run -p rocketmq-controller --example controller_metrics_example
cargo run -p rocketmq-controller --example cli_usage -- -c ./controller-node1.toml -pFocused checks for this crate:
cargo test -p rocketmq-controller --lib
cargo test -p rocketmq-controller --tests --no-run
cargo test -p rocketmq-controller --examples --no-run
cargo test -p rocketmq-controller --test controller_failover_sloSelect additional checks for this crate from the repository root:
cargo fmt -p rocketmq-controller -- --check
cargo clippy -p rocketmq-controller --no-deps -- -D warningsController metadata consensus and message payload durability are separate boundaries. A metadata leader election alone
does not prove message RPO=0. A strict RPO=0 result is accepted only when the tested broker path used synchronous local
flush, completed the required replica acknowledgements, retained clean election, recovered every message for which the
producer received PutOk, and kept confirmOffset monotonic and no greater than the legal in-sync acknowledgement.
The qualification API records one failover as an ordered timeline:
- T0 fault injected
- T1 Controller leader elected
- T2 broker master elected
- T3 store write authority granted
- T4 NameServer route converged
- T5 producer first succeeded
FailoverQualificationReport emits versioned JSON evidence and explicitly rejects incomplete or unsafe runs. It does
not turn the configured election timeout into a measured RTO and does not claim a percentile from a single run. Use the
end-to-end fault harness to collect repeated samples before publishing p50/p95/p99 or an SLO. The pinned OpenRaft alpha
dependency must continue to pass linearizability, storage-fault, multi-node, and failover qualification gates before a
release is promoted.
cargo bench -p rocketmq-controller --bench controller_benchThis benchmark measures deterministic Controller hot paths: replicated heartbeat application, Raft request encoding, and qualification evidence recording. Leader election and end-to-end RTO are integration/fault tests, not Criterion microbenchmarks.
Licensed under the Apache License, Version 2.0.