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
| Role | Holds | On the data path? |
|---|---|---|
nanocached-node | Cache entries (memory-bounded, LRU) | Yes — clients connect directly |
nanocached-discovery | Soft-state membership registry | No — consulted only at bootstrap and on refresh |
| SDK client | The node list, R, and open connections | Is 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:
- A node that is asked for a key it does not own (stale client
routing) answers
W; the client refreshes its node list and retries once. - A write whose primary just died fails its connection, which triggers the same refresh-and-retry — so writes recover as soon as discovery drops the dead node (bounded by the liveness timeout, seconds).
- Every heartbeat acknowledgment carries the current member list, so
the surviving nodes' own ring view follows that eviction within one
heartbeat interval too — a key re-homed from the dead node is served
by its new owner rather than answered
W. - That same ring-view update triggers re-replication:
rendezvous hashing's removal property means an eviction (or any leave)
can only ever promote a node into a key's top-R, never leave a
slot vacant, but nothing had sent that promoted node a copy — so at
R ≥ 2 the cluster silently ran one replica short after the
first failure, and a second failure on the sole remaining
copy lost the key. Once a survivor notices the drop, the highest-ranked
surviving old owner of each affected key streams it to the newly
promoted owner(s) as a put-if-absent handoff (
U's trailingAtoken — see Protocol), so a re-replication racing a newer client write can never regress it. This runs independently of, and is ordered against, a join in progress: a join's sweep of its own dead-copy marks waits for any in-flight re-replication to finish first (so it can't reclaim a survivor's only copy of a key re-replication hasn't delivered yet), and a re-replication that starts while a join is running waits for the join to finish first, since the join changes the ring again.
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:
- The joining node registers; discovery tells the existing nodes to migrate the keys the newcomer now owns.
- While migration runs, existing nodes keep serving those keys and forward concurrent writes to the joiner, so no update is lost.
- When every node reports completion, discovery promotes the joiner into the served node list; clients pick it up on their next refresh.
- 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:
- Every node and client is configured with the same replica list. Nodes heartbeat to all replicas; clients try their configured addresses in order.
- Losing any replica — including the first-listed primary — costs neither cache traffic nor client bootstrap. Only joins pause until the primary is back, because exactly one orchestrator keeps handoffs simple.
- A restarted replica answers
B(busy) during a startup grace while live members re-announce, so clients never bootstrap from a half-recovered list; SDKs skip a busy replica like a dead one. A join requested during that grace is held until it ends, for the same reason: its handoff must come from every member, not from the few that have re-announced so far.
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 |