How RocketMQ-Rust Ensures Memory Safety: Ownership, Borrowing, and Safe Abstractions

RocketMQ-Rust guarantees memory safety by leveraging Rust's ownership model and confining all unsafe operations to well-tested abstractions like ArcMut and MappedBuffer, preventing data races, dangling pointers, and buffer overflows at compile time.

RocketMQ-Rust is a Rust implementation of Apache RocketMQ that prioritizes memory safety without sacrificing performance. Unlike traditional message queue implementations in C++, RocketMQ-Rust utilizes Rust's compile-time guarantees to eliminate entire classes of memory vulnerabilities. The codebase maintains safety through strict ownership discipline, reference-counted interior mutability, and runtime bounds checking on memory-mapped files.

Compile-Time Safety Through Rust's Ownership Model

Strict Ownership and Borrowing in Core Data Structures

All core data structures in RocketMQ-Rust—such as messages, buffers, and configuration objects—adhere to Rust's single-ownership principle. In rocketmq-common/src/lib.rs, the compiler guarantees that no two mutable references coexist, preventing data races at compile time. When data must be shared across threads, the codebase uses thread-safe containers like Arc and RwLock rather than raw pointers.

Safe Abstractions for Shared Mutable State

ArcMut: Reference-Counted Interior Mutability

The ArcMut<T> type provides a safe wrapper around SyncUnsafeCell, enabling shared mutable access without exposing unsafe operations to callers. Internally, ArcMut isolates all unsafe pointer dereferencing inside methods that perform explicit bounds checks and atomic position updates. The public API—including new, mut_from_ref, as_ref, and as_mut—remains completely safe to call.

use rocketmq::ArcMut;

// Create an ArcMut holding a vector
let shared_vec = ArcMut::new(vec![1, 2, 3]);

// Clone to share across threads
let cloned = shared_vec.clone();

// Mutate safely from one thread
{
    let mut_ref = cloned.mut_from_ref(); // safe because we own the only mutable reference here
    mut_ref.push(4);
}

// Read from another thread
let read_ref: &Vec<i32> = shared_vec.as_ref();
assert_eq!(read_ref, &vec![1, 2, 3, 4]);

Source: [rocketmq/src/arc_mut.rs](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq/src/arc_mut.rs)

Memory-Mapped File Safety

Runtime Bounds Checking in MappedBuffer

The MappedBuffer and DefaultMappedFile types expose read/write operations that first verify the requested offset and length against the file size. Any attempt to access memory out of bounds returns a typed error (MappedFileError::OutOfBounds), ensuring that the underlying MmapMut slice is never accessed unsafely.

use rocketmq_store::log_file::mapped_file::MappedBuffer;
use std::sync::Arc;
use parking_lot::RwLock;
use memmap2::MmapMut;
use tempfile::tempfile;

// Create a temporary file and map it
let file = tempfile().unwrap();
file.set_len(4096).unwrap();
let mmap = unsafe { MmapMut::map_mut(&file).unwrap() };
let mmap_arc = Arc::new(RwLock::new(mmap));

// Build a buffer covering the whole file
let buf = MappedBuffer::new(mmap_arc.clone(), 0, 4096).unwrap();

// Write data (checked at runtime)
buf.write(0, b"RocketMQ").unwrap();

// Zero-copy read – returns a `Bytes` that shares the mmap
let bytes = buf.read_zero_copy(0..9).unwrap();
assert_eq!(&bytes[..], b"RocketMQ");

Source: [rocketmq-store/src/log_file/mapped_file/mapped_buffer.rs](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/log_file/mapped_file/mapped_buffer.rs)

Atomic Position Tracking in DefaultMappedFile

Write, commit, and flush positions are stored in AtomicI32 and AtomicU64 types. All updates use fetch_add and store operations with appropriate memory orderings (Acquire, Release, AcqRel). This removes the need for external locks when tracking how much of the file has been written or flushed, while remaining free of data races.

use rocketmq_store::log_file::mapped_file::DefaultMappedFile;
use rocketmq_common::common::message::message_ext_broker_inner::MessageExtBrokerInner;
use rocketmq_store::base::append_message_callback::AppendMessageCallback;

