EPIC-023: Documentation du modèle de concurrence pour utilisateurs avancés et contributeurs.
User-facing write throughput guidance: see
docs/guides/WRITE_CONCURRENCY.mdfor the single-writer-per-collection model, batching patterns, and the Community/Enterprise split. This document covers the internal lock ordering and concurrency primitives used across the engine.
VelesDB utilise un modèle de concurrence basé sur:
- Sharding: Partitionnement des données pour réduire la contention
- RwLock: Lecture parallèle, écriture exclusive (parking_lot)
- Lock-free atomics: For compteurs, métriques, HNSW entry-point reads, and CSR snapshot swap
- ArcSwap: Lock-free CSR snapshot reads for graph traversal (zero contention on reads)
- Lock ordering: Ordre déterministe pour prévenir les deadlocks
Everything below this section describes concurrency inside one process —
because one process is all there can be. Database::open_impl
(crates/velesdb-core/src/database/mod.rs) takes an exclusive flock on
<data_dir>/velesdb.lock at open and holds it for the Database's entire
lifetime (the RAII _lock_file field releases it on drop, including on
crash, via the OS). A second process — another daemon, an embedded binding,
a CLI pointed at the same directory — fails at open with
Error::DatabaseLocked before it can reach any read, write, or
read-modify-write. The daemon surfaces this as actionable stderr guidance
after a bounded retry (velesdb-memory's daemon_startup.rs).
This is the invariant the rest of the model leans on. Notably, the
working-context index in velesdb-memory needs no cross-process
compare-and-swap precisely because no second process can hold the store
(#1958 was closed on that proof). Two real-process tests pin it:
crates/velesdb-memory/tests/http_lock_contention.rs— a second HTTP daemon on a held store exits non-zero with the lock guidance;crates/velesdb-memory/tests/working_index_two_daemons_process.rs— a contender fails fast mid-save-loop, and no index entry is lost across contention or the process handoff.
If single-writer-at-open is ever relaxed, those tests fail first, and every "intra-process is all there is" claim in this document must be re-derived.
Known limits: flock is advisory (a process bypassing Database::open
is not stopped by it), and its semantics on network filesystems (NFS, some
FUSE mounts) are weaker than on local disks — the same caveat any
flock-based scheme carries. Deleting velesdb.lock while a holder is live
breaks the exclusion for future openers; nothing in-tree does this.
┌─────────────────────────────────────────────────────────────────┐
│ ConcurrentEdgeStore │
├─────────┬─────────┬─────────┬─────────┬─────────┬───────────────┤
│ Shard 0 │ Shard 1 │ Shard 2 │ Shard 3 │ ... │ Shard N-1 │
│ RwLock │ RwLock │ RwLock │ RwLock │ │ RwLock │
└─────────┴─────────┴─────────┴─────────┴─────────┴───────────────┘
│
Shard = hash(node_id) % num_shards
Default shards: 256 (configurable via with_shards() or with_estimated_edges())
Shard selection:
- Small graphs (< 1K edges): 1 shard
- Medium graphs (1K-64K): 16-64 shards
- Large graphs (64K-1M): 64-128 shards
- Very large graphs (> 1M): 256 shards
| Component | Lock Type | Contention | Notes |
|---|---|---|---|
| EdgeStore shards | parking_lot::RwLock |
Low | Per-shard, fine-grained |
| HNSW layers | parking_lot::RwLock |
Medium | Global, read-heavy |
| HNSW neighbors | parking_lot::RwLock |
Medium | Per-node |
| PropertyIndex | parking_lot::RwLock |
Low | Per-property |
| HNSW entry point | AtomicUsize + promotion Mutex |
None on reads | Lock-free reads; the first claim and each promotion write under the promotion lock |
| HNSW max layer | AtomicUsize + promotion Mutex |
None on reads | Moves with the entry point, under the same lock |
| Metrics counters | AtomicU64 |
None | Lock-free |
| Edge ID registry | RwLock<HashMap> |
Low | Global, for existence checks |
| CsrSnapshot | ArcSwap<Arc<CsrSnapshot>> |
None | Lock-free reads via atomic swap; lazy rebuild on dirty flag |
| Quantizer (RaBitQ/SQ8) | parking_lot::RwLock |
None (after training) | Write-once then read-only |
| Code store (RaBitQ/SQ8) | parking_lot::RwLock |
Low | Write per insert (held for one append) |
| Training buffer (RaBitQ/SQ8) | parking_lot::Mutex |
Low | Pre-training only |
| MmapStorage (compaction) | parking_lot::RwLock |
High (during compaction) | Exclusive write lock for full compaction duration |
These types are safe to share across threads and can be moved between threads:
// Safe to share and send
Collection: Send + Sync
HnswIndex: Send + Sync
ConcurrentEdgeStore: Send + Sync
ConcurrentNodeStore: Send + Sync
Database: Send + SyncThese types contain non-thread-safe internal state:
// Must stay on creation thread
GraphTraversal: !Send // Contains references
QueryCursor: !Send // Iterator state
BfsIterator: !Send // Traversal stateVelesDB uses compile-time assertions to verify thread safety:
// In ConcurrentEdgeStore
const _: () = {
const fn assert_send_sync<T: Send + Sync>() {}
assert_send_sync::<ConcurrentEdgeStore>();
};When multiple locks are needed, acquire them in this order:
1. edge_ids (global registry)
2. shards[0]
3. shards[1]
4. ...
5. shards[N-1]
The rank ordinals below mirror the typed LockRank newtype defined in
crates/velesdb-core/src/lock_rank.rs, which derives a total ordering over its
ordinals. That newtype and its public assert_lock_order helper are a typed
vocabulary for the intended order, not a wired runtime check:
assert_lock_order has zero production call sites — it is exercised only by
lock_rank_tests.rs and tests/concurrency_lock_order.rs. No production
binary, debug or release, calls it, so it never fires on a real acquisition.
Enforcement in practice, by build and by tier:
- Production (release) binaries have no runtime lock-order check at all. This is zero-overhead by design — no thread-local stack and no atomic on the hot search path.
- The HNSW tier has a debug-only, warn-only, partial tracker — a
private mechanism separate from
assert_lock_order. Therecord_lock_acquirefunction incrates/velesdb-core/src/index/hnsw/native/graph/locking.rsis#[cfg(debug_assertions)]; on an out-of-order acquisition it increments an atomic violation counter and emits atracing::warn!— it never panics. It is also partial: only theGpuVectorsSnapshot,EntryPointPromotion,VectorsandLayersranks are ever recorded.Neighborsis#[allow(dead_code)]with norecord_lock_acquirecall site, so 1 of the 5 core ranks is untracked even in debug. - The collection tier —
Collection's own field order (config,vector_storage,payload_storage, ...; the=== LOCK ORDERING ===block incrates/velesdb-core/src/collection/types.rs, and the "Collection-level lock order" section below) — is convention-only, backed by no assertion of any kind. This is the tier whose vector/payload pair actually deadlocked (2026-07); its lock-order correctness rests on code review plus the regression tests, not on any runtime mechanism.
The ordinal table below is the human-readable record of the intended order; it
is not enforced by a compiled assertion at acquisition time. Two LockRank
types carry it in code, with distinct roles (#2013): the public registry
(crate::lock_rank::LockRank) owns the numbers and the premium reservation,
and the private mechanism (HnswLockRank,
index/hnsw/native/graph/locking.rs) is what the debug tracker actually
records. Their lock-step is not a promise: the private enum's discriminants
are defined from the registry's constants, so any divergence is a compile
error.
For HNSW index operations that touch the GPU snapshot cache, the entry-point promotion lock, vector storage, graph layers, and neighbor lists, the global lock acquisition order is:
gpu_vectors_snapshot (rank 5) → promotion (rank 8) → vectors (rank 10)
→ layers (rank 20) → neighbors (rank 30)
(Rank 15 was the PDX block-columnar layout, removed with the unwired PDX machinery — the ordinal is retired, not reassigned.)
| Lock | Rank | Registry constant (= HnswLockRank discriminant) |
Component | Notes |
|---|---|---|---|---|
gpu_vectors_snapshot |
5 | LockRank::GPU_VECTORS_SNAPSHOT |
GPU flat-vector snapshot cache (Mutex) |
Acquired before vectors in the GPU path (gpu feature); writers release vectors before reacquiring it to invalidate |
promotion |
8 | LockRank::ENTRY_POINT_PROMOTION |
HNSW entry-point moves (Mutex<()>) |
Taken by the first claim, each promotion and reorder_for_locality's renumbering; the anchor reparenting it guards then takes vectors, layers and list locks, one at a time |
vectors |
10 | LockRank::VECTORS |
ContiguousVectors (single vector store since PERF1) |
Acquired first among the core HNSW locks in upsert and search paths |
layers |
20 | LockRank::LAYERS |
HNSW layer structure (RwLock) |
Global graph topology |
neighbors |
30 | LockRank::NEIGHBORS |
Per-node neighbor lists (RwLock) |
Fine-grained, acquired last |
The base layer's per-node reachability anchors (AtomicU32, #2259) are no
lock of their own: an anchor is written under the list lock of the node it
names (rank 30), with the layers write lock held (reorder_for_locality), or,
when a loaded graph rebuilds them, under the layers read lock and each
parent's list read lock while nothing else runs. Clearing a promoted entry
point's anchor, under the promotion lock, only unprotects an edge.
Rule: Never acquire a lower-rank lock while holding a higher-rank lock.
For example, acquiring vectors while holding neighbors is forbidden. The
typed assert_lock_order(previously_held, about_to_acquire) helper expresses
this rule but is not wired into any acquisition path (see the enforcement
note above). In debug builds the HNSW tier's record_lock_acquire will warn
(never panic) on a violation — and only among the tracked ranks
(GpuVectorsSnapshot / EntryPointPromotion / Vectors / Layers).
Lock order is not enough for HnswIndex::inner. Ordering rules keep two
threads from taking the same pair of locks in opposite directions; this one
keeps a single lock from deadlocking against itself through a thread pool:
A thread holding
HnswIndex::innernever waits on the global rayon pool.
parking_lot's RwLock is deliberately not read-re-entrant under a waiting
writer: a pending write() blocks every new read(), so that a stream of
readers cannot starve a writer. That fairness is what closes the cycle of
#2343 when the rule above is broken:
- a
vacuumasks forinner.write()and waits for the current readers; - a batch operation holds
inner.read()and waits for its rayon jobs; - those jobs run on the global pool, where the per-query search and
rerank_candidates_simdtakeinner.read()— new readers, so the pending writer parks them; - the parked workers are the ones step 2 is waiting for.
Nobody advances, and no ordering discipline would have prevented it: only one lock is involved.
Both operations that hold the guard and join on rayon therefore submit to a
dedicated pool, graph_pool in index/hnsw/index/batch.rs:
HnswIndex::insert_batch_parallel's place phase and HnswIndex::link_placed's
connect phase. Their jobs reach the graph's arena, layers and entry point and
never HnswIndex itself, so no job on that pool can take inner; rayon does
not steal work across pools, so the isolation is structural rather than a
timing assumption.
Two checks hold the rule, because neither is sufficient alone:
crates/velesdb-core/src/index/hnsw/index/global_pool_tests.rsparks every global rayon worker, then requires each guard holder to finish anyway. It is deterministic: an operation that joins on the global pool cannot finish, whatever the machine's timing. Before the fix it timed out at 60 s; after it, both holders finish in 0.05 s.scripts/check_hnsw_rayon_pool.py(CI joblint) refuses a new rayon submission in any non-test file underindex/hnsw/that is neither inside an.install(closure — on any dedicated pool, not onlygraph_pool— nor listed with the reason it cannot close the cycle. It cannot decide by itself whether a submission runs under a held guard — the one that deadlocked reachedpar_iterthree calls down, in another file — so it requires the reason to be written rather than inferred.
A wall-clock bound on a racing test would not do instead: under heavy load a batch that is merely slow and one that is deadlocked print the same thing. Telling them apart needs the stack, not the clock — which is why the rule is held by a structural test and a guard, not by a timeout.
The inclusive ordinal range [40, 59] is reserved for premium-owned lock
classes (cluster state, tenant store, server-level locks). Core never
assigns a rank at or above 40, leaving that band exclusively for premium so
it can order its own locks relative to core without collision. The bounds are
exposed as LockRank::PREMIUM_MIN (40) and LockRank::PREMIUM_MAX (59), and
LockRank::premium(value) constructs a premium rank, returning None for any
value outside the reserved range. Because all core ranks (5–30) sit strictly
below the premium band, the union of core and premium ranks forms one
authoritative total ordering shared across both engines: premium locks are
always acquired after the core locks whose data they wrap.
| Range | Owner | Ordinals |
|---|---|---|
| Core | velesdb-core |
5, 8, 10, 20, 30 (15 retired) |
| Premium (reserved) | velesdb-premium |
40–59 inclusive |
When an edge spans two shards (source in shard A, target in shard B):
// ✅ CORRECT: Ascending order
let (first_idx, second_idx) = if source_shard < target_shard {
(source_shard, target_shard)
} else {
(target_shard, source_shard)
};
let mut first = shards[first_idx].write();
let mut second = shards[second_idx].write();// ❌ WRONG: May cause deadlock
let mut source = shards[source_shard].write();
let mut target = shards[target_shard].write(); // DEADLOCK if another thread holds target first!Deleting a graph node now cascades to its edges: the node-delete path
(collection/core/crud_read_delete.rs) calls
ConcurrentEdgeStore::remove_node_edges(id) so both outgoing and incoming
edges are removed, leaving no dangling edges pointing at a phantom node (#900).
The cascade follows the global lock order: it acquires edge_ids first,
then the affected shards in ascending index order, using a BTreeSet whose
iteration is already sorted:
let mut shards_to_clean: BTreeSet<usize> = BTreeSet::new();
// BTreeSet iteration is already sorted ascending
for &idx in &shards_to_clean {
guards.push(shards[idx].write());
}The snapshot invalidation / debounced rebuild (see below) runs after the
edge_ids write lock is released, never while holding it, to avoid deadlocking
against the downstream rebuild lock acquisition.
Separate from the HNSW LockRank total ordering above, Collection's own
fields (config, vector_storage, payload_storage, the quantization
caches, edge_wal_lock, the property/range/label indexes, ...) follow a
documentation-enforced ascending order recorded as a plain comment block at
the top of crates/velesdb-core/src/collection/types.rs (=== LOCK ORDERING ===). It is not backed by a typed LockRank-style assertion —
this section only narrates one addition to it: position 3b,
edge_wal_lock, sitting right after payload_storage (3).
edge_wal_lock has two acquisition patterns:
add_edge/add_edges_batch(collection/core/graph_api.rs) holdpayload_storage's read guard from the endpoint referential-integrity check through the end of the edge write, and acquireedge_wal_lockwhile that read guard is still held — so the acquisition order for this path ispayload_storage(3) → edge_wal_lock(3b).remove_edgeand the delete-cascade path (collection/core/crud_read_delete.rs'scascade_delete_node_edges) acquireedge_wal_lockalone, with no other collection lock held — the delete path has already released itspayload_storagewrite guard by the time it reaches the cascade.
This closes a race left open by the original #1442 fix: add_edge used to
release payload_storage's read guard immediately after checking that both
endpoints exist, then separately acquire edge_wal_lock to write the edge.
A concurrent delete() of an endpoint could land in that window — its own
write guard acquisition would not contend with anything, since the reader
had already let go — and finish removing the node (and, if the edge
happened to already exist, cascading its removal) before add_edge ever
wrote it. The result was a "phantom" edge: present in the edge store,
invisible to all_node_ids()/MATCH, since both resolve nodes from the
payload store.
Holding payload_storage(3) across edge_wal_lock(3b) fixes this: a
concurrent delete() needs payload_storage's write guard, which
cannot be acquired while add_edge holds the read guard, so the delete
blocks until the edge is fully durable (WAL + edge-store apply). By the
time the delete proceeds, the edge already exists, so its own cascade sees
and removes it — no phantom, no compensation logic, no rollback. The same
guard is held once for the whole batch in add_edges_batch, so "an
endpoint disappears mid-batch" is impossible by construction.
The order stays acyclic: neither remove_edge nor the delete cascade ever
acquires payload_storage while holding edge_wal_lock, so there is no
path back from 3b to 3.
Residual latency: writers to payload_storage (upserts, deletes) can
now stall behind an in-flight edge write for up to one fsync (the edge
WAL append) — typically 0.05–5 ms on an SSD, consistent with the
single-writer-per-collection model this section belongs under (see
guides/WRITE_CONCURRENCY.md).
Known limitation (accepted): this only protects edges written through
Collection::add_edge/add_edges_batch. Edges loaded from a pre-existing
WAL/snapshot at Collection::open (replay) are trusted as-is and never
re-validated — replay intentionally bypasses referential-integrity
validation so legitimate edge-only databases created before the #1442 fix
keep their data. velesdb-cli graph doctor <collection> audits/repairs
such legacy phantom edges as an explicit, opt-in tool (read-only report by
default, --purge/--stub to fix); see
#1469 and
guides/GRAPH_PATTERNS.md.
The codec-generic QuantizedPrecisionHnsw<D, C> (RaBitQPrecisionHnsw and
Sq8PrecisionHnsw are aliases over it) uses interior mutability for its
trained quantizer, encoded vector store, and pre-training buffer — the
contract below is written once in quantized_precision.rs and holds for
every codec:
| Field | Type | Access Pattern |
|---|---|---|
quantizer |
RwLock<Option<Arc<C::Quantizer>>> |
Write-locked once during training, then read-only |
store |
RwLock<Option<C::Store>> |
Write during insert (push encoded vector) |
training_buffer |
Mutex<Vec<Vec<f32>>> |
Write during pre-training inserts |
install_gate |
RwLock<()> |
Read across each insert; write for training/install |
During train_codec(), locks are acquired and released in this order:
install_gate.write() → quantizer.write() → store.write() → training_buffer.lock()
Locks are released between acquisitions (not held simultaneously). The
training function uses a double-check locking pattern: it first checks
quantizer under a read lock, and if training is needed, re-checks
under a write lock to prevent duplicate training from concurrent threads.
install_trained_quantizer() (quantizer restore at open / TRAIN QUANTIZER live install) follows the same quantizer → store → training_buffer order. It holds quantizer.write() for the whole
re-encode so concurrent inserts cannot interleave store pushes, and it
reads the inner.vectors snapshot and releases it before taking
store.write() — never waiting on the store lock while holding
vectors (a search thread holds store.read() while acquiring
inner.vectors.read()). The install gate (write) is taken first: every
in-flight insert holds it for read across its whole body, so the rebuild's
snapshot can never miss an insert that already passed the trained check —
the store is positional (entry N = node N) and one missed push would shift
every subsequent encoding onto the wrong node.
When training completes, the store is set BEFORE the quantizer:
store.write() ← set encoded vectors
quantizer.write() ← set trained quantizer
This ordering prevents an inconsistent snapshot where a search thread sees
a trained quantizer but an empty store. A search thread that acquires
quantizer.read() and sees Some(...) is guaranteed that
store.read() also contains Some(...) with all pre-training
vectors already encoded.
| Pattern | Impact | Acceptable? |
|---|---|---|
rabitq_store.write() serializes post-training inserts |
One short hold per push (a single encoded-vector append) | Yes: store push is a trivial Vec::push operation |
| Training blocks all inserts while it runs | One-time event per index lifetime | Yes: training runs once when the buffer threshold is reached |
reorder_for_locality() is offline-only |
Takes &self but must not run during concurrent search |
Yes: only called during explicit maintenance, not on the hot path |
- Post-training inserts hold
rabitq_store.write()for the minimum duration needed to push a single encoded vector (oneVec::push). - Training is amortized: it runs once when the training buffer reaches its threshold, then never again for the lifetime of the index.
reorder_for_locality()is documented as offline-only and is not exposed through any concurrent API path.
Readers load the entry point and the max layer, AtomicUsize fields of
NativeHnsw, without a lock:
entry_point: AtomicUsize, // NO_ENTRY_POINT (usize::MAX) when empty
max_layer: AtomicUsize, // Current maximum layerWriters take the promotion lock (promotion: Mutex<()>, rank 8): the two
cases below, and the renumbering reorder_for_locality applies.
- Empty index:
claim_first_entry_pointCASesentry_pointfromNO_ENTRY_POINTto the new node ID and setsmax_layer, under the lock; only one thread wins. A graph that already has an entry point returns before taking the lock. A single insert that saw an empty graph but loses the claim connects through the winner, and a batch consumes its first node only when that node won (#2259). - Layer promotion:
promote_entry_pointreturns before taking the lock when the node's layer does not exceedmax_layer, the common case. Under the lock it checks again, stores the newmax_layer, swaps in the new entry point (swap,AcqRel) and callsreparent_entry_point, which, under the vectors and layers read locks, clears the new root's anchor and links the previous entry point into the tree below it: into its list, or one below it when that list is full of protected entries. Two promotions that overlapped without the lock could reparent in the wrong order, the later one's clearing the anchor the earlier one had just given its node, and strand the tree below it (#2259).
Transient inconsistency window: Between the max_layer store and the
swap, a concurrent reader may see the new max_layer with the old
entry_point. This is safe: search_layer_single
finds no neighbour of the old EP at a layer where it has no edges, and
returns it unchanged: a no-op greedy descent.
Entry-point promotion is rare (O(log_M(N)) times per index lifetime), so the lock is almost never contended.
EdgeStore and ConcurrentEdgeStore maintain an optional CsrSnapshot --
a Compressed Sparse Row representation of outgoing edges for zero-copy
neighbor access during BFS/DFS traversal.
Lifecycle:
- Build:
build_read_snapshot()materializes edges into contiguous arrays (targets: Vec<u64>,edge_ids: Vec<u64>) with aFxHashMap<u64, (offset, len)>index. Auto-built after loading from disk and afterflush(). - Read:
csr_snapshot().neighbors(source_id)returns&[u64]-- a zero-copy slice into the contiguous target array. - Invalidate: Every write operation (
add_edge,remove_edge,remove_node_edges) setscsr_snapshot = Noneand increments a pending-write counter. Subsequent reads fall back to per-shard edge lookup until the snapshot is rebuilt.
Debounced rebuild (#905): The actual O(N+E) CSR rebuild (which clones every
edge into a fresh EdgeStore) is debounced rather than run on the first
read after any write. A mutation only flips the dirty flag and bumps
pending_writes; the rebuild is deferred until the accumulated write count
reaches CSR_REBUILD_WRITE_THRESHOLD (64). Batch writes count their full size
toward the threshold. While dirty-but-below-threshold, reads remain correct by
falling back to per-shard lookup — debouncing trades a slightly slower fallback
read for avoiding a full rebuild on every interleaved read/write. A completed
rebuild clears the dirty flag and resets the counter to 0.
Thread safety (ConcurrentEdgeStore): The read snapshot is not stored
under the shard RwLocks. It lives in a lock-free ArcSwap<CsrSnapshot>
(csr_snapshot) paired with an AtomicBool dirty flag (csr_dirty) and an
AtomicU64 write counter (pending_writes). Readers load() the current
Arc<CsrSnapshot> with zero contention. Every mutation (add_edge,
remove_edge, remove_node_edges) sets csr_dirty and bumps pending_writes;
the rebuild is deferred (issue #905 debounce) and stores a fresh snapshot into
the ArcSwap. Two rebuild paths clear the dirty flag under different
guarantees, and both are correct:
build_read_snapshot()(snapshot.rs) holdsedge_ids.read()for the whole build and across the flag clear — it walks every edge under that one read guard,stores the snapshot, and only afterwards resetspending_writesto 0 and clearscsr_dirty, all while still holding the guard. No writer can interleave, so a blanket reset is sound.ensure_csr_fresh()(query.rs) holds no such lock. It therefore clears the dirty flag first (csr_dirty.swap(false, AcqRel)), snapshots the observedpending_writesbefore reading the shards, rebuilds, then subtracts only that observed count (fetch_updatewithsaturating_sub). A concurrentfetch_addlanding during the rebuild is preserved rather than clobbered, so the next reader still observes the snapshot as due for rebuild.
The deferred rebuild path acquires edge_ids read-only (never write) and
releases per-shard read locks promptly, so it never violates the
edge_ids → shards ordering and the caller must not hold an edge_ids write
lock across it.
RwLockallows multiple concurrent readers- Sharding distributes reads across independent locks
- Recommendation: Default 256 shards optimal for most workloads
- Writers block readers on same shard
- Cross-shard writes require 2 locks
- Recommendation:
- Use batch inserts to amortize lock overhead
- Consider
with_estimated_edges()to optimize shard count
- Uses "Read-Copy-Drop" pattern to minimize lock duration:
// ✅ CORRECT: Copy data, drop lock immediately
let neighbors: Vec<u64> = {
let guard = shard.read();
guard.get_outgoing(node).iter().map(|e| e.target()).collect()
}; // Guard dropped here
for neighbor in neighbors {
// Process without holding lock
}-
CsrSnapshot fast-path: When a CSR read snapshot is available (built after load or after
build_read_snapshot()), BFS/DFS reads neighbors via a contiguous&[u64]slice instead of per-shard edge lookup. Falls back to the shard-based path when the snapshot is invalidated by writes. -
Parent-pointer path reconstruction: BFS traversal uses a
FxHashMap<u64, (u64, u64)>parent-pointer map instead of cloning path vectors at every edge expansion. Paths are reconstructed on-demand viareconstruct_path()only when a result is emitted, avoiding O(depth) allocations per expansion step.
-
Cross-shard operations hold multiple locks:
- Edge spanning 2 shards requires 2 locks + edge_ids lock
- Mitigation: Lock ordering prevents deadlocks
-
Large traversals can block writers:
- BFS/DFS with many nodes may hold locks longer
- Mitigation: Read-Copy-Drop pattern releases locks quickly
-
An HNSW vacuum holds the index write lock for its swap (#2262):
- The rebuild runs beside searches and writes, which carry on against the old graph. The writes made meanwhile are then copied into the new graph without the write lock, in at most four rounds, each taking the writes made during the one before, until 64 or fewer are left or four have run
- Under the write lock, searches and writes wait while the swap copies, one insert at a time, every id mapped but not yet carried when the lock is granted, rebuilds the mapping of every live id and drops the old graph. The rounds and their threshold say when the catch-up stops trying, not how much is left: a write in flight — a whole batch — maps its ids after the last round looked and before the lock is granted, so what the swap copies has no upper bound. #2335 tracks bounding it
- Nothing runs on rayon under that write lock: rayon workers running batch searches park on the index lock, so a parallel insert waiting for them there never finishes. The swap once did, and the vacuum hung with every search
reorder_for_localityand a running vacuum wait for each other: one maintenance lock, which writers and saves never take, serializes the two. A save made during the rebuild saves the old graph; the swap waits for its dump like for any read guard- Saves of one index wait for each other on a lock of their own, never on the maintenance lock: two saves into one directory rewrote the same files under one generation
- The swap holds the index write lock for the mapping rebuild only, never for a copy: a publishing lock, taken shared by every mapping publication and exclusively by the vacuum for its last catch-up, freezes the remainder so it is settled outside the graph guard. Taken before the index lock, always, so the orders cannot cross; searches never take it. Unsealed, a write landing between the last catch-up round and the grant of the write guard was copied under that guard with no upper bound — measured at a whole 8 000-id working set in one swap (#2335)
- The rebuild uses the index's own M,
ef_constructionand alpha, read from its graph, not the defaults for its dimension - A save holds the graph's vector read lock across its vectors file and its graph file, so an insert waits for both to be written: released between the two, a node pushed in between and linked into saved nodes left a save the next load refused
- Mitigation: vacuum while writes are light; incremental updates remain preferred over a full rebuild
-
No transactional semantics:
- Operations are atomic per-operation, not per-batch
- Mitigation: Use flush() for durability checkpoints
-
Enlarged crash recovery window during batch upsert:
- The 3-phase upsert pipeline (
batch_store_all->per_point_updates->bulk_index_or_defer) writes vectors and payloads to storage before inserting into the HNSW graph. A crash between Phase 1 and Phase 3 leaves vectors in storage but missing from the HNSW index. - Mitigation: On
Collection::open(), gap detection comparesstorage.ids()againstindex.mappingsand re-indexes any missing vectors. See HNSW Crash Recovery for the full recovery architecture and SOUNDNESS.md for the slot allocation invariants.
- The 3-phase upsert pipeline (
-
Dimension shards appropriately:
// For 100K edges let store = ConcurrentEdgeStore::with_estimated_edges(100_000);
-
Prefer batch operations:
// ✅ Better: One lock acquisition collection.upsert(vec![point1, point2, point3])?; // ❌ Worse: Three lock acquisitions collection.upsert(vec![point1])?; collection.upsert(vec![point2])?; collection.upsert(vec![point3])?;
-
Limit traversal depth:
// Always specify max_depth to prevent runaway traversals let nodes = store.traverse_bfs(start, 5); // Max 5 hops
-
Follow lock ordering strictly:
- Document lock order in new concurrent structures
- Use BTreeSet/BTreeMap for automatic ordering
-
Use Read-Copy-Drop pattern:
- Never hold locks while processing data
- Copy what you need, release lock, then process
-
Add compile-time Send+Sync checks:
const _: () = { const fn assert_send_sync<T: Send + Sync>() {} assert_send_sync::<YourNewConcurrentType>(); };
-
Write Loom tests for new concurrent code:
#[cfg(loom)] #[test] fn test_your_concurrent_operation() { loom::model(|| { // Test concurrent access patterns }); }
HNSW graph persistence is intentionally deferred: Collection::flush() only
saves the HNSW graph to disk when inserts_since_last_hnsw_save exceeds
HNSW_SAVE_THRESHOLD (10 000 inserts). This amortizes the cost of
serializing the full HNSW graph (metadata, mappings, vectors, and graph
structure) across many write operations instead of paying it on every flush.
The trade-off is an enlarged crash recovery window: if the process crashes between a vector storage write and the next HNSW save, the HNSW index on disk will be missing those vectors. Two complementary recovery layers ensure no data is lost.
┌────────────────────────────────────────────────────────────────────────┐
│ Collection::open() │
│ │
│ 0. load_config() — config.json │
│ │
│ 1. MmapStorage::new() │
│ ├─ Load vectors.idx (ID → offset mapping) │
│ ├─ Replay vectors.wal → restore writes since last flush_index() │
│ ├─ Record WAL-touched ids (drained by step 7, pass 3) │
│ └─ Truncate WAL after successful replay │
│ │
│ 2. LogPayloadStorage::new() │
│ └─ Load payloads.snapshot, replay payloads.log past it │
│ │
│ 3. load_or_create_hnsw() │
│ ├─ sweep_stale_arenas(): drop orphan arena files, before any load │
│ ├─ Gate: native_meta.bin present? (commit point, written LAST) │
│ ├─ Load native_hnsw.graph/.vectors/.gen + native_mappings.bin — │
│ │ all generation-stamped (#617); a legacy native_vectors.bin │
│ │ is generation-checked then discarded (PERF1) │
│ └─ Load failure or config mismatch → empty index (rebuild below) │
│ │
│ 4. load_bm25_index() │
│ └─ bm25.snapshot + bm25.wal replay; payload rebuild if no snapshot │
│ │
│ 5. Property index, label index (rebuilt from payloads), range index, │
│ edge_store.bin, named sparse indexes │
│ └─ each sparse snapshot, then its WAL replayed over it │
│ │
│ 6. reconcile_point_count() │
│ └─ Set config.point_count = storage.len() (authoritative source) │
│ │
│ 7. recover_index_state() — reconciliation passes │
│ ├─ Pass 1 (gap): recover_hnsw_gap │
│ │ ├─ Early exit: if storage.len() == hnsw.len() → no gap │
│ │ ├─ find_gap_ids: storage.ids() \ index.mappings │
│ │ ├─ retrieve_valid_vectors: load from mmap, validate dimension │
│ │ └─ reindex_vectors: insert_batch_parallel into HNSW │
│ ├─ Pass 2 (orphans): ids in index.mappings \ storage → remove │
│ ├─ Pass 3 (stale): WAL-touched ids on both sides — re-upsert when │
│ │ the indexed vector ≠ storage (storage is the source of truth) │
│ ├─ Pass 4 (unlinked): mapped ids whose node has no layer-0 link │
│ │ (entry point exempt) — re-upsert onto a fresh, linked node │
│ └─ Any pass mutated the index → index.save() before open returns │
│ (the WAL was truncated; the delta has no other witness) │
│ │
│ 8. restore_auto_reindex_from_config(), │
│ restore_secondary_indexes_from_config() │
│ │
│ 9. run_post_open_hooks() │
│ ├─ reindex edge properties from edge_store.bin │
│ ├─ replay edges.wal over edge_store.bin │
│ └─ restore_persisted_quantizers(): trained quantizer artifacts │
└────────────────────────────────────────────────────────────────────────┘
Module: crates/velesdb-core/src/storage/mmap/wal_replay.rs
Every MmapStorage::store() and store_batch() call writes a CRC32-framed
entry to vectors.wal before updating the mmap and in-memory index. The
WAL format uses a single-byte opcode prefix:
| Op | Name | Frame Layout |
|---|---|---|
0x01 |
Store | [op:1B][id:8B LE][len:4B LE][data:N B][crc32:4B LE] |
0x02 |
Delete | [op:1B][id:8B LE][crc32:4B LE] |
On MmapStorage::new(), the constructor calls replay_wal_to_index() which:
- Opens
vectors.waland validates it uses the CRC32-framed format (legacy pre-#317 WAL files without CRC are detected and skipped). - Reads entries sequentially, verifying each CRC32 checksum. A CRC mismatch or truncated entry indicates a crash mid-write; replay stops at the corruption boundary (all prior valid entries are applied).
- For store entries: writes the vector data into the mmap at the correct offset and updates the sharded index.
- For delete entries: removes the ID from the sharded index.
- Truncates the WAL file to zero after successful replay, preventing double-replay on the next startup.
This layer recovers vectors that were written to the WAL but not yet
persisted to vectors.idx (the index file is only written by
flush_index() or flush_full(), not by the fast flush() path).
Module: crates/velesdb-core/src/collection/core/recovery.rs
After storage is fully reconstructed (Layer 1) and the persisted HNSW
index is loaded (or an empty one built when the load fails),
Collection::open() calls run_crash_recovery(), which runs the passes
below against the storage state, in order:
Pass 1 — gap (recover_hnsw_gap):
-
Early exit heuristic: If
storage.len() == 0orstorage.len() == hnsw.len(), returns 0 (no gap). This check is O(1) and avoids a full scan in the common case. -
Gap ID detection (
find_gap_ids): Iterates all IDs instorage.ids()and filters those not present inindex.mappings.contains(id). This is O(storage_count) with O(1) per-ID lookup in the sharded mappings. -
Vector retrieval (
retrieve_valid_vectors): Loads each gap vector from mmap storage, skipping entries with mismatched dimension (corruption) or missing data (concurrent deletion betweenids()andretrieve()calls). -
Re-indexing (
reindex_vectors): Batch-inserts all valid gap vectors into the HNSW graph viainsert_batch_parallel. The re-index uses the same parallel rayon-based insertion as normal upserts.
Pass 2 — orphans (remove_orphan_ids): ids present in
index.mappings but absent from storage (a delete reached the vector
WAL but not the next HNSW save) are soft-deleted from the index so the
tombstone cannot resurface in search results.
Pass 3 — stale (reindex_stale_wal_ids): for every id touched by
the Layer-1 WAL replay that is present on both sides, the indexed
sidecar vector is compared against the storage bytes; on mismatch the
storage value is re-upserted (an upsert landed in the WAL after the
last HNSW save). An index loaded without sidecar vector storage cannot
be compared — when WAL-touched ids overlap its mappings it is replaced
by an empty index and fully rebuilt by pass 1
(rebuild_if_unverifiable).
Pass 4 — unlinked (relink_unlinked_ids): a mapping is not proof of
graph membership. upsert_bulk's V2 path places each vector, maps its id
to the slot it got, and leaves linking that slot to the AsyncIndexBuilder, so
a save that races it persists mappings for nodes nothing links to (#2246).
Every mapped id whose node has an empty layer-0 list — the entry point
excepted, since a graph's first node has nothing to link to — is
re-upserted onto a fresh, linked node, and its old slot is left behind as
a tombstone.
When any pass mutated the index, Collection::open() re-saves it
before returning: the vector WAL was truncated during replay, so
without a fresh save the reconciled delta would be undetectable after
the next crash. For the same reason, compact_vector_storage() (which
also truncates the WAL) re-saves the HNSW index after compaction.
Three distinct write paths can leave vectors in storage but absent from HNSW:
| Gap Source | Mechanism | Typical Window |
|---|---|---|
| Normal insert gap | batch_store_all writes vectors before bulk_index_or_defer inserts into HNSW. Crash between Phase 1 and Phase 3 of the 3-phase upsert pipeline. |
Duration of Phase 2 (secondary indexes, quantization, text indexing) |
| Deferred indexer gap | DeferredIndexer buffers vectors in memory (up to merge_threshold, default 1 024) before batch-merging into HNSW. Crash before merge loses the buffer. |
Up to merge_threshold vectors (memory-only, not WAL-protected) |
| Delta buffer gap | DeltaBuffer accumulates vectors during background HNSW rebuild. Crash before deactivate_and_drain loses the buffer. |
Duration of the rebuild operation |
All three gaps are recovered by the same recover_hnsw_gap mechanism
because the recovery compares the final storage state against the HNSW
mappings, regardless of how the gap originated.
If a crash occurs between an HNSW delete and the corresponding storage
delete being persisted, a previously deleted vector may appear in storage
but not in HNSW. This is indistinguishable from an insert gap. Recovery
will re-index the deleted vector, effectively "resurrecting" it. This is
an intentional trade-off: resurrecting a deleted vector is preferable to
silently losing an inserted one. The window for this scenario is very
small (within a single delete() call).
Recovery latency grows with the number of gap vectors. No run measures it; what dominates at each size:
| Gap Size | Dominant Cost |
|---|---|
| 0 (no gap) | O(1) early exit heuristic |
| 1–100 vectors | Storage retrieval + HNSW insert |
| 100–1 000 vectors | Parallel HNSW batch insert |
| 1 000–10 000 vectors | Parallel HNSW batch insert (rayon) |
| > 10 000 vectors | Proportional to gap size; mitigated by HNSW_SAVE_THRESHOLD |
The HNSW_SAVE_THRESHOLD (10 000) bounds the maximum gap size in practice:
flush() forces an HNSW save after 10 000 inserts, so the worst-case
recovery inserts at most ~10 000 vectors into the graph. A graceful
shutdown via flush_full() saves the HNSW graph unconditionally, reducing
the gap to zero for planned restarts.
| Parameter | Default | Location | Effect |
|---|---|---|---|
HNSW_SAVE_THRESHOLD |
10 000 | Collection::flush() in flush.rs |
Maximum inserts before flush() forces an HNSW save. Lower values reduce worst-case recovery time but increase flush latency. |
DurabilityMode |
Fsync |
MmapStorage |
Controls WAL write behavior. Fsync: full durability. FlushOnly: user-space flush only (faster, risk of OS-crash data loss). None: no WAL writes (for bulk import; no WAL replay possible). |
DeferredIndexerConfig.merge_threshold |
1 024 | collection.streaming.deferred |
Number of buffered vectors before deferred merge into HNSW. Larger values increase the deferred indexer gap window. |
DeferredIndexerConfig.max_buffer_age_ms |
5 000 | collection.streaming.deferred |
Maximum age of buffered vectors before a time-based merge. Provides a time bound on the deferred gap. |
| Method | WAL fsync | mmap flush | vectors.idx | HNSW save | Use Case |
|---|---|---|---|---|---|
Collection::flush() |
Yes | Yes | No | Only if > 10K inserts | Normal operation, periodic durability |
Collection::flush_full() |
Yes | Yes | Yes | Always | Graceful shutdown, before compaction |
MmapStorage::flush() |
Yes | Yes | No | N/A | Storage-level fast barrier |
MmapStorage::flush_full() |
Yes | Yes | Yes | N/A | Storage-level complete barrier |
HNSW index persistence uses atomic write-tmp-fsync-rename for crash safety.
Each save writes six files, every one stamped with the same monotonic
generation: u64 (#617) so a crash between two renames is detected on
load. native_meta.bin is written LAST — its generation is the
authoritative commit point that load_sidecars checks the other
artefacts against, and its presence is the gate load_or_create_hnsw
uses to attempt a load at all.
| File | Contents | Format |
|---|---|---|
native_hnsw.vectors |
Vector data in NodeId order |
Custom binary via file_dump |
native_hnsw.graph |
Graph structure (layers, neighbors) + params incl. VAMANA alpha (header v2) |
Custom binary via file_dump |
native_hnsw.gen |
Graph generation marker | postcard-serialized u64 |
native_mappings.bin |
id_to_idx, idx_to_id, next_idx, generation |
postcard-serialized HashMaps |
native_vectors.bin (legacy, pre-PERF1) |
Vec<(internal_idx, Vec<f32>)>, generation |
Read for the generation check only, payload discarded; deleted on the next save |
native_meta.bin |
Dimension, metric, vector storage flag, storage mode, generation | postcard-serialized tuple |
Decision: the standalone
hnsw_delta_walmodule has been removed fromvelesdb-core. O(delta) fast / warm-standby recovery becomes a premium concern, built by consuming theWalCursorseam (crates/velesdb-core/src/storage/wal_cursor.rs) rather than a core-owned delta-WAL. This disposition is recorded here, alongside the WAL shippability contract, because both concern the WAL/recovery boundary.
Core previously carried a standalone storage/hnsw_delta_wal.rs module
that logged incremental graph mutations (edge add/remove, entry-point
changes) with the intent of enabling O(delta) graph recovery instead of a
full O(N*M) rebuild. It was never wired into the open / flush /
recovery path — Collection::open() and the recovery flow described above
never read or wrote it.
Why it was not wired, and was removed instead of activated:
- Dead-in-core infrastructure. Core recovery correctness is already
fully provided by the reconciliation passes
(
collection/core/recovery.rs) against storage as the single source of truth. The delta WAL added nothing to correctness; the persisted graph load plus reconciliation is sufficient. - No speculative infrastructure. An unexercised WAL format that must be kept byte-compatible with a graph format it never observes is exactly the silent-drift hazard the disposition set out to eliminate. Wiring it would make the HNSW graph a partial source of truth alongside storage, re-opening soundness invariants for a performance gain that only matters at cluster scale.
- The O(delta) benefit is an enterprise concern. Fast failover and
warm-standby recovery belong to premium's clustering/replication story,
not the local-first core. Premium can build it by consuming the
WalCursorseam — the clean, supported extension point — which unifies "shippable WAL" and "delta recovery" behind one open seam rather than two half-built mechanisms. Core never depends on premium; premium consumes the cursor.
On-disk / migration impact: none. Current versions never wrote a delta-WAL file, so the removal is a pure source deletion with no on-disk artifact to migrate and no format-compatibility work required.
Recovery behavior is validated by the following test suite:
| Test | File | Scenario |
|---|---|---|
test_no_gap_returns_zero |
recovery_tests.rs |
No gap: storage and HNSW counts match |
test_empty_collection_no_recovery |
recovery_tests.rs |
Empty collection: early exit |
test_crash_gap_detected_and_recovered |
recovery_tests.rs |
Simulated gap: 2 vectors in storage but not HNSW |
test_gap_recovery_on_collection_reopen |
recovery_tests.rs |
End-to-end: create, gap, flush, drop, reopen, verify search |
test_metadata_only_skips_recovery |
recovery_tests.rs |
Metadata-only collections skip recovery |
| WAL replay tests | wal_recovery_tests.rs |
CRC validation, legacy format detection, truncation |
Storage compaction holds the MmapStorage write lock for the entire
duration of the operation. This is enforced at two levels:
-
Synchronous path (
MmapStorage::compact(&mut self)): The method takes&mut self, so the caller must already hold an exclusive reference. No concurrent reads or writes are possible while compaction runs. -
Asynchronous path (
compact_async(storage: Arc<RwLock<MmapStorage>>)): Acquiresstorage.write()inside aspawn_blockingtask and holds the write guard for the full compaction cycle. All readers and writers on the sameRwLockare blocked until the guard is dropped.
compact_async()
├─ spawn_blocking
│ ├─ storage.write() ← exclusive lock acquired
│ ├─ MmapStorage::compact() ← rewrite active vectors to .tmp
│ │ ├─ build temp file
│ │ ├─ copy active vectors
│ │ ├─ atomic_replace(.tmp → .dat)
│ │ └─ rebuild index + flush
│ └─ drop(guard) ← exclusive lock released
On large collections (>1M vectors), compaction rewrites the entire active vector set to a new file and atomically replaces the original. This can block all reads and writes for seconds, depending on disk throughput and vector dimensionality. This is an intentional correctness-over-performance trade-off: holding the exclusive lock prevents readers from observing a partially rewritten file and writers from appending to a file that is about to be replaced.
recover_compaction_artifacts() runs automatically during
MmapStorage::new() to repair any interrupted compaction. The recovery
logic inspects leftover intermediate files:
| State on Disk | Interpretation | Recovery Action |
|---|---|---|
.bak exists, original missing |
Crash after rename-to-backup, before new file swap | Restore .bak as original |
.bak exists, original exists |
Compaction completed, backup not yet cleaned up | Remove .bak |
.tmp exists |
Incomplete compaction (temp file never swapped in) | Remove .tmp |
This ensures the storage directory is always in a consistent state before the mmap file is opened, regardless of when the previous process crashed.
Module: crates/velesdb-core/src/storage/compaction.rs
(recover_compaction_artifacts, atomic_replace)
Copy-on-write compaction that allows concurrent reads during the rewrite phase is planned for the enterprise edition. The current exclusive-lock design is the baseline for correctness validation.
These validate hand-written loom models of the lock ordering
(tests/loom_tests.rs, src/storage/loom_tests.rs), not the production
parking_lot/dashmap types directly. The tests gate on cfg(loom); the
crate's build.rs bridges the loom Cargo feature to that cfg, so the feature
flag alone activates them — RUSTFLAGS="--cfg loom" is not required locally
(CI still sets it explicitly, which is harmless — the cfg is idempotent). A bare
cargo test --features loom (without the bridge) would compile loom but run
zero tests, which is the trap the build.rs removes.
These are the commands CI runs (quality-deep.yml), with the preemption bound
it uses; the RUSTFLAGS is redundant locally and harmless.
# Integration models (tests/loom_tests.rs)
LOOM_MAX_PREEMPTIONS=3 RUSTFLAGS="--cfg loom" cargo test -p velesdb-core --features loom,persistence --test loom_tests -- --test-threads=1
# Storage models (src/storage/loom_tests.rs, a unit-test target)
LOOM_MAX_PREEMPTIONS=3 RUSTFLAGS="--cfg loom" cargo test -p velesdb-core --features loom,persistence storage::loom -- --test-threads=1# Run stress tests with multiple threads
cargo test --test stress_concurrency_tests -- --test-threads=1For how an insert gets its slot — one allocator, the arena, with the mapping following it under the index read guard — see SOUNDNESS.md: HNSW Slot Allocation.
- Rust Atomics and Locks (Mara Bos)
- The Rustonomicon - Concurrency
- parking_lot documentation
- Loom crate
Last updated: 2026-08-09 · Applies to: velesdb-core 6.0.0 (this revision: corrected the lock-order enforcement section — assert_lock_order is unwired in production, the HNSW tracker is debug-only, warn-only and partial, the collection tier is convention-only — and the CSR snapshot thread-safety section to describe the lock-free ArcSwap + dirty-flag protocol; previous revision noted: HNSW persisted-graph reload at open; storage compaction concurrency)