Architecture

The clients do the clustering. Nodes store data, discovery servers track membership, and every routing decision happens in the SDK — so the data path has no extra hop.

The three roles

RoleHoldsOn the data path?
nanocached-nodeCache entries (memory-bounded, LRU)Yes — clients connect directly
nanocached-discoverySoft-state membership registryNo — consulted only at bootstrap and on refresh
SDK clientThe node list, R, and open connectionsIs the data path

An SDK bootstraps by connecting to any configured address. The server's handshake reply says whether it is a node or a discovery server; against a discovery server the client fetches the node list (which carries the replication factor R) and then talks to nodes directly.

Routing: rendezvous hashing

For each key, every node gets a score — fmix64(fnv1a(node-name) XOR fnv1a(key)) — and the nodes ranked by descending score are the key's owner list. The top-R nodes hold the key. A namespaced key hashes be32(namespace-length) || namespace || key in place of the bare key, so one namespace's entries still spread over every node and two namespaces' identical key names land on different ones; the default namespace keeps the bare form so existing placements survive an upgrade.

Rendezvous (highest-random-weight) hashing has a property that matters here: adding or removing a node never reorders the surviving nodes in any key's ranking. A new node takes exactly its share of keys and nothing else moves; a dead node's keys fall to the next owner that already holds a replica.

All six SDKs and the server implement this pipeline bit-for-bit and assert the same published test vectors, so a polyglot deployment routes identically from every language.

Replication

Writes fan out from the client to all R owners; the primary's result decides the operation, and a dead replica never fails a write. Reads go to the primary and fail over to the next owner only when the holder is unreachable.

Two self-healing behaviors close the failure gaps:

These healers share one consistency model: each acts on a snapshot of membership taken when it starts, and no two of them coordinate. A planned leave (the SIGTERM decommission — see Deployment) computes its handoffs the same way, from one roster fetch as the drain begins. Two membership changes landing inside one drain window — two nodes scaled in at once, or an eviction while a third drains — can therefore leave keys transiently below R: each leaver picks entrants from a roster that still lists the other, so a key whose owners included both comes out of the double-leave one handoff short. Nothing regresses (handoffs are put-if-absent) and reads keep landing on the surviving copies; a key that does lose its last copy is a cache miss, not corruption. There is no anti-entropy sweep: the gap closes per key as a later ring change re-ranks it, a client rewrites it, or its TTL turns it over — and it never opens at all when scale-in retires one node at a time.

Liveness is binary, though: a node that is slow but still answering is never evicted, and every request that touches it waits out its round trip — with R copies on N nodes, about R/N of them. The SDKs' hedged reads (opt-in) send a read to the next owner once the primary has been silent for a configured interval and take the first answer, so reads route around a slow owner; writes cannot, since every copy must be written.

Node identity is a random per-process name, decoupled from the network address — routing hashes the name, connections use the address — so a node restarting on the same address is correctly treated as a new member.

Joining a node: staged handoff

