Skip to content

Latest commit

 

History

21 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

KV Store

A distributed key-value store that keeps your data available and consistent while machines fail underneath it.

Writes are replicated across a cluster of nodes before being acknowledged, so losing a node loses no acknowledged data. The cluster elects a new leader automatically when one goes down, and clients follow the change without reconfiguration. Reads and writes are linearizable: every client sees a single, consistent ordering of operations, never a stale value from a node that has fallen behind.

Storage is an LSM engine written from scratch — a write-ahead log, in-memory memtable, on-disk SSTables with bloom filters, and background compaction — so the dataset is not bounded by memory and survives a crash.

Features

  • Replicated PUT / GET / DELETE over a simple binary protocol
  • Automatic leader election and failover; clients follow redirects on their own
  • Linearizable reads and writes, verified by a linearizability checker
  • Durable through crashes: consensus state and stored data are both fsynced
  • Log compaction via snapshots, so a long-running cluster doesn't grow forever
  • Nodes that fall far behind catch up by snapshot transfer rather than full replay
  • Retried writes are deduplicated server-side, so a client retry is not applied twice

Build

docker compose run --rm dev
# inside the container:
mkdir -p build_linux && cd build_linux
cmake -DENABLE_ASAN=ON .. && make -j4

Running a cluster

# in three separate shells (or backgrounded), from build_linux/:
./kv_store 1
./kv_store 2
./kv_store 3

Client usage

./kv_client put foo bar     # OK
./kv_client get foo         # bar
./kv_client del foo         # OK
./kv_client get foo         # (not found)

The client takes the full cluster membership and finds the leader itself, following redirects as leadership moves — any node will do as a starting point. A C++ client library is in include/raft_client.h; tests/chaos/raft_client.py is an equivalent Python client speaking the same wire protocol.

Each node listens on a Raft peer port (800<id>) and a client-facing port (900<id>), bound to its own loopback alias (127.0.0.<id+1>). Each node keeps its data in node_data_<id>/, so a whole cluster can run from one directory.

Cluster size

Defaults to 3 nodes. Both binaries take --num-nodes N for a different-sized cluster on the same localhost layout, or --cluster for explicit membership (real hosts, non-sequential ports):

./kv_store 1 --num-nodes 5
./kv_client get foo --num-nodes 5

./kv_store 1 --cluster 1:10.0.0.1:8001:9001,2:10.0.0.2:8001:9001,3:10.0.0.3:8001:9001

Every node and client in a cluster must be given the same membership.

Tests

All commands assume you're inside the dev container (docker compose run --rm dev) and have built into build_linux/.

# unit tests: storage engine + consensus state persistence
./build_linux/run_tests

# linearizability checker's own tests (no cluster needed)
cd tests/chaos && python3 test_linearizability_checker.py

Durability under faults — sustained writes while nodes are killed and network-partitioned, then verifies every acknowledged write is still readable:

python3 tests/chaos/chaos_test.py \
  --bin-dir build_linux --duration 25 --fault-interval 4

Partitions use real iptables rules and need NET_ADMIN (already granted to the dev service). Add --no-partitions to run kill-only without it.

Linearizability — concurrent clients hammer a small shared keyspace while faults are injected; the recorded history is then checked for any ordering no correct system could have produced:

python3 tests/chaos/concurrent_lin_test.py \
  --bin-dir build_linux --duration 25 --workers 8 --num-keys 6 \
  --with-chaos --fault-interval 5

Benchmark — throughput and latency percentiles for a configurable read/write mix:

python3 tests/chaos/benchmark.py \
  --bin-dir build_linux --duration 20 --concurrency 8 --num-keys 1000 --read-ratio 0.8

The test scripts accept the same --num-nodes / --cluster flags and pass them through to the processes they start. Each cleans up its own node_data_*/ directories first, so runs don't leak state into each other.

Performance

Reads are answered from the leader's state machine once a majority has confirmed it is still leader, rather than being replicated through the log. That keeps reads linearizable while avoiding a log append and an fsync each, and one confirmation round covers every read waiting on it. On a read-heavy workload it's worth roughly 20–30% throughput over log-replicated reads.

Indicative numbers for 3 nodes on loopback inside Docker, 4 concurrent clients, 80% reads: ~60–90 ops/sec, p50 45–60 ms, p99 under 130 ms. Absolute figures move a lot with host load, so the honest measurement is the A/B, not any single number — --log-reads runs the cluster the other way, and the benchmark takes the same flag:

python3 tests/chaos/benchmark.py --bin-dir build_linux --duration 10 --read-ratio 0.8
python3 tests/chaos/benchmark.py --bin-dir build_linux --duration 10 --read-ratio 0.8 --log-reads

Latency is dominated by the replication round-trip and the event loop's poll interval rather than by storage.

Limitations

  • Cluster membership is fixed at startup — no online reconfiguration
  • One event loop per node: consensus is deliberately single-threaded, so throughput comes from batching rather than parallelism
  • Tested against process crashes, not power loss

About

A distributed key value store built for strong consistency and durability. It survives node crashes via automatic leader election and write replication while ensuring clients never read stale data. Powered by a custom LSM tree storage engine, it safely persists datasets larger than RAM and recovers seamlessly from failures.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages