Distributed Systems: Cluster Mode¶
Pion implements Redis Cluster protocol compatibility, enabling use with cluster-aware clients such as
Valkey GLIDE, redis-py cluster mode, Lettuce, and Jedis.
Two cluster modes are available: single-node cluster (one Pion instance owns all 16,384 hash slots)
and multi-node cluster (multiple Pion instances with gossip health monitoring and WAL replication).
Multi-node is a preview. Slot migration, gossip, failover and replication exist and are tested, but production multi-node deployment is not supported yet. Replication covers worker 0's keyspace only, so run replicated servers with
-w 1.
Implemented Commands¶
CLUSTER INFO¶
Returns cluster health metadata in key:value\r\n format (Redis 7 compatible):
CLUSTER INFO
→ cluster_enabled:1
cluster_state:ok
cluster_slots_assigned:16384
cluster_slots_ok:16384
cluster_slots_pfail:<live pfail count>
cluster_slots_fail:<live fail count>
cluster_known_nodes:<1 + peer_count>
cluster_size:<1 + peer_count>
...
cluster_slots_pfail and cluster_slots_fail reflect live gossip health state updated by
the background PING thread.
CLUSTER MYID¶
Returns the node's 40-character hex node ID (derived from FNV-1a of host:port):
CLUSTER KEYSLOT \<key>¶
Returns the CRC16 hash slot for a key (0–16383):
Hash tags ({...}) are supported: only the substring inside {} is hashed, matching Redis Cluster spec.
CLUSTER NODES¶
Returns one line per node with live health status:
CLUSTER NODES
→ "<id> <host>:<port>@<bus-port> myself,master - 0 0 1 connected 0-16383"
"<id> <peer-host>:<peer-port>@<bus-port> master - 0 0 1 connected|pfail|fail <slots>"
Peer health values: connected (online), pfail (possible failure — gossip PING timeout),
fail (confirmed failure — repeated PING failures).
CLUSTER SLOTS / CLUSTER SHARDS¶
Returns slot-to-node mapping. Same format as previous; health reflected via peer entries.
CLUSTER MEET \<host> \<port>¶
Registers a new peer at runtime (no restart required):
- Adds peer to
ClusterState.peers[]with evenly distributed slot range - Gossip starts health-probing the new peer immediately
cluster_epochis incremented
CLUSTER FORGET \<node-id>¶
Removes a peer by its 40-char hex node ID:
CLUSTER REPLICATE \<node-id>¶
Marks this node as a replica of the given peer. Used in conjunction with
--cluster-replica / --cluster-primary-host at startup, or applied dynamically:
CLUSTER FAILOVER [FORCE]¶
Promotes a replica to primary (clears is_replica flag, increments cluster_epoch):
CLUSTER FAILOVER → +OK (graceful — waits for replication sync)
CLUSTER FAILOVER FORCE → +OK (forced — skips health checks, immediate promotion)
The FORCE variant bypasses gossip health verification and promotes immediately. Use when the primary is unreachable and graceful failover would hang.
CLUSTER SETSLOT \<slot> \<subcommand> [node-id]¶
Manages slot ownership during migration:
CLUSTER SETSLOT 1234 IMPORTING <source-node-id> → +OK
CLUSTER SETSLOT 1234 MIGRATING <target-node-id> → +OK
CLUSTER SETSLOT 1234 STABLE → +OK
CLUSTER SETSLOT 1234 NODE <node-id> → +OK
Used for zero-downtime slot migration between nodes. IMPORTING/MIGRATING set transitional state; NODE finalizes ownership; STABLE cancels migration.
CLUSTER RESET¶
Resets replica state (promotes to standalone primary):
HELLO [protover]¶
RESP2/RESP3 negotiation (required by Valkey GLIDE). Pion always responds with proto:2
to prevent GLIDE from enabling RESP3 features:
CRC16 Hash Slot Computation (src/network/cluster.mojo)¶
ClusterState implements the full Redis CRC16 polynomial (0x1021) using a 256-entry
precomputed lookup table. The keyslot() method handles hash tags per spec:
if '{' found before '}' and '}' comes after '{':
hash only the bytes between '{' and '}'
else:
hash the entire key
owns_slot(slot) returns True for all slots in single-node mode (owns all 16,384).
MOVED Redirect¶
When a cluster-aware client sends a command for a key that maps to a slot owned by a peer, Pion
returns a MOVED redirect:
In single-node mode all slots are owned locally so MOVED is never generated for a correctly configured
client. MOVED responses are generated when peer_nodes is configured with multiple hosts and the
target slot maps to a peer.
Gossip: Peer Health Monitoring¶
Pion uses lightweight TCP-based health probing (not SWIM UDP) via a background C pthread
(PionGossipBlock in src/ffi/uring_wrap.c).
Mechanism:
- Every
gossip_ping_ms(default 1000ms), the gossip thread iterates all peers - For each peer: opens a TCP connection to
peer_host:peer_port, sendsPING\r\n, expects+PONGwithin 250ms - On success:
peer_health[i] = 0(online), resets consecutive fail counter - On failure: increments consecutive fail counter
>= pfail_thresholdconsecutive failures (default 5) →peer_health[i] = 1(pfail)>= fail_thresholdconsecutive failures (default 15) →peer_health[i] = 2(fail)
Implementation: peer_health is an InlineArray[UInt8, 16] inside the heap-allocated
ClusterState. The gossip thread writes directly through a pointer, eliminating any
per-tick C call from Mojo. Single-byte writes are inherently atomic on x86_64 and ARM64.
Architecture: Only worker 0 starts the gossip pthread. All workers read the same
ClusterState health array.
Auto-Failover¶
Worker 0 monitors gossip health state and automatically promotes a replica to primary when the
current primary is detected as FAIL (confirmed failure after fail_threshold consecutive PING
failures, default 15 = ~15 seconds of unreachability).
Mechanism:
- Worker 0 checks
peer_health[]each event-loop tick - When a peer transitions to
FAILstate, worker 0 triggers automatic promotion - The replica clears
is_replica, incrementscluster_epoch, and begins accepting writes - Gossip continues monitoring — if the old primary recovers, it must be manually reconfigured
Performance: Measured failover time is 463ms from primary failure detection to replica
accepting writes (tested in tests/test_failover.py). The ~5s detection window is dominated
by the gossip fail_threshold (15 consecutive 1-second PINGs). Timing may vary on different hardware.
Testing:
This test starts a primary and replica, kills the primary, and verifies the replica auto-promotes and begins serving reads/writes within the expected timeframe.
Manual override: CLUSTER FAILOVER (graceful) or CLUSTER FAILOVER FORCE (immediate, skips
health checks) can be issued at any time to trigger manual promotion without waiting for auto-detection.
INFO REPLICATION¶
Returns replication status in key:value\r\n format (Redis-compatible):
INFO REPLICATION
→ # Replication
role:master (or role:slave if --cluster-replica)
connected_slaves:0
master_repl_offset:0
In replica mode, additional fields report the primary host/port and replication lag. The role
field reflects the current state — after CLUSTER FAILOVER, a former replica reports role:master.
WAL Replication¶
Primary→replica streaming of Write-Ahead Log entries via a background C pthread.
Primary side (PionReplPrimary):
- Listens on port + 10000 (e.g., port 1974 → replication port 11974)
- Accepts replica connections, validates handshake (PION-REPL 1.0\r\n → +OK\r\n)
- Background thread polls wal.tail_offset every 1ms and streams new bytes to each replica
- Handles replica reconnects transparently; restarts from offset 0 on reconnect
Replica side (PionReplReplicaBlock):
- Connects to primary_host:primary_repl_port (primary's port + 10000)
- Receives WAL bytes into a 4MB ring buffer (C-side, protected by mutex)
- Mojo drains the ring each event-loop tick (kqueue and io_uring paths) via
pion_repl_replica_drain() — zero-copy into a pre-allocated drain buffer
- apply_wal_entries() parses and applies entries: SET → keyspace.set(), DEL → keyspace.remove_generic()
WAL entry format (same as local WAL, see src/io/wal.mojo):
[4B entry_len LE][1B cmd_id][4B key_len LE][key bytes][4B val_len LE][val bytes]
cmd_id: 1=SET, 2=DEL
Since the WAL does not wrap in Phase 1, replication is a linear byte stream from offset 0.
Enabling Cluster Mode¶
Single-node cluster (GLIDE-compatible)¶
./pion-server --cluster --cluster-host 127.0.0.1 --cluster-nodes 127.0.0.1:1974
# Verify
redis-cli -p 1974 CLUSTER INFO | head -3
redis-cli -p 1974 CLUSTER MYID
redis-cli -p 1974 CLUSTER KEYSLOT mykey
redis-cli -p 1974 HELLO 2
Multi-node cluster (primary + replica)¶
# Primary node (port 1974, replication listener on 11974)
./pion-server --cluster --cluster-host 10.0.0.1 \
--cluster-nodes "10.0.0.1:1974,10.0.0.2:1975"
# Replica node (connects to primary for WAL streaming)
./pion-server --cluster --cluster-host 10.0.0.2 -p 1975 \
--cluster-nodes "10.0.0.1:1974,10.0.0.2:1975" \
--cluster-replica \
--cluster-primary-host 10.0.0.1 \
--cluster-primary-port 1974
Runtime peer management¶
# Add peer at runtime (no restart)
redis-cli -p 1974 CLUSTER MEET 10.0.0.3 1976
# Remove peer
redis-cli -p 1974 CLUSTER FORGET <node-id>
# Promote replica to primary
redis-cli -p 1975 CLUSTER FAILOVER
Configuration flags¶
--cluster Enable cluster mode
--cluster-host <ip> Advertised IP in CLUSTER NODES output
--cluster-nodes <list> Comma-separated "host:port,..." including self
--cluster-replica This node replicates from primary
--cluster-primary-host <h> Primary's advertised host
--cluster-primary-port <p> Primary's main port (replication port = p + 10000)
--gossip-ping-ms <ms> Gossip PING interval (default: 1000ms)
ClusterConfig (src/common/config.mojo):
ClusterConfig:
enabled: false (opt-in via --cluster)
my_host: "127.0.0.1"
peer_nodes: ""
is_replica: false
primary_host: ""
primary_port: 0
gossip_ping_ms: 1000
pfail_threshold: 5 (consecutive PING failures → pfail)
fail_threshold: 15 (consecutive PING failures → fail)
Implementation (src/ffi/uring_wrap.c)¶
All gossip and replication logic runs as C pthreads to avoid blocking Mojo's event loop (macOS Mojo green threads don't yield at OS syscalls; pthreads are full OS threads).
| C struct | Role |
|---|---|
PionGossipBlock |
Gossip pthread state; writes peer_health[] directly |
PionReplPrimary |
Primary listener: accept + stream thread |
PionReplReplicaBlock |
Replica receiver thread + 4MB ring buffer |
Mojo calls per event-loop tick (both kqueue and io_uring paths):
- pion_repl_replica_drain(handle, buf, max) → drain ring into apply buffer
Client Compatibility¶
| Client | Status |
|---|---|
redis-cli (single-node) |
✅ All 9 CLUSTER subcommands verified |
redis-py cluster mode |
✅ CLUSTER SLOTS/SHARDS format correct |
| Valkey GLIDE | ✅ HELLO 2 + single-node mode |
| Lettuce / Jedis | ✅ CLUSTER NODES format with live health |
Architecture Notes¶
Worker 0 exclusivity: Only worker 0 starts gossip and replication pthreads. This is necessary because Pion caps workers at 4 P-cores on macOS (main.mojo policy) — all threads are needed for the event loop. pthreads bypass this limit (they are scheduled by the OS independently).
Shared-nothing consistency: In shared-nothing mode, each worker has its own keyspace.
Replication only applies to worker 0's keyspace. Use -w 1 (the default) for fully replicated
workloads; -w N > 1 refuses to start without --independent-workers.
macOS cluster testing: Use redis-cli in cluster mode or single-node GLIDE.
Full multi-node testing requires two separate Pion instances on separate ports/IPs.
Cluster Protocols¶
1. SWIM Gossip Protocol¶
Replaced basic TCP PING with UDP SWIM protocol (pion_gossip_start_swim() in uring_wrap.c).
Features: - Random target selection — each round probes a random peer (not sequential) - Indirect ping — if direct ping fails, asks K=3 random delegates to probe the target - State exchange — SWIM messages carry node_id, epoch, role, slot_count, pfail bitmap - Suspicion protocol — SUSPECT/ALIVE/CONFIRM messages for graceful failure detection - TCP fallback — each round also does a TCP PING for reliability - Quorum FAIL promotion — pfail→fail requires majority of nodes reporting PFAIL
Ports: SWIM UDP = main port + 20000 (e.g., 21974 for port 1974).
Usage: Enabled by default when --cluster is set. The GossipManager.start(use_swim=True) parameter
controls SWIM vs legacy TCP. Legacy TCP is still available with start(use_swim=False).
2. Raft Metadata Consensus¶
Lightweight Raft for cluster metadata — leader election and topology change commits. NOT for data replication (WAL handles that).
Features: - Leader election — randomized timeout (150-300ms), RequestVote RPC, majority vote - Heartbeat — leader sends AppendEntries every 50ms to maintain authority - Term tracking — monotonic term counter prevents stale leaders - Log replication — topology changes committed via Raft log (slot migration, failover) - Split-brain prevention — only Raft leader can commit topology changes
Ports: Raft TCP = main port + 30000 (e.g., 31974 for port 1974).
C API:
void* pion_raft_create(const char* node_id, int raft_port);
void pion_raft_set_peer(void* block, int i, const char* host, int port);
int pion_raft_start(void* block);
int pion_raft_is_leader(void* block);
int pion_raft_get_state(void* block); // 0=Follower, 1=Candidate, 2=Leader
3. vNode Migration (Slot Migration)¶
Full slot migration with CLUSTER SETSLOT + MIGRATE + ASK/ASKING protocol.
Commands:
- CLUSTER SETSLOT <slot> MIGRATING <node-id> — mark slot for export
- CLUSTER SETSLOT <slot> IMPORTING <node-id> — mark slot for import
- CLUSTER SETSLOT <slot> NODE <node-id> — finalize slot ownership
- CLUSTER SETSLOT <slot> STABLE — cancel migration
- MIGRATE host port key|"" db timeout [COPY] [REPLACE] [KEYS key ...] — push keys to target
- ASKING — per-fd flag for importing slot acceptance
- DUMP key — serialize key to binary format
- RESTORE key ttl payload [REPLACE] — deserialize and store
- CLUSTER GETKEYSINSLOT slot count — find keys in a slot
- CLUSTER COUNTKEYSINSLOT slot — count keys in a slot
ASK/MOVED redirect flow:
1. Client sends command for migrating slot
2. If key exists locally → serve it
3. If key not found → -ASK <slot> <target> redirect
4. Client sends ASKING to target, then retries command
5. Target serves the command (ASKING flag clears after one use)
MIGRATE internals: Uses dump_key() to serialize, TCP to target, RESTORE to deserialize.
Supports COPY (keep local), REPLACE (overwrite existing), multi-key via KEYS keyword.
4. Cross-Worker WAL Aggregation¶
Workers 1..N forward WAL entries to worker 0 via per-worker ring buffers (4MB each). Worker 0 drains all rings each event-loop tick and applies entries to the master WAL.
C API:
void* pion_wal_agg_create(int worker_count);
int pion_wal_agg_push(void* block, int worker_id, const uint8_t* entry, int entry_len);
int pion_wal_agg_drain(void* block, int worker_id, uint8_t* out_buf, int max_bytes);
void pion_wal_agg_stop(void* block);
Ring buffer: 4MB per worker, mutex-protected, supports wrap-around. Worker 0 can drain entries from any worker's ring for replication.
5. Replica-Reads¶
Per-connection READONLY/READWRITE enforcement on replica nodes.
Commands:
- READONLY — enable read-only mode on this connection (allows reads on replicas)
- READWRITE — restore default mode (reject reads on replicas)
Write rejection:
- Non-READONLY connections on replicas: -MOVED to primary
- READONLY connections on replicas: reads pass through, writes get -READONLY error
- Write commands identified by exclusion: GET/MGET/PING/EXISTS/HGET/LLEN/LRANGE/DBSIZE/
GETBIT/BITCOUNT/PFCOUNT/CLUSTER/CONFIG/INFO/HELLO/READONLY/READWRITE/ASKING/QUIT/ECHO
are reads; everything else is a write.
CLUSTER NODES output: Replicas now show myself,slave instead of myself,master.