A distributed, in-memory key-value store built to make distributed-systems behavior tangible. Shardis implements sharding, leader–follower replication, durable writes, failover, live slot migration, pub/sub, and cluster observability using real processes and WebSocket connections.
It is intentionally a small-scale learning and demonstration system—not a Redis replacement. The goal is readable implementation, honest trade-offs, and reproducible failure testing.
- Dashboard: shardis-benchmarks.vercel.app
- Nodes: leader
/healthz· follower/healthz
Free Render services sleep after idle traffic and can take about a minute to wake. The dashboard retries automatically and shows a waking state during that interval.
- Distributed data path — Redis CRC16 hash slots, hash tags, cross-shard
MOVEDredirects, andASKredirects while slots migrate. - Durability — append-only logging with
fsyncbefore acknowledgement, snapshots, compaction, replay, and crash recovery. - Replication and failover — leader–follower replication, heartbeats, runtime membership, gossip-backed redirect updates, and an opt-in Raft-lite mode.
- Real operations — Docker Compose topology, health checks, JSON metrics, structured logs, CLI tooling, and a Next.js cluster dashboard.
- Resilience testing — unit tests plus real multi-process integration tests, including
SIGKILLrecovery and Compose smoke tests in CI.
┌─────────────────────┐
│ Dashboard / CLI │
│ WebSocket clients │
└──────────┬──────────┘
│
MOVED / ASK │ GET · SET · DEL · PUB/SUB
▼
┌───────────────────────────────────────────────────┐
│ Shardis cluster │
│ │
│ Shard A Shard B Shard C │
│ a1 (leader) ──► a2 b1 (leader) ──► b2 c1 ──► c2│
│ │ │ │ │
│ AOF + snapshot AOF + snapshot AOF + snapshot
└───────────────────────────────────────────────────┘
The local environment runs three shards and six nodes. Any node can accept a client connection; when a key belongs elsewhere, the client is redirected to the active leader for that shard.
- Node.js 20 or newer
- pnpm 9
- Docker Desktop and Docker Compose (for the full cluster)
pnpm install
pnpm build
pnpm testdocker compose up -d --build
curl http://localhost:7001/healthzThe six nodes are exposed on ports 7001–7006. Stop the cluster and remove its local volumes with:
docker compose down -vpnpm --filter @shardis/cli build
node packages/cli/dist/shardis-cli.js --url ws://localhost:7001/ws SET hello world
node packages/cli/dist/shardis-cli.js --url ws://localhost:7001/ws GET helloThe CLI follows MOVED redirects automatically. To use the compact binary protocol instead of JSON:
node packages/cli/dist/shardis-cli.js --binary --url ws://localhost:7001/ws SET hello worldpnpm --filter @shardis/dashboard devOpen http://localhost:3000. With the Compose cluster running, the dashboard shows topology, node health, replication state, events, and an in-browser console.
| Area | Included behavior |
|---|---|
| Storage | In-memory keys, TTLs, active and lazy expiry, and LRU eviction |
| Persistence | AOF writes before acknowledgement, snapshots, compaction, and recovery replay |
| Routing | 16,384 CRC16 hash slots, hash tags, MOVED redirects, and runtime slot ownership |
| Resharding | Slot-by-slot MIGRATING / IMPORTING handoff with ASK redirects and resume support |
| Replication | Leader streaming, follower acknowledgements, lag visibility, full resync, and membership relay |
| Failover | Deterministic promotion by default; opt-in Raft-lite elections and majority commit tracking |
| Messaging | Local sharded pub/sub and optional cluster-scoped publish relay |
| Protocols | JSON over WebSocket by default, plus an opt-in compact binary protocol |
| Observability | /healthz, /metrics, /topology, structured logs, CLI, and dashboard |
Shardis is designed to be observed under failure. Kill the leader of shard A, then inspect its follower:
docker compose kill -s SIGKILL node-a1
curl http://localhost:7002/healthz
docker compose up -d node-a1In deterministic mode, the lowest live follower id promotes. The returning node reconnects and completes a full resynchronization. The test suite also covers follower outages, stale redirects, dynamic joins, graceful shutdown, and crash recovery.
Client requests are WebSocket messages using commands such as GET, SET, DEL, EXPIRE, TTL, SUBSCRIBE, UNSUBSCRIBE, and PUBLISH.
| Endpoint | Purpose |
|---|---|
GET /healthz |
Node identity, shard, role, and uptime |
GET /metrics |
Operation, eviction, connection, and replication-lag metrics |
GET /topology |
The topology loaded by that node; used for dashboard bootstrap |
WS /ws |
WebSocket client and peer protocol endpoint |
By default, pub/sub is sharded: publishes reach subscribers connected to the receiving node. Use scope: "cluster" (CLI: PUBLISH channel message --cluster) to relay an event once to other shards.
All node configuration is environment-driven. Copy .env.example when running a node outside Compose.
| Variable | Description | Default |
|---|---|---|
NODE_ID |
Stable node identifier | node-a1 |
ROLE |
Boot role hint | leader |
SHARD_ID |
Node's shard | shard-a |
CLUSTER_CONFIG_PATH |
Static topology file | ./cluster.config.local.json |
PORT |
HTTP and WebSocket port | 7000 |
DATA_DIR |
AOF and snapshot directory | ./data/<NODE_ID> |
FAILOVER_MODE |
deterministic or raft |
deterministic |
JOIN_URL / NODE_URL |
Dynamic-follower join configuration | unset / derived |
PUBLIC_DEMO / DEMO_WRITE_KEY |
Demo write protection | false / unset |
CLUSTER_SECRET |
Shared peer authentication secret | unset |
MAX_CONNECTIONS_PER_IP |
Concurrent connection cap per source IP | 20 |
Benchmarks are recorded with a timestamp and commit hash in docs/benchmarks.md. Current local measurements include:
| Scenario | Result |
|---|---|
| Throughput | 2,260 ops/sec with 10 clients over 5 seconds |
| Scaling sweep | 18,123 ops/sec peak at 5 clients in a short single-node sweep |
| Failover recovery | 2,846 ms total with a 3,000 ms heartbeat timeout |
| Replication lag | 4.5 ms average; 7 ms p95 across 50 samples |
| Range recompute | ~50% of keys move when changing 3 shards to 4 |
These are local, controlled-environment measurements—not capacity guarantees. Reproduce them with the benchmark package after building the node:
pnpm --filter @shardis/benchmarks bench:<name>Shardis is not intended to be exposed directly to the public internet.
- Nodes speak plain
ws://internally. Terminate TLS at a reverse proxy or managed platform and expose onlywss://to clients. PUBLIC_DEMO=truecan requireDEMO_WRITE_KEYfor mutating operations; this is a demo safeguard, not user authentication or tenant isolation.- The node applies request, connection, key-size, and value-size limits. Peer authentication is available through
CLUSTER_SECRET. - Per-key ACLs, encryption at rest, mTLS, off-node backups, and full production hardening are deliberately out of scope.
Please review SECURITY.md before reporting a vulnerability or deploying an internet-facing instance.
The reduced two-node Render demo is intentionally a preview environment, not an always-on service. Free Render web services sleep after 15 minutes without traffic and can take roughly a minute to wake. The dashboard retries nodes automatically and displays a waking state during that interval. We deliberately do not run synthetic keep-alive traffic: keeping two services active all month would exceed a free workspace's shared instance-hour budget. Use local Docker Compose for the full, durable six-node cluster.
- Default failover is not consensus. Deterministic promotion is useful for controlled environments but can have a split-brain window during a network partition.
- Raft-lite is opt-in. It adds elections, terms, vote and log checks, and majority commit tracking, but membership changes do not use Raft joint consensus.
- Persistence is local. A disk failure can lose that node's AOF and snapshot unless operators copy them elsewhere.
- This is deliberately modest in scale. The project prioritizes visibility and correctness under documented scenarios over throughput tuning or broad production guarantees.
packages/
node/ Store engine, persistence, protocol, replication, routing
cli/ Redirect-aware WebSocket client
dashboard/ Next.js cluster dashboard
benchmarks/ Throughput, failover, lag, and reshard benchmark runners
docs/
architecture.md System design and decisions
benchmarks.md Dated benchmark history
deployment.md Render and Vercel deployment guide
ci-cd.md CI/CD and release operations
pnpm build
pnpm test
pnpm lint
pnpm audit --audit-level moderateGitHub Actions builds and tests every package, audits dependencies, runs the real-process node integration suite, and starts the complete Compose topology for an end-to-end CLI write/read smoke test. Additional workflows run CodeQL, Gitleaks, scheduled chaos checks, benchmarks, and container-image releases.
