RustFS Deep Dive: Metadata & Chunk Server Design Analysis
RustFS Deep Dive: Metadata & Chunk Server Design Analysis
Change Plan for POSIX Filesystem & Block Storage Workloads
Part 1 — Current Metadata Architecture (What Exists)
1.1 The filemeta Crate (10K lines)
RustFS stores all metadata in xl.meta files co-located with data on each disk. There is no separate metadata server. The on-disk format uses a binary header:
Header: [b'X', b'L', b'2', b' '] (XL_FILE_HEADER)
Version: 3 (XL_META_VERSION)
Format: MessagePack-encoded payload
The central metadata structure is ObjectInfo:
// Reconstructed from debug logs and DeepWiki analysis
// Source: crates/ecstore/src/store_api.rs, crates/filemeta/src/fileinfo.rs
struct ObjectInfo {
bucket: String,
name: String, // "key" — full object path
storage_class: Option<String>, // "STANDARD", "REDUCED_REDUNDANCY", etc.
mod_time: Option<DateTime>,
size: u64,
actual_size: u64, // pre-compression size
is_dir: bool, // synthetic directory marker
user_defined: HashMap<String, String>, // x-amz-meta-* headers
parity_blocks: u32,
data_blocks: u32,
version_id: Option<String>,
delete_marker: bool,
transitioned_object: TransitionedObject,
restore_ongoing: bool,
restore_expires: Option<DateTime>,
user_tags: String,
parts: Vec<ObjectPartInfo>, // multipart upload parts
is_latest: bool,
content_type: Option<String>,
content_encoding: Option<String>,
expires: Option<DateTime>,
num_versions: u32,
successor_mod_time: Option<DateTime>,
put_object_reader: Option<...>,
etag: Option<String>,
// ... internal fields
}
Each object part carries its own erasure info:
struct ObjectPartInfo {
etag: String,
number: u32, // part number (1-based)
size: u64,
actual_size: u64,
mod_time: Option<DateTime>,
index: Option<Vec<u8>>,
checksums: Option<ChecksumInfo>,
error: Option<String>,
}
struct ErasureInfo {
algorithm: String, // "rs-vandermonde"
data_blocks: u32, // e.g., 4
parity_blocks: u32, // e.g., 4
block_size: u64, // fixed 1 MiB (1048576)
distribution: Vec<u32>, // shard-to-disk mapping per part
checksums: Vec<ChecksumInfo>, // per-block bitrot checksums
// index: shard index within the set
}
struct ChecksumInfo {
part_number: u32,
algorithm: BitrotAlgorithm, // HighwayHash256S
hash: Vec<u8>,
}
Critical observations:
-
No inode concept. Objects are identified by
(bucket, name)string pairs, not numeric inode IDs. There is no inode table, no inode allocation bitmap. -
No POSIX attributes. No
uid,gid,mode,nlink,atime,ctime,mtime(onlymod_timewhich is S3'sLast-Modified). No extended attributes beyond S3'sx-amz-meta-*. -
Directories are synthetic. The
is_dirflag is a convenience marker for the S3 "prefix" illusion. There is no actual directory tree data structure —ListObjectsscans and filters by prefix. -
Metadata is embedded with data. Each
xl.metafile lives alongside the erasure-coded shards on every disk in the set. To read metadata, you must reach quorum on the disk set (read N copies, pick the latest version). -
No open-file table. There is no concept of
open()returning a file descriptor. Every S3 GET is stateless. -
Version-linked list. Versioning chains objects by
version_id+successor_mod_time. This is S3 versioning, not POSIX versioning.
1.2 The ecstore Crate — Storage Engine (87K lines)
The monolithic heart of RustFS. Key subsystems:
ecstore/src/
├── store.rs # ECStore: top-level coordinator
├── store_api.rs # Public API surface (PutObject, GetObject, etc.)
├── store_list_objects.rs # ListObjects implementation
├── sets.rs # ErasureSets: groups of SetDisks
├── set_disk.rs # SetDisks: single erasure set operations (2500+ lines)
├── pools.rs # EndpointServerPools: multi-pool management
├── disk/
│ ├── local.rs # LocalDisk: physical disk abstraction
│ ├── error_reduce.rs # Quorum error handling
│ └── ...
├── erasure_coding/
│ ├── decode.rs # Reed-Solomon decode + reconstruct
│ ├── bitrot.rs # Bitrot reader/writer
│ ├── bitrot_verify.rs # Checksum verification
│ └── ...
├── bucket/
│ ├── replication/ # Cross-site bucket replication
│ ├── lifecycle/ # Object lifecycle management
│ └── ...
├── rpc/
│ ├── http_auth.rs # HMAC-SHA256 inter-node auth
│ └── ... # gRPC/tonic service definitions
├── config/ # ECStore configuration
├── data_usage.rs # Capacity tracking
├── rebalance.rs # Disk rebalancing
└── ...
Object-to-disk mapping flow:
bucket + object_key
│
▼
SipHash(deployment_id, bucket, key) → set_index (which erasure set)
│
▼
CRC32(key) → distribution[N] (which disk in the set gets which shard)
│
▼
Disk path: /{volume}/{bucket}/{key_hash}/xl.meta + part.N files
The SetDisks coordinator (the most important struct) handles:
- Write: split data → EC encode → parallel write shards → quorum check → write xl.meta
- Read: parallel read xl.meta from N disks → quorum pick → parallel read shards → EC decode
- Delete: parallel delete across disk set
- Heal: detect missing/corrupt shards → reconstruct from parity
Inter-node RPC:
Endpoint: /rustfs/rpc/read_file_stream
Params: disk, volume, path, offset, length
Auth: HMAC-SHA256 with shared secret
Transport: gRPC/tonic (with HTTP fallback for file streaming)
Also via protobuf: /node_service.NodeService/*
The RPC already supports offset + length parameters for reading partial file streams — this is critical and reusable for byte-range POSIX reads.
1.3 The lock Crate (7.1K lines)
crates/lock/src/
├── lock_manager.rs # LockManager: coordinates namespace locks
├── namespace_lock.rs # NamespaceLock / NamespaceLockGuard
└── ...
Locking is at object granularity: (bucket, key) pairs. Used inside SetDisks operations via NamespaceLock (referenced at set_disk.rs:80).
What's missing for POSIX:
- No byte-range locks (only whole-object)
- No shared/exclusive differentiation at byte-range level
- No distributed lease protocol (no heartbeat, no timeout-based revocation)
- No deadlock detection across nodes
1.4 The io-core Crate (6.5K lines)
Provides zero-copy I/O primitives:
- Buffer pool with slab allocation
- Direct I/O support (O_DIRECT alignment)
- I/O scheduling and backpressure
- Async read/write with tokio
1.5 The rio Crate (6.9K lines)
Composable reader pipeline:
Input stream → Encrypt → Compress → Hash → Limit → Output
Reusable as-is for the POSIX write path.
Part 2 — Gap Analysis by Data Structure
2.1 Metadata Model: Object vs. Inode
| Concept | RustFS (S3 Object Model) | POSIX FS (Needed) | Block Device (Needed) |
|---|---|---|---|
| Identity | (bucket, key) string pair | inode_id (u64) | (volume_id, LBA) |
| Namespace | Flat key space with / delimiter | Hierarchical directory tree | Flat linear address space |
| Attributes | S3 metadata headers | uid, gid, mode, nlink, size, timestamps (a/m/c/birth) | Volume size, block_size |
| Size semantics | Immutable after PUT | Grows/shrinks with write/truncate | Fixed at creation (or grow) |
| Link model | None | Hard links (nlink), symlinks | N/A |
| Directory | Synthetic (ListObjects prefix scan) | Real directory entries (dirent) | N/A |
| Open state | Stateless (each GET independent) | File descriptor + offset + flags | Session/connection state |
| Locking | Object-level namespace lock | Byte-range flock/fcntl + leases | N/A (handled by FS above) |
2.2 Data Path: Immutable Object vs. Mutable File
S3 PUT (current):
Client sends entire object → server receives all bytes →
EC encode → write all shards atomically → write xl.meta
POSIX write (needed):
Client opens file → seeks to offset → writes partial bytes →
Only affected chunks updated → metadata (size, mtime) updated in-place
Block write (needed):
Client writes to LBA offset → only affected 4K block written →
Write barrier/flush support required
The fundamental mismatch: RustFS treats objects as immutable blobs written atomically, while POSIX and block I/O require in-place mutation of arbitrary byte ranges.
2.3 Metadata Lookup: Prefix Scan vs. Directory Lookup
S3 ListObjects (current):
Input: bucket + prefix + delimiter
Action: Scan ALL keys, filter by prefix ← O(N) per bucket
Return: Matching keys + common prefixes
POSIX readdir (needed):
Input: directory inode_id
Action: Lookup directory's child entries ← O(1) per directory
Return: List of (name, inode_id, type)
POSIX stat (needed):
Input: inode_id
Action: Single key lookup ← O(1)
Return: Full attribute struct
ListObjects scales poorly as a directory substitute — a bucket with 10M objects means scanning 10M keys to list a single subdirectory.
Part 3 — Concrete Change Plan
Phase A: Foundation Layer — New Core Data Structures
A.1 New Crate: crates/inode/ — Inode System
// crates/inode/src/types.rs
/// Unique identifier for every file, directory, symlink
/// Uses a 64-bit space: upper 16 bits = MDS partition, lower 48 bits = local counter
pub type InodeId = u64;
pub const ROOT_INODE: InodeId = 1;
pub const INVALID_INODE: InodeId = 0;
/// POSIX inode attributes — stored in MDS
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct InodeAttr {
pub ino: InodeId,
pub size: u64,
pub blocks: u64, // 512-byte blocks (for stat)
pub atime: Timespec,
pub mtime: Timespec,
pub ctime: Timespec,
pub birthtime: Timespec,
pub kind: InodeKind, // RegularFile | Directory | Symlink | BlockDevice
pub perm: u16, // rwxrwxrwx
pub nlink: u32,
pub uid: u32,
pub gid: u32,
pub rdev: u32, // for device files
pub flags: u32,
pub xattrs: BTreeMap<String, Vec<u8>>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub enum InodeKind {
RegularFile,
Directory,
Symlink(String), // target path
BlockDevice { volume_id: VolumeId },
}
/// Chunk location: where a region of a file's data lives
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct ChunkLocation {
pub chunk_id: ChunkId, // globally unique
pub oss_node: NodeId, // which storage server
pub set_index: u32, // which erasure set
pub version: u64, // for COW
}
/// Maps file content to storage locations
/// Key insight: this replaces S3's flat `ObjectPartInfo` list
/// with a range-indexed structure that supports byte-range I/O
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct ChunkMap {
pub inode: InodeId,
pub chunk_size: u64, // e.g., 4 MiB
/// BTreeMap<chunk_index, Vec<ChunkLocation>>
/// chunk_index = file_offset / chunk_size
/// Vec because of replication factor or EC distribution
pub chunks: BTreeMap<u64, Vec<ChunkLocation>>,
}
impl ChunkMap {
/// Given a byte range [offset, offset+length), return which chunks to read
pub fn chunks_for_range(&self, offset: u64, length: u64) -> Vec<(u64, &[ChunkLocation])> {
let start_idx = offset / self.chunk_size;
let end_idx = (offset + length - 1) / self.chunk_size;
self.chunks.range(start_idx..=end_idx)
.map(|(idx, locs)| (*idx, locs.as_slice()))
.collect()
}
}
Why this design:
InodeIduses partitioned 64-bit space to support future MDS sharding (DNE)ChunkMapis a BTreeMap for efficient range queries — critical forfseek + freadChunkLocationdecouples the "where is the data" from "what is the data" — the same chunk storage can serve POSIX, S3, and block workloadsInodeKind::BlockDeviceallows block volumes to appear as special files
A.2 New Crate: crates/dentry/ — Directory Entry System
// crates/dentry/src/lib.rs
/// A single directory entry
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct DirEntry {
pub name: String, // filename component (not full path)
pub inode: InodeId,
pub kind: InodeKind, // cached for readdir optimization
}
/// Directory content: stored as a sorted list per directory inode
/// For large directories (>10K entries), split into sub-pages
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct DirectoryData {
pub parent: InodeId,
pub entries: BTreeMap<String, DirEntry>, // sorted by name
}
impl DirectoryData {
pub fn lookup(&self, name: &str) -> Option<&DirEntry> {
self.entries.get(name)
}
pub fn insert(&mut self, name: String, entry: DirEntry) -> Option<DirEntry> {
self.entries.insert(name, entry)
}
pub fn remove(&mut self, name: &str) -> Option<DirEntry> {
self.entries.remove(name)
}
pub fn list(&self, offset: usize, count: usize) -> Vec<&DirEntry> {
self.entries.values().skip(offset).take(count).collect()
}
}
Why BTreeMap, not HashMap: readdir must return entries in a stable order (required by POSIX telldir/seekdir). BTreeMap gives sorted iteration for free.
A.3 Extend Existing lock Crate — Byte-Range Locks + Leases
// crates/lock/src/posix_lock.rs — NEW FILE
/// POSIX byte-range lock
#[derive(Clone, Debug)]
pub struct ByteRangeLock {
pub inode: InodeId,
pub owner: LockOwner, // (client_id, pid)
pub kind: LockKind, // Shared | Exclusive
pub start: u64, // byte offset
pub end: u64, // byte offset (u64::MAX = EOF)
}
#[derive(Clone, Debug)]
pub enum LockKind {
Shared, // F_RDLCK — many readers
Exclusive, // F_WRLCK — single writer
}
/// Client lease for caching correctness (like Lustre's LDLM)
#[derive(Clone, Debug)]
pub struct Lease {
pub inode: InodeId,
pub client_id: ClientId,
pub kind: LeaseKind,
pub granted_at: Instant,
pub expires_at: Instant, // must heartbeat to renew
}
#[derive(Clone, Debug)]
pub enum LeaseKind {
ReadLease, // client can cache reads
WriteLease, // client has exclusive write + read cache
LayoutLease(ChunkRange), // client knows chunk layout (won't change)
}
/// Extended LockManager with byte-range and lease support
pub struct PosixLockManager {
/// Byte-range locks: BTreeMap indexed by (inode, start_offset)
range_locks: RwLock<BTreeMap<(InodeId, u64), Vec<ByteRangeLock>>>,
/// Active leases: HashMap indexed by (inode, client_id)
leases: RwLock<HashMap<(InodeId, ClientId), Lease>>,
/// Existing namespace lock (reused for directory ops)
namespace_lock: NamespaceLock,
}
impl PosixLockManager {
/// Check if a byte-range lock can be granted
pub fn can_lock(&self, req: &ByteRangeLock) -> bool {
// Check for overlapping conflicting locks
// Shared locks compatible with shared; exclusive conflicts with all
// Interval-tree overlap detection
todo!()
}
/// Revoke leases when a conflicting operation arrives
/// (e.g., client A has ReadLease, client B wants to write)
pub async fn revoke_conflicting_leases(
&self,
inode: InodeId,
kind: LeaseKind,
) -> Result<(), LockError> {
// Send revocation callback to affected clients
// Wait for acknowledgment or timeout
todo!()
}
}
Design rationale — lease model inspired by Lustre's LDLM:
- ReadLease: client can trust its page cache for reads. Revoked when any client writes.
- WriteLease: only one client gets this. Can cache writes locally, flush on lease revocation.
- LayoutLease: client caches the chunk map. Revoked when file is re-striped or extended.
This is what makes multi-reader AI workloads fast: hundreds of GPU nodes get ReadLeases on training data files, and since nobody writes during training, the leases are never revoked → every read hits client-side cache after first fetch.
Phase B: Metadata Server (rustfs-mds)
B.1 New Crate: crates/mds/ — Metadata Server
// crates/mds/src/store.rs
/// Metadata storage backend trait — allows swappable implementations
#[async_trait]
pub trait MetadataStore: Send + Sync {
// Inode operations
async fn get_inode(&self, ino: InodeId) -> Result<InodeAttr>;
async fn set_inode(&self, attr: &InodeAttr) -> Result<()>;
async fn alloc_inode(&self) -> Result<InodeId>;
async fn free_inode(&self, ino: InodeId) -> Result<()>;
// Directory operations
async fn get_directory(&self, ino: InodeId) -> Result<DirectoryData>;
async fn set_directory(&self, ino: InodeId, dir: &DirectoryData) -> Result<()>;
// Chunk map operations
async fn get_chunk_map(&self, ino: InodeId) -> Result<ChunkMap>;
async fn set_chunk_map(&self, ino: InodeId, map: &ChunkMap) -> Result<()>;
async fn update_chunks(
&self,
ino: InodeId,
updates: &[(u64, Vec<ChunkLocation>)], // (chunk_index, new_locations)
) -> Result<()>;
// Path resolution (full path → inode, walking directory tree)
async fn resolve_path(&self, path: &str) -> Result<InodeId>;
// Transactions
async fn begin_tx(&self) -> Result<TxId>;
async fn commit_tx(&self, tx: TxId) -> Result<()>;
async fn abort_tx(&self, tx: TxId) -> Result<()>;
}
B.2 Storage Backend Options
// crates/mds/src/backend/rocksdb.rs
/// Single-node MDS using RocksDB for metadata persistence
///
/// Key schema:
/// i:{inode_id} → InodeAttr (msgpack)
/// d:{inode_id} → DirectoryData (msgpack)
/// c:{inode_id} → ChunkMap (msgpack)
/// p:{parent_inode}:{name} → child InodeId (for fast lookup)
/// n:next_inode → u64 (inode counter)
///
/// Column families:
/// "inodes" — inode attributes
/// "dirs" — directory entries
/// "chunks" — chunk maps
/// "path_index" — parent+name → inode index
/// "wal" — write-ahead log for transactions
pub struct RocksDbMetadataStore {
db: Arc<rocksdb::DB>,
}
// crates/mds/src/backend/raft.rs
/// HA MDS cluster using Raft consensus (via openraft)
///
/// All metadata mutations go through Raft log.
/// Reads can be served from any replica (with stale-read option)
/// or from leader only (strong consistency).
pub struct RaftMetadataStore {
raft: Arc<openraft::Raft<MdsTypeConfig>>,
store: Arc<RocksDbMetadataStore>, // local state machine
}
Why RocksDB: LSM-tree gives excellent write throughput for metadata-heavy workloads (creating millions of small files). Read-optimized with bloom filters. Same choice as JuiceFS, CephFS's RocksDB option, and TiKV.
B.3 MDS gRPC Service Definition
// crates/mds/proto/mds.proto
service MetadataService {
// Path operations
rpc Lookup(LookupReq) returns (LookupResp); // parent + name → inode + attr
rpc Resolve(ResolveReq) returns (ResolveResp); // full path → inode + attr
// Inode attribute operations
rpc GetAttr(GetAttrReq) returns (GetAttrResp);
rpc SetAttr(SetAttrReq) returns (SetAttrResp);
// Directory operations
rpc ReadDir(ReadDirReq) returns (stream ReadDirResp);
rpc MkDir(MkDirReq) returns (MkDirResp);
rpc RmDir(RmDirReq) returns (RmDirResp);
// File operations
rpc Create(CreateReq) returns (CreateResp); // create + open
rpc Unlink(UnlinkReq) returns (UnlinkResp);
rpc Rename(RenameReq) returns (RenameResp);
rpc Link(LinkReq) returns (LinkResp); // hard link
rpc Symlink(SymlinkReq) returns (SymlinkResp);
// Chunk map operations (for data path)
rpc GetChunkMap(GetChunkMapReq) returns (GetChunkMapResp);
rpc AllocChunks(AllocChunksReq) returns (AllocChunksResp);
rpc CommitChunks(CommitChunksReq) returns (CommitChunksResp);
// Lock operations
rpc AcquireLock(LockReq) returns (LockResp);
rpc ReleaseLock(UnlockReq) returns (UnlockResp);
rpc AcquireLease(LeaseReq) returns (LeaseResp);
rpc RenewLease(LeaseRenewReq) returns (LeaseRenewResp);
// Session management (client registration)
rpc OpenSession(SessionReq) returns (SessionResp);
rpc Heartbeat(HeartbeatReq) returns (HeartbeatResp);
rpc CloseSession(CloseReq) returns (CloseResp);
// Block volume operations
rpc CreateVolume(CreateVolumeReq) returns (CreateVolumeResp);
rpc DeleteVolume(DeleteVolumeReq) returns (DeleteVolumeResp);
rpc ResizeVolume(ResizeVolumeReq) returns (ResizeVolumeResp);
rpc SnapshotVolume(SnapshotReq) returns (SnapshotResp);
rpc GetVolumeMap(GetVolumeMapReq) returns (GetVolumeMapResp);
}
B.4 MDS Namespace Partitioning (Distributed Namespace — DNE)
For scale beyond a single MDS (millions of files, hundreds of clients):
Strategy: Subtree partitioning (like Lustre DNE phase 1)
/
├── dataset/ → MDS-0 (AI training data, high read metadata traffic)
│ ├── imagenet/
│ └── llm-tokens/
├── checkpoints/ → MDS-1 (model checkpoints, high write traffic)
│ ├── epoch-001/
│ └── epoch-002/
├── scratch/ → MDS-2 (temporary files, high create/delete rate)
└── volumes/ → MDS-3 (block device metadata)
├── vm-disk-001
└── vm-disk-002
Routing: path prefix → MDS partition
Cross-MDS rename: two-phase commit between source and target MDS
Future: directory-level hash partitioning (DNE phase 2), where entries within a single large directory are distributed across multiple MDS nodes.
Phase C: Chunk Server Adaptation (rustfs-oss)
C.1 New Chunk-Level gRPC Service (Alongside Existing S3 API)
// crates/ecstore/proto/chunk_service.proto
service ChunkService {
// Write a chunk (new or overwrite)
rpc WriteChunk(stream WriteChunkReq) returns (WriteChunkResp);
// Read byte range from a chunk
rpc ReadChunk(ReadChunkReq) returns (stream ReadChunkResp);
// Read multiple chunks in parallel (batch read)
rpc ReadChunks(ReadChunksReq) returns (stream ReadChunksResp);
// Delete chunks
rpc DeleteChunks(DeleteChunksReq) returns (DeleteChunksResp);
// Chunk metadata
rpc StatChunk(StatChunkReq) returns (StatChunkResp);
// Allocate space for new chunks (thin provisioning)
rpc AllocateChunks(AllocateReq) returns (AllocateResp);
// Flush: ensure chunk data is durable on disk
rpc FlushChunk(FlushReq) returns (FlushResp);
// Truncate: resize a chunk (for file truncation)
rpc TruncateChunk(TruncateReq) returns (TruncateResp);
}
message ReadChunkReq {
string chunk_id = 1;
uint64 offset = 2; // byte offset within chunk
uint64 length = 3; // bytes to read
bool direct_io = 4; // bypass OS page cache (for O_DIRECT)
}
message WriteChunkReq {
string chunk_id = 1;
uint64 offset = 2; // byte offset within chunk (for partial writes!)
bytes data = 3;
bool sync = 4; // fsync after write
}
Critical change: partial chunk writes. The existing ecstore always writes complete EC-encoded objects. For POSIX/block, we need to write arbitrary byte ranges within a chunk. Two strategies:
Strategy 1: Chunk = Plain file on local disk (no EC per-chunk)
- Each chunk is a regular file on a single disk
- Replication factor (e.g., 3 copies) provides durability
- Supports in-place partial writes trivially
- Used by: CephFS (RADOS objects), JuiceFS (S3 objects with append)
Strategy 2: Chunk = EC-encoded across disk set (current model adapted)
- Partial write requires: read-modify-write of affected EC block
- 1 MiB EC block: a 4K write requires reading 1 MiB, modifying 4K, re-encoding, writing back
- Write amplification: 256x for 4K writes → unacceptable for block device
- Used by: only when full-stripe writes are common (AI data ingest)
Recommended hybrid approach:
Data path selection based on workload:
POSIX files (AI training data):
- Write path: EC striping (current model) — data is written once in large sequential writes
- Read path: parallel stripe reads — data is read many times from multiple clients
- Chunk size: 4 MiB (configurable)
- EC is perfect here because the pattern is write-once-read-many
Block volumes (VM disks):
- Write path: 3-way replication, NOT EC — random 4K writes demand in-place mutation
- Read path: read from nearest replica
- Chunk size: 4 MiB, but internal block granularity = 4K (sub-chunk writes)
- Uses a journaling/WAL layer for write ordering guarantees
POSIX files (random write workloads):
- Write path: replication with COW — write new chunk version, update chunk map
- Read path: read from latest chunk version
- Chunk size: 1 MiB
- COW avoids read-modify-write amplification
C.2 Changes to Existing ecstore Crate
The existing 87K-line monolith needs decomposition before adding POSIX/block support:
CURRENT (monolith):
ecstore/
├── Everything in one crate (87K lines)
PROPOSED (decomposed):
crates/
├── ecstore-core/ # ~20K: EC encoding/decoding, shard math, quorum logic
│ ├── erasure_coding/
│ ├── bitrot.rs
│ └── quorum.rs
│
├── ecstore-disk/ # ~15K: LocalDisk, disk health, I/O scheduling
│ ├── local.rs
│ ├── health.rs
│ └── scheduler.rs
│
├── ecstore-pool/ # ~10K: EndpointServerPools, set management, rebalance
│ ├── pools.rs
│ ├── sets.rs
│ └── rebalance.rs
│
├── ecstore-s3/ # ~25K: S3 API object operations (SetDisks current logic)
│ ├── set_disk.rs # Existing S3 PUT/GET/DELETE logic
│ ├── store_api.rs
│ └── store_list_objects.rs
│
├── ecstore-chunk/ # ~8K: NEW — chunk-level operations for POSIX/block
│ ├── chunk_service.rs # gRPC ChunkService implementation
│ ├── chunk_store.rs # Local chunk storage manager
│ └── chunk_replicate.rs
│
├── ecstore-block/ # ~10K: NEW — block volume extent management
│ ├── volume.rs
│ ├── extent_map.rs
│ ├── cow.rs
│ └── journal.rs
│
└── ecstore-bucket/ # ~12K: Bucket management, replication, lifecycle
├── bucket/
└── lifecycle/
C.3 Concrete ecstore Interface Changes
// crates/ecstore-chunk/src/chunk_service.rs
/// Adapts existing ecstore primitives for chunk-level operations
pub struct ChunkServiceImpl {
/// Reuse existing disk pool for chunk storage
pools: Arc<EndpointServerPools>,
/// Reuse existing I/O pipeline
io_core: Arc<IoCore>,
}
impl ChunkServiceImpl {
/// Write partial data to a chunk
///
/// For EC-mode chunks: performs read-modify-write on the affected EC block
/// For replica-mode chunks: writes directly to the chunk file
pub async fn write_chunk(
&self,
chunk_id: &ChunkId,
offset: u64,
data: &[u8],
mode: ChunkMode,
) -> Result<()> {
match mode {
ChunkMode::Replicated { factor } => {
// Write to all replicas in parallel
// Uses existing io-core zero-copy primitives
let disks = self.pools.get_replica_disks(chunk_id, factor)?;
let writes = disks.iter().map(|disk| {
disk.write_at(chunk_id.path(), offset, data)
});
futures::future::try_join_all(writes).await?;
Ok(())
}
ChunkMode::ErasureCoded { data_blocks, parity_blocks } => {
if data.len() as u64 >= EC_BLOCK_SIZE {
// Full-block write: encode and write directly
self.write_ec_full_block(chunk_id, offset, data, data_blocks, parity_blocks).await
} else {
// Partial write: read-modify-write
self.write_ec_partial(chunk_id, offset, data, data_blocks, parity_blocks).await
}
}
}
}
/// Read byte range from a chunk
///
/// Reuses existing RPC: /rustfs/rpc/read_file_stream?offset=X&length=Y
pub async fn read_chunk(
&self,
chunk_id: &ChunkId,
offset: u64,
length: u64,
) -> Result<Bytes> {
// Determine which disk(s) have this chunk
let disks = self.pools.get_chunk_disks(chunk_id)?;
// Read from nearest/fastest disk (or parallel EC decode)
let reader = self.io_core.create_reader(disks, offset, length);
reader.read_all().await
}
}
Phase D: Block Storage Layer
D.1 New Crate: crates/block-volume/
// crates/block-volume/src/volume.rs
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Volume {
pub id: VolumeId,
pub name: String,
pub size: u64, // logical size in bytes
pub block_size: u32, // typically 4096
pub chunk_size: u64, // typically 4 MiB
pub provisioning: Provisioning, // Thin | Thick
pub state: VolumeState, // Available | InUse | Snapshotting
pub snapshots: Vec<SnapshotId>,
pub created_at: DateTime,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub enum Provisioning {
Thin, // allocate chunks on first write
Thick, // pre-allocate all chunks
}
/// Extent map: maps volume LBA ranges to storage chunks
/// Uses a B-tree for O(log N) lookup of any LBA
///
/// For thin provisioning: entries only exist for written regions
/// Unwritten regions return zeroes
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct ExtentMap {
pub volume_id: VolumeId,
pub chunk_size: u64,
/// BTreeMap<extent_index, ExtentEntry>
/// extent_index = lba_offset / chunk_size
pub extents: BTreeMap<u64, ExtentEntry>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct ExtentEntry {
pub chunk_id: ChunkId,
pub oss_node: NodeId,
pub version: u64, // for COW snapshots
pub dirty: bool, // has unflushed writes
}
/// COW snapshot: shares extent map with parent, diverges on write
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Snapshot {
pub id: SnapshotId,
pub parent: VolumeId,
pub created_at: DateTime,
pub extent_map: ExtentMap, // frozen copy of parent's map at snapshot time
}
D.2 Write Path with COW for Snapshots
impl ExtentMap {
/// Write to a volume with COW semantics
///
/// If this extent was inherited from a snapshot (version < current_version),
/// allocate a new chunk and copy-on-write
pub async fn write_block(
&mut self,
lba: u64,
data: &[u8],
chunk_service: &ChunkServiceImpl,
mds: &MetadataServiceClient,
) -> Result<()> {
let extent_idx = lba / self.chunk_size;
let offset_in_chunk = lba % self.chunk_size;
match self.extents.get(&extent_idx) {
Some(entry) if entry.version == self.current_version => {
// Same version: write in-place (no COW needed)
chunk_service.write_chunk(
&entry.chunk_id,
offset_in_chunk,
data,
ChunkMode::Replicated { factor: 3 },
).await?;
}
Some(entry) => {
// Inherited from snapshot: COW
// 1. Allocate new chunk
let new_chunk = mds.alloc_chunks(1).await?[0].clone();
// 2. Copy old content to new chunk
let old_data = chunk_service.read_chunk(
&entry.chunk_id, 0, self.chunk_size
).await?;
chunk_service.write_chunk(
&new_chunk.chunk_id, 0, &old_data,
ChunkMode::Replicated { factor: 3 },
).await?;
// 3. Apply the new write
chunk_service.write_chunk(
&new_chunk.chunk_id, offset_in_chunk, data,
ChunkMode::Replicated { factor: 3 },
).await?;
// 4. Update extent map
self.extents.insert(extent_idx, ExtentEntry {
chunk_id: new_chunk.chunk_id,
oss_node: new_chunk.oss_node,
version: self.current_version,
dirty: true,
});
}
None => {
// Thin provisioned: first write, allocate new chunk
let new_chunk = mds.alloc_chunks(1).await?[0].clone();
// Zero-fill then write (or use fallocate + punch hole)
chunk_service.write_chunk(
&new_chunk.chunk_id, offset_in_chunk, data,
ChunkMode::Replicated { factor: 3 },
).await?;
self.extents.insert(extent_idx, ExtentEntry {
chunk_id: new_chunk.chunk_id,
oss_node: new_chunk.oss_node,
version: self.current_version,
dirty: true,
});
}
}
Ok(())
}
}
D.3 Block Export: iSCSI Target
// crates/iscsi-target/src/lib.rs
/// Maps iSCSI SCSI commands to volume operations
pub struct RustfsIscsiTarget {
volumes: Arc<RwLock<HashMap<Lun, Volume>>>,
chunk_service: Arc<ChunkServiceImpl>,
mds: Arc<MetadataServiceClient>,
}
impl RustfsIscsiTarget {
/// Handle SCSI READ(16) — translate LBA to chunk read
async fn handle_read(&self, lun: Lun, lba: u64, transfer_len: u32) -> Result<Vec<u8>> {
let volume = self.volumes.read().await.get(&lun)?;
let byte_offset = lba * volume.block_size as u64;
let byte_length = transfer_len as u64 * volume.block_size as u64;
// Resolve extent map → chunk locations
let chunks = volume.extent_map.chunks_for_range(byte_offset, byte_length);
// Parallel read from chunk servers
let mut result = Vec::with_capacity(byte_length as usize);
for (chunk_idx, locations) in chunks {
let chunk_offset = if chunk_idx == byte_offset / volume.chunk_size {
byte_offset % volume.chunk_size
} else { 0 };
let data = self.chunk_service.read_chunk(
&locations[0].chunk_id, chunk_offset, /* length */
).await?;
result.extend_from_slice(&data);
}
Ok(result)
}
/// Handle SCSI WRITE(16) — translate LBA to chunk write
async fn handle_write(&self, lun: Lun, lba: u64, data: &[u8]) -> Result<()> {
let mut volume = self.volumes.write().await.get_mut(&lun)?;
volume.extent_map.write_block(
lba * volume.block_size as u64,
data,
&self.chunk_service,
&self.mds,
).await
}
/// Handle SYNCHRONIZE_CACHE — flush dirty chunks
async fn handle_flush(&self, lun: Lun) -> Result<()> {
let volume = self.volumes.read().await.get(&lun)?;
let dirty_chunks: Vec<_> = volume.extent_map.extents.values()
.filter(|e| e.dirty)
.collect();
for entry in dirty_chunks {
self.chunk_service.flush_chunk(&entry.chunk_id).await?;
}
Ok(())
}
}
Phase E: FUSE Client for POSIX
E.1 New Crate: crates/fuse-client/
// crates/fuse-client/src/lib.rs
use fuser::{Filesystem, Request, ReplyAttr, ReplyData, ReplyEntry, ReplyDirectory};
pub struct RustfsFuse {
mds: MetadataServiceClient, // gRPC to MDS
chunk_clients: Vec<ChunkServiceClient>, // gRPC to OSS nodes
open_files: RwLock<HashMap<u64, OpenFile>>, // fh → OpenFile
read_ahead: ReadAheadEngine,
lease_manager: ClientLeaseManager,
}
struct OpenFile {
inode: InodeId,
flags: u32, // O_RDONLY, O_WRONLY, O_RDWR, O_DIRECT
offset: AtomicU64, // current seek position
chunk_map: Arc<ChunkMap>, // cached from MDS
lease: Option<Lease>, // read or write lease
}
impl Filesystem for RustfsFuse {
fn lookup(&mut self, _req: &Request, parent: u64, name: &OsStr, reply: ReplyEntry) {
// MDS RPC: Lookup(parent_inode, name) → (child_inode, attr)
let resp = self.mds.lookup(parent, name.to_str().unwrap());
match resp {
Ok(r) => reply.entry(&TTL, &to_fuse_attr(&r.attr), r.generation),
Err(_) => reply.error(ENOENT),
}
}
fn read(
&mut self, _req: &Request, ino: u64, fh: u64,
offset: i64, size: u32, _flags: i32, _lock: Option<u64>,
reply: ReplyData,
) {
let file = self.open_files.read().get(&fh);
let chunk_map = &file.chunk_map;
// Determine which chunks cover [offset, offset+size)
let chunks = chunk_map.chunks_for_range(offset as u64, size as u64);
// Parallel fetch from OSS nodes (the KEY optimization for AI workloads)
let mut buf = vec![0u8; size as usize];
let fetches = chunks.iter().map(|(idx, locs)| {
let oss = &self.chunk_clients[locs[0].oss_node];
let chunk_offset = /* compute offset within this chunk */;
let chunk_length = /* compute how many bytes from this chunk */;
oss.read_chunk(locs[0].chunk_id, chunk_offset, chunk_length)
});
let results = futures::future::join_all(fetches);
// Assemble results into buf...
// Trigger read-ahead for next chunks
self.read_ahead.on_read(ino, offset as u64, size as u64, chunk_map);
reply.data(&buf);
}
fn write(
&mut self, _req: &Request, ino: u64, fh: u64,
offset: i64, data: &[u8], _write_flags: u32, _flags: i32,
_lock: Option<u64>, reply: fuser::ReplyWrite,
) {
let file = self.open_files.write().get_mut(&fh);
// Ensure we have a write lease
if file.lease.is_none() || !file.lease.as_ref().unwrap().is_write() {
file.lease = Some(self.mds.acquire_lease(ino, LeaseKind::WriteLease));
}
// Determine affected chunks
let chunk_map = &mut file.chunk_map;
let start_chunk = offset as u64 / chunk_map.chunk_size;
let end_chunk = (offset as u64 + data.len() as u64 - 1) / chunk_map.chunk_size;
for chunk_idx in start_chunk..=end_chunk {
let chunk_offset = /* offset within this chunk */;
let chunk_data = /* slice of data for this chunk */;
if let Some(locs) = chunk_map.chunks.get(&chunk_idx) {
// Existing chunk: write in-place (or COW if needed)
self.chunk_clients[locs[0].oss_node].write_chunk(
locs[0].chunk_id, chunk_offset, chunk_data
);
} else {
// New chunk: allocate from MDS, then write
let new_loc = self.mds.alloc_chunks(1);
self.chunk_clients[new_loc.oss_node].write_chunk(
new_loc.chunk_id, chunk_offset, chunk_data
);
chunk_map.chunks.insert(chunk_idx, vec![new_loc]);
}
}
// Update file size in MDS if extended
let new_size = (offset as u64 + data.len() as u64).max(file.attr.size);
if new_size > file.attr.size {
self.mds.set_attr(ino, SetAttrReq { size: Some(new_size), .. });
}
reply.written(data.len() as u32);
}
}
E.2 Read-Ahead Engine (Critical for AI Workloads)
// crates/fuse-client/src/readahead.rs
pub struct ReadAheadEngine {
/// Per-file tracking of access patterns
trackers: RwLock<HashMap<InodeId, AccessTracker>>,
/// Background prefetch task queue
prefetch_queue: mpsc::Sender<PrefetchTask>,
/// Maximum read-ahead window (default: 32 MiB, tunable)
max_readahead: u64,
}
struct AccessTracker {
last_offset: u64,
last_size: u64,
stride: Option<u64>, // detected stride pattern
sequential_count: u32, // consecutive sequential reads
readahead_window: u64, // current window size (grows/shrinks)
}
impl ReadAheadEngine {
pub fn on_read(&self, ino: InodeId, offset: u64, size: u64, chunk_map: &ChunkMap) {
let mut trackers = self.trackers.write();
let tracker = trackers.entry(ino).or_default();
// Detect access pattern
let is_sequential = offset == tracker.last_offset + tracker.last_size;
let is_strided = !is_sequential
&& tracker.stride.is_some()
&& offset == tracker.last_offset + tracker.stride.unwrap();
if is_sequential || is_strided {
tracker.sequential_count += 1;
// Exponentially grow read-ahead window (like Linux VFS)
tracker.readahead_window = (tracker.readahead_window * 2)
.min(self.max_readahead);
} else {
// Random access: shrink window
tracker.sequential_count = 0;
tracker.readahead_window = 0;
}
// Issue prefetch for the window ahead
if tracker.readahead_window > 0 {
let prefetch_offset = offset + size as u64;
let prefetch_length = tracker.readahead_window;
let chunks = chunk_map.chunks_for_range(prefetch_offset, prefetch_length);
for (idx, locs) in chunks {
self.prefetch_queue.send(PrefetchTask {
ino, chunk_idx: idx, locations: locs.to_vec(),
});
}
}
tracker.last_offset = offset;
tracker.last_size = size as u64;
if !is_sequential && tracker.sequential_count == 0 {
// Try to detect stride
let diff = offset.wrapping_sub(tracker.last_offset);
tracker.stride = Some(diff);
}
}
}
Why this matters for AI: PyTorch DataLoader workers read sequentially through their assigned shard, then jump to the next shard. The stride detector catches this pattern and prefetches the next shard's chunks before the reader gets there. With 4 MiB chunks and 32 MiB read-ahead, each worker has 8 chunks pre-fetched — enough to keep GPU feeding continuous.
Phase F: S3-POSIX Unified Namespace (Bidirectional Access)
The S3 API and POSIX mount should see the same data:
S3 bucket "training-data" ←→ POSIX mount at /mnt/rustfs/training-data/
S3 PUT training-data/imagenet/shard-0001.tar
↓
MDS: create inode, allocate chunks, build chunk map
↓
POSIX: ls /mnt/rustfs/training-data/imagenet/ → shows shard-0001.tar
POSIX: echo "hello" > /mnt/rustfs/training-data/notes.txt
↓
MDS: create inode, allocate chunk, write data
↓
S3 GET training-data/notes.txt → returns "hello"
Implementation: unified metadata path
// Modify existing S3 PUT handler to go through MDS
// BEFORE (current):
// S3 PUT → app/object_usecase → storage/ecfs → ecstore → disk
// AFTER (unified):
// S3 PUT → app/object_usecase → MDS (create inode + chunk map)
// → ecstore-chunk (write chunk data)
// → MDS (commit chunks, update size)
// Modify existing S3 GET handler similarly:
// S3 GET → app/object_usecase → MDS (resolve path → inode → chunk map)
// → ecstore-chunk (read chunks)
The existing ecstore/store_api.rs ObjectInfo struct gets a new optional field:
struct ObjectInfo {
// ... existing S3 fields ...
/// If present, this object is backed by the POSIX metadata system.
/// The inode contains authoritative size, timestamps, and chunk layout.
/// S3 metadata fields above are derived from the inode for compatibility.
posix_inode: Option<InodeId>,
}
Part 4 — Summary: Change Inventory
New Crates
| Crate | Est. Lines | Purpose |
|---|---|---|
crates/inode/ | 2-3K | Inode types, ChunkMap, InodeAttr |
crates/dentry/ | 1-2K | Directory entry types and operations |
crates/mds/ | 15-20K | Metadata server (RocksDB + Raft + gRPC service) |
crates/ecstore-chunk/ | 5-8K | Chunk-level gRPC service on storage nodes |
crates/block-volume/ | 8-12K | Volume manager, extent map, COW, journal |
crates/iscsi-target/ | 6-10K | iSCSI SCSI command handler |
crates/nbd-server/ | 2-3K | NBD protocol server (simplest block export) |
crates/fuse-client/ | 10-15K | FUSE client with read-ahead, leases, caching |
Modified Crates
| Crate | Change | Impact |
|---|---|---|
crates/lock/ | Add PosixLockManager, byte-range locks, leases | +3-4K lines |
crates/ecstore/ | Decompose into 5-6 sub-crates; add chunk-mode writes | Major refactor |
crates/filemeta/ | Add InodeId bridge to ObjectInfo; extend xl.meta format | +1-2K lines |
crates/protos/ | Add MDS and ChunkService protobuf definitions | +500 lines |
rustfs/src/storage/ecfs.rs | Route through MDS for unified namespace | ~500 line modification |
rustfs/src/app/ | Add POSIX-aware object use cases | +1-2K lines |
New Binaries
| Binary | Purpose |
|---|---|
rustfs-mds | Standalone metadata server |
rustfs-fuse | POSIX FUSE client |
rustfs-iscsi | iSCSI target daemon |
rustfs-nbd | NBD server daemon |
Unchanged (Fully Reused)
| Crate | Why |
|---|---|
io-core | Zero-copy I/O, buffer pool, direct I/O — used by all new paths |
rio | Encrypt → compress → hash pipeline — used for chunk data pipeline |
crypto / kms | Encryption at rest works for POSIX and block data too |
iam / policy | Access control applies to POSIX operations via MDS |
obs / metrics | Observability covers all new subsystems |
heal / scanner | Data integrity applies to chunks regardless of access method |
notify / audit | Event notifications for POSIX operations |
Part 5 — Implementation Priority
IMMEDIATE (enables AI POSIX workload):
inode/ + dentry/ + mds/ (RocksDB single-node) + fuse-client/
+ chunk-level read in ecstore (reuse existing RPC offset+length)
+ read-ahead engine
Result: mount, read training data with fseek/fread from many GPUs
NEXT (enables block storage):
block-volume/ + nbd-server/ (simplest block export)
+ COW snapshots + extent map
+ lock crate extensions (byte-range)
Result: boot a VM from RustFS block volume
THEN (production hardening):
HA MDS (Raft) + MDS sharding (DNE)
+ iscsi-target/ (production VMs)
+ ecstore decomposition (the big refactor)
+ unified S3-POSIX namespace
+ RDMA transport for chunk reads
LATER (performance parity with Lustre):
Kernel FUSE bypass (io_uring passthrough or kernel module)
+ GPUDirect Storage integration
+ NVMe-oF target
+ client-side write-back caching with lease coherency