Skip to main content
Blog

Replacing Our Message Broker with... S3?

Rivet 3.0's only dependency is S3. Here's how we replaced NATS with it.

Replacing Our Message Broker with... S3?

Rivet is a fast, high-density, and scalable orchestrator for agentic workloads such as agent orchestration, workflows, and realtime apps. The primitive at the center of Rivet is an Actor, frequently used as an open-source alternative to Durable Objects.

We’re putting the finishing touches on Rivet 3.0, which has three goals:

  • S3 as the only dependency, making scaling cheaper and easier to operate.
  • No performance sacrifice compared to the SSD-based databases behind Rivet 2.0, like Postgres and FoundationDB.
  • As easy as possible to run, whether you’re self-hosting Rivet or embedding it inside your own application.

That obviously means removing those databases. It also means replacing our message broker, NATS, which has been the backbone between all of our services.

Replacing NATS with S3 sounds absurd, and for most workloads it would be. Here’s how we did it.

What is NATS?

NATS is a lightweight, open-source message broker, and it’s the message broker Rivet has used between all of its nodes: sending messages between nodes, broadcasting events to the whole cluster, and generally anything that helps us scale beyond a single machine.

We use it in two ways:

  • Request/reply: One node sends a request to a specific subscriber, such as the node running a given Actor, and waits for a reply.
  • Broadcast: One node publishes an event, and every node subscribed to that subject receives it.
REQUEST / REPLYNode ANATSNode Brequestrequestreplyreplyone request, one replyBROADCASTNode ANATSNode BNode Cpublishone publish, every subscriber receives it

Before going any further, it’s worth saying that NATS has been phenomenal. It’s arguably the lowest-maintenance piece of our stack:

  • We’ve never had an outage caused by NATS.
  • Its CPU and memory usage is laughably low.
  • It’s incredibly easy to set up and run.
  • We’ve hit exactly two bugs, both in the Rust client.

If you need a message broker, use NATS.

So why replace it?

We want Rivet to be as easy as possible to run, both when self-hosting and when embedding Rivet as a dependency inside your own application. Every extra service is one more thing to deploy, monitor, and upgrade. Removing NATS makes S3 our only dependency.

There’s a performance motivation too. Every message through NATS travels from the sender to NATS and then from NATS to the receiver.

THROUGH NATS2 hopsNode ANATSNode Bextra hop + copyDIRECT1 hopNode ANode B

Talking directly instead:

  • Removes a network hop from every request, which directly cuts latency inside Rivet.
  • Removes a data copy through the broker, which makes throughput cheaper.

Taking a step back: our actual workload

S3 writes take tens of milliseconds and are billed per request. To be clear, S3 is not a generic replacement for NATS.

But here’s the thing: we don’t need all of NATS. When we looked at every place Rivet uses it, it came down to two patterns:

  • Frequent, latency-sensitive request/reply between nodes.
  • Infrequent broadcasts to other nodes that are fine taking a moment to arrive.

We don’t rely on most of the situations where NATS shines, and these two patterns have very different requirements, so we solved them separately.

The naive approach: fan out to every node

The naive way to get rid of NATS is to have each node talk to every other node directly. For example, when a node needs to broadcast something like actor.wake, it sends it to every other Rivet node, and nodes without a matching subscriber ignore it.

3 NODES2 requestsABC10 NODES9 requests

With 3 nodes, this is fine: each broadcast costs 2 requests. With 10 nodes, each broadcast costs 9 requests. The cost of a single broadcast grows with the size of the cluster, which is O(n) per broadcast. If every node broadcasts at a steady rate, total cluster traffic is n × (n − 1), or O(n²): 6 messages for 3 nodes, 90 for 10 nodes, and 9,900 for 100 nodes.

This is fan-out amplification: one logical event costs work proportional to the cluster size. It’s doubly bad because it gets worse precisely when you scale. Adding nodes should increase capacity, but here each new node also makes every broadcast more expensive for every other node.

NATS solves this with a hub. A publisher sends one message to its NATS server, and the NATS cluster takes care of delivering it to every subscriber.

Node ANATS CLUSTERNATS serverNATS serverNATS server1 messageNode A’s cost stays at one message no matter how many nodes subscribe

The publisher’s cost is constant regardless of cluster size. That’s a big part of why brokers exist.

Messages each node sends per broadcastFanoutNATS00252550507575100100Rivet nodesmessagesFanout, 3 nodes: 2 messagesFanout, 10 nodes: 9 messagesFanout, 25 nodes: 24 messagesFanout, 50 nodes: 49 messagesFanout, 75 nodes: 74 messagesFanout, 100 nodes: 99 messagesFanout: 99NATS, 3 nodes: 1 messagesNATS, 10 nodes: 1 messagesNATS, 25 nodes: 1 messagesNATS, 50 nodes: 1 messagesNATS, 75 nodes: 1 messagesNATS, 100 nodes: 1 messagesNATS: 1

The alternative: gossip

Another option we considered was gossip. Instead of sending a broadcast to every node, a node tells a few random peers, and those peers each tell a few more. The message spreads like a rumor, reaching the whole cluster in O(log n) rounds while each node only does a constant amount of work.

ROUND 1ROUND 2Node ANode BNode CNode DNode ENode FNode Gduplicate,dropped