Adding a node must not drop the keys it takes ownership of, so joins are orchestrated by discovery in stages:

  1. The joining node registers; discovery tells the existing nodes to migrate the keys the newcomer now owns.
  2. While migration runs, existing nodes keep serving those keys and forward concurrent writes to the joiner, so no update is lost.
  3. When every node reports completion, discovery promotes the joiner into the served node list; clients pick it up on their next refresh.
  4. A source keeps the copies it handed off, marked dead, until discovery confirms the join (the joiner shows up in its heartbeat acknowledgment, or in the roster of the next join's migrate request); an abandoned join restores them instead, so no key is lost to a half-finished handoff.

Migrated-away entries are swept only after the join is confirmed, and the forwarding window stays open past completion to cover clients still routing with the older node list. Joins are serialized — one handoff at a time, cluster-wide — so a fleet started at once simply queues; a source that has finished one handoff accepts the next join's request while it is still forwarding for the previous one.

Discovery without a single point of failure

The registry is deliberately soft state: nodes announce themselves and heartbeat continuously, so a discovery server's knowledge is always rebuildable from the nodes alone. That allows N independent replicas with zero coordination — no consensus, no leader election, no shared storage:

Pair this with a supervisor that restarts a dead replica (systemd, Kubernetes) and two replicas make discovery effectively always available.

Membership carries a second, per-node secret alongside the shared one: each node registers with a random membership token, and the frames that skip a node's wrong-node check — a join's M, a decommission's or re-replication's U/u — must present the receiver's token, so holding the cluster-wide shared secret alone doesn't let a client forge them. The flip side is deliberate: any registered node can fetch every member's token from discovery (the T frame), because a draining node must authorize handoffs to every peer, and a survivor re-replicating after an eviction must authorize them to owners it never met at registration time. Compromising one node therefore yields enough to impersonate any node in the cluster — the tokens narrow who can speak the internal frames from "anything holding the shared secret" down to "cluster members", not further. Size the blast radius accordingly: keep the data and operations ports internal (see Deployment's security groups), and treat a node-level compromise as a cluster-level event.

The proxy tier (optional)

Proxyless client-side sharding holds one connection per node per client process — a node-side connection count that grows with fleet size, capped by each node's connection limit. For large fleets, nanocached-proxy terminates many client connections and routes on their behalf: it fetches the roster from discovery, performs the same rendezvous-hash routing the SDKs do (namespaced keys included), fans writes out to all R owners, retries a stale-view W after a roster refresh, and fans c/F to every member. To a client it looks like a single node that owns every key, so any SDK pointed at one address — and any bare protocol client — works against it unchanged, with no ring view or discovery client of its own. Proxies are stateless and scale horizontally: each registers itself with discovery, and a client in proxy mode asks discovery for the proxy list and picks one — the only address anything ever needs is discovery's (a LB/VIP in front of the proxies works too). Clients that prefer cluster-direct connections keep working as before. Every client connection shares one tagged, pipelined backend connection per node, so the node-side connection count is the proxy count — not the fleet size — with bounded per-client and per-backend in-flight windows providing fairness and backpressure.

One documented limit of “looks like a single node”: the proxy enqueues every request on its backend connection in the order the client sent it, but a request the proxy has to retry (a stale-view W after a roster refresh, a replica fallback) is re-enqueued afterwards, outside that order. A client that pipelines a G for a key and, without awaiting its reply, an S for the same key on the same connection can — if the G lands on a stale view during a ring change — receive the value its own S just wrote. That is the same outcome an SDK’s own W retry produces in cluster-direct mode and is not a consistency violation (the write was already on the wire when the read was answered), so it is left as is rather than serialising same-key requests per connection. A client that needs the read to strictly precede the write awaits the G reply before sending the S.

Known limits

A ring change stalls a node for a few tens of milliseconds per million keys. Re-replication, a staged handoff, and a decommission each begin by taking a snapshot of every live key the node holds, and that walk runs on the node’s single cache task, so requests that arrive during it wait. Measured in release builds at about 25 ms per million live keys (5 ms at 200k, 52 ms at 2M); at the default --max-memory of 256 MiB a node holds on the order of 2–3M small entries, so one membership change costs a single pause of a few tens of milliseconds — the same order as the expiry sweep’s own periodic refill, and only on membership changes, never on the request path. Bounding it would need an ordered per-namespace key index that every write maintains; that standing cost on the hot path was judged not worth removing a rare, bounded pause. Size the node count so a single node’s keyspace stays in the low millions if even that pause matters to you.

Design decisions

The main choices, in summary:

Decision
Client-side hashing with a lightweight discovery server
Shared-secret authentication via environment variable
TLS via rustls, required once configured
Server type in the auth response → one unified connect
Staged node join with orchestrated data handoff
Node identity decoupled from network address
Discovery HA via soft-state replicas and announces
Client-side replication via rendezvous hashing