In the previous installment of Module of the Week, we learned how Effect Cluster uses the actor model to distribute work across many nodes. We also met DadBot, a bot that joins your chats and interprets everything said as an invitation for a dad joke, whether you want one or not.
In this post, we are going to define the Dad entity, then follow a message from Slack across the cluster and into the right Dad’s mailbox. Along the way, we’ll see how the runners agree on where each Dad lives without a central coordinator.
By the end of this post, you’ll know:
- How to define an entity and send it messages
- How a message finds its entity across a cluster
- How runners agree on who owns which shard, with nobody in charge
DadBot
Before we dive in, let’s take a look at the architecture of our DadBot application.
On each node, we run one copy of the DadBot app. The app has two parts. An HTTP server handles DadBot’s webhook, and a cluster runner which hosts Dad entities. Every chat gets its own Dad entity which lives on exactly one runner.
Slack
- #engineering Sebastian: Why is this test flaky?
- #random
- #general
- …and a million more
Node 1
Runner 1
HTTP server
POST /webhook
Cluster
- #random
- #design
- #lunch
Node 2
Runner 2
HTTP server
POST /webhook
Cluster
- #general
- #support
- #backend
Node 3
Runner 3
HTTP server
POST /webhook
Cluster
- #engineering
- #help
- #oncall
When someone posts in a chat, Slack calls the webhook through a load balancer, which passes the request to the HTTP server on one of the nodes. The webhook handler has no idea which node the Dad for that chat lives on, so it hands the message to the cluster runner on its own node to route it from there.
So how does the runner find the right Dad, when he could be on any node in the cluster?
We’ll get there, but first, we need a Dad.
Defining an entity
An entity has two halves. Its protocol is the list of messages it understands, and its behavior is what it does in response to a message.
The protocol
Dad’s protocol is short. He reads every message in the chat, and when the moment is right, he replies with a pun.
import { Schema } from "effect"import { Entity } from "effect/cluster"import { Rpc } from "effect/rpc"
// "Dad" is the entity type. There's one definition, but one Dad per chat.export const Dad = Entity.make("Dad", [ // The only message Dad understands: someone posted in his chat Rpc.make("NewMessage", { // Who said it, and what they said payload: { from: Schema.String, text: Schema.String }, // No `success` schema, so the sender gets nothing back. Dad replies // in the chat himself, when he's ready. }),])Entity.make takes a name, "Dad", and a list of RPCs. Each RPC is one type of message Dad can handle. Right now there’s only NewMessage, which carries who said something and what they said. There’s no success schema, because Dad doesn’t answer the sender. When he’s ready, he posts to the chat himself.
The payload is a schema. Cluster uses them to validate every message and turn it into bytes before it crosses the network.
Notice that we haven’t defined any logic for our entity yet, we’ve only specified the messages that it can handle.
The behavior
With Dad’s protocol defined, we can write a handler for each of his messages. We do that with Dad.toLayer, which takes an effect that sets up Dad’s state and returns his handlers.
export const DadLayer = Dad.toLayer( // This effect runs once every time a Dad starts up, so every chat's Dad // gets his own copy of everything below Effect.gen(function* () { // Every Dad runs the same code, so he checks his address to find out // who he is: // - `address.entityType`: "Dad" // - `address.entityId`: the chat, like "#engineering" // - `address.shardId`: the shard he lives in const address = yield* Entity.CurrentAddress const punGenerator = yield* PunGenerator const chat = yield* Chat
// Dad's state. Only this Dad can see it, so it needs no locks. // Every joke he's told in this chat: const told = new Set<string>() // Messages he hasn't answered yet: const unread = yield* Queue.unbounded<{ from: string; text: string }>()
yield* Effect.log(`Dad for ${address.entityId} has entered the chat`)
// Wait for one message, then keep collecting until nobody has posted // for 2.5 seconds const nextBatch = Effect.gen(function* () { const batch = [yield* Queue.take(unread)] // waits for the first message while (true) { // Wait up to 2.5 seconds for another one const next = yield* Queue.take(unread).pipe( Effect.timeoutOption("2500 millis"), ) // Silence. Everyone's stopped typing. Dad clears his throat. // This is the moment he's been waiting for. if (Option.isNone(next)) return batch batch.push(next.value) } })
// The joke loop, running in the background for as long as Dad is alive yield* Effect.gen(function* () { const messages = yield* nextBatch // Ask for a pun about the whole batch that he hasn't told yet const joke = yield* punGenerator.generate(messages, { told }) told.add(joke) // Post it to the chat. Nothing goes back to whoever sent the messages. yield* chat.reply(address.entityId, joke) }).pipe( Effect.forever, // then wait for the next batch, forever Effect.forkScoped, // in the background, stopping when Dad stops )
// One handler per message in the protocol. NewMessage just queues the // message and returns, so the sender never waits for a joke. return { NewMessage: ({ payload }) => Effect.asVoid(Queue.offer(unread, payload)), } }),)Let’s walk through it. First, Dad accesses three services:
CurrentAddress: tells him which chat he is inPunGenerator: where his puns come fromChat: lets him post in the chat
Then he sets up some state. The told set holds every joke he’s told, and the unread queue tracks messages he hasn’t answered yet.
The NewMessage handler doesn’t tell jokes. It enqueues the message into unread and returns.
The jokes come from a loop running in the background. It waits until the chat has been quiet for 2.5 seconds, turns everything in unread into one pun, and posts it. Just as we described in Part 1, Dad lets the conversation settle before he delivers the punchline.
Sending a message
To send Dad a message, we ask for his client.
Here’s what DadBot’s webhook handler does when Sebastian’s message comes in:
const handleMessage = Effect.gen(function* () { // A function from an entity ID to a client for that entity const makeDad = yield* Dad.client
// The Dad for #engineering. We say *which* Dad, never *where* he is. const dad = makeDad("#engineering")
// Cluster finds the runner that owns this Dad and delivers the message yield* dad.NewMessage( { from: "Sebastian", text: "Why is this test flaky?" }, // Fire and forget: don't wait for Dad's handler to finish { discard: true }, )})Dad.client gives you a function from an entity ID to a client. You never say where Dad is. You say which Dad, and cluster works out the rest.
discard: true makes the call fire and forget. The webhook sends the message and moves on, without waiting for Dad’s handler to reply. Dad’s joke goes to the chat, not back to the webhook.
Running the cluster
You don’t need a fleet of machines to try this. TestRunner.layer gives you a whole cluster in one process, with in-memory storage. Let’s have Sebastian and Mattia post in #engineering a second apart, and Tim post in #random.
const program = Effect.gen(function* () { const makeDad = yield* Dad.client
// Someone posts in a chat: log it, then send it to that chat's Dad const post = (chat: string, from: string, text: string) => Effect.gen(function* () { yield* Effect.log(`${chat} · ${from}: ${text}`) yield* makeDad(chat).NewMessage({ from, text }, { discard: true }) })
// Two messages in #engineering, a second apart. Dad should answer // them both with one joke once the chat goes quiet. yield* post("#engineering", "Sebastian", "Why is this test flaky?") yield* Effect.sleep("1 second") yield* post("#engineering", "Mattia", "It's flaky for me too")
// A different chat, so a different Dad yield* post("#random", "Tim", "Is CI down?")
// Give the Dads time to think before the program exits yield* Effect.sleep("10 seconds")})
// The services Dad depends onconst ServicesLayer = Layer.mergeAll(ChatLayer, PunGeneratorLayer)
const MainLayer = DadLayer.pipe( Layer.provide(ServicesLayer), // A whole cluster in one process, with in-memory storage Layer.provideMerge(TestRunner.layer),)
program.pipe(Effect.provide(MainLayer), Effect.runPromise)Here’s what it printed:
[11:34:06.306] INFO: #engineering · Sebastian: Why is this test flaky?[11:34:06.411] INFO: Dad for #engineering has entered the chat[11:34:07.415] INFO: #engineering · Mattia: It's flaky for me too[11:34:07.416] INFO: #random · Tim: Is CI down?[11:34:07.417] INFO: Dad for #random has entered the chat[11:34:11.921] INFO: #engineering · DadBot: I'd tell you a UDP joke, but you might not get it.[11:34:11.921] INFO: #random · DadBot: I'd tell you a UDP joke, but you might not get it.Sebastian’s message woke Dad up, but he didn’t answer it. He waited. Mattia chimed in a second later, and once #engineering had been quiet for 2.5 seconds, Dad answered them both with one joke. That’s exactly what Dad did in Part 1, and we didn’t write a single line of coordination code.
Meanwhile, Tim’s message woke a brand new Dad in #random, with his own empty told. Different chat, different Dad. He’s allowed to tell the UDP joke again, because nobody in #random has heard it.
Following a message
So far, we’ve only seen Dad from the outside. We call dad.NewMessage(...), and the message somehow reaches him.
Our cluster has three runners. Sebastian just posted in #engineering, and the load balancer handed his message to Runner 1. The Dad for #engineering could be on any of the three, so how does Runner 1 find him?
Who’s in the cluster
First, Runner 1 needs to know which other runners exist. You might expect a central coordinator that keeps the list, or runners gossiping their presence to one another. Cluster does neither.
Instead, the list lives in the shared database. Every runner adds itself when it starts, and checks back regularly to see who else is around.
Leaving works the same way. A runner that shuts down cleanly deletes itself from the list on the way out. A runner that crashes can’t, which is what the optional RunnerHealth service is for. Plug one in, and cluster uses it to check on runners and mark crashed ones as unhealthy. We’ll dig into how that works in a later post.
Entity to shard
The first step in routing a message is working out which shard its entity belongs to. That part is easy. Cluster hashes the entity ID and takes the remainder:
// Same ID, same hash, same shard, on every runner, every timeshard = (hash("#engineering") % shardsPerGroup) + 1Cluster defaults shardsPerGroup to 300, so the Dad entity for #engineering lives in shard 276, and he always will. The entity ID never changes and neither does the number of shards in the cluster, so the computed shard never changes either. Any runner can work it out on its own, with nothing to look up and nothing to store.
That’s the whole point of shards. Instead of keeping track of millions of Dads, the cluster only has to keep track of 300 shards.
300 shards? Back in my day, we had one shard and we were grateful.
Shard to runner
Next, cluster needs to work out which runner owns shard 276. This part is harder, because runners come and go. You deploy, autoscaling kicks in, or someone trips over a power cord.
A remainder worked for shards, so let’s try it for runners. With three runners, shard 276 goes to runner (276 % 3) + 1. Change the number of runners below from 3 to 4 and watch what happens to the shards.
0 shards moved Didn't move
Adding a single runner moved three quarters of the shards. A shard only stays put when (shard % 3) + 1 and (shard % 4) + 1 match, which is true for one shard in four.
Every move a shard makes has a cost. When a shard changes runners, its entities must shut down on the old one and then restart on the new one.
Consistent hashing
The fix is an idea from 1997, called consistent hashing.
Imagine a line holding every hash cluster can produce, from about minus a billion to plus a billion. Runners and shards both get hashed onto it, runners by their address and shards by their number. Each shard then goes to whichever runner landed closest.
Notice what happens at the ends of our line. Runner 2 has the lowest point, so every shard below it is closest to Runner 2. That’s 181 of the 300 shards, which is quite an imbalance. They don’t get assigned to Runner 2 for any good reason, though, there’s just no other runners below that point on the line.
The fix is to bend the line into a ring, so the two ends meet. The lowest and highest hashes now meet at the top of the ring, so a shard past the highest hash wraps around to the lowest. Dad’s shard lands just before Runner 2’s point, so Runner 2 keeps him.
Color every shard by the runner it’s closest to, and the ring splits into one arc per runner. Each arc reaches halfway to the next runner’s point on either side. Runner 1’s arc now crosses the wrap around point at the top of the ring and it picks up 81 of the shards Runner 2 previously had on the line.
When a fourth runner joins, it drops its point on the ring and takes only the shards that are now closer to it than to anyone else. Every other shard stays exactly where it was.
0 shards moved Didn't move
Instead of moving 225 shards like we did with our modulo-based strategy, we now move only 93 shards. Our Dad for #engineering is one of them this time, since Runner 4 lands right next to his shard.
There’s still a problem. With one point each, the arcs come out uneven. In the diagram above, Runner 1 owns nearly half the ring, purely by chance.
The fix is to give each runner lots of points, called virtual nodes. Cluster gives each runner 128 by hashing host:port:1, host:port:2, and so on. Drag the Points per runner slider below and watch the bars even out.
Shards per runner
- Runner 1 144
- Runner 2 100
- Runner 3 56
Here’s everything together. Add or remove runners, change the points per runner, and see how many shards move each time.
Shards per runner · click to add or remove
But what if we have runners that have different resource constraints? Perhaps we have one runner with twice the memory of the others and we want to allow it to be assigned twice the number of shards as other runners. Cluster allows us to assign a weight to runners. A weight of 2 doubles the runners points on the ring from 128 to 256.
Try picking different weights for Runner 1 below and seeing how that affects the distribution of shards.
0 shards added to Runner 1 shards removed from Runner 1 Didn't move
Shards per runner
- Runner 1 95
- Runner 2 106
- Runner 3 99
Let’s go back to where we started. In Shard to runner, going from three runners to four with modulo moved three quarters of the shards. Below is the same diagram, except Effect’s HashRing assigns the shards this time.
Try adding a fourth runner and see how many shards move. Then flip to Modulo to see the result using the modulo strategy.
0 shards moved Didn't move
With the hash ring, 85 shards move, not 225. All of them go to the new runner, and Runners 1, 2 and 3 never pass shards among themselves. That’s close to the best you could hope for. A fourth runner needs a quarter of the 300 shards, so at least 75 have to move.
Bounded loads
The hash ring keeps moves small, but even with 128 points per runner, it doesn’t keep the shards perfectly even. Spread 300 shards across 20 runners, and one runner might get 7 while another gets 22.
To even things out, cluster borrows an idea from a 2016 paper on consistent hashing with bounded loads. Every runner gets a cap, equal to its fair share of the shards. Shards still go to the closest runner first, but once a runner is full, the shards that wanted it go to the nearest runner with room instead.
Here are six runners with 128 points each. Switch on Cap at fair share and watch which shards move.
25 shards over the line
With a cap in place, every runner gets exactly its fair share. The trade-off is more shard movement. When a runner joins, every cap shrinks, and when one leaves, every cap grows. Either way, some shards move even though their nearest runner hasn’t changed. Cluster ends up moving about twice as many shards as it strictly needs to, in exchange for a perfectly even split.
Agreeing on an owner
There’s no central server deciding who owns what shard. Each runner reads the list of runners from the database and runs the same hash ring math itself, so they all come up with the same assignment… eventually.
Runners reread the list from the database at slightly different moments. Right after a new runner joins, one runner might have the new list while another still has the old one, and both can think they own shard 276. If both started a Dad for #engineering, we’d have two dads in one chat.
To prevent that, a runner has to lock a shard in the database before it starts any of its entities. Only one runner can hold a shard’s lock at a time, and the lock expires if its owner stops renewing it.
Click Runner 1 below to start it, and watch shard 276 change hands.
Click a runner to start or stop it
For a moment during the handoff, nobody owns shard 276. Messages for it aren’t lost, though. The sender keeps retrying until the new owner takes the lock. What can never happen is two runners owning the shard at once.
If Runner 2 crashes instead of handing off cleanly, nobody is left to release the lock. The new owner has to wait for it to expire, configured by the shardLockExpiration, and Dad for #engineering goes quiet in the meantime. Like a real dad when you ask who ate the leftovers.
Delivering a message
With all that in place, sending a message is almost boring. Hash the entity ID to get the shard, look up who owns the shard, and send the message straight there. There’s no router in the middle, no broker, and no queue between runners. Every runner can do the math, so every runner is its own router.
Every runner is its own router? I’ve been my own GPS for years. Never once asked for directions.
Here’s the whole trip for Sebastian’s message, from webhook to punchline.
Slack #engineering
Runner 1
Owns 100 shards
hash("#engineering") % 300 + 1 = ? shard ? →
? Runner 2
Owns 100 shards. Not on this trip
Runner 3
Owns 100 shards, including 276
Shard 276 Lock
Dad · #engineering
- State
- Not running
- Mailbox
- –
Wrapping up
That’s it for this week! We’ve covered how to define an entity and send it messages, how a message finds its entity using shards and a hash ring, and how runners agree on who owns each shard without anyone in charge. All of the code from this post is on GitHub, if you’d like to run it yourself.
There’s one thing we haven’t talked about yet. Everything so far is fast, but none of it is durable. If a runner crashes, any messages waiting for Dad are gone, and so is his told set. When his shard moves, he wakes up on a new runner with no idea what he’s already said.
Next time, we’ll look at what happens when a runner dies, how cluster can persist messages so they survive a crash, and how to give Dad a memory that survives a move, so he never tells the same joke twice, no matter which runner he wakes up on.
Until then, Happy Effecting!