beanstalkd-rs: A Drop-in beanstalkd in Rust, with a Raft Cluster
beanstalkd is a work queue written in C: a single binary with a text protocol you can type into telnet, and a handful of commands (put, reserve, delete, release, bury, kick). It does very little, which is why a lot of job systems still run on it. In production you eventually run into its limits. It uses one thread, it has no replication, TLS needs something like stunnel in front of it, and when the binlog hits a disk error it turns the binlog off and keeps serving.
Today I tagged beanstalkd-rs 0.5.0, its first release. It reimplements beanstalkd 1.13 in Rust with the same protocol byte for byte, so existing clients work without changes. It also adds TLS, Prometheus metrics, a stricter write-ahead log, and an optional 3- or 5-node Raft cluster that looks like a single beanstalkd server to clients.
I started on September 25 and released on October 10: about 90 commits and 77k lines of Rust, tests included. Below are the benchmark results, how compatibility is tested, the internal design, the cluster, and the places where it is still slower or heavier than the original.
Performance
All numbers are from the README. Each cell is operations per second against beanstalkd built with -O2, both servers on default settings. The machine is Linux 7.0 in an OrbStack VM on an Apple M6, with the servers on 6 cores and the load generator on the other 6, over loopback. Each cell is the median of 5 alternated runs with 16-byte job bodies.
| Workload | beanstalkd | beanstalkd-rs | |
|---|---|---|---|
| put-reserve-delete, 1 connection | 56,176 | 53,943 | 0.96× |
| put-reserve-delete, 10 connections | 339,283 | 468,310 | 1.38× |
| put-reserve-delete, 100 connections | 392,713 | 516,336 | 1.31× |
| 100 connections, 4 KiB bodies | 362,247 | 463,416 | 1.28× |
| 100 connections, 16 pipelined | 490,247 | 643,204 | 1.31× |
| producers and consumers, 100 connections | 422,349 | 500,891 | 1.19× |
with a binlog (-b, default fsync), 100 connections |
267,391 | 594,148 | 2.22× |
Without a binlog, CPU per operation at 10 and 100 connections is 10–24% lower (at 100 connections, 1.94 µs against 2.55). With -b, the thread that owns all queue state (described below) runs on its own OS thread so that file writes and fsync never block a tokio worker, and the default goes up to two workers. On native macOS on the same machine, the speedup is 1.2–1.4× at 10 and 100 connections and 1.9× with a binlog.
For the cluster, I ran 3 nodes on the same machine with every operation committed by a majority before the reply. At 100 connections with the data on tmpfs, it does about 316k ops/s through the leader and 236k through a follower. On a real disk, each commit has to wait for an fsync on a majority of nodes, so the disk's sync rate is the ceiling: 48k on this VM's volume.
Quick start
# Docker (linux/amd64 and linux/arm64)
docker run -d -p 127.0.0.1:11300:11300 -v bstk-data:/data \
shonhen/beanstalkd-rs -l 0.0.0.0 -p 11300 -b /data
# or from source (Rust 1.98+)
cargo build --release --locked -p bstk-server
./target/release/beanstalkd-rs -l 127.0.0.1 -p 11300 -b ./binlog
The flags are the same as the original's, -l -p -z -b -f -F -s -V -v, with the same defaults. The exception is -u. Switching users is left to systemd, and the release archive includes a unit file for that. Prebuilt archives are available for Linux x86_64 and aarch64 (glibc or static musl) and for macOS aarch64.
Testing compatibility
The rule is that when protocol.txt and the C source disagree, the source wins. beanstalkd-rs does what prot.c actually does, and every such case is recorded in docs/COMPAT.md.
There are three kinds of tests:
- Differential tests.
scripts/build-ref.shbuilds the C beanstalkd from a pinned commit. The same script runs against both servers, and the replies are compared byte for byte, masking only fields that always differ (pid, uptime, rusage, hostname, version, time-left). There are 189 cases. - Real clients. Python's greenstalk and Go's go-beanstalk each run a full flow against both servers.
- An engine oracle. Before I rewrote the engine's data structures for speed, I froze a copy of the old engine (
bstk-engine-oracle). A property test feeds both engines the same(time, message)sequences, down to deadlines that tie to the nanosecond, and their replies and stats must match at every step.
These tests found a bug in the original. On Linux, beanstalkd rounds its epoll timeout down to whole milliseconds and can wake up to 1 ms early. If the clock it reads then equals a job's TTR deadline exactly, the timer is lost, and the job stays reserved past its TTR until that connection does more I/O. It happened in about 3% of runs. beanstalkd-rs expires the job on time, and the build script patches the reference copy so the differential tests don't fail at random (COMPAT D12/D14).
There are only a few intentional differences, and each has a reason. The one that matters most is D5. If writing the binlog, an fsync or a compaction fails, beanstalkd-rs exits. The original disables its binlog without telling anyone and keeps replying INSERTED to writes it can no longer save.
One engine thread, I/O on the other cores
beanstalkd's semantics are single-threaded. A reserve can watch several tubes and must return the most urgent job across all of them, so splitting the state into shards would change the protocol. I kept that and used the other cores for I/O:
- The engine (
bstk-engine) is a deterministic state machine. It never reads the clock or uses randomness. The caller passes innow, and the pid, hostname and rusage come in through a trait. The samenowand the same sequence of messages always produce the same output. Without that, neither the differential tests nor Raft replication would work. - One actor owns the engine, so there are no locks. Parsing, encoding and socket I/O run on tokio workers. Each connection has at most one command in flight to the actor, so pipelined commands are answered in order.
- Timers are indexed. An early profile showed the cost per operation growing with the number of tubes and connections, because
tickscanned all of them after every message. Delayed jobs, connection deadlines and paused tubes now each have aBTreeSetindex, sotickreturns immediately when nothing is due. Scale tests with 100k and 1M buried or reserved jobs catch anything that turns quadratic. - I picked the thread counts by benchmarking. Fewer tokio workers means fewer cross-thread wake-ups per command. The default is one worker for plain TCP, and two with TLS, a binlog or a cluster, because a single worker maxed out a core in those modes.
--threads Noverrides it.
The whole workspace sets unsafe_code = "forbid", and unwrap() is not allowed outside tests.
The write-ahead log
With -b DIR, every state change is written to the log at the same points where the original writes its binlog, and a reply is sent only after its change is in the log. The file format is my own and can't read the original's; that was never a goal. Records carry a CRC-32C and live in preallocated segment files. Compaction copies live jobs out of old segments and then deletes those segments. The design doc shows why replay ends up with the same live jobs whether a crash happens before, during or after a compaction step. If the last segment ends in a torn record, replay cuts it off there, which is the cost of -f/-F after a power loss, as in the original. Corruption anywhere earlier stops the server from starting, because that is real data loss.
Every decoder that reads untrusted bytes is fuzzed: the protocol parser, the WAL reader, the Raft wire format and snapshots. Fuzzing found two u64::MAX cases, one in segment numbers and one in job ids, where the counter could wrap around and reuse ids. Both values are now rejected.
The cluster
With a [cluster] section in the config, 3 or 5 nodes replicate the queue over mutual TLS using openraft 0.9.
Every input to the engine is a Raft log entry: client commands, connects and disconnects, and timer ticks. Each node applies the committed entries to its own engine, so every node holds the same state, connections and reservations included. No reply is sent before a majority has committed the entry.
A client can connect to any node. A follower forwards its clients' commands to the leader, then sends the replies itself once it has applied the entries. When the leader changes, no replies are lost, and a reserve waiting on a surviving node keeps waiting.
When a node stops responding, the leader proposes a DropNode entry that disconnects that node's clients. Their reserved jobs go back to ready, as they would after a disconnect on a single server.
beanstalkd-rs cluster add | promote | remove | set-addr changes membership while the cluster is serving. It can grow a cluster from 1 to 3 nodes or from 3 to 5, replace a node or a disk, or remove the leader. Only one voter can change at a time, only a learner that has caught up can be promoted, and node ids are never reused. A removed node shuts itself down once a majority of the remaining voters confirm the removal.
The part that took the most work was a node that comes back with an empty disk. Before it was wiped, it may have acknowledged log entries and voted in elections, and it no longer remembers either. If it simply rejoined, Raft's safety guarantees could break. So before it starts, such a node asks the current voters for their state, adopts the highest vote they report, and stays out of elections until it has caught up. DESIGN §8 has the full argument and the openraft source lines it depends on.
Two chaos harnesses test all of this. One runs in-process on a simulated network. The other runs real server processes behind proxies that can be paused. Both record every client operation through crashes, network partitions, disk wipes and membership changes. Afterwards, each job's history is checked for linearizability against a model of a single server. Acceptance required 1,000 in-process seeds and 20 multi-process runs per scenario. They found two real bugs, and both now have regression tests.
Known weaknesses
- With one connection, beanstalkd-rs is 4% slower than the original on Linux and the same on macOS. Each request is one round trip, so latency decides that number.
- With a binlog at 10 connections, it is 0.96× the original on Linux. With a binlog it also uses 1.0–1.4× the original's CPU per operation, so the 2.22× in the table comes at a CPU cost.
- Each small job uses 249 bytes against the original's 219, 13% more. At 4 KiB the two are the same.
- In a cluster at one connection, every log entry costs an
fdatasyncon every node, plus an extra AppendEntries round in openraft 0.9 that only carries the commit index. I ported the code to openraft 0.10 on a separate branch to see whether it helped. The port took about an hour and the tests passed, but CPU usage didn't improve, so I'm staying on 0.9.25 until 0.10.0 is released.--threads 1lowers the cluster's CPU per operation by about 30% at one connection, but costs 19–32% of throughput at 100–300 connections. It's documented as an option for low-traffic clusters, not the default. - The binlog format is different from the original's, so beanstalkd-rs can't read an existing beanstalkd binlog directory. Drain the old server's queue before you switch.
The development plan
The repo includes docs/PLAN.md, which covers phases P0 to P9. Each phase lists the facts I checked before planning, the acceptance criteria, and at the end, the benchmark section or CI run that shows each criterion was met. When a target wasn't reached, the plan says "not met" and records the measured values and the reason. I did the same in last month's simpleconf re-measurement: every number in a README should be reproducible from a run someone can repeat.
Links
- GitHub: shaunlee/beanstalkd-rs (MIT, like beanstalkd)
- Docker Hub: shonhen/beanstalkd-rs
- Docs: OPERATIONS · DESIGN · COMPAT · BENCH
If you already run beanstalkd, you can try it by pointing one worker pool at a beanstalkd-rs instance. If you find behavior that differs from the original and isn't listed in COMPAT.md, please open an issue.

