Skip to content

Latest commit

 

History

History
1142 lines (926 loc) · 57.4 KB

File metadata and controls

1142 lines (926 loc) · 57.4 KB

VelesDB Concurrency Model

EPIC-023: Documentation du modèle de concurrence pour utilisateurs avancés et contributeurs.

User-facing write throughput guidance: see docs/guides/WRITE_CONCURRENCY.md for 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.

Overview

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

Inter-Process Exclusion (the outermost boundary)

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.

Architecture

Sharding Strategy

┌─────────────────────────────────────────────────────────────────┐
│                    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

Lock Types

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

Thread Safety Guarantees

Send + Sync Types

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 + Sync

!Send Types (Single-Thread Only)

These types contain non-thread-safe internal state:

// Must stay on creation thread
GraphTraversal: !Send  // Contains references
QueryCursor: !Send     // Iterator state
BfsIterator: !Send     // Traversal state

Compile-Time Verification

VelesDB uses compile-time assertions to verify thread safety:

// In ConcurrentEdgeStore
const _: () = {
    const fn assert_send_sync<T: Send + Sync>() {}
    assert_send_sync::<ConcurrentEdgeStore>();
};

Lock Ordering (Deadlock Prevention)

Rule: Always Acquire Locks in Ascending Order

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]

Global Lock Order (HNSW + Storage)

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. The record_lock_acquire function in crates/velesdb-core/src/index/hnsw/native/graph/locking.rs is #[cfg(debug_assertions)]; on an out-of-order acquisition it increments an atomic violation counter and emits a tracing::warn! — it never panics. It is also partial: only the GpuVectorsSnapshot, EntryPointPromotion, Vectors and Layers ranks are ever recorded. Neighbors is #[allow(dead_code)] with no record_lock_acquire call 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 in crates/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).

The HNSW Index Guard and rayon: a liveness rule, not an ordering one

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::inner never 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:

  1. a vacuum asks for inner.write() and waits for the current readers;
  2. a batch operation holds inner.read() and waits for its rayon jobs;
  3. those jobs run on the global pool, where the per-query search and rerank_candidates_simd take inner.read() — new readers, so the pending writer parks them;
  4. 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.rs parks 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 job lint) refuses a new rayon submission in any non-test file under index/hnsw/ that is neither inside an .install( closure — on any dedicated pool, not only graph_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 reached par_iter three 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.

Reserved Premium Rank Range [40, 59]

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

Cross-Shard Operations

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!

Cascade Delete (remove_node_edges)

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.

Collection-level lock order

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) hold payload_storage's read guard from the endpoint referential-integrity check through the end of the edge write, and acquire edge_wal_lock while that read guard is still held — so the acquisition order for this path is payload_storage(3) → edge_wal_lock(3b).
  • remove_edge and the delete-cascade path (collection/core/crud_read_delete.rs's cascade_delete_node_edges) acquire edge_wal_lock alone, with no other collection lock held — the delete path has already released its payload_storage write 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.

Quantized-Precision Backend Interior Mutability (RaBitQ, SQ8)

Lock Layout

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

Training Lock Order

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.

Store-Before-Quantizer Ordering

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.

RaBitQ Contention Analysis

Known Contention Patterns

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

Mitigation

  • Post-training inserts hold rabitq_store.write() for the minimum duration needed to push a single encoded vector (one Vec::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.

HNSW Entry-Point Promotion

Entry-Point Updates

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 layer

Writers take the promotion lock (promotion: Mutex<()>, rank 8): the two cases below, and the renumbering reorder_for_locality applies.

  1. Empty index: claim_first_entry_point CASes entry_point from NO_ENTRY_POINT to the new node ID and sets max_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).
  2. Layer promotion: promote_entry_point returns before taking the lock when the node's layer does not exceed max_layer, the common case. Under the lock it checks again, stores the new max_layer, swaps in the new entry point (swap, AcqRel) and calls reparent_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.

CsrSnapshot Invalidation Pattern

Graph Edge CSR Read Snapshot

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:

  1. Build: build_read_snapshot() materializes edges into contiguous arrays (targets: Vec<u64>, edge_ids: Vec<u64>) with a FxHashMap<u64, (offset, len)> index. Auto-built after loading from disk and after flush().
  2. Read: csr_snapshot().neighbors(source_id) returns &[u64] -- a zero-copy slice into the contiguous target array.
  3. Invalidate: Every write operation (add_edge, remove_edge, remove_node_edges) sets csr_snapshot = None and 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) holds edge_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 resets pending_writes to 0 and clears csr_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 observed pending_writes before reading the shards, rebuilds, then subtracts only that observed count (fetch_update with saturating_sub). A concurrent fetch_add landing 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.

