Writing a simple Go async message queue server
Marton Trencseni - Wed 10 June 2026 - Programming
Introduction
After writing async message queue servers in Python, Javascript and C++, and then one more in Rust, I wrote one final version, in Go. Of all the languages in this series, Go is arguably the most purpose-built for exactly this task. It was designed at Google for networked, concurrent server software, with goroutines and channels as first-class citizens. So if there is a language that should make a small async message queue server feel natural, it is this one.
As before, the same checklist applies — this is a good toy precisely because it touches the parts that matter:
- it has sockets
- it has concurrency
- it has shared mutable state
- it has parsing and validation
This post walks through the Go implementation, with attention to the Go-specific parts, and in particular to the one place where Go pushed back. The full code is on Github.

What the server does
The Go version follows the same exact protocol as the previous implementations, and passes exactly the same unit tests.
Clients connect over TCP and send JSON commands, one per line:
{"command":"subscribe","topic":"news"}
{"command":"send","topic":"news","msg":"hello","delivery":"all"}
{"command":"unsubscribe","topic":"news"}
The server supports subscribe, unsubscribe, and send. Each topic has a set of subscribers, a small cache of recent messages, and a monotonically increasing message index. Messages can be delivered in two modes: all (send to every subscriber) or one (send to one random subscriber). Messages delivered with all may also be cached, so that a client subscribing later can catch up.
The data structures
Let's start with the state.
type Topic struct {
subscribers map[int64]*Outbound
cache []map[string]interface{}
nextIndex int64
}
type SharedState struct {
mu sync.Mutex
topics map[string]*Topic
cacheSize int
nextID int64
}
Compared to the Rust version, the first thing to notice is what is missing. In Rust the shared state was an Arc<Mutex<SharedState>> — an atomically reference-counted, lock-protected heap object, with the sharing made painfully explicit. In Go, that machinery is invisible. Go is garbage collected, so I just pass a *SharedState pointer to every goroutine and the runtime keeps the object alive as long as anyone holds a reference. There is no Arc, because there is nothing to count by hand.
What remains explicit is the lock. The sync.Mutex lives right inside SharedState as a field, and I take it whenever I touch the topic map:
func (s *SharedState) allocID() int64 {
s.mu.Lock()
defer s.mu.Unlock()
id := s.nextID
s.nextID++
return id
}
The defer s.mu.Unlock() is the Go idiom for scope-based unlocking. It is not as automatic as Rust's guard object — I have to remember to write the defer — but it achieves the same thing: the lock is released when the function returns, no matter which path it takes.
Goroutines instead of tasks
The accept loop is about as simple as it gets:
for {
conn, err := listener.Accept()
if err != nil {
fmt.Fprintf(os.Stderr, "Accept error: %v\n", err)
continue
}
go handleClient(conn, state)
}
The go keyword starts a new goroutine. A goroutine is a lightweight, runtime-scheduled thread; starting tens of thousands of them is cheap and normal, which is why the same 10,000-connection test that the other implementations pass works here without any special effort. In the Rust version this was tokio::spawn; in Python it was asyncio. Here it is a language keyword, which tells you something about Go's priorities.
JSON values as dynamic maps
Like every other implementation in this series, the Go version treats JSON dynamically rather than defining typed structs. The Go equivalent of "some JSON thing" is interface{}, and a JSON object decodes into a map[string]interface{}:
var obj map[string]interface{}
if err := json.Unmarshal(lineBytes, &obj); err != nil {
return true, sendFailure(w, "Could not parse json")
}
The cost of this dynamism is type assertions everywhere. To read a field, I have to assert its type and check the result:
if _, ok := obj["command"].(string); !ok {
return false
}
There is one Go-specific trap worth calling out: encoding/json decodes all JSON numbers into float64, never int. So validating that last_seen is an integer means asserting float64 and then checking it has no fractional part:
if v, ok := obj["last_seen"]; ok {
f, ok := v.(float64)
if !ok {
return false
}
if f != float64(int64(f)) {
return false
}
}
And when I stamp a message with its index, the index goes in as a float64, which marshals back to 0, 1, 2… exactly as the protocol and the unit tests expect.
The interesting part: Go has no unbounded channel
This is the part of the project where Go actually made me stop and think, and it is the most interesting thing in the whole implementation.
In the Rust version, each subscriber had an mpsc::UnboundedSender<Value>. When a send arrives, the sender pushes the message into every subscriber's channel and moves on. The channel is unbounded, so the push never blocks and never drops; a slow client just accumulates a backlog.
My first instinct in Go was to give each session a buffered channel and deliver into it. But Go channels are not unbounded — they are either unbuffered or have a fixed capacity, by design. That leaves two bad options for delivery:
- Block on a full channel. This deadlocks. Imagine client A sending to client B while B is simultaneously sending to A; each
handleSendblocks trying to deliver into the other's full channel, holding things up, and the whole thing wedges. - Drop on a full channel (a
selectwith adefaultcase). No deadlock, but now a slow subscriber silently loses messages — a weaker guarantee than every other implementation in the series.
Go's omission of unbounded channels is deliberate: the language wants you to think about backpressure. But here I specifically want the Rust semantics — never block, never drop — so I built an unbounded queue by hand, out of a mutex, a condition variable, and a slice:
type Outbound struct {
mu sync.Mutex
cond *sync.Cond
queue [][]byte
closed bool
done chan struct{}
}
func (o *Outbound) push(msg []byte) {
o.mu.Lock()
if !o.closed {
o.queue = append(o.queue, msg)
o.cond.Signal()
}
o.mu.Unlock()
}
push only ever appends to a slice, so it never blocks and never drops — exactly the property I wanted. A dedicated writer goroutine drains the queue, sleeping on the condition variable when it is empty:
func (o *Outbound) run(w *bufio.Writer) {
defer close(o.done)
for {
o.mu.Lock()
for len(o.queue) == 0 && !o.closed {
o.cond.Wait()
}
if len(o.queue) == 0 && o.closed {
o.mu.Unlock()
return
}
batch := o.queue
o.queue = nil
o.mu.Unlock()
for _, msg := range batch {
if _, err := w.Write(msg); err != nil {
return
}
w.Write([]byte("\r\n"))
}
w.Flush()
}
}
sync.Cond is one of the lower-level, less fashionable corners of the Go standard library — most Go concurrency is supposed to be done with channels — but this is exactly the textbook use case for it: a consumer that must wait efficiently until a producer makes work available. The cond.Wait() atomically releases the lock and sleeps, and re-acquires the lock when woken by cond.Signal(). The for loop around Wait() (rather than an if) is the standard guard against spurious wakeups.
The payoff is that this writer goroutine is the only goroutine that ever writes to the socket. Command responses and messages delivered from other clients all funnel through the same queue, so there is no chance of two goroutines interleaving bytes on the wire.
A blocking read loop, no select needed
Because the writer goroutine owns all output, the rest of the connection handler can be refreshingly boring. There is no select multiplexing reads against deliveries — the read side is just a plain blocking loop:
func handleClient(conn net.Conn, state *SharedState) {
sess := newSession(state)
w := bufio.NewWriter(conn)
go sess.out.run(w)
defer func() {
sess.cleanup() // remove from topics: no further deliveries enqueued
sess.out.close() // let the writer flush what remains and exit
<-sess.out.done // wait for the writer before closing the socket
conn.Close()
}()
reader := bufio.NewReaderSize(conn, 128*1024)
for {
line, err := reader.ReadBytes('\n')
if len(line) > 0 {
if !sess.processLine(line) {
return
}
}
if err != nil {
return
}
}
}
The defer block is worth a second look. Go's defers run in LIFO order when the function returns, which gives a clean, ordered teardown: first remove the session from every topic so no new messages get enqueued, then signal the writer to flush whatever is left and exit, then wait for it via the done channel, and only then close the socket. This is the same scope-based cleanup discipline that Rust gets from Drop, except in Go it is written out by hand — which, as with the mutex, is the recurring theme: Go gives you the structured-cleanup pattern, but you have to spell it out.
Delivery and lock ordering
Because push can never block, I can afford to deliver while still holding the state lock, which keeps per-topic ordering well defined:
if delivery == "all" {
for _, out := range topic.subscribers {
out.push(msgBytes)
}
} else {
if len(topic.subscribers) > 0 {
keys := make([]int64, 0, len(topic.subscribers))
for k := range topic.subscribers {
keys = append(keys, k)
}
chosen := keys[rand.Intn(len(keys))]
topic.subscribers[chosen].push(msgBytes)
doCache = false
}
}
The lock order is always state-lock first, then the per-session Outbound lock, and never the reverse. The writer goroutine, which holds the Outbound lock inside run, never reaches for the state lock. Keeping that ordering consistent everywhere is what guarantees there is no lock cycle, and therefore no deadlock — the same property I would have gotten for free from Rust's unbounded channel, but here reasoned out explicitly.
Lines of code
Updating the table from the summary post with Rust and Go:
| Language | LOC |
|---|---|
| Python | 202 |
| JavaScript | 311 |
| C++ | 343 |
| Rust | 334 |
| Go | 474 |
Go comes out as the longest, though the comparison is a little unfair: a good chunk of those lines are the hand-rolled Outbound queue and the comments explaining the locking, neither of which the other versions need (Python gets a collections.deque, Rust gets an unbounded channel). Strip the comments and blank lines and it is closer to 400. Still, Python remains the clear winner on brevity, and nothing here has dislodged it.
Conclusion
Go was, as expected, the most frictionless language in this series for the shape of the problem. Goroutines make "one lightweight thread per connection" the obvious and idiomatic design, the standard library has everything needed (net, bufio, encoding/json, sync) with no third-party dependencies at all, and the 10k-connection test passes without a second thought. For a networked concurrent server, the language really does feel like it was built for the job — because it was.
The one genuinely interesting wrinkle was the unbounded channel, or rather its absence. Go's opinion that channels should be bounded is a reasonable one — it nudges you toward thinking about backpressure — but when I specifically wanted the never-block, never-drop semantics of the other implementations, I had to step down to sync.Mutex and sync.Cond and build the queue myself. That is the Go experience in miniature: the common path is short and pleasant, and when you step off it, the lower-level primitives are right there, but you are now writing and reasoning about the concurrency yourself.
Set against Rust, the contrast is clean. Rust forces you to make ownership, sharing, mutability, and lifetimes explicit everywhere, all the time, at compile time. Go hides almost all of that behind the garbage collector and cheap goroutines, and only asks you to be explicit about the things it cannot infer — which locks protect which data, and in which order you take them. Rust makes the compiler your reviewer; Go trusts you, and hands you a race detector for when that trust is misplaced. For a weekend project like this one, Go's bargain is the more comfortable one. For a large system where the invariants have to hold under a team of contributors, I can see the appeal of Rust's insistence.
With Python, JavaScript, C++, Rust and now Go all speaking the same wire protocol and passing the same unit tests, this little message queue has turned into a surprisingly good lens for comparing languages. I think this is where I finally let the series rest.