RustFS Unified Storage — Low-Level Design (LLD)
RustFS Unified Storage — Low-Level Design (LLD)
Document: RUSTFS-LLD-001
Status: DRAFT
Version: 0.1.0
Companion: RUSTFS-FEAT-001 (Feature Spec), RUSTFS-LLD analysis docs
Audience: Systems engineers (implementation-ready)
Scope: I/O stack, kernel bypass, metadata indexing/cache, transport,
error handling, rebalancing, protocol compatibility
0. Reading Guide
This LLD assumes the data structures and component topology from the Feature Spec (MDS / DSS / Client tiers, inode model, chunk maps, leases). It drills into the mechanisms that determine whether the system actually hits its latency and bandwidth targets. Each section ends with a decision (the chosen approach) and fallbacks (what to use if the primary approach is unavailable on a given deployment).
1. The I/O Stack (Bottom to Top)
1.1 Layered I/O Model
┌────────────────────────────────────────────────────────────────┐
│ APPLICATION (PyTorch, QEMU, S3 SDK) │
└────────────────────────────────────────────────────────────────┘
│ syscalls (read/write/fseek) │ SCSI/NVMe │ HTTP
▼ ▼ ▼
┌────────────────────────────────────────────────────────────────┐
│ ACCESS LAYER │
│ ┌─────────────┐ ┌──────────────┐ ┌─────────────────┐ │
│ │ FUSE / CUSE │ │ iSCSI/NVMe-oF│ │ HTTP/S3 handler │ │
│ │ + io_uring │ │ target │ │ (hyper) │ │
│ │ passthrough │ │ │ │ │ │
│ └─────────────┘ └──────────────┘ └─────────────────┘ │
└────────────────────────────────────────────────────────────────┘
│
▼
┌────────────────────────────────────────────────────────────────┐
│ CLIENT CORE (rustfs-client lib) │
│ metadata cache · chunk cache · lease mgr · read-ahead │
│ request router · retry/backoff · circuit breaker │
└────────────────────────────────────────────────────────────────┘
│ control plane (MDS) │ data plane (DSS)
▼ ▼
┌─────────────────────┐ ┌─────────────────────────────────────┐
│ TRANSPORT: gRPC │ │ TRANSPORT: tiered │
│ (tonic over HTTP/2) │ │ 1. RDMA verbs (preferred) │
│ for metadata ops │ │ 2. io_uring + TCP (zero-copy) │
│ │ │ 3. plain gRPC (fallback) │
└─────────────────────┘ └─────────────────────────────────────┘
│ │
▼ ▼
┌─────────────────────┐ ┌─────────────────────────────────────┐
│ MDS │ │ DSS chunk service │
│ RocksDB engine │ │ io_uring disk I/O │
│ Raft replication │ │ buffer pool (reuse io-core) │
└─────────────────────┘ └─────────────────────────────────────┘
│ │
▼ ▼
NVMe SSD (metadata) NVMe SSD (chunk data)
+ optional SPDK userspace driver
1.2 Why Three Data-Plane Transports
Transport Latency Bandwidth CPU cost Requirement
──────────────── ──────── ────────── ───────── ─────────────────────
RDMA verbs ~2 µs 200 Gbps very low RDMA NIC (CX-5/6/7,
(RoCEv2/IB) BlueField DPU)
io_uring + TCP ~15 µs 100 Gbps low Linux 5.10+, any NIC
(zero-copy send) (best general default)
gRPC/tonic ~50 µs 40-80 Gbps moderate universal fallback
Decision: runtime transport negotiation. Client and DSS exchange capabilities at session open; the best mutually-supported transport is selected per-connection. Metadata always uses gRPC (small messages, latency dominated by RocksDB/Raft, not transport).
2. io_uring Integration
2.1 Where io_uring Is Used
Location Purpose Crate
──────────────────── ───────────────────────────────── ──────────────────
DSS disk I/O chunk read/write to NVMe chunk-service
DSS network I/O zero-copy TCP send of chunk data chunk-service
FUSE client /dev/fuse processing via io_uring fuse-client
(FUSE-over-io_uring, Linux 6.14+)
Client network I/O batched chunk fetch over TCP rustfs-client
WAL (block volumes) ordered durable appends iscsi-target
2.2 DSS Disk I/O Engine
// crate: rustfs-chunk-service
// file: crates/chunk-service/src/iouring_engine.rs
/// io_uring-based disk I/O engine for chunk storage.
///
/// Rust ecosystem choice: `tokio-uring` (preferred for tokio integration)
/// or `io-uring` (raw, lower-level, used when we need full control of SQE flags).
///
/// We use a hybrid: tokio-uring for the async runtime integration on the
/// network side, and a dedicated raw io-uring ring per disk for the storage
/// side (so disk I/O doesn't compete with network I/O in the same ring).
pub struct DiskIoEngine {
/// One ring per physical disk — avoids head-of-line blocking across disks
/// and lets each disk's queue depth be tuned independently.
rings: Vec<DiskRing>,
}
struct DiskRing {
disk_id: u32,
ring: io_uring::IoUring,
/// Registered fixed buffers for zero-copy (avoids per-op buffer pinning).
/// Reuses the existing rustfs-io-core buffer pool.
buffers: FixedBufferPool,
/// Registered file descriptors for the chunk files on this disk.
fixed_fds: FdTable,
queue_depth: u32, // typically 128-256 per NVMe
}
impl DiskIoEngine {
/// Read a chunk range using io_uring.
///
/// Uses IORING_OP_READ_FIXED with a registered buffer for zero-copy.
/// For O_DIRECT, the buffer must be aligned to the device block size.
pub async fn read_chunk(
&self,
disk_id: u32,
fd: RawFd,
offset: u64,
length: u64,
direct_io: bool,
) -> Result<FixedBuf> {
let ring = &self.rings[disk_id as usize];
let buf = ring.buffers.acquire(length as usize).await;
let sqe = if direct_io {
// O_DIRECT path: aligned, bypasses page cache
opcode::ReadFixed::new(Fixed(fd), buf.as_mut_ptr(), length as u32, buf.index())
.offset(offset)
.build()
} else {
// Buffered path: uses page cache (good for hot data re-reads)
opcode::Read::new(Fd(fd), buf.as_mut_ptr(), length as u32)
.offset(offset)
.build()
};
// Submit and await completion
let cqe = ring.submit_and_wait(sqe).await?;
if cqe.result() < 0 {
return Err(Error::from_errno(-cqe.result()));
}
buf.set_len(cqe.result() as usize);
Ok(buf)
}
/// Write with optional O_DSYNC for durability.
///
/// For block volumes (sync=true), chains IORING_OP_WRITE_FIXED
/// followed by IORING_OP_FSYNC using IOSQE_IO_LINK so they execute
/// in order without a separate submission round trip.
pub async fn write_chunk(
&self,
disk_id: u32,
fd: RawFd,
offset: u64,
data: &[u8],
sync: bool,
) -> Result<usize> {
let ring = &self.rings[disk_id as usize];
let buf = ring.buffers.acquire_with_data(data).await;
if sync {
// Linked write + fdatasync (single submission, ordered execution)
let write_sqe = opcode::WriteFixed::new(
Fixed(fd), buf.as_ptr(), data.len() as u32, buf.index())
.offset(offset)
.flags(squeue::Flags::IO_LINK) // link to next op
.build();
let fsync_sqe = opcode::Fsync::new(Fd(fd))
.flags(types::FsyncFlags::DATASYNC)
.build();
ring.submit_linked(&[write_sqe, fsync_sqe]).await?;
} else {
let write_sqe = opcode::WriteFixed::new(
Fixed(fd), buf.as_ptr(), data.len() as u32, buf.index())
.offset(offset)
.build();
ring.submit_and_wait(write_sqe).await?;
}
Ok(data.len())
}
}
Decision: raw io-uring crate with one ring per disk, registered fixed
buffers, registered FDs. Queue depth 128–256 per NVMe. Linked write+fsync for
synchronous block writes.
Fallback: when io_uring is unavailable (kernel < 5.10, or restricted
container), fall back to tokio::fs with a thread pool. The DiskIoEngine
trait abstracts this so callers don't change.
2.3 FUSE over io_uring
// crate: rustfs-fuse
// file: crates/fuse-client/src/uring_fuse.rs
/// FUSE request processing over io_uring (Linux 6.14+ FUSE-over-io_uring).
///
/// Classic FUSE: each request is a read() on /dev/fuse, reply is a write().
/// Two syscalls + two context switches per op → high latency floor.
///
/// FUSE-over-io_uring: requests and replies flow through io_uring SQ/CQ rings
/// shared with the kernel. Eliminates per-request syscalls; enables batching.
///
/// Measured improvement (upstream): up to 2-3x throughput on metadata-heavy
/// workloads, ~30-50% latency reduction on small reads.
pub struct UringFuseSession {
fuse_fd: RawFd,
ring: io_uring::IoUring,
/// Per-CPU request rings for scalability across cores.
nr_queues: u32,
}
impl UringFuseSession {
/// Register FUSE rings with the kernel via FUSE_DEV_IOC_CLONE +
/// the io_uring registration ioctl.
pub fn setup(mountpoint: &Path, nr_queues: u32) -> Result<Self> {
// 1. mount with FUSE, get /dev/fuse fd
// 2. negotiate FUSE_OVER_IO_URING capability in FUSE_INIT
// 3. register io_uring rings for FUSE request/response
todo!()
}
}
Decision: detect FUSE-over-io_uring at mount time (FUSE_INIT capability
negotiation). Use it when available (kernel ≥ 6.14). The fuser crate base is
kept for the operation handlers; only the transport between kernel and userspace
changes.
Fallback: classic /dev/fuse read/write loop with a multi-threaded worker
pool (the default fuser behavior). Still acceptable; loses ~30% on small ops.
2.4 Network Zero-Copy Send (DSS → Client)
// crate: rustfs-chunk-service
// file: crates/chunk-service/src/uring_net.rs
/// Zero-copy chunk transfer over TCP using io_uring.
///
/// For a chunk read served over TCP:
/// 1. Chunk data is read from disk into a registered buffer (READ_FIXED)
/// 2. The SAME buffer is sent on the socket (SEND_ZC — zero-copy send)
/// 3. No memcpy between disk buffer and network buffer
///
/// IORING_OP_SEND_ZC (Linux 6.0+) avoids the kernel copying userspace data
/// into socket buffers; the NIC DMAs directly from the registered buffer.
pub async fn send_chunk_zerocopy(
ring: &IoUring,
socket_fd: RawFd,
buf: &FixedBuf,
) -> Result<usize> {
let sqe = opcode::SendZc::new(Fd(socket_fd), buf.as_ptr(), buf.len() as u32)
.build();
let cqe = ring.submit_and_wait(sqe).await?;
// SEND_ZC posts two CQEs: one for submission, one for buffer release.
// Must wait for the IORING_CQE_F_NOTIF flag before reusing the buffer.
Ok(cqe.result() as usize)
}
Decision: IORING_OP_SEND_ZC for chunk payloads ≥ 16 KiB (zero-copy has
fixed setup cost; small payloads use regular send). Combined with READ_FIXED
the chunk never gets copied in DSS userspace.
3. Kernel Bypass (RDMA & SPDK)
3.1 RDMA Data Plane
CLIENT DSS
┌──────────────────────┐ ┌──────────────────────┐
│ rustfs-client │ │ chunk-service │
│ │ │ │
│ registered MR ◄─────┼──────┼─► registered MR │
│ (chunk cache) │ │ (buffer pool) │
│ │ │ │ │ │
│ ▼ │ │ ▼ │
│ RDMA verbs (libibverbs via rust bindings) │
│ │ │ │ │ │
└───────┼──────────────┘ └────────┼──────────────┘
│ RoCEv2 / InfiniBand │
└───────────────────────────────┘
RDMA NIC ↔ RDMA NIC
3.1.1 RDMA Read Flow (Client Pulls Chunk)
// crate: rustfs-rdma
// file: crates/rdma/src/client.rs
/// RDMA-based chunk read. The client issues an RDMA READ to pull chunk data
/// directly from the DSS's memory into its own registered buffer.
///
/// Control flow (small, over RC QP send/recv):
/// 1. Client → DSS: "I want chunk X, range [off,len]" (send)
/// 2. DSS: reads chunk from disk (io_uring) into a registered MR
/// 3. DSS → Client: "data is at addr A, rkey R, len L" (send)
/// 4. Client: RDMA READ from (A, R) into local buffer ← no DSS CPU involved
/// 5. Client → DSS: "done, release MR" (send)
///
/// For hot data already in DSS memory, steps 2 is skipped (cache hit).
///
/// Rust binding choice: `async-rdma` crate, or direct FFI to libibverbs
/// via `rdma-sys`. We wrap in our own safe abstraction.
pub struct RdmaChunkClient {
pd: ProtectionDomain,
cq: CompletionQueue,
qps: HashMap<NodeId, QueuePair>, // one RC QP per DSS node
mr_pool: MemoryRegionPool, // pre-registered chunk buffers
}
impl RdmaChunkClient {
pub async fn read_chunk_rdma(
&self,
node: NodeId,
chunk_id: &ChunkId,
offset: u64,
length: u64,
) -> Result<RdmaBuf> {
let qp = self.qps.get(&node).ok_or(Error::NoConnection)?;
// 1. Send read request (control message)
let req = ChunkReadControl { chunk_id, offset, length };
qp.post_send(&serialize(&req)?).await?;
// 2. Receive remote memory descriptor
let desc: RemoteMemoryDesc = qp.post_recv().await?;
// 3. RDMA READ: pull data directly from DSS memory
let local_buf = self.mr_pool.acquire(length as usize);
qp.post_read(
local_buf.as_mut_ptr(),
local_buf.lkey(),
desc.remote_addr,
desc.rkey,
length as u32,
).await?;
// 4. Acknowledge (lets DSS release its MR)
qp.post_send(&serialize(&ChunkReadAck { token: desc.token })?).await?;
Ok(local_buf)
}
}
Decision: RDMA READ-based pull model for reads (DSS CPU not involved in the
data copy), RDMA WRITE-based push for writes. One Reliable Connection (RC) Queue
Pair per client↔DSS pair. Pre-registered memory regions on both sides drawn from
the existing io-core buffer pool, registered once at startup.
Fallback: if no RDMA NIC, transport negotiation drops to io_uring+TCP. The
ChunkTransport trait makes this transparent to the chunk service logic.
3.1.2 GPUDirect Storage (Phase 5)
// crate: rustfs-rdma
// file: crates/rdma/src/gds.rs
/// GPUDirect Storage: read chunk data directly into GPU HBM, bypassing
/// the CPU and host RAM entirely.
///
/// Path: NVMe/network → GPU memory (DMA), no bounce buffer through CPU.
///
/// Uses NVIDIA cuFile API (libcufile) + DMA-BUF. The GPU buffer is registered
/// as an RDMA MR (via nvidia-peermem or DMA-BUF), then RDMA READ targets it.
///
/// For AI training, this means the DataLoader can read a batch straight into
/// the GPU tensor's backing memory — no CPU copy, no host RAM staging.
pub async fn read_chunk_to_gpu(
&self,
node: NodeId,
chunk_id: &ChunkId,
gpu_ptr: CuDevicePtr, // registered GPU memory
gpu_mr: &MemoryRegion,
offset: u64,
length: u64,
) -> Result<()> {
let qp = self.qps.get(&node)?;
let req = ChunkReadControl { chunk_id, offset, length };
qp.post_send(&serialize(&req)?).await?;
let desc: RemoteMemoryDesc = qp.post_recv().await?;
// RDMA READ directly into GPU memory region
qp.post_read(gpu_ptr.as_ptr(), gpu_mr.lkey(),
desc.remote_addr, desc.rkey, length as u32).await?;
Ok(())
}
3.2 SPDK Userspace NVMe (Optional DSS Backend)
For DSS nodes that want maximum disk throughput and lowest latency,
SPDK provides a userspace NVMe driver — the kernel block layer is bypassed
entirely.
Tradeoff:
+ Polled I/O, no interrupts, no context switches → ~10 µs → ~2 µs
+ Dedicated CPU cores poll NVMe submission/completion queues
- Disk is exclusively owned by the SPDK process (no other access)
- Higher idle CPU usage (polling burns a core)
Decision: SPDK backend is OPT-IN per DSS node (config flag).
Default backend = io_uring (good enough for most, no core dedication).
SPDK backend for latency-critical block storage nodes.
Rust integration: FFI to SPDK C library (spdk-rs bindings or custom FFI).
SPDK runs its own reactor; we bridge to tokio via a channel-based shim.
// crate: rustfs-chunk-service
// file: crates/chunk-service/src/spdk_backend.rs
/// SPDK userspace NVMe backend (opt-in).
///
/// SPDK owns dedicated poller cores. Chunk I/O requests are passed from
/// tokio tasks to SPDK reactor cores via lock-free SPSC rings.
pub struct SpdkDiskEngine {
/// Handle to SPDK bdev (block device abstraction)
bdev: SpdkBdev,
/// Channels to submit I/O to SPDK reactor cores
submit_rings: Vec<spsc::Producer<SpdkIoRequest>>,
}
// The DiskIoEngine trait is implemented for both IoUringDiskEngine
// and SpdkDiskEngine — selected at DSS startup based on config.
4. Metadata Indexing, Query, and Cache
4.1 MDS Storage Engine Layout
RocksDB tuning for metadata workload (many small KV, point lookups + range scans):
Column Family Key Pattern Read pattern Tuning
─────────────── ─────────────────────── ─────────────── ──────────────────
inodes BE(ino) point lookup bloom filter,
block cache 40%
dentries BE(parent) ++ name point + prefix prefix bloom,
scan (readdir) block cache 30%
chunkmaps BE(ino) point lookup block cache 20%,
(large values) larger block size
volumes vol_id point lookup small CF
leases BE(ino) ++ client point + scan in-memory only
(no persistence)
locks BE(ino) ++ range range scan in-memory only
Global RocksDB options:
- WAL on separate NVMe (or NVDIMM) for low-latency commits
- Compaction: leveled, with separate threads
- Memtable: 256 MB, 4 memtables
- Block cache: shared, sized to ~50% of MDS node RAM
- Direct I/O for compaction (avoid double caching)
4.2 Inode Number Allocation
// crate: rustfs-mds
// file: crates/mds/src/alloc.rs
/// Inode allocator. Hands out InodeIds in batches to avoid hot RocksDB key.
///
/// Each MDS partition owns a 48-bit local space. The allocator grabs ranges
/// of 10,000 inodes at a time from the persistent counter, then hands them
/// out from memory. This makes inode creation a memory operation, not a
/// RocksDB write, except once per 10,000 creates.
pub struct InodeAllocator {
partition: u16,
next: AtomicU64, // next inode to hand out (in-memory)
range_end: AtomicU64, // end of current reserved range
persist_lock: Mutex<()>, // serialize range reservation
store: Arc<dyn MetadataStore>,
}
impl InodeAllocator {
pub async fn alloc(&self) -> Result<InodeId> {
loop {
let n = self.next.fetch_add(1, Ordering::Relaxed);
if n < self.range_end.load(Ordering::Acquire) {
return Ok(self.make_inode_id(n));
}
// Reserve a new range (rare, ~once per 10K allocs)
self.reserve_range().await?;
}
}
fn make_inode_id(&self, local: u64) -> InodeId {
((self.partition as u64) << 48) | (local & 0xFFFF_FFFF_FFFF)
}
}
4.3 Path Resolution Cache (MDS-Side)
// crate: rustfs-mds
// file: crates/mds/src/path_cache.rs
/// Path resolution cache: full path → inode, avoiding repeated dentry walks.
///
/// "/buckets/training-data/imagenet/shard-0001.tar" requires 4 dentry lookups
/// to resolve from root. This cache short-circuits hot paths.
///
/// Invalidation: any rename/unlink/rmdir on a path component invalidates all
/// cached descendants. Implemented with a generation counter per directory.
pub struct PathCache {
/// Full path → (inode, dir_generation_snapshot)
cache: moka::sync::Cache<String, (InodeId, u64)>,
/// Per-directory generation counters; bumped on any mutation
generations: DashMap<InodeId, u64>,
}
impl PathCache {
pub fn lookup(&self, path: &str) -> Option<InodeId> {
let (ino, gen_snapshot) = self.cache.get(path)?;
// Validate: has the parent directory changed since we cached?
let parent = self.parent_inode_of(ino)?;
let current_gen = self.generations.get(&parent).map(|g| *g).unwrap_or(0);
if current_gen == gen_snapshot {
Some(ino)
} else {
self.cache.invalidate(path);
None
}
}
/// Called on any namespace mutation under `dir`.
pub fn invalidate_dir(&self, dir: InodeId) {
self.generations.entry(dir).and_modify(|g| *g += 1).or_insert(1);
}
}
Decision: moka (concurrent, TinyLFU eviction) for the path cache and
inode/dentry caches. Generation-counter validation avoids cache-wide flushes on
single mutations.
4.4 Client-Side Metadata Cache + Lease Coherency
// crate: rustfs-fuse
// file: crates/fuse-client/src/meta_cache.rs
/// Client metadata cache, coherent via leases.
///
/// Without leases: caches are TTL-based (stale window = TTL).
/// With leases: cache is valid until the MDS revokes the lease, so the
/// attribute timeout can be effectively infinite (strong coherency).
///
/// This is what lets `stat()` on a hot file be a 0.01ms memory op instead of
/// a 1ms MDS round trip — and still be correct.
pub struct ClientMetaCache {
inodes: moka::sync::Cache<InodeId, CachedInode>,
dentries: moka::sync::Cache<(InodeId, String), DirEntry>,
/// Inodes for which we hold a ReadMeta lease → cache never goes stale
leased: DashMap<InodeId, Lease>,
}
struct CachedInode {
attr: InodeAttr,
valid_until: Instant, // ignored if a lease is held
has_lease: bool,
}
impl ClientMetaCache {
pub fn get_attr(&self, ino: InodeId) -> Option<InodeAttr> {
let cached = self.inodes.get(&ino)?;
if cached.has_lease || cached.valid_until > Instant::now() {
Some(cached.attr)
} else {
None
}
}
/// Called by the lease revocation handler (MDS → client callback).
pub fn on_lease_revoked(&self, ino: InodeId, kind: LeaseKind) {
match kind {
LeaseKind::ReadMeta => { self.inodes.invalidate(&ino); }
LeaseKind::ReadData | LeaseKind::Layout => { /* handled by data cache */ }
LeaseKind::WriteData => { /* flush dirty + invalidate */ }
}
self.leased.remove(&ino);
}
}
4.5 Chunk Map Query Optimization
Problem: a 1 TB file with 4 MiB chunks has 262,144 chunk entries.
Fetching the whole chunk map for a small read is wasteful.
Solution: ranged chunk map queries + client-side partial caching.
GetChunkMap(ino, offset, length) returns ONLY the chunks covering that range.
Client caches the chunks it has seen in a per-inode BTreeMap.
A Layout lease guarantees the cached chunk locations stay valid.
For sequential reads, read-ahead pre-fetches the NEXT range's chunk map
along with the data, so the chunk map lookup is never on the critical path
after the first read.
// Client-side partial chunk map with lease-backed validity
pub struct PartialChunkMap {
ino: InodeId,
chunk_size: u64,
known: BTreeMap<u64, Vec<ChunkLocation>>, // chunk_index → locations
layout_lease: Option<Lease>,
fully_cached: bool,
}
impl PartialChunkMap {
/// Returns chunks if fully present locally, else the missing ranges to fetch.
pub fn resolve_or_miss(&self, offset: u64, length: u64)
-> ResolveResult
{
let start = offset / self.chunk_size;
let end = (offset + length - 1) / self.chunk_size;
let mut missing = Vec::new();
for idx in start..=end {
if !self.known.contains_key(&idx) {
missing.push(idx);
}
}
if missing.is_empty() {
ResolveResult::Hit(self.collect_range(start, end))
} else {
ResolveResult::Miss(coalesce_ranges(missing))
}
}
}
5. Client ↔ Chunk Server Communication
5.1 Connection Model
Control plane (client ↔ MDS):
- Long-lived gRPC/HTTP2 connection, multiplexed streams
- One connection per client to each MDS (pooled)
- Bidirectional stream for lease revocation callbacks
Data plane (client ↔ DSS):
- Connection pool per DSS node
- Transport: RDMA QP / io_uring TCP / gRPC (negotiated)
- For RDMA: one RC QP per (client, DSS) pair, kept warm
- For TCP: N connections per DSS (N = nr CPU cores), load-balanced
// crate: rustfs-client
// file: crates/client/src/dss_pool.rs
/// Connection pool to all DSS nodes. Handles transport selection,
/// health tracking, and load balancing across connections.
pub struct DssConnectionPool {
nodes: DashMap<NodeId, DssNode>,
transport: TransportKind, // negotiated at startup
}
struct DssNode {
node_id: NodeId,
endpoint: SocketAddr,
conns: Vec<ChunkTransport>, // multiple conns for parallelism
health: NodeHealth,
latency_ewma: AtomicU64, // exponentially-weighted moving avg (ns)
inflight: AtomicU32, // current outstanding requests
}
/// Unified transport abstraction — all three transports implement this.
#[async_trait]
pub trait ChunkTransport: Send + Sync {
async fn read_chunk(&self, req: ReadChunkReq) -> Result<Bytes>;
async fn write_chunk(&self, req: WriteChunkReq) -> Result<WriteChunkResp>;
async fn read_batch(&self, reqs: Vec<ReadChunkReq>) -> Result<Vec<Bytes>>;
}
5.2 Request Batching and Coalescing
// crate: rustfs-client
// file: crates/client/src/batcher.rs
/// Batches multiple chunk reads destined for the same DSS node into a single
/// ReadChunkBatch RPC, reducing round trips.
///
/// A 64 MiB read = 16 chunks. If 4 chunks live on DSS-3, they go in one batch
/// RPC instead of 4 separate RPCs.
///
/// Batching window: 50 µs (configurable). Requests arriving within the window
/// to the same DSS are coalesced. Trades a tiny latency increase for far fewer
/// round trips under high concurrency.
pub struct ChunkReadBatcher {
pending: DashMap<NodeId, Vec<PendingRead>>,
window: Duration,
}
struct PendingRead {
req: ReadChunkReq,
respond: oneshot::Sender<Result<Bytes>>,
}
impl ChunkReadBatcher {
pub async fn submit(&self, node: NodeId, req: ReadChunkReq)
-> oneshot::Receiver<Result<Bytes>>
{
let (tx, rx) = oneshot::channel();
self.pending.entry(node).or_default().push(PendingRead { req, respond: tx });
// A background flusher fires every `window` or when batch hits max size
rx
}
}
5.3 Load Balancing Across Replicas/Shards
For a chunk with R replicas (or for EC shards), pick where to read from:
Priority order:
1. Locality: replica on the client's own node (DSS co-located with compute)
2. Health: skip nodes marked degraded/healing
3. Latency: lowest latency_ewma
4. Load: lowest inflight count (avoid hotspotting one node)
Implementation: power-of-two-choices — sample 2 candidate replicas,
pick the one with lower (latency_ewma × (1 + inflight)). Avoids the
herding problem of always picking the single "best" node.
fn pick_replica(&self, locs: &[ChunkLocation]) -> &ChunkLocation {
// Filter healthy
let healthy: Vec<_> = locs.iter()
.filter(|l| self.node_health(l.node_id).is_serving())
.collect();
// Locality shortcut
if let Some(local) = healthy.iter().find(|l| l.node_id == self.local_node) {
return local;
}
// Power of two choices
let a = healthy.choose(&mut rng);
let b = healthy.choose(&mut rng);
[a, b].into_iter().flatten()
.min_by_key(|l| self.score(l.node_id))
.unwrap()
}
fn score(&self, node: NodeId) -> u64 {
let lat = self.node(node).latency_ewma.load(Relaxed);
let inflight = self.node(node).inflight.load(Relaxed) as u64;
lat * (1 + inflight)
}
6. Error Handling and Processing
6.1 Error Taxonomy
// crate: rustfs-common
// file: crates/common/src/errors.rs
#[derive(Debug, thiserror::Error)]
pub enum StorageError {
// --- Transient (retry) ---
#[error("DSS node {0} temporarily unavailable")]
NodeUnavailable(NodeId),
#[error("request timed out after {0:?}")]
Timeout(Duration),
#[error("insufficient quorum: got {got}, need {need}")]
InsufficientQuorum { got: u32, need: u32 },
#[error("lease revoked during operation")]
LeaseRevoked,
// --- Data integrity (heal) ---
#[error("bitrot detected on chunk {0}")]
BitrotDetected(ChunkId),
#[error("chunk {0} corrupt, reconstructing")]
ChunkCorrupt(ChunkId),
// --- Permanent (fail fast) ---
#[error("inode {0} not found")]
InodeNotFound(InodeId),
#[error("path not found: {0}")]
PathNotFound(String),
#[error("permission denied")]
PermissionDenied,
#[error("no space left on device")]
NoSpace,
// --- Consistency (escalate) ---
#[error("MDS split-brain suspected")]
SplitBrain,
#[error("chunk map / data mismatch for inode {0}")]
MetadataDataMismatch(InodeId),
}
impl StorageError {
pub fn class(&self) -> ErrorClass {
match self {
Self::NodeUnavailable(_) | Self::Timeout(_)
| Self::InsufficientQuorum {..} | Self::LeaseRevoked
=> ErrorClass::Transient,
Self::BitrotDetected(_) | Self::ChunkCorrupt(_)
=> ErrorClass::Heal,
Self::InodeNotFound(_) | Self::PathNotFound(_)
| Self::PermissionDenied | Self::NoSpace
=> ErrorClass::Permanent,
Self::SplitBrain | Self::MetadataDataMismatch(_)
=> ErrorClass::Critical,
}
}
}
pub enum ErrorClass { Transient, Heal, Permanent, Critical }
6.2 Retry, Backoff, and Hedging
// crate: rustfs-client
// file: crates/client/src/retry.rs
/// Retry policy with jittered exponential backoff and request hedging.
///
/// Hedging: if a read hasn't responded within p99 latency, fire a second
/// request to a DIFFERENT replica. Whichever returns first wins; the other
/// is cancelled. This caps tail latency caused by a single slow/busy node.
pub struct RetryPolicy {
max_attempts: u32, // default 3
base_backoff: Duration, // default 10ms
max_backoff: Duration, // default 1s
hedge_after: Duration, // default = node p99 latency
}
impl RetryPolicy {
pub async fn execute_read<F, Fut>(&self, locs: &[ChunkLocation], op: F)
-> Result<Bytes>
where F: Fn(&ChunkLocation) -> Fut, Fut: Future<Output = Result<Bytes>>
{
let mut attempt = 0;
loop {
let primary = pick_replica(locs);
let hedge_timer = tokio::time::sleep(self.hedge_after);
tokio::select! {
result = op(primary) => {
match result {
Ok(data) => return Ok(data),
Err(e) if e.class() == ErrorClass::Transient => {
attempt += 1;
if attempt >= self.max_attempts {
return Err(e);
}
self.backoff(attempt).await;
}
Err(e) => return Err(e), // permanent: fail fast
}
}
_ = hedge_timer => {
// Primary is slow — fire hedge to a different replica
let secondary = pick_different_replica(locs, primary);
return tokio::select! {
r = op(primary) => r,
r = op(secondary) => r,
};
}
}
}
}
fn backoff(&self, attempt: u32) -> impl Future<Output=()> {
let base = self.base_backoff * 2u32.pow(attempt - 1);
let capped = base.min(self.max_backoff);
let jitter = rand::random::<f64>() * capped.as_secs_f64() * 0.5;
tokio::time::sleep(capped + Duration::from_secs_f64(jitter))
}
}
6.3 Circuit Breaker (Per DSS Node)
// crate: rustfs-client
// file: crates/client/src/circuit.rs
/// Per-node circuit breaker. Stops sending to a node that's failing,
/// gives it time to recover, avoids cascading timeouts.
///
/// States:
/// Closed: normal operation, count failures
/// Open: node is failing, reject immediately, route elsewhere
/// HalfOpen: probe with a few requests; if they succeed, close
pub struct CircuitBreaker {
state: AtomicCell<CircuitState>,
failure_count: AtomicU32,
failure_thresh: u32, // default 5 consecutive failures
open_duration: Duration, // default 5s before half-open probe
opened_at: AtomicCell<Option<Instant>>,
}
enum CircuitState { Closed, Open, HalfOpen }
impl CircuitBreaker {
pub fn allow_request(&self) -> bool {
match self.state.load() {
CircuitState::Closed => true,
CircuitState::Open => {
if self.opened_at.load().unwrap().elapsed() > self.open_duration {
self.state.store(CircuitState::HalfOpen);
true // allow a probe
} else {
false // reject, route to another replica
}
}
CircuitState::HalfOpen => true,
}
}
pub fn record_success(&self) {
self.failure_count.store(0, Relaxed);
self.state.store(CircuitState::Closed);
}
pub fn record_failure(&self) {
let n = self.failure_count.fetch_add(1, Relaxed) + 1;
if n >= self.failure_thresh {
self.state.store(CircuitState::Open);
self.opened_at.store(Some(Instant::now()));
}
}
}
6.4 Data Integrity Errors → Healing
On bitrot/corrupt detection during a read:
1. Client receives ChunkCorrupt error from DSS (DSS verified checksum on read)
2. For EC chunks: client reads other shards, reconstructs locally,
continues serving the application WITHOUT waiting for heal
3. For replicated chunks: client reads a different replica, continues
4. Client/DSS reports the corrupt chunk to the MDS heal queue (async)
5. HealManager (existing rustfs-heal crate) reconstructs the bad shard/replica
in the background and updates the chunk map
The application read NEVER fails due to a single corrupt chunk as long as
quorum/replicas remain. This reuses the existing rustfs heal + bitrot machinery.
async fn read_with_heal_fallback(&self, entry: &ChunkRangeEntry) -> Result<Bytes> {
match self.read_primary(entry).await {
Ok(data) => Ok(data),
Err(StorageError::ChunkCorrupt(id)) | Err(StorageError::BitrotDetected(id)) => {
// Report to heal queue (fire-and-forget)
self.heal_reporter.report_corrupt(&id);
// Reconstruct from remaining shards/replicas
match entry.locations.len() {
_ if self.is_ec(entry) => self.ec_reconstruct(entry).await,
_ => self.read_alternate_replica(entry).await,
}
}
Err(e) => Err(e),
}
}
6.5 MDS Failure Handling
MDS leader failure (Raft):
- Followers detect missing heartbeat → election → new leader (typ. < 2s)
- Clients retry metadata ops against new leader (discovered via Raft membership)
- In-flight leases survive (part of Raft state machine, replicated)
- Locks survive (also in Raft state machine)
MDS partition unreachable (sharded namespace):
- Only the subtree owned by that partition is affected
- Other namespace subtrees keep serving
- Clients cache "MDS partition down" and return EIO for affected paths
until partition recovers
7. Rebalancing
7.1 When Rebalancing Triggers
Trigger Action
─────────────────────────────────── ─────────────────────────────────────
New DSS node added migrate chunks to fill new capacity
DSS node decommissioned evacuate all chunks before removal
Disk usage skew > threshold (e.g. move chunks from hot to cold nodes
20% imbalance across nodes)
Hot chunk detected (read hotspot) add extra replica on a cooler node
Storage policy change re-encode affected chunks
7.2 Chunk Placement Algorithm
// crate: rustfs-mds
// file: crates/mds/src/placement.rs
/// Chunk placement: decides which DSS nodes/disks store a chunk's
/// replicas or EC shards.
///
/// Goals:
/// - Even capacity distribution (weighted by disk size)
/// - Failure-domain awareness (don't put all replicas in one rack)
/// - Stability (adding a node moves minimal data — uses rendezvous hashing)
///
/// Algorithm: Weighted Rendezvous Hashing (HRW) with failure domains.
/// For each chunk, compute weight(node) = -ln(hash01) / node_capacity_weight
/// Pick top-R nodes, but enforce failure-domain diversity.
///
/// Why HRW not consistent hashing rings: HRW gives better balance and the
/// minimal-disruption property without the complexity of virtual nodes.
pub struct ChunkPlacer {
nodes: Vec<NodeWeight>,
failure_domains: HashMap<NodeId, FailureDomain>, // rack/zone
}
struct NodeWeight {
node_id: NodeId,
capacity: u64,
used: u64,
weight: f64, // capacity-based, adjusted for current usage
}
impl ChunkPlacer {
pub fn place(&self, chunk_id: &ChunkId, replicas: u32) -> Vec<NodeId> {
let mut scored: Vec<(f64, NodeId, FailureDomain)> = self.nodes.iter()
.map(|n| {
let h = hash_to_unit_interval(chunk_id, n.node_id);
let score = -h.ln() * n.weight; // weighted rendezvous
(score, n.node_id, self.failure_domains[&n.node_id])
})
.collect();
scored.sort_by(|a, b| b.0.partial_cmp(&a.0).unwrap());
// Pick top-R with failure-domain diversity
let mut chosen = Vec::new();
let mut used_domains = HashSet::new();
for (_, node, domain) in scored {
if chosen.len() == replicas as usize { break; }
// Prefer spreading across domains; relax if not enough domains
if !used_domains.contains(&domain) || chosen.len() + remaining_domains < replicas {
chosen.push(node);
used_domains.insert(domain);
}
}
chosen
}
}
7.3 Online Rebalancing Without Disruption
// crate: rustfs-mds
// file: crates/mds/src/rebalance.rs
/// Background rebalancer. Moves chunks to restore balance without blocking
/// client I/O.
///
/// Migration is COPY-THEN-SWITCH:
/// 1. Copy chunk to the new target node (background, rate-limited)
/// 2. Verify checksum on the new copy
/// 3. Atomically update the chunk map in MDS (add new location)
/// 4. Remove the old location from the chunk map
/// 5. Delete the old chunk data
///
/// During migration, reads can use either location (chunk map has both).
/// Writes during migration go to both until step 3, then to the new one.
///
/// Rate limiting: rebalancing yields to client I/O. Uses a token-bucket
/// limiter on migration bandwidth (default: cap at 10% of node bandwidth,
/// burst higher when client load is low).
pub struct Rebalancer {
mds: Arc<dyn MetadataStore>,
placer: ChunkPlacer,
rate_limiter: TokenBucket,
in_progress: DashMap<ChunkId, MigrationState>,
}
impl Rebalancer {
pub async fn migrate_chunk(&self, chunk_id: &ChunkId, to: NodeId) -> Result<()> {
// 1. Copy (rate-limited)
self.rate_limiter.acquire(chunk_size).await;
let data = self.read_chunk(chunk_id).await?;
let new_loc = self.write_chunk_to(to, &data).await?;
// 2. Verify
if self.checksum(&new_loc).await? != self.expected_checksum(chunk_id) {
self.delete_chunk(&new_loc).await?;
return Err(StorageError::ChunkCorrupt(chunk_id.clone()));
}
// 3-4. Atomic chunk map swap (transactional in MDS)
self.mds.swap_chunk_location(chunk_id, /*old*/, new_loc).await?;
// 5. Delete old
self.delete_old_chunk(chunk_id).await?;
Ok(())
}
}
7.4 Node Decommission and Failure Recovery
Graceful decommission (admin-initiated):
1. Mark node as "draining" in MDS (no new chunks placed)
2. Rebalancer evacuates all chunks to other nodes (copy-then-switch)
3. When empty, node is removed from cluster membership
Node failure (unplanned):
1. MDS detects via missed heartbeats (DSS → MDS health ping every 5s)
2. Node marked "failed"; its replicas/shards now under-replicated
3. HealManager re-replicates affected chunks from surviving copies
to restore the replication/parity factor
4. Client reads/writes route around the failed node immediately
(circuit breaker + chunk map has other locations)
8. Protocol Compatibility
8.1 S3 Compatibility (Preserve Existing)
The existing S3 API surface is preserved 100%. The only change is the backend:
S3 handler → UnifiedObjectLayer → MDS + DSS (instead of → ecstore directly)
Compatibility requirements:
- ETag semantics: MD5 for single-part, "{md5}-{N}" for multipart
- Multipart upload: parts staged as chunks, assembled on CompleteMultipartUpload
- Versioning: maps to inode version chain (S3 version_id ↔ inode generation)
- Conditional requests (If-Match, If-None-Match): checked against inode etag
- Range requests: map directly to chunk map resolve_range()
- List with prefix/delimiter: MDS readdir + prefix filter (now O(1) per dir)
- Object tagging, ACLs, bucket policies: stored as inode xattrs / S3 metadata
Edge cases that need explicit handling:
- S3 allows "keys" with "/" that don't correspond to real directories.
→ MDS auto-creates intermediate directories (mkdir -p semantics) on PUT.
- S3 objects can have the same name as a "directory" prefix.
→ Resolved by the bucket_root mapping; a key "a/b" always creates dir "a".
- Empty "directory marker" objects (key ending in "/"):
→ Mapped to an actual empty directory inode.
8.2 POSIX Compatibility
POSIX semantics supported:
✓ open/close/read/write/lseek/pread/pwrite
✓ stat/fstat/lstat with correct st_size, st_blocks, timestamps
✓ mkdir/rmdir/readdir/opendir (with telldir/seekdir cookies)
✓ rename (atomic within a single MDS partition; 2PC cross-partition)
✓ link (hard link — increments nlink, two dentries → one inode)
✓ symlink/readlink
✓ chmod/chown/utimes
✓ truncate/ftruncate
✓ fsync/fdatasync (flush dirty chunks + commit to MDS)
✓ flock (advisory whole-file) and fcntl byte-range locks
✓ getxattr/setxattr/listxattr/removexattr
✓ mmap (read-only fully; read-write via FUSE writeback cache)
✓ O_DIRECT (bypasses client cache → direct DSS chunk I/O)
POSIX semantics with caveats:
~ Strict atime updates: relaxed by default (relatime semantics) for
performance; strict atime is opt-in (costly: a write per read).
~ Close-to-open consistency guaranteed; stronger consistency requires
disabling client caching (or relies on lease coherency, which we provide).
Not supported (by design):
✗ Device files other than the block-volume special files
✗ Named pipes / sockets in the filesystem (use real IPC)
8.3 Block Protocol Compatibility
iSCSI (RFC 7143):
✓ SCSI READ(10/16), WRITE(10/16), VERIFY, SYNCHRONIZE_CACHE
✓ INQUIRY, READ_CAPACITY(10/16), REPORT_LUNS, TEST_UNIT_READY
✓ Persistent Reservations (PR) — maps to MDS-based volume locks
(needed for clustered VM hosts / shared disk)
✓ UNMAP / WRITE_SAME (thin provisioning TRIM)
✓ Multi-path (MPIO): volume exported from multiple iSCSI targets,
initiator load-balances; MDS coordinates consistency
NVMe-oF (NVMe over Fabrics):
✓ NVMe Read, Write, Flush, Dataset Management (TRIM)
✓ Transport: TCP (universal) and RDMA (RoCEv2/IB, low latency)
✓ Namespaces map 1:1 to volumes
✓ Reservations (NVMe Reservation commands) for shared access
Hypervisor integration:
✓ QEMU/KVM: via iSCSI initiator or NVMe-oF, or libvirt storage pool
✓ VMware: via iSCSI (VMFS datastore) — needs VAAI primitives (future)
✓ Kubernetes: CSI driver
- block mode (RWO): one volume per pod, iSCSI/NVMe-oF attach
- filesystem mode (RWX): FUSE mount shared across pods
8.4 Kubernetes CSI Driver
// crate: rustfs-csi
// file: crates/csi/src/lib.rs
/// CSI driver implementing the Container Storage Interface.
///
/// Exposes two storage classes:
/// - "rustfs-block" (RWO): provisions a Volume, attaches via NVMe-oF/iSCSI
/// - "rustfs-fs" (RWX): provisions a directory, mounts via FUSE
///
/// Implements:
/// - Controller service: CreateVolume, DeleteVolume, ControllerPublishVolume,
/// CreateSnapshot, ExpandVolume
/// - Node service: NodeStageVolume, NodePublishVolume (mount), NodeUnpublish
pub struct RustfsCsiDriver {
mds: MdsClient,
}
// CSI CreateVolume → MDS CreateVolume (block) or MkDir (filesystem)
// CSI NodePublishVolume (block) → attach NVMe-oF namespace, present /dev node
// CSI NodePublishVolume (fs) → rustfs-mount the subtree at target path
// CSI CreateSnapshot → MDS SnapshotVolume (COW)
8.5 Protocol Interop Matrix
Written via S3 Written via POSIX Written via Block
───────────────── ────────────── ──────────────── ─────────────────
Readable via S3 ✓ native ✓ (if under ✗ (block volumes
/buckets/) not in S3 ns)
Readable via POSIX ✓ (appears in ✓ native ✓ (as /volumes/
/buckets/dir) special file)
Readable via Block ✗ ✗ ✓ native
Key rule: file/object data is bidirectionally accessible (same inode,
two views). Block volumes are a distinct inode kind, accessible as a raw
device via block protocols or as a special file via POSIX, but NOT as an
S3 object (block semantics don't map to object semantics).
9. Consolidated Tech Stack
Layer Technology Rust crate / library
───────────────── ────────────────────────── ────────────────────────────
Async runtime tokio + tokio-uring tokio, tokio-uring
Disk I/O io_uring (raw) io-uring; SPDK FFI (opt-in)
Network (data) RDMA verbs / io_uring SEND_ZC rdma-sys/async-rdma; io-uring
Network (control) gRPC over HTTP/2 tonic, prost
Metadata store RocksDB rocksdb (rust-rocksdb)
Consensus Raft openraft
Caches TinyLFU concurrent cache moka
Concurrent maps sharded lock-free maps dashmap
FUSE FUSE (+ io_uring passthru) fuser (+ custom uring layer)
Erasure coding Reed-Solomon (existing) rustfs-ecstore-core (reuse)
Serialization MessagePack (metadata), rmp-serde; prost (wire)
Protobuf (wire)
Checksums HighwayHash (existing) rustfs-checksums (reuse)
Block targets iSCSI / NVMe-oF custom + SPDK nvmf (opt-in)
Buffer pool slab + fixed io_uring bufs rustfs-io-core (reuse/extend)
Compression existing pipeline rustfs-rio (reuse)
Encryption SSE-S3/KMS (existing) rustfs-crypto/kms (reuse)
Observability OpenTelemetry + Prometheus rustfs-obs/metrics (reuse)
10. Critical Path Latency Budgets
10.1 POSIX read(), warm cache, 4 MiB
Step Budget Mechanism
───────────────────────────────────────── ──────── ─────────────────────────
FUSE request in (io_uring) 1 µs FUSE-over-io_uring
Inode cache hit (lease-backed) 0.1 µs moka in-memory
Chunk map cache hit (layout lease) 0.5 µs PartialChunkMap
Pick replica (power-of-two) 0.2 µs in-memory
RDMA READ 4 MiB ~20 µs 200 Gbps NIC
Assemble + FUSE reply (io_uring) 2 µs zero-copy
───────────────────────────────────────── ────────
TOTAL (RDMA, uncached data) ~24 µs
If data is in client chunk cache: ~3 µs (no network at all)
10.2 POSIX read(), cold, fallback TCP, 4 MiB
FUSE request in 5 µs classic /dev/fuse
Inode/chunkmap lookup (MDS round trip) 1 ms gRPC + RocksDB
Pick replica 0.2 µs
io_uring TCP read 4 MiB ~400 µs 100 Gbps
Assemble + reply 10 µs
───────────────────────────────────────── ────────
TOTAL ~1.4 ms (dominated by MDS lookup)
Optimization: with read-ahead + layout lease, the MDS lookup is off the
critical path for sequential access → drops to ~420 µs.
10.3 Block write 4 KiB, replicated x3, synchronous
iSCSI WRITE command parse 2 µs
WAL append (io_uring linked write+fsync) ~30 µs NVMe WAL device
Write to 3 replicas (parallel) ~30 µs max of 3, io_uring O_DSYNC
MDS chunk map update (only if new chunk) 0 or 1ms thin provision first-write
───────────────────────────────────────── ────────
TOTAL (existing chunk) ~60 µs
TOTAL (first write to new chunk) ~1 ms (one-time MDS alloc)
11. Implementation Sequencing (Mapped to LLD Components)
Sprint Component Depends on
────── ──────────────────────────────────────────── ──────────────────
1-2 io-core buffer pool extension (fixed bufs) (existing io-core)
1-2 DiskIoEngine (io_uring backend) buffer pool
3-4 chunk-service gRPC + io_uring read/write DiskIoEngine
3-4 MDS RocksDB engine + inode allocator (rustfs-inode types)
5-6 MDS namespace service + path cache RocksDB engine
5-6 rustfs-client core + DssConnectionPool chunk-service
7-8 FUSE client (classic) + read/write paths client core, MDS
7-8 retry/backoff/circuit breaker client core
9-10 read-ahead + lease manager (client + MDS) FUSE, MDS
9-10 error→heal integration (reuse rustfs-heal)
11-12 ChunkPlacer + Rebalancer MDS, chunk-service
11-12 S3 UnifiedObjectLayer MDS, client core
13-14 block VolumeService + extent map + COW MDS, chunk-service
13-14 iSCSI target + WAL VolumeService
15-16 RDMA transport (data plane) client, chunk-service
15-16 FUSE-over-io_uring FUSE client
17-18 Raft HA MDS MDS
17-18 CSI driver MDS, FUSE, iSCSI
19-20 NVMe-oF target VolumeService
19-20 SPDK backend (opt-in) DiskIoEngine trait
21+ GPUDirect Storage, MDS sharding (DNE) RDMA, Raft MDS
12. Risks and Mitigations
Risk Mitigation
────────────────────────────────────────── ────────────────────────────────
io_uring kernel version fragmentation trait-abstracted DiskIoEngine with
tokio::fs fallback; feature-detect
RDMA not available on commodity clusters transport negotiation → io_uring TCP
FUSE overhead caps POSIX throughput FUSE-over-io_uring; later kernel
module or CUSE if needed
MDS becomes metadata bottleneck path cache + client lease caching +
MDS sharding (DNE)
RocksDB compaction stalls metadata ops WAL on separate device; rate-limit
compaction; monitor write stalls
Rebalancing competes with client I/O token-bucket rate limiting; yield
to client load
EC partial-write amplification policy routing: EC for write-once,
replication for random-write
Cross-MDS rename complexity 2PC; or constrain renames within a
partition where possible
Split-brain in Raft proper quorum config (odd nodes);
fencing tokens on lease grants
This LLD is the implementation contract for the I/O, transport, metadata, error-handling, rebalancing, and protocol subsystems. It should be reviewed against the Feature Spec (RUSTFS-FEAT-001) before each implementation sprint and updated as decisions in Section 12 (open questions) are resolved.