// Create a mapped file with a numeric filename
let file = DefaultMappedFile::new(
    CheetahString::from("/tmp/00000000000000000000"),
    8 * 1024 * 1024, // 8 MiB
);

// Prepare a message
let mut msg = MessageExtBrokerInner::default();
msg.body = b"Hello".to_vec();

// Append safely – positions updated atomically
let result = file.append_message(&mut msg, &callback, &PutMessageContext::default());
assert_eq!(result.status, AppendMessageStatus::Ok);

Source: [rocketmq-store/src/log_file/mapped_file/default_mapped_file_impl.rs](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/log_file/mapped_file/default_mapped_file_impl.rs)

Isolating Unsafe Operations

Platform-Specific FFI Wrappers

Functions that call OS APIs such as mlock, munlock, and VirtualLock are isolated in the ffi.rs module. These platform-specific unsafe calls are wrapped in safe Rust methods that return Result types, ensuring that unsafe code never propagates to callers.

Source: [rocketmq-store/src/utils/ffi.rs](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/utils/ffi.rs)

Verification Through Testing

The project maintains comprehensive unit and integration tests that exercise every public function, particularly edge cases like out-of-bounds writes, zero-copy reads, and concurrent access. The test suite runs in continuous integration, guaranteeing that any change breaking safety contracts is caught early.

Source: [rocketmq-store/tests/commitlog_load_tests.rs](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/tests/commitlog_load_tests.rs)

Summary

RocketMQ-Rust achieves memory safety through a defense-in-depth strategy that combines compile-time guarantees with runtime verification:

  • Compile-time guarantees: Rust's ownership and borrowing rules prevent data races and dangling pointers before the code runs.
  • Safe abstractions: Types like ArcMut encapsulate interior mutability, exposing only safe APIs while isolating unsafe blocks.
  • Runtime bounds checking: MappedBuffer validates all offsets and lengths before accessing memory-mapped regions, returning typed errors for violations.
  • Atomic synchronization: DefaultMappedFile uses AtomicI32 and AtomicU64 with explicit memory orderings to track file positions without locks.
  • FFI isolation: Platform-specific unsafe calls are confined to ffi.rs and wrapped in safe functions returning Result.

Frequently Asked Questions

Does RocketMQ-Rust use any unsafe code?

Yes, but it is minimal and strictly contained. Unsafe code appears only in low-level abstractions like ArcMut (which wraps SyncUnsafeCell) and memory-mapped file operations. Every unsafe block is wrapped by a safe API that enforces invariants through bounds checking and reference counting, ensuring callers never interact with raw pointers directly.

How does RocketMQ-Rust prevent data races during concurrent message processing?

The codebase prevents data races through Rust's compile-time ownership rules combined with atomic operations. For shared state, ArcMut ensures that mutable references cannot coexist, while DefaultMappedFile uses AtomicI32 and AtomicU64 with Acquire/Release memory orderings to update write positions. Additionally, RwLock protects memory-mapped regions, allowing multiple readers or a single writer without undefined behavior.

What prevents buffer overflows when reading from memory-mapped files?

All access to memory-mapped files goes through MappedBuffer, which performs explicit runtime bounds checking before every read or write operation. Methods like write and read_zero_copy verify that the requested offset plus length does not exceed the file size, returning MappedFileError::OutOfBounds if the check fails. This ensures the underlying MmapMut slice is never accessed with invalid indices.

How does ArcMut differ from standard Rust smart pointers?

ArcMut is a specialized abstraction that provides reference-counted interior mutability specifically designed for RocketMQ-Rust's concurrent needs. Unlike Arc<RwLock<T>> which uses locking, ArcMut wraps SyncUnsafeCell to allow mutable access through reference counting, but isolates all unsafe pointer dereferencing inside methods that enforce safety invariants. The public API (new, mut_from_ref, as_ref, as_mut) is entirely safe, unlike raw UnsafeCell usage which requires unsafe blocks at every call site.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →