Architecture
Design and consistency model. See vs Redis for how this compares to a single-writer system, and Go Backend for the package-level breakdown.
Peer-to-peer, no coordinator
Every node runs a full MeshStore — there is no leader, no coordinator process, and no single node whose failure takes down writes. Reads and writes are served locally by whichever node receives them; consistency across nodes is achieved after the fact by gossip, not enforced before the fact by consensus. This is the central structural trade-off of the whole system: availability and write latency are prioritized over strong consistency, in the same family as Cassandra or Riak rather than a single-writer system like standalone Redis.
Conflict resolution: CRDTs + Lamport clocks
State that can be modified from multiple nodes concurrently is represented as a Conflict-free Replicated Data Type, chosen so that merging two nodes’ versions of the same value is always well-defined and always converges to the same result regardless of merge order:
GCounter(grow-only counter) — backs rate limiting. Each node tracks its own increments; the logical value is the sum across all nodes. Merging twoGCounters is just taking the per-node max, so merging is commutative, associative, and idempotent — safe to apply out of order or more than once.GSet(grow-only set) — backs replay protection. A nonce is either in the set or not; merging is a set union.- Ordering, where it matters (e.g. sorted set score updates, stream entries), uses Lamport clocks — a logical counter that establishes a causal
happens-beforeordering across nodes without requiring synchronized wall-clock time.
Everything built on top of these primitives inherits the same convergence guarantee: apply updates in any order, on any subset of nodes, and every node ends up in the same state once gossip has propagated. What you give up is strong consistency — a read immediately after a write on a different node can briefly see stale state until the next gossip round.
Gossip and state sync
coordination/gossip.go implements real peer-to-peer gossip, riding on the same HTTP API every SDK uses rather than a separate wire protocol: each node periodically picks a random peer, GETs its /internal/state, and merges the result into local state via MeshStore.MergeState, which applies each primitive’s own CRDT merge (GCounter.Merge, GSet.Merge) so the merge is commutative, associative, and idempotent regardless of round order. A node joins a cluster with -join <peer-http-addr>, which does a one-time POST /internal/peers/join handshake that also returns the peer’s own peer list, so joining through a single existing member discovers the rest of an already-formed cluster.
This replaces what used to be here: for most of this project’s history, performGossip was a literal no-op placeholder (its own comment said as much), and there was no way to even join two real node processes into a cluster. That’s now fixed for the three original CRDT-backed primitives — verified against three separately launched OS processes, not just tests, including a concurrent-write scenario and correct GCounter aggregation across the cluster.
Current limitation, stated plainly: Persistence is the one feature group that doesn’t fit this pattern at all — its WAL/snapshot files are a per-node durability mechanism for the state above, not data with a CRDT of its own to replicate. Metrics is a partial exception (see below) rather than a clean yes/no. Ranking itself is genuinely stateless (nothing persists between calls), but a new registry built on top of it — named, reusable ranking configs — does replicate; see below. Everything else — the original three primitives (rate limiting, replay protection, cache) plus, as of this session, Sorted Sets, Streams, Pipelines, Search, Job Queues, Pub/Sub, Transactions, WASM Scripting, and named Ranking configs — replicates across nodes today.
The real scalability ceiling: this is a full-state transfer every round, not a delta. GetState() returns the entire replicated state on every gossip round, regardless of how much actually changed since the last one — bandwidth and serialization cost scale with total data volume, not with the size of the actual change. Fixing that properly means per-key delta tracking across all thirteen replicated primitives (each needs a monotonic version a peer can say “I already have everything up to”), which is a genuinely large redesign with real correctness risk to already-working, extensively-tested replication – out of scope for a single pass. What ships instead is a real, measured, safe partial mitigation: handleInternalState now gzip-compresses its response whenever the requester’s Accept-Encoding says it can decode gzip, which every client this codebase makes can (gossip’s own client, PeerManager’s health-check client, the join-request client), automatically, via Go’s net/http.Transport – no client-side changes were needed at all. Measured against a representative populated state (2000 cache keys, 500 sorted-set members, 500 stream entries): ~476KB of JSON compresses to ~39KB, about 8% of the original; verified live against two real processes with 1000 populated cache keys, /internal/state dropped from 165,374 to 18,056 bytes (~89%), with gossip still converging correctly afterward. This doesn’t fix the O(total data) shape of the problem, but it substantially reduces the actual bytes crossing the network every round, which is often the practical bottleneck.
There used to be a second thing living in this file’s history worth noting: coordination/state_sync.go (StateSync/MerkleTree) looked, on the surface, like it might already be the delta-sync infrastructure described above. It wasn’t, and has been removed. It was never instantiated anywhere outside its own unit tests; its merge logic covered only 3 of the 13 primitives that now replicate and used a blind overwrite instead of any real CRDT rule (weaker than the real MergeState already live); and its “Merkle tree” wasn’t actually usable for diffing at all – buildMerkleTree iterated a Go map (non-deterministic order) to build the tree, so it couldn’t produce a structure comparable across two different nodes even in principle, and no tree-diff/comparison traversal existed anywhere in the file, only a single whole-tree root-hash equality check. Wiring it in would have been a regression, not the scalability fix it appeared to be.
Sorted Sets was the first of the ten feature groups added after the original three primitives to gain replication, and the easiest: SortedSet.Merge already implemented a real (score, timestamp, node) CRDT conflict resolution for local testing purposes, so extending it to gossip was a matter of wiring a wire-format snapshot (SortedSet.Snapshot/MergeSnapshot, including tombstoned members so deletes propagate, not just additions) through GetState/MergeState, not designing new conflict-resolution logic.
Streams was the second, and needed one real fix first: entry IDs were <timestamp>-<sequence>, where sequence is a plain per-Stream-instance counter — two different nodes’ independent Stream objects for the same stream name both start at 0 and increment with no coordination between them, so two nodes producing entries in the same millisecond (ordinary under real write volume) could produce the identical ID for two different entries. Harmless single-node, but fatal for a union-merge, since the ID is exactly the identity a merge uses to decide “have I seen this already”. Fixed by adding the node to the ID (<timestamp>-<sequence>-<node>) — safe because nothing in the codebase parses an ID’s structure, every use is an opaque string lookup, confirmed via a repo-wide search before making the change. With globally-unique IDs, merging is a plain set union (entries are immutable once appended, unlike Sorted Sets or Cache, so there’s no per-entry conflict to resolve) — the one thing that needs care is that Range/GetFirst/GetLast are positional over the entry slice, not ID-derived, so merged entries are inserted in the correct chronological position, not just appended.
Pipelines was the third, and structurally close to Cache: a named registry with full-value replacement, no per-key merge within one pipeline’s own contents. Pipeline.Created was already re-stamped on every registration (not just the first), so it doubled as an LWW version with only one new field needed (Node, the tiebreaker). One real gap, documented rather than closed: DeletePipeline is a hard local delete with no tombstone, so a deleted pipeline gets silently re-introduced by the next gossip round from any peer that still has it — unlike Sorted Sets, pipeline deletion does not replicate.
Search was the fourth, and the same shape as Pipelines: Document gained a Node field alongside its existing (but previously never actually set) Timestamp, giving each indexed document a (Timestamp, Node) LWW-register version. Before wiring replication, a real pre-existing bug was fixed: IndexDocument never removed a document’s previous BM25 contribution before re-indexing the same ID, so calling it twice for the same ID — an ordinary update, or, once gossip repeatedly re-indexes a merged document, a routine occurrence — silently double-counted docCount, duplicated postings, and inflated totalTerms, corrupting every document’s BM25 score (docCount feeds IDF, totalTerms/docCount is avgDocLen, both are direct terms in the scoring formula). Fixed by de-indexing a document’s old contribution before re-indexing it, and building MergeSnapshot on top of IndexDocument itself (not a raw map write) so an adopted peer document is always correctly re-indexed. Same gap as Pipelines, not closed: DeleteDocument is a hard local delete with no tombstone, so a deleted document is silently re-introduced by the next gossip round from any peer that still has it. Verified live against two real server processes in both directions.
Job Queues was the fifth, and needed a small cleanup first: Job carried a mu sync.RWMutex field that was never actually locked or unlocked anywhere (confirmed by grep before touching it) — dead weight that would have blocked the value-copy Snapshot() needs (Go’s vet flags copying a struct containing a mutex), so it was removed and replaced with the Node field every other feature’s LWW version needed. JobQueue.Snapshot()/MergeSnapshot() replicate its append-only job log via a (UpdatedAt, Node) LWW-register comparison per job ID, placing an adopted job into the correct derived list (PendingJobs/ProcessingJobs) for its merged status. This feature’s gap is more consequential than the others’ documented gaps: it has no delete to worry about, but merging state doesn’t provide exclusive claims across nodes. JobQueue.GetNextJob’s CAS-like check only prevents two workers on the same node from claiming a job — two different nodes can each claim the same pending job locally before a gossip round tells either about the other’s claim, so a job can be processed more than once across the cluster. Job processing is at-least-once across the mesh, not exactly-once. Verified live against two real server processes in both directions: enqueued on node1, converged pending to node2, claimed on node2, and the claim (status and worker ID) converged back to node1.
Pub/Sub was the sixth, and needed the same real fix Streams needed: Message.ID was topic+nanotime only, which two node processes publishing to the same topic in the same nanosecond could collide on, silently dropping one message from a union-merge – fixed by folding the node ID into it. Deliberately scoped narrower than the other five, though: PubSubBroker.Snapshot()/MergeSnapshot() converge each topic’s message history (a set union by Message.ID, re-sorted by timestamp and trimmed to the existing history cap) across nodes, so GetMessageHistory/GetTopics/GetStats see the full cluster picture — but a Subscriber’s Channel is in-process memory with no serialization story across a gossip round, so merged messages are not pushed into live subscriber channels. Subscribe and Poll for a given subscriber ID must land on the same node, a pre-existing constraint gossip does not change either way.
Transactions was the seventh, and needed a real durability fix first, unrelated to replication itself: CommitTransaction applied a transaction’s queued Set operations directly to ms.cache but never called ms.persistence.LogOperation, unlike the plain Set() path — so a committed transaction’s writes were silently lost on the next process restart even though an equivalent direct Set to the same key would have survived it. Fixed by logging with the same convention Set() and MergeState’s cache adoption already use. Transaction gained UpdatedAt/Node as its LWW-register version (Created is stamped once and never changes, so it can’t double as a version the way Pipeline’s Created could). Transaction.mu is actively used (unlike Job’s dead mutex, which could just be deleted), so Snapshot/MergeSnapshot copy its fields individually into fresh Transaction values rather than dereferencing the original. Same at-least-once-style gap as Job Queues: this converges transaction metadata (status, queued operations) but not cross-node atomicity — AddTransactionOperation and CommitTransaction for one transaction ID must land on the same node, or the second call can fail with “transaction not found” until the next gossip round. A committed transaction’s actual Set effects don’t depend on this at all — they already replicate via Cache’s own existing LWW merge, since they land in the same ms.cache map Cache gossip already reads. Verified live against two real server processes, including killing and restarting one node alone to confirm the WAL fix.
WASM Scripting was the eighth, and structurally different from every prior feature: there is no compiled artifact to gossip. wazero.CompiledModule can’t cross a process boundary and isn’t safe to trust as opaque bytes from the wire, so only a script’s Go source and version metadata (Node added alongside the existing Compiled timestamp) replicate, and adopting a peer’s version means actually invoking the TinyGo compiler locally — the same real, multi-second cost Compile’s own doc comment already describes, not a cheap struct swap. WasmEngine can be nil on a node where TinyGo wasn’t found at startup — both GetState and MergeState guard on that, and a cluster mixing nodes with and without TinyGo simply never converges WASM scripts onto the nodes lacking it. Verified live against two real server processes: node2, which had never seen the script’s name before, independently recompiled the peer’s source and served correct output, proving it wasn’t just copying a cached artifact.
Security fix, found by a targeted review after this session’s replication work: MergeSnapshot originally used the same (Compiled, Node) LWW-register comparison every other feature’s merge uses, letting a peer’s “newer” claimed version overwrite an existing local script. That’s fine for inert data (a cache value, a job payload) that does nothing until a separately-authenticated caller acts on it — but WASM’s merge invokes the compiler as a direct side effect of the merge itself, purely via the gossip path (/internal/state, gated only by X-Cluster-Secret, a materially weaker and separately-distributed credential from the X-API-Key that gates script compilation everywhere else). Allowing overwrite meant a node holding only the cluster secret could silently swap what an existing, previously-trusted script name resolves to, by gossiping a fabricated future Compiled timestamp — and a later /script/execute call by a legitimate API-key holder would then run different code than they registered. Fixed by making MergeSnapshot insert-only: a peer’s script is compiled and adopted only for a name this node has never seen before; an existing name can never be replaced via gossip, no matter what version metadata a peer claims, only through the local, API-key-gated Compile path. The accepted tradeoff: a script update no longer propagates via gossip, only its first appearance does — an operator updating a script must still do so directly on every node. Verified live: registered a script under a name on one node, then registered a different, later-timestamped script under the identical name on a peer (simulating the hijack) — confirmed the first node kept running its own original version.
Metrics was the ninth, and the odd one out: rather than a registry of named entries with an LWW-register version, its 13 monotonic counters (consume/seen/get/set totals, hits/misses, evictions, gossip in/out/errors) are a natural fit for the grow-only-counter CRDT already used for rate limiting (core.GCounter) — each node reports only its own current value, and merging takes the max per (metric, node) pair, so a count can only grow, never regress or double-count regardless of gossip order or repeated merges of the same state. Latency percentiles (p50/p99) are deliberately excluded: computing them independently on different nodes’ sample sets and combining the numbers afterward doesn’t produce a meaningful cluster percentile — a real statistical limitation, the same reason Prometheus itself aggregates raw histogram buckets rather than combining pre-computed percentiles, and this project’s ring-buffer sampling doesn’t preserve buckets. The merge is inlined into MeshStore.MergeState rather than routed through core.GCounter itself, since GCounter.Increment takes a mutex per call and Metrics’ existing lock-free atomic counters are on the hot path of every single Consume/Seen/Get/Set — converting them would have been a real performance regression for a nice-to-have observability feature. A new GetClusterMetrics (GET /metrics/cluster) reports each counter’s cluster-wide sum plus a per-node breakdown; the existing GetMetrics//metrics is untouched and still reports only local activity. Verified live against two real server processes: 3 consume calls on one node and 1 on the other converged to a consume_total of 4 with the correct per-node split, while each node’s own local /metrics kept reporting only its own count.
Ranking configs was the tenth, and a genuinely new feature rather than a fix: Rank itself has no state to replicate (a fresh Ranker is built from caller-supplied arguments on every call, nothing persists), so there was nothing to point replication at. What does exist now is ranking.Registry, a named (strategy, boosts) pair a client registers once via RegisterRankingConfig and references by name via RankWithConfig on every later call — the same shape Pipelines gives operation sequences, right down to sharing its exact gap: DeleteRankingConfig is a hard local delete with no tombstone, so a deleted config is silently re-introduced by the next gossip round from any peer that still has it. Verified live against two real server processes: registered a boosted config on one node, confirmed the other – which had never seen it – correctly ranked the boosted item first after gossip converged.
Persistence is the one feature group that remains entirely out of scope for gossip replication itself, for a fundamentally different reason than “not designed yet” — see the limitation paragraph above. Persistence’s own snapshot/restore mechanism, separately, now covers all ten of the other feature groups’ local durability, ranking configs included — see Storage layout below.
Cache is a real per-key LWW-register CRDT: each entry carries the wall-clock time it was written and its writer’s node ID, and MergeState adopts a peer’s entry for a key only when it’s strictly newer than the local one (the writer’s node ID breaks an exact tie) — two nodes concurrently writing the same key converge to whichever write actually happened later, on every node, regardless of gossip order. This also required separating “when this node wrote a WAL entry” from “the cache entry’s actual CRDT version”: a value learned via gossip carries its original writer’s timestamp for versioning purposes, which can be earlier than this node’s own most recent snapshot, while the WAL entry recording it is still logged at “now” for the WAL’s own recovery-cutoff bookkeeping — conflating the two made a gossip-learned value vanish on the receiving node’s own next restart, reverting it to that node’s older local write (confirmed live against two real server processes before and after the fix).
A bounded mitigation for the Job Queue / Transaction race window, not a fix. Both features’ documented gaps (above, and in the Job Queues/Transactions replication paragraphs) come from the same root cause: the only path to cross-node visibility was waiting for the next periodic pull-gossip round (5s by default), so two nodes could act on the same job or transaction before either learned about the other. Fixing this for real means consensus — genuinely out of scope. What’s real and shipped: handleInternalState now also accepts POST (gzip-aware, same as GET) as a push target, and GossipCoordinator.PushState sends this node’s current state to one peer (the same healthy-preferring selection performGossip’s pull already uses) instead of waiting to be pulled from — one peer, not a fan-out, so it costs the same as one ordinary gossip round. MeshStore.ClaimJob and CommitTransaction trigger this push on success, debounced and rate-limited (a size-1 buffered channel drained by one consumer loop enforcing a 200ms minimum gap) so a burst of concurrent claims can’t turn this into a flood of full-state requests — a test confirmed 200 concurrent triggers collapsing into 2 actual pushes. Verified live: a job claim reached the peer in 45ms against a 5s default interval, roughly a 100x reduction in the race window — not an elimination of it, since two nodes claiming at nearly the same instant can still both push before either arrives.
This closes yet another instance of the “looks wired, does nothing” pattern found repeatedly this session: queue.JobManager.broadcastJobState was already being called from ClaimJob/CompleteJob, with its own comment reading “This simulates gossip protocol… For now, it’s a no-op.” It stays a no-op in queue itself (that package has no business depending on coordination/HTTP — wrong layer for it), but its evident intent is now actually implemented, correctly layered, at the MeshStore/GossipCoordinator level instead.
Cluster membership and failure detection
coordination/peer_manager.go and coordination/failure_detector.go (the latter now removed) were both, until this session, exactly the kind of code this document has flagged before: real-looking, fully unit-tested in isolation, and never instantiated anywhere outside their own constructors. PeerManager.performHealthCheck’s own comment admitted it – “Simulate health check - in production, this would be actual HTTP/gRPC calls” – and api/health.go’s HealthChecker.GetReadinessStatus() hardcoded every check to true regardless of any real state.
PeerManager is now wired into GossipCoordinator as the single source of truth for peer health (FailureDetector was removed rather than also wired in – it tracked the identical failure/recovery concept with a different, no-more-capable API, and using both would just mean two overlapping subsystems fed the same signal). Two real signals now feed it: performGossip records a genuine success or failure against PeerManager for every actual HTTP round it makes to a peer, and performHealthCheck independently does a real GET <peer>/health on its own fixed schedule (the same unauthenticated endpoint load balancers already use) – this second, independent check matters because gossip only ever contacts one randomly chosen peer per round, so with more than a couple of peers, relying on gossip alone to notice a failure (or a recovery) would leave most peers unchecked for arbitrarily long stretches. performGossip also now prefers a peer PeerManager currently considers healthy, falling back to trying any known peer if none are healthy, so a total outage doesn’t permanently prevent ever attempting a peer again once it might have recovered.
Wiring this in immediately surfaced a real bug: PeerManager’s constructor eagerly created its health-check time.Ticker, which panics on a non-positive interval – a legitimate pre-existing pattern in this codebase (a GossipCoordinator built for tests that only exercise HTTP handlers and never call Start, e.g. api/auth_test.go). Fixed by creating the ticker lazily inside Start(), the same “only ever ticks if Start is actually called” shape gossipLoop’s own ticker already has.
api/health.go’s HealthChecker is real now too, split into the two checks a real orchestrator (Kubernetes or otherwise) actually wants distinguished: liveness (GET /livez) stays healthy as long as the process is up and responding at all, deliberately not factoring in peer connectivity – an orchestrator that kills and restarts every node it can’t currently reach its peers from would turn an ordinary network partition into a much worse, self-inflicted cascading outage. Readiness (GET /readyz) is the one that’s actually derived from PeerManager’s real tracking: a standalone node with no configured peers is always ready (there’s no cluster connectivity to lose), and a node with peers is ready as long as at least one is currently healthy, only flipping to not-ready (HTTP 503) when every single configured peer is simultaneously unreachable. Both endpoints, like /health, stay unauthenticated – an orchestrator probing a node’s liveness can’t be expected to know a secret either. A new GET /peers/health also exposes the same per-peer detail (success/failure counts, response time, healthy/unhealthy) that /peers (addresses only) never did.
Verified live against two real server processes: confirmed both nodes’ peer health started healthy, killed one, confirmed the survivor marked it unhealthy and /readyz correctly flipped to 503 while /livez stayed 200, restarted the killed node, and confirmed the survivor detected the recovery and /readyz returned to 200 – all driven by the independent health-check loop, not by gossip happening to pick that peer.
Storage layout
Each node’s live state is entirely in-memory (Go maps and the CRDT types above), guarded by sync.RWMutex. Every successful rate-limit consume, replay-protection check, and cache write logs to a write-ahead log automatically — no explicit call needed — and a new process recovers automatically on startup: it loads the most recent snapshot (if any) and replays every WAL entry logged after it, reconstructing state exactly as it was before the process stopped, including a hard crash (SIGKILL, no graceful shutdown), which is the case this has actually been verified against. create_snapshot remains available to explicitly compact the WAL into a point-in-time snapshot — see API Reference: Persistence. The WAL also now compacts itself automatically: MeshStore.autoSnapshotLoop calls the equivalent of create_snapshot every defaultSnapshotInterval (5 minutes) — this closes a real bug a sustained 10-minute load test found, where PersistenceEngine.snapshotInterval was stored and even reported via GetPersistenceStats but nothing ever actually triggered a periodic snapshot, so the WAL grew without bound (confirmed live at ~213MB over 10 minutes with zero snapshots taken) unless an operator called create_snapshot explicitly.
Job Queues also had a real memory-growth bug, found the same way. JobQueue.maxAge (“Clean up old jobs”, per its own original field comment) was stored at construction and never read anywhere, so every job ever created — including ones completed or failed long ago — stayed in memory (Jobs/JobIndex) for the process’s entire lifetime, independent of whether jobs were being claimed and completed promptly. Fixed: the existing periodic cleanup now also evicts terminal-state jobs (Completed/Failed/Cancelled only — Pending/Processing jobs are never evicted just for being old, since they’re still active) older than maxAge from the in-memory log.
Both fixes came out of a real sustained-load test: 10 minutes, ~2.3M requests, 30 concurrent workers, a mixed workload across every major feature, against one real server process — the kind of test a short benchmark or a 20-second load-test run structurally can’t perform. Goroutine count stayed perfectly flat at ~70 for the entire run (no leak there), which is exactly the useful negative result this kind of test is also good for. /debug/pprof/* is now wired in permanently — it’s what made watching real goroutine counts during the test possible in the first place, and is useful operational infrastructure independent of this one test. (It was initially gated by the same API key as every other SDK-facing endpoint; a later security review found that too weak and moved it to the cluster-secret tier – see below.)
A separate multi-node test ran cmd/loadtest concurrently against all three nodes of a real cluster at once (not one node in isolation), while gossip actively converged in the background: ~2M total requests, zero errors, and — checked precisely, not just “close enough” — every node’s /metrics/cluster converged to the exact same per-node breakdown, matching each loadtest instance’s own reported count exactly. No correctness issues found under real concurrent multi-node write pressure.
Snapshot/restore now covers all nine gossip-replicated feature groups added this session, not just the original three. This was a real gap: none of Sorted Sets, Streams, Pipelines, Search, Job Queues, Pub/Sub, Transactions, WASM Scripts, or Metrics were ever WAL-logged, so a lone node with no peers (or an entire cluster restarting at once, with no peer left to re-gossip from) silently lost all nine on every restart — invisible in the common multi-node case, since a peer repopulates the restarted node on the next gossip round, which is exactly why it went unnoticed until specifically investigated. persistence.Snapshot gained one field per feature group, each exactly the output of that feature’s own Snapshot(); RestoreFromLatestSnapshot feeds each one back through that same feature’s own MergeSnapshot() — the identical operation gossip performs when learning from a peer, just sourced from disk. WASM Scripts are included at a real cost: restoring recompiles every persisted script via TinyGo (the same multi-second cost WasmEngine.MergeSnapshot’s doc comment describes), added to that node’s own startup time, accepted because the alternative is permanently losing scripts on a lone node’s restart. Building the test for this also caught a real bug: stampLocalMetrics (which feeds Metrics into GetState/CreateSnapshot/GetClusterMetrics) unconditionally overwrote its own node’s slot with its live counter value, so a freshly-restarted process (whose in-memory counter always starts at 0) would immediately regress a just-restored historical high-water mark back to 0 on the very next call — violating the “a GCounter slot only grows” invariant every peer’s merge depends on. Fixed to take the max, like MergeState already does for a peer’s reported value.
Named ranking configs (added after this paragraph was first written) are covered by the same snapshot/restore extension as the other nine. Rank itself still has no per-call state to lose. Persistence’s own mechanism (the WAL and snapshot files themselves) is naturally outside its own scope.
Sandboxed scripting
WASM Scripting deliberately does not run untrusted code in-process. A script is compiled server-side by an external TinyGo process into a WASI WebAssembly module, then executed inside a dedicated wazero sandbox with an enforced memory limit and execution timeout, isolated from the host Go process’s memory and from other running scripts. See vs Redis: Scripting for why this exists instead of an embedded Lua VM.
HTTP as the transport
All access — from every one of the 7 SDKs — goes through a single HTTP API (api/http.go). There is no persistent-connection or pipelining protocol like Redis’s RESP; every operation is one HTTP request/response round-trip. This is simple and easy to reason about, but it puts a real floor under per-operation latency compared to a protocol designed to keep a connection open and pipeline commands — see vs Redis: Performance.
TLS is opt-in, and off by default. Every request – SDK traffic, gossip’s /internal/state pulls, PeerManager’s health-check pings – has always been plaintext HTTP, protected only by the X-API-Key/X-Cluster-Secret headers, which authenticate but do not encrypt: anyone with visibility into the network path can read or tamper with traffic in flight. Three flags turn this on, independently: -tls-cert/-tls-key make a node’s own HTTP server TLS-only (api.HTTPServer.StartTLS, wrapping ListenAndServeTLS – there’s no mixed plaintext/TLS mode, since gossip rides the same port SDKs use); -tls-ca makes a node’s outgoing gossip/health-check/join requests use https:// and verify each peer’s certificate against that CA (GossipCoordinator.SetTLSConfig, which propagates the identical tls.Config to PeerManager so health checks get the same real verification, not a weaker path). A real deployment sets all three on every node, signed by one shared cluster CA. Verified live with two real processes and openssl-generated certs: a plain HTTP request to the TLS-only listener is rejected, HTTPS with the correct CA converges gossip normally, HTTPS without CA trust fails certificate verification outright (not silently accepted), and an actual cache value set on one node was confirmed to replicate to the other entirely over the encrypted channel.
When all three flags are set together, TLS is mutual, not just server-side. Without this, a node’s HTTP surface could be reached by anyone able to complete a TLS handshake at all — encryption was real, but authorization of who was connecting still rested entirely on the cluster secret, sent as a plain header inside that encrypted channel, not on anything TLS itself verified. With -tls-cert/-tls-key/-tls-ca all set, this node’s outgoing requests (gossip, health checks, joining) present its own certificate, and this node’s server (api.HTTPServer.StartTLSWithConfig, configured with ClientAuth: tls.RequireAndVerifyClientCert) requires and verifies every incoming connection’s client certificate against the same shared CA before serving anything — /internal/state included — rejecting the TLS handshake itself for a caller that can’t present a cluster-issued certificate, before the cluster-secret header is ever even read. Verified with TestMutualTLSRequiresValidClientCert (api/mtls_test.go): a real, self-signed test CA and leaf certificates generated entirely in-process (no openssl dependency), confirming a server requiring client certs accepts a certificate signed by the trusted CA and rejects both no certificate and one signed by an untrusted CA.
Certificates now hot-reload, too. They were previously loaded once at startup (tls.LoadX509KeyPair) and never looked at again – fine for a long-lived, manually-managed cert, but a real deployment using short-lived certificates from an automated PKI would need to restart every node on every rotation, defeating much of the point of short-lived certs. coordination.CertReloader watches the cert/key file pair’s mtimes and re-parses only when something actually changed (one stat() per handshake in steady state, not a re-parse); a reload that fails outright (e.g. the files are mid-rewrite, briefly mismatched) isn’t surfaced as an error – it just keeps serving the last successfully loaded certificate rather than disrupting an in-flight handshake over a transient, self-correcting race. main.go now routes all TLS, not just mutual TLS, through StartTLSWithConfig with GetCertificate wired to a CertReloader (the plain StartTLS/ListenAndServeTLS path loads its cert file once and structurally can’t support rotation). Verified live: a real running process was presenting CN=node-v1; its cert/key files were replaced on disk (same paths, no restart) with a CN=node-v2 pair signed by the same CA; the process began presenting CN=node-v2 within seconds, no restart or explicit reload call involved.
A broader security review pass across everything built this session (the closest honest proxy available for external audit — not a substitute for one) found one real, newly-introduced issue: /debug/pprof/cmdline returns this process’s raw os.Args verbatim, including -cluster-secret if an operator passed it as a CLI flag instead of the recommended TOLLMESH_CLUSTER_SECRET env var. That combination previously required local process-list (ps) access to observe; pprof, gated only by the same X-API-Key as every other SDK-facing endpoint, made it reachable over the network by anyone holding that materially weaker, more widely distributed credential. /peers/health had a related, lower-severity gap: it discloses reachability/timing for hosts a cluster-secret holder chose to add as peers, to anyone holding only the API key. A related claim from the same pass — that the new real peer health checks constitute an SSRF vulnerability — was investigated and downgraded: choosing the probe’s target already requires the cluster secret (via the pre-existing /internal/peers/join gate), so an API-key-only holder can only read coarse healthy/unhealthy metadata about a target someone else already chose, never pick it themselves; gossip’s own pre-existing pull already crosses the identical trust boundary with far richer content (the entire replicated state, not a health boolean). Fixed by extending authMiddleware’s existing cluster-secret gate to also cover /debug/pprof/* and /peers/health, even though neither lives under /internal/. Verified live: both correctly reject a request bearing only the (correct) API key and accept one bearing the cluster secret, while ordinary endpoints are unaffected.