This gives gossip attractive scaling properties. A node only ever talks to a few peers at a time, and to make sure every node eventually hears a broadcast, each message is retransmitted about log n times. In HashiCorp’s memberlist, that’s 4 * ceil(log10(n + 1)) messages per node, so per-node cost grows as O(log n) and total cluster traffic as O(n log n), instead of the O(n) per node of direct fanout.

Messages each node sends per broadcastNATSGossip04812162002505007501,000Rivet nodesmessagesNATS, 3 nodes: 1 messagesNATS, 10 nodes: 1 messagesNATS, 100 nodes: 1 messagesNATS, 1000 nodes: 1 messagesNATS: 1Gossip, 3 nodes: 4 messagesGossip, 9 nodes: 4 messagesGossip, 10 nodes: 8 messagesGossip, 99 nodes: 8 messagesGossip, 100 nodes: 12 messagesGossip, 999 nodes: 12 messagesGossip: 12

HashiCorp’s stack is built on this. Consul and Nomad use Serf, which is built on their memberlist library, an implementation of the SWIM gossip protocol. Gossip is a good fit when broadcasts need to arrive quickly but not instantly, which describes ours well.

We didn’t go with it for two reasons:

  • Complexity: Gossip means owning peer forwarding, deduplication, repairing messages that a node missed, and correct behavior while nodes join and leave. Every one of those is a source of subtle bugs.
  • State: Every node has to remember which messages it has already seen so it can drop duplicates, like Node E above. With S3, a reader keeps one position per publisher, reads what’s new, and forgets it.

Solving request/reply

Request/reply turns out to be a simple problem, because we already have a list of every Rivet node and how to send requests to it.

The catch is that it requires changing logic outside of our messaging layer: anywhere we need request/reply, the caller now has to know which node owns the subscriber.

For example, a worker running an Actor used to blindly listen on a topic for that Actor, like actor.123, and anyone who wanted to talk to the Actor sent a request to that topic:

client.request("actor.123", payload);

Now we store the ID of the worker running it on the Actor’s record in S3, and when we need to talk to the Actor, we send the request straight to that worker’s internal address:

client.request(actor.workerId, "actor.123", payload);
ACTOR RECORDactor_id: 123worker_id: node-bNode ANode BActor 1231. who owns Actor 123?2. requestreplyone hop, no broker in between

The result is that request/reply is now one direct hop between two nodes, the theoretical minimum, with a little more care about where requests and replies come from.

Solving broadcasting

Rivet’s broadcasts are:

  • Low frequency: Events are published occasionally, not per request.
  • Latency-tolerant: It’s fine if delivery takes a moment.
  • Reach every node: Every node that subscribes must eventually see every event.

Examples include cache purges and configuration changes. None of this logic is on the hot path for low-latency behavior such as Actor scheduling and networking, which makes it a good fit for S3.

One batched log per node

Each node writes its broadcasts to its own log in S3. Instead of one write per event, it batches everything published in a 100 ms window into a single write, then advances the log’s head to the new batch using head fencing.

S3node A logbatch 1cache.purgeconfig.updatebatch 2cache.purgenode B logbatch 1pool.updatebatch 2cache.purgecache.purgenode C logbatch 1config.updatebatch 2pool.updateNode A1 batch / 100 msNode B1 batch / 100 msNode C1 batch / 100 ms

Every node then polls every other node’s log every 250 ms.

Each log is two kinds of objects in S3, a small head that records the log’s newest batch and the immutable batches themselves:

heads/
  node-a = 2
  node-b = 2
batches/
  node-a/
    1 = [cache.purge, config.update]
    2 = [cache.purge]
  node-b/
    1 = [pool.update]
    2 = [cache.purge, cache.purge]

A poll works like this:

  1. List the heads: One S3 list request on heads/ returns every head with its version fingerprint (ETag). If a head’s fingerprint is the same as on the last poll, that log hasn’t changed, so the node skips it.
  2. Read the changed heads: For each head whose fingerprint changed, the node reads it to find that log’s newest batch.
  3. Fetch the new batches: The node downloads every batch between the last one it read and the newest one, such as batches/node-a/2.
  4. Deliver locally: Each node has many local subscribers, so it filters the messages by subject and delivers each one only to the subscribers that care about it.

As an unintended side effect of this architecture, broadcasts are durable since they’re stored in S3, even though they don’t need to be.

Overall, this makes broadcasts cheap, scalable, and simple for our specific use case, though it doesn’t fit most NATS workloads.

Why not a log per subject?

With S3, the expensive part is the number of write operations, not the bytes written. A log per subject would turn one write per node into one write per active subject. With high-cardinality subjects, that’s a lot of writes happening at once.

It also makes reading harder. There’s no cheap way for a reader to know which subject logs exist or have changed, so it would have to discover and poll a growing set of logs.

Why not one big shared log?

Writing to a log requires a conditional write to its head (head fencing) so writers can’t clobber each other. That’s cheap with a single writer per log. With one shared log, every node would compete to update the same head, causing constant retries that get worse as the cluster grows. Per-node logs cost readers a bit more polling, which is fine for our workload.

Looking ahead to Rivet 3.0

As we wrap up the finishing touches on Rivet 3.0, we’re excited to share more about its inner workings.

We believe S3 is finally at the point where, with the right engineering decisions, it is the best durable backend for most storage by providing fast, cheap, and scalable storage.

Learn more