This documentation is part of the "Projects with Books" initiative at zenOSmosis.
The source code for this project is available on GitHub.
Concurrency and Thread Safety
Loading…
Concurrency and Thread Safety
Relevant source files
- src/main.rs
- src/utils/format_bytes.rs
- tests/concurrency_tests.rs
- tests/mmap_and_zero_copy_tests.rs
- tests/parallel_iterator_tests.rs
The DataStore engine is designed for high-performance concurrent access, utilizing a layered locking strategy that balances data integrity with low-latency reads. By combining atomic offsets, read-write locks for metadata, and mutex-protected memory remapping, the engine allows multiple threads to read and write simultaneously with minimal contention.
Layered Locking Strategy
The core of the thread safety model lies in DataStore, which manages several synchronized components to ensure that append-only operations do not corrupt the in-memory index or the memory-mapped view used by readers.
Architecture Overview
| Component | Type | Purpose |
|---|---|---|
file | Arc<RwLock<BufWriter<File>>> | Synchronizes disk appends to the storage file. src/storage_engine/data_store.rs28 |
key_indexer | Arc<RwLock<KeyIndexer>> | Protects the hash-to-offset mapping for lookups. src/storage_engine/data_store.rs31 |
mmap | Arc<Mutex<Arc<Mmap>>> | Ensures safe, atomic remapping of the memory-mapped file view. src/storage_engine/data_store.rs29 |
tail_offset | AtomicU64 | Provides a lock-free “high-water mark” for the valid file end. src/storage_engine/data_store.rs30 |
Sources: src/storage_engine/data_store.rs:27-33
Data Flow and Synchronization
The following diagram illustrates how different locks are acquired during a write operation to maintain consistency between the physical file, the memory map, and the index.
Write Operation Synchronization Flow
graph TD
subgraph "DataStore::write [src/storage_engine/data_store.rs]"
A["Start Write"] --> B["Acquire file.write()
RwLock"]
B --> C["Append Entry to Disk via BufWriter"]
C --> D["Acquire mmap Mutex"]
D --> E["Remap File to Memory (new Mmap)"]
E --> F["Update tail_offset (AtomicU64)"]
F --> G["Release mmap Mutex"]
G --> H["Acquire key_indexer.write()
RwLock"]
H --> I["Update KeyIndexer with Hash/Offset"]
I --> J["Release key_indexer lock"]
J --> K["Release file lock"]
end
Sources: src/storage_engine/data_store.rs:28-33 src/storage_engine/data_store.rs:176-215
Key Implementation Details
Atomic Tail Offset
The tail_offset is an AtomicU64 that tracks the end of the valid data chain. It is updated using Ordering::SeqCst (Sequential Consistency) in the re_map_and_update_index function to ensure that once a write is finalized, all subsequent readers see the updated file size immediately. This prevents readers from attempting to access data that hasn’t been fully committed to the memory map yet. src/storage_engine/data_store.rs30 src/storage_engine/data_store.rs209
Safe Remapping via Mutex
Because Mmap views in Rust are generally considered unsafe to change while references to them exist, DataStore wraps the Arc<Mmap> in a Mutex. When a write extends the file, re_map_and_update_index is called:
- It acquires the
Mutexto prevent concurrent remapping attempts. src/storage_engine/data_store.rs198 - It creates a new
Mmapspanning the new file size viaMmapOptions. src/storage_engine/data_store.rs:173-174 src/storage_engine/data_store.rs203 - It replaces the old
Arc<Mmap>with a new one. src/storage_engine/data_store.rs208
Existing EntryHandle objects hold their own Arc<Mmap>, meaning they can continue to read from the “old” memory segment even after the DataStore has remapped to a larger one, ensuring zero-copy read stability. src/storage_engine/data_store.rs29
KeyIndexer Concurrency
The KeyIndexer uses an RwLock to allow multiple concurrent readers to perform lookups (e.g., via read or exists) while blocking them only during the brief moment the index is updated after a successful write. src/storage_engine/data_store.rs31 src/storage_engine/data_store.rs:212-214
Parallel Processing and Rayon Integration
When the parallel feature is enabled, the engine integrates with Rayon to accelerate compute-intensive batch operations and iteration.
Parallel Entry Iteration
The par_iter_entries method (available via the parallel feature) allows for parallel processing of all entries in the store. This is highly effective for tasks like compact(), where the engine must check the “liveness” of every entry against the current index. tests/parallel_iterator_tests.rs:35-38
Code Entity Mapping: Parallel Iteration
Sources: src/storage_engine/data_store.rs:23-24 tests/parallel_iterator_tests.rs:3-6 tests/parallel_iterator_tests.rs:35-38
Batch Operations
The batch_write and batch_read functions utilize the internal locking mechanisms to perform bulk updates efficiently:
batch_write: Acquires thefilewrite lock once for the entire batch, significantly reducing the overhead of repeated lock acquisition and remapping. src/storage_engine/data_store.rs:253-270batch_read: Can be executed concurrently by multiple threads as it only requires read access to theKeyIndexerand theMmap. src/storage_engine/data_store.rs:313-328
Concurrency Testing
The robustness of this locking strategy is verified in tests/concurrency_tests.rs, which includes stress tests for various thread-safety scenarios:
concurrent_write_test: 16 threads writing unique keys simultaneously to ensure theRwLockandAtomicU64correctly serialize appends without data loss. tests/concurrency_tests.rs:113-161interleaved_read_write_test: Usestokio::sync::Notifyto coordinate threads, verifying that a reader sees a value immediately after a writer releases its locks. tests/concurrency_tests.rs:165-214concurrent_slow_streamed_write_test: Simulates high-latency IO (e.g., network streams) using aSlowReaderto ensure that long-runningwrite_streamoperations do not deadlock or corrupt the file structure. tests/concurrency_tests.rs:16-109
Thread Safety Entity Map
Sources: src/storage_engine/data_store.rs:28-33 tests/concurrency_tests.rs:165-214
Dismiss
Refresh this wiki
Enter email to refresh