Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

GitHub

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

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

ComponentTypePurpose
fileArc<RwLock<BufWriter<File>>>Synchronizes disk appends to the storage file. src/storage_engine/data_store.rs28
key_indexerArc<RwLock<KeyIndexer>>Protects the hash-to-offset mapping for lookups. src/storage_engine/data_store.rs31
mmapArc<Mutex<Arc<Mmap>>>Ensures safe, atomic remapping of the memory-mapped file view. src/storage_engine/data_store.rs29
tail_offsetAtomicU64Provides 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:

  1. It acquires the Mutex to prevent concurrent remapping attempts. src/storage_engine/data_store.rs198
  2. It creates a new Mmap spanning the new file size via MmapOptions. src/storage_engine/data_store.rs:173-174 src/storage_engine/data_store.rs203
  3. 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:

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 the RwLock and AtomicU64 correctly serialize appends without data loss. tests/concurrency_tests.rs:113-161
  • interleaved_read_write_test : Uses tokio::sync::Notify to coordinate threads, verifying that a reader sees a value immediately after a writer releases its locks. tests/concurrency_tests.rs:165-214
  • concurrent_slow_streamed_write_test : Simulates high-latency IO (e.g., network streams) using a SlowReader to ensure that long-running write_stream operations 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