Performance vs Safety Tradeoffs

Read-Heavy Workloads

  • RwLock allows multiple concurrent readers
  • Sharding distributes reads across independent locks
  • Recommendation: Default 256 shards optimal for most workloads

Write-Heavy 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

Graph Traversal

  • 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 via reconstruct_path() only when a result is emitted, avoiding O(depth) allocations per expansion step.

Known Limitations

  1. Cross-shard operations hold multiple locks:

    • Edge spanning 2 shards requires 2 locks + edge_ids lock
    • Mitigation: Lock ordering prevents deadlocks
  2. Large traversals can block writers:

    • BFS/DFS with many nodes may hold locks longer
    • Mitigation: Read-Copy-Drop pattern releases locks quickly
  3. 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_locality and 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_construction and 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
  4. No transactional semantics:

    • Operations are atomic per-operation, not per-batch
    • Mitigation: Use flush() for durability checkpoints
  5. 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 compares storage.ids() against index.mappings and re-indexes any missing vectors. See HNSW Crash Recovery for the full recovery architecture and SOUNDNESS.md for the slot allocation invariants.

Best Practices

For Users

  1. Dimension shards appropriately:

    // For 100K edges
    let store = ConcurrentEdgeStore::with_estimated_edges(100_000);
  2. 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])?;
  3. Limit traversal depth:

    // Always specify max_depth to prevent runaway traversals
    let nodes = store.traverse_bfs(start, 5);  // Max 5 hops

For Contributors

  1. Follow lock ordering strictly:

    • Document lock order in new concurrent structures
    • Use BTreeSet/BTreeMap for automatic ordering
  2. Use Read-Copy-Drop pattern:

    • Never hold locks while processing data
    • Copy what you need, release lock, then process
  3. Add compile-time Send+Sync checks:

    const _: () = {
        const fn assert_send_sync<T: Send + Sync>() {}
        assert_send_sync::<YourNewConcurrentType>();
    };
  4. Write Loom tests for new concurrent code:

    #[cfg(loom)]
    #[test]
    fn test_your_concurrent_operation() {
        loom::model(|| {
            // Test concurrent access patterns
        });
    }

HNSW Crash Recovery

Problem Statement

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.

Recovery Architecture

┌────────────────────────────────────────────────────────────────────────┐
│                        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     │
└────────────────────────────────────────────────────────────────────────┘

Layer 1: Vector Storage WAL Replay

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:

  1. Opens vectors.wal and validates it uses the CRC32-framed format (legacy pre-#317 WAL files without CRC are detected and skipped).
  2. 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).
  3. For store entries: writes the vector data into the mmap at the correct offset and updates the sharded index.
  4. For delete entries: removes the ID from the sharded index.
  5. 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).

Layer 2: HNSW Reconciliation Passes

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):

  1. Early exit heuristic: If storage.len() == 0 or storage.len() == hnsw.len(), returns 0 (no gap). This check is O(1) and avoids a full scan in the common case.

  2. Gap ID detection (find_gap_ids): Iterates all IDs in storage.ids() and filters those not present in index.mappings.contains(id). This is O(storage_count) with O(1) per-ID lookup in the sharded mappings.

  3. Vector retrieval (retrieve_valid_vectors): Loads each gap vector from mmap storage, skipping entries with mismatched dimension (corruption) or missing data (concurrent deletion between ids() and retrieve() calls).

  4. Re-indexing (reindex_vectors): Batch-inserts all valid gap vectors into the HNSW graph via insert_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.

Gap Sources

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.

Known Limitation: Delete-Insert Ambiguity

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).

Startup Latency Impact

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.

Configuration Knobs

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.

Flush Variants

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

Persistence Format

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

HNSW Delta WAL — Removed from Core (Disposition)

Decision: the standalone hnsw_delta_wal module has been removed from velesdb-core. O(delta) fast / warm-standby recovery becomes a premium concern, built by consuming the WalCursor seam (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 WalCursor seam — 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.

Test Coverage

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 Concurrency

Exclusive Lock Scope

Storage compaction holds the MmapStorage write lock for the entire duration of the operation. This is enforced at two levels:

  1. 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.

  2. Asynchronous path (compact_async(storage: Arc<RwLock<MmapStorage>>)): Acquires storage.write() inside a spawn_blocking task and holds the write guard for the full compaction cycle. All readers and writers on the same RwLock are 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

Latency Impact

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.

Crash Recovery

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)

Future Roadmap

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.

Testing Concurrency

Running Loom Tests

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

Stress Testing

# Run stress tests with multiple threads
cargo test --test stress_concurrency_tests -- --test-threads=1

HNSW Slot Allocation

For 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.

References


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)