Instalment 4 · Course 1 (Go) · Milestones 9–12
The colony gets metrics, a profiler, a live picture in your browser, a rewrite guided by an actual profile, and finally a world that lives in a separate process you can kill while ants are talking to it.
Every snippet compiled, vetted, race-tested and run. Every number measured. The container has one core, so parallel speedups are understated and I say so where it matters. The finished project is 2,961 lines of Go.
Counters, gauges and latency histograms for everything the colony does, exposed over HTTP in a format a scraper understands, with pprof mounted alongside so you can profile a live process without restarting it.
Atomic counters versus owner-owned counters, cheap histograms without allocation, the Prometheus text format, net/http servers and graceful shutdown, mounting net/http/pprof deliberately, and reading a CPU profile.
We already have two kinds of counter, and the distinction is worth making explicit before adding twenty more.
| Kind | Written by | Implementation | Example |
|---|---|---|---|
| Ledger | Only the owner goroutine | A plain int field | delivered, carrying, lost |
| Metric | Any goroutine | atomic.Int64 | timeouts, restarts, round-trip latency |
The ledger must be exactly right, because the conservation invariant depends on it, and it is cheap to keep right because one goroutine owns it. Metrics are written from everywhere and need to be approximately right and very cheap. Do not merge these two ideas: a ledger behind atomics invites someone to read two counters that were never consistent with each other, and a metric behind the owner adds a message round trip to every observation.
For histograms, the naive approach (keep every observation, sort, take a percentile) allocates unboundedly and is useless in a hot path. The standard answer is bucketing, and there is a trick that makes it nearly free: use powers of two as bucket boundaries, and find the bucket with a single bit operation.
Most stacks reach for a metrics client library immediately: import a Prometheus SDK, register each metric with a global registry, and let the library own an HTTP handler you never read the internals of. It is quick to start and gives you a lot for free — histograms with configurable buckets, a battle-tested text encoder, client-side aggregation.
The idiomatic Go answer used here skips the dependency: a plain struct of atomic fields, a hand-written WriteProm that Fprintfs the exposition format directly, and pprof mounted by hand on a private mux. Zero third-party code, a metrics endpoint you can read start to finish in one file, and no framework magic deciding what gets exposed. The honest cost: you re-derive small things (percentile-bucket math, HELP-line formatting) that a mature client library already got right and tested against real scrape edge cases, and if this project ever needed histograms with configurable, non-power-of-two buckets, the library would save real time.
// Counter is a monotonically increasing number.
type Counter struct{ v atomic.Int64 }
func (c *Counter) Inc() { c.v.Add(1) }
func (c *Counter) Add(n int64) { c.v.Add(n) }
func (c *Counter) Value() int64 { return c.v.Load() }
// Gauge is a number that goes up and down.
type Gauge struct{ v atomic.Int64 }
func (g *Gauge) Set(n int64) { g.v.Store(n) }
func (g *Gauge) Value() int64 { return g.v.Load() }
Wrapping atomic.Int64 in a named type with two methods looks like pointless ceremony until you notice what it prevents: nobody can accidentally read a counter with = instead of .Load(), and the type name documents whether a number is expected to go down. go vet also refuses to let you copy a struct containing an atomic, which catches a whole class of "my metrics stopped moving" bugs.
// Histogram records durations in power-of-two microsecond buckets. Bucket i
// holds observations in [2^(i-1), 2^i) microseconds, which covers 1µs to
// about 4 seconds in 24 buckets for the cost of one atomic add per
// observation and no allocation at all.
type Histogram struct {
buckets [24]atomic.Int64
count atomic.Int64
sumUS atomic.Int64
}
func (h *Histogram) Observe(d time.Duration) {
us := d.Microseconds()
if us < 0 {
us = 0
}
h.count.Add(1)
h.sumUS.Add(us)
h.buckets[bucketFor(us)].Add(1)
}
func bucketFor(us int64) int {
if us <= 0 {
return 0
}
// bits.Len64 gives the position of the highest set bit, which is
// floor(log2)+1: exactly the bucket index we want.
i := bits.Len64(uint64(us))
if i > 23 {
i = 23
}
return i
}
math/bits.Len64 compiles to a single CPU instruction on every architecture Go supports. So one observation is three atomic adds and one instruction, with no branching on bucket boundaries, no allocation, and no lock. The whole histogram is 26 machine words and can live inline in a struct.
// Quantile returns the upper edge of the bucket holding the qth quantile.
// It is approximate by construction: the answer is only ever accurate to a
// factor of two, which is usually enough to tell "fine" from "on fire".
func (h *Histogram) Quantile(q float64) time.Duration {
total := h.count.Load()
if total == 0 {
return 0
}
target := int64(q * float64(total))
seen := int64(0)
for i := range h.buckets {
seen += h.buckets[i].Load()
if seen >= target {
return time.Duration(int64(1)<<uint(i)) * time.Microsecond
}
}
return time.Duration(int64(1)<<23) * time.Microsecond
}
Be honest in the doc comment about accuracy. A p99 of "8 ms" here means "somewhere between 4 and 8 ms", and that is fine for alerting and useless for a latency SLO quoted to three decimal places. Real systems use HDR histograms or t-digests when they need better; the point of showing the cheap version is that you now know what the expensive ones are buying.
Note also that Quantile reads 24 atomics one at a time while other goroutines are writing, so the snapshot it returns never existed at any single instant. For a monitoring endpoint that is completely acceptable and worth knowing.
// Set is every metric the colony publishes. Named fields rather than a map
// of strings: the compiler checks the names and no lookup happens in the
// hot path.
type Set struct {
Delivered Counter
FailedPickups Counter
Lost Counter
Restarts Counter
Dropped Counter
Timeouts Counter
Requests Counter
AntsAlive Gauge
QueueDepth Gauge
RoundTrip Histogram
Decide Histogram
}
Most metrics libraries hand you counter("antfarm_requests_total").Inc(), which needs a map lookup, a string hash and usually a mutex or a sharded cache on every observation. A struct of named fields costs a pointer offset, is checked by the compiler, and shows up in your editor's autocomplete. You lose dynamic metric names, which for an application (as opposed to a library) you rarely want anyway.
// WriteProm renders the set in the Prometheus text exposition format, which
// is simple enough to produce by hand and is what most scrapers speak.
func (s *Set) WriteProm(w io.Writer) {
counter := func(name, help string, c *Counter) {
fmt.Fprintf(w, "# HELP %s %s\n# TYPE %s counter\n%s %d\n", name, help, name, name, c.Value())
}
...
counter("antfarm_food_delivered_total", "Food units delivered to the nest.", &s.Delivered)
gauge("antfarm_ants_alive", "Ants currently running.", &s.AntsAlive)
hist("antfarm_round_trip", &s.RoundTrip)
}
Local closures as helpers inside a function is a very Go thing to do. They capture w, they are not visible outside, and they turn eleven repetitive blocks into eleven readable lines. The naming convention is worth copying: _total suffix for counters, base units (seconds, bytes) rather than milliseconds, and a # HELP line that says what the number means.
// Handler builds the observability endpoints. We mount pprof explicitly
// rather than relying on its init() registering itself on the default mux,
// because exposing profiles by accident is a real security problem.
func (s *Set) Handler() http.Handler {
mux := http.NewServeMux()
mux.HandleFunc("GET /healthz", func(w http.ResponseWriter, r *http.Request) {
fmt.Fprintln(w, "ok")
})
mux.HandleFunc("GET /metrics", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/plain; version=0.0.4; charset=utf-8")
s.WriteProm(w)
})
mux.HandleFunc("GET /debug/pprof/", pprof.Index)
mux.HandleFunc("GET /debug/pprof/cmdline", pprof.Cmdline)
mux.HandleFunc("GET /debug/pprof/profile", pprof.Profile)
mux.HandleFunc("GET /debug/pprof/symbol", pprof.Symbol)
mux.HandleFunc("GET /debug/pprof/trace", pprof.Trace)
return mux
}
Almost every Go tutorial tells you to write import _ "net/http/pprof". The underscore means "import for side effects only", and the side effect is that the package's init() registers its handlers on http.DefaultServeMux. If anything in your program then serves DefaultServeMux on a public port, you have published heap profiles, goroutine stacks, command-line arguments and a CPU profiler to the internet. This has caused real incidents.
Mounting the handlers yourself on your own mux, as above, makes the exposure a deliberate decision. In production you would bind this listener to localhost or an internal interface and reach it through a tunnel.
Note the Go 1.22 method patterns: "GET /metrics" matches only GET requests, and "GET /debug/pprof/" with a trailing slash is a subtree match. Before 1.22 you checked r.Method by hand.
The engine gets an embedded metrics.Set, and the interesting instrumentation is two lines in the request path:
func (e *Engine) handle(r request) {
e.M.Requests.Inc()
e.M.QueueDepth.Set(int64(len(e.reqs)))
...
}
len() on a channel returns how many values are buffered in it right now. It is a genuinely useful gauge (queue depth is the single best early warning that a consumer is falling behind) and it is a terrible basis for logic, because by the time you act on it the number has changed. Measure with it, never branch on it.
func (c *antClient) roundTrip(ctx context.Context, r request) (response, error) {
start := time.Now()
defer func() { c.e.M.RoundTrip.Observe(time.Since(start)) }()
...
}
And the HTTP server in main.go, with the shutdown discipline from Milestone 7:
// serve starts an HTTP server and returns a function that shuts it down
// without dropping in-flight requests.
func serve(addr string, h http.Handler, name string) func() {
srv := &http.Server{
Addr: addr,
Handler: h,
ReadHeaderTimeout: 5 * time.Second, // never accept a slow-loris header
}
go func() {
if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
log.Printf("antfarm: %s server: %v", name, err)
}
}()
return func() {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
if err := srv.Shutdown(ctx); err != nil {
log.Printf("antfarm: %s server shutdown: %v", name, err)
}
}
}
ListenAndServe always returns a non-nil error, and it returns http.ErrServerClosed when you shut it down on purpose, so the errors.Is check is how you avoid logging a scary message during a clean exit. srv.Shutdown stops accepting, waits for in-flight handlers, and respects the context as a deadline. ReadHeaderTimeout is there because the zero-value http.Server has no timeouts at all, which means one slow client can hold a connection open indefinitely.
$ antfarm -ants 1500 -grid 96x96 -food 8 -food-per-source 400 \
-duration 6s -metrics 127.0.0.1:9090 -view 127.0.0.1:8080 \
-chaos crash=0.0005,drop=0.002
colony: 1500 ants, 96x96 grid, 8 food sources of 400, behaviour trail, chaos {...}
metrics: http://127.0.0.1:9090/metrics pprof: http://127.0.0.1:9090/debug/pprof/
view: http://127.0.0.1:8080/
$ curl -s 127.0.0.1:9090/metrics
# HELP antfarm_food_delivered_total Food units delivered to the nest.
# TYPE antfarm_food_delivered_total counter
antfarm_food_delivered_total 955
# HELP antfarm_ant_restarts_total Ants restarted after a crash.
# TYPE antfarm_ant_restarts_total counter
antfarm_ant_restarts_total 792
# HELP antfarm_messages_dropped_total Requests discarded by chaos injection.
# TYPE antfarm_messages_dropped_total counter
antfarm_messages_dropped_total 6141
# HELP antfarm_request_timeouts_total Requests that never got a reply.
# TYPE antfarm_request_timeouts_total counter
antfarm_request_timeouts_total 5865
# HELP antfarm_requests_total Requests handled by the world owner.
# TYPE antfarm_requests_total counter
antfarm_requests_total 3034578
# HELP antfarm_ants_alive Ants currently running.
# TYPE antfarm_ants_alive gauge
antfarm_ants_alive 1500
# HELP antfarm_queue_depth Requests waiting for the world owner.
# TYPE antfarm_queue_depth gauge
antfarm_queue_depth 0
# TYPE antfarm_round_trip_seconds summary
antfarm_round_trip_seconds{quantile="0.5"} 0.000002
antfarm_round_trip_seconds{quantile="0.9"} 0.001024
antfarm_round_trip_seconds{quantile="0.99"} 0.008192
antfarm_round_trip_seconds_count 3033172
Three million requests in about four seconds, and look at the latency distribution: a median of 2 µs, a p90 of 1 ms, a p99 of 8 ms. That is a spread of four thousand to one between the typical case and the bad case, and it is the signature of a queue. When the owner is free your request is answered immediately; when 1,500 ants arrive at once you wait behind them. Averages would have hidden this completely: the mean here is around 100 µs, a number that describes nobody's experience. Always look at the tail.
$ curl -o cpu.prof "127.0.0.1:9090/debug/pprof/profile?seconds=10"
$ go tool pprof -top -nodecount=14 cpu.prof
File: antfarm
Type: cpu
Duration: 10.10s, Total samples = 9970ms (98.69%)
flat flat% sum% cum cum%
1030ms 10.33% 10.33% 3230ms 32.40% runtime.selectgo
1000ms 10.03% 20.36% 1000ms 10.03% runtime.nanotime (partial-inline)
820ms 8.22% 28.59% 820ms 8.22% runtime.lock2
720ms 7.22% 35.81% 790ms 7.92% runtime.unlock2
560ms 5.62% 41.42% 560ms 5.62% time.Now
350ms 3.51% 44.93% 520ms 5.22% runtime.casgstatus
320ms 3.21% 48.14% 320ms 3.21% runtime.duffcopy
220ms 2.21% 50.35% 800ms 8.02% runtime.sellock
200ms 2.01% 54.56% 6120ms 61.38% sim.(*Engine).runAnt
180ms 1.81% 56.37% 770ms 7.72% runtime.mallocgc
120ms 1.20% 60.68% 450ms 4.51% sim.(*antClient).roundTrip.func1
That is a real profile of our real program, and it is damning in the most useful way. Read the columns first: flat is time spent in that function itself, cum is time in it plus everything it called. So runAnt has 2% flat and 61% cumulative: it does nothing itself and everything beneath it.
Now read the names. selectgo (32% cumulative) is the runtime implementing select. lock2, unlock2 and sellock are the locks the runtime uses inside channels. nanotime and time.Now together are 15%. casgstatus is goroutine state transitions, which is scheduling. Add it up: the overwhelming majority of the CPU is channel machinery, timers and scheduling, and almost none of it is ant simulation.
Not one line of Decide, senseAt or applyAction appears in the top fourteen. Our program is not computing; it is communicating about computing.
time.Now at 5.6% flat, plus a large share of nanotime's 10%, is substantially the RoundTrip histogram we added twenty minutes ago: two clock reads per round trip, three million round trips. Observability is not free, and measuring at microsecond granularity in a path that takes microseconds is measuring the measurement.
This is not an argument against instrumenting. It is an argument for knowing the cost and sampling when it is too high: observe one round trip in every hundred and multiply, which loses nothing statistically and cuts the overhead by 99%. Exercise 9 asks you to do it.
time.Now has fallen. Then answer: what does sampling do to _count, and how should the exposition handle that?curl -o heap.prof 127.0.0.1:9090/debug/pprof/heap, then go tool pprof -top heap.prof. Identify the largest allocation site in the colony and explain why it exists.curl "127.0.0.1:9090/debug/pprof/goroutine?debug=1" and find the line where 1,500 ants are blocked. This is the single most useful debugging technique in Go and it takes ten seconds.1. Sampling. The cheapest correct approach is a shared atomic counter and modulo:
type SampledHistogram struct {
Histogram
n atomic.Int64
rate int64 // observe one in every rate
}
func (h *SampledHistogram) Maybe(start time.Time) {
if h.n.Add(1)%h.rate != 0 {
return
}
h.Observe(time.Since(start))
}
But note the trap: the caller still has to call time.Now() to produce start, so the expensive part happens anyway. The version that actually saves time decides before reading the clock, which means the sampling decision must happen at the top of roundTrip:
var start time.Time
sampled := c.e.M.RoundTrip.Should() // one atomic add, no clock
if sampled {
start = time.Now()
}
defer func() {
if sampled {
c.e.M.RoundTrip.Observe(time.Since(start))
}
}()
This is a small but perfect example of why you profile rather than guess: the obvious refactor moves no work at all, because the cost was in the clock read, not in the histogram.
_count now undercounts by a factor of rate, so either multiply it on the way out (and document that it is estimated) or publish the raw request counter separately and let the dashboard divide. Silently exposing a sampled count as though it were exact is how dashboards start lying.
2. Heap profile. The largest live allocation is the grids: a 128×128 []int plus a []float64, allocated once. The largest rate of allocation is the frame snapshots from Milestone 10, which build three slices per frame. Everything in the ant path allocates nothing, which is what 0 allocs/op in the Milestone 2 benchmark predicted.
3. Goroutine profile. You will see something like 1500 @ ... sim.(*antClient).roundTrip followed by runtime.selectgo, which tells you where they are parked and, with debug=2, how long they have been there. When a Go service hangs in production, this endpoint is almost always how you find out why, and it works on a process that is otherwise unresponsive.
Remove the ReadHeaderTimeout: 5 * time.Second line from serve's http.Server and rebuild. Open a raw connection to the metrics port (nc 127.0.0.1 9090, or a small Go program that dials and never writes) and leave it sitting there; it stays open indefinitely, because the zero-value http.Server has no timeouts of any kind. Open a few hundred of these in a loop and watch the process happily hold every one open forever. Put the line back, repeat, and confirm each connection is dropped after five seconds. That single struct field is the entire difference between a server that resists a slow-loris attack and one that does not.
import _ "net/http/pprof" plus a public DefaultServeMux. Profiles on the internet.Counter. The copy has its own value and stops tracking. go vet catches it.len(ch). It is stale the instant you read it. Fine to publish, wrong to decide with.http.Server with no timeouts. The zero value has none, which is a denial-of-service waiting to happen.bits.Len64 produce a histogram bucket, and why is that cheap?import _ "net/http/pprof"?runAnt is 2% flat and 61% cumulative. Explain both numbers.time.Now in the profile, and what would you do about it?This milestone leans on three things the standard library gives you for nothing: sync/atomic for counters that need no lock, net/http/pprof for a production-grade profiler built into every binary, and a runtime that starts 1,500 goroutines in milliseconds so instrumenting them barely shows up next to the cost of the workload itself. Attaching a CPU profiler to a live process with two lines of wiring, no agent, no restart, no separate tool to install, is a genuine and unusual convenience.
The honest counterweight: Go's atomics give you a counter and a gauge for free, but nothing richer. A proper HDR histogram, client-side aggregation across processes, or exemplars linking a slow trace to its exact metric bucket are things a mature client library (or a language with one, like Java's Micrometer) hands you out of the box; here we built a deliberately approximate histogram by hand and said so in the doc comment. And pprof's profile format, while excellent for CPU and heap, is Go-specific — you cannot point the same tool at a Python worker in the same fleet. For a single Go binary this milestone is close to the best case for the standard library; for a polyglot system you would reach for OpenTelemetry regardless of language and treat Go's native tools as a local debugging aid rather than the whole observability story.
See the colony. An ANSI-redrawn terminal view for a quick look, and a browser page fed by server-sent events for a good one. Neither may ever slow the simulation down.
Aggregating at the source, snapshot values as an interface between subsystems, ANSI escape codes, server-sent events, http.Flusher, request contexts as the client's lifetime, and dropping frames as a policy.
The naive viewer asks for the whole world and draws it. At 512×512 with pheromone that is half a megabyte per frame, serialised to JSON, sixty times a second, and all of it built on the owner's goroutine while 50,000 ants wait. The viewer would become the bottleneck of the simulation it is supposed to observe.
So the rule is: aggregate at the source. The owner builds a small, fixed-size Frame, downsampled to whatever resolution the client asked for, and that is all that ever leaves the simulation.
owner goroutine viewer(s)
─────────────── ─────────
holds 512x512 grids ──► Frame: 96x96 buckets, ~3 slices
builds one Frame terminal renderer → ANSI text
per request SSE handler → JSON over HTTP
(either, both, or none)
A second rule follows from the first: the viewer may never block the simulation. If a frame cannot be produced in time, the viewer skips it. A dropped frame is invisible to a human; a stalled colony is not.
A typical "add a live view" implementation queries the primary data store directly, once per connected client, at whatever resolution the client asked for — that is how most admin dashboards and BI tools are built, and it works fine when the store is a database built to handle read concurrency.
Here the "store" is 50,000 goroutines mid-flight and a single owner goroutine that is also the bottleneck of the entire simulation, so querying it per client at full resolution would make the viewer the most expensive client the system has. The idiomatic move is to aggregate once, at the source, into a small fixed-size Frame, and let every viewer — terminal, browser tab, or a hundred browser tabs — share that same cheap snapshot. It is more code up front than "just query it", and it is the only version that does not fall over the moment someone opens a second tab.
// Frame is a downsampled picture of the colony, built by the owner and
// small enough to send many times a second. Aggregating at the source
// rather than shipping the whole grid is what keeps the viewer cheap.
type Frame struct {
Cols int `json:"cols"`
Rows int `json:"rows"`
Food []int `json:"food"` // Cols*Rows, summed per block
Pher []float64 `json:"pher"` // Cols*Rows, summed per block
Ants []int `json:"ants"` // Cols*Rows, ants per block
Nest [2]int `json:"nest"`
Stats Stats `json:"stats"`
}
func (e *Engine) frameNow(cols int) Frame {
if cols <= 0 || cols > e.cfg.Width {
cols = min(64, e.cfg.Width)
}
rows := max(1, e.cfg.Height*cols/e.cfg.Width)
f := Frame{ /* ... allocate cols*rows slices ... */ }
for y := range e.cfg.Height {
for x := range e.cfg.Width {
p := world.Position{X: x, Y: y}
food, pher := e.world.Food.At(p), e.world.Pher.At(p)
if food == 0 && pher == 0 {
continue
}
i := (y*rows/e.cfg.Height)*cols + x*cols/e.cfg.Width
f.Food[i] += food
f.Pher[i] += pher
}
}
for _, p := range e.positions {
f.Ants[(p.Y*rows/e.cfg.Height)*cols+p.X*cols/e.cfg.Width]++
}
f.Stats = e.statsNow()
return f
}
Two things to notice.
The struct tags (`json:"cols"`) control JSON field names. They are a string literal attached to a field, read by encoding/json through reflection at run time. go vet checks them for you, which is how I caught a duplicate tag while writing this: struct field Rows repeats json tag "cols". A typo in a string that silently produces the wrong wire format is exactly the kind of thing that should be checked, and it is.
The continue on empty cells is not micro-optimisation, it is the difference between scanning and doing work: most of a colony's grid is empty most of the time, so the loop touches 262,144 cells but does arithmetic on very few. The scan itself is the remaining cost, and Exercise 10 asks what to do about it.
Where do ant positions come from? The owner does not own ants, so it cannot read their positions. But it sees every move request, so it can remember them:
case ActMove:
np := e.world.Clamp(r.pos.Add(r.act.Dir))
if r.id >= 0 && r.id < len(e.positions) {
e.positions[r.id] = np
}
A slice indexed by ant ID, written only by the owner. The viewer gets ant positions with no extra messages and no synchronisation, because the information was already flowing past. Deriving a view from a message stream you already have is the cheapest kind of observability, and it generalises: event sourcing is this idea taken seriously.
// shades map a density to a character, densest last.
var shades = []rune{' ', '.', ':', '-', '=', '+', '*', '#', '%', '@'}
// Terminal redraws the colony in place using ANSI escape codes.
func Terminal(ctx context.Context, w io.Writer, e *sim.Engine, cols int, every time.Duration) error {
ticker := time.NewTicker(every)
defer ticker.Stop()
fmt.Fprint(w, "\x1b[?25l") // hide the cursor
defer fmt.Fprint(w, "\x1b[?25h") // and always put it back
for {
select {
case <-ctx.Done():
return nil
case <-ticker.C:
f, err := e.Frame(ctx, cols)
if err != nil {
continue // a dropped frame is not worth stopping for
}
fmt.Fprint(w, "\x1b[H\x1b[2J") // home, then clear
fmt.Fprint(w, Render(f))
}
}
}
ANSI escape codes are a tiny language your terminal speaks: \x1b is the escape character, [H moves the cursor to the top left, [2J clears the screen, [?25l and [?25h hide and show the cursor. No library needed for this much. The defer that restores the cursor matters: a program that exits with the cursor hidden leaves the user's terminal broken, and "put back what you changed" applies to terminals as much as to files.
Render is a separate pure function taking a Frame and returning a string, so it can be tested without a terminal, a simulation or a goroutine. Splitting "produce the bytes" from "write the bytes to a live device" is a habit worth having; it also makes the renderer reusable by the web view if you ever want an ASCII mode.
func Render(f sim.Frame) string {
var b strings.Builder
b.Grow(f.Rows*(f.Cols+1) + 256)
...
}
strings.Builder with Grow is the right way to assemble a string in Go: building with s += ... in a loop is O(n²) because strings are immutable and every concatenation copies. Grow pre-sizes the buffer so the builder never reallocates.
For the browser we need a push channel. The options are WebSockets (bidirectional, needs a library or a hand-rolled handshake), long polling (awkward), or server-sent events (one-way, plain HTTP, built into every browser as EventSource, about fifteen lines of server code). We only push, so SSE wins easily.
func stream(w http.ResponseWriter, r *http.Request, e *sim.Engine, cols int, every time.Duration) {
flusher, ok := w.(http.Flusher)
if !ok {
http.Error(w, "streaming unsupported", http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
// r.Context() is cancelled when the browser tab closes, which is how
// this goroutine learns to stop.
ctx, cancel := context.WithCancel(r.Context())
defer cancel()
ticker := time.NewTicker(every)
defer ticker.Stop()
enc := json.NewEncoder(w)
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
f, err := e.Frame(ctx, cols)
if err != nil {
continue // drop this frame rather than stall the simulation
}
fmt.Fprint(w, "data: ")
if err := enc.Encode(f); err != nil {
return
}
fmt.Fprint(w, "\n")
flusher.Flush()
}
}
}
data: then a payload then a blank line. json.Encoder.Encode conveniently appends a newline, so the extra Fprint(w, "\n") completes the pair.w.(http.Flusher) is a type assertion: http.ResponseWriter is an interface, and the concrete value behind it may or may not support flushing. Without Flush, Go buffers the response and the browser receives nothing until the handler returns, which for a stream is never. This is the classic SSE bug.r.Context() is cancelled when the client disconnects. That is how this goroutine finds out the tab was closed. Without it, every closed tab would leave a goroutine ticking forever, and a page that the user opens and closes fifty times leaks fifty goroutines.continue on error is the "never block the simulation" rule. If Frame times out because the owner is saturated, we skip that frame and try again in 100 ms. The browser client is 25 lines of dependency-free JavaScript in a Go string constant: an EventSource, a canvas, and a loop writing RGBA pixels straight into an ImageData. Ants go in the red channel, food in green, pheromone in blue and a little green, and the canvas is scaled up with image-rendering: pixelated. No build step, no framework, and the binary still has nothing to ship alongside it.
frameNow takes at 512×512 by adding a histogram, then serve two browser tabs at 60 fps and watch the round-trip p99 for ants. Quantify the harm the viewer does.Frame to every client that asks in between. Second: keep running per-block totals updated incrementally as food and pheromone change, so building a frame is a copy rather than a scan. Compare the complexity of the two.1. Expect several milliseconds per frame at 512×512, because the scan is a quarter of a million cells. Two clients at 60 fps is 120 scans per second, which can consume a large fraction of the owner's time, and you will see the ant round-trip p99 climb. The viewer, added to observe the system, has changed the system.
2a. Caching. Roughly ten lines: keep lastFrame Frame and lastFrameAt time.Time in the owner, and serve the cached copy if it is fresh enough. It decouples client count from cost entirely, which is the whole problem, and it is the right first fix.
case reqFrame:
if time.Since(e.lastFrameAt) < e.cfg.FrameInterval {
r.reply <- response{frame: e.lastFrame}
return
}
e.lastFrame = e.frameNow(r.id)
e.lastFrameAt = time.Now()
r.reply <- response{frame: e.lastFrame}
One subtlety: the cached Frame contains slices, and now several goroutines hold copies of the struct pointing at the same backing arrays. Since nothing mutates a frame after it is built, that is safe, but it is safe by convention rather than by construction, and it deserves a comment. If a client ever wanted to modify a frame in place, this would become a race.
2b. Incremental totals. Keep blockFood []int and blockPher []float64 updated in applyAction, so the frame is a copy of two small slices. Frame cost becomes proportional to the frame, not the world. The price: every mutation site must remember to update the totals, the block mapping is now duplicated in two places, and any bug produces a view that drifts slowly away from reality. This is the classic denormalisation trade, and the honest answer for this project is to cache first and only denormalise if the profile still says so.
3. Pause. The tempting place is a paused atomic.Bool that every ant checks, which works and spreads the concept over 50,000 goroutines. The tidier answer is a single paused flag owned by the owner: when paused, it stops answering reqSense and reqAct and just buffers, so ants block naturally in their round trips with no new code at all. It also gives you the right semantics for free (a paused colony's clients are waiting, not spinning) and pausing becomes a property of the world rather than an agreement among ants. Watch out for one thing: ants will start timing out, so pausing should also suppress the timeout counter, or your metrics will report an incident every time someone hits the button.
Comment out the flusher.Flush() call in stream and rebuild. Open the browser view: the connection succeeds, curl -N -i 127.0.0.1:8080/stream even shows the right Content-Type: text/event-stream header, but no frame ever arrives, because Go's net/http buffers the response until either the buffer fills or the handler returns — and a stream handler never returns. The canvas just sits black. Put Flush() back and watch frames appear within the tick period. This is the single most common bug in a hand-rolled SSE handler, and it produces no error anywhere: the connection is open, the headers are correct, and the browser is simply waiting for bytes that Go is holding onto.
Flush(). The stream appears dead; the browser shows nothing forever.r.Context(). A goroutine per closed tab, leaked permanently. ResponseWriter from two goroutines. It is not safe for concurrent use. One handler, one writer.+= in the render loop. Quadratic.http.Flusher do and what happens without it?Render a separate pure function?Two standard-library pieces carry this milestone: net/http's Flusher interface makes server-sent events a fifteen-line handler with no dependency, and goroutines make "one goroutine renders ANSI, one goroutine streams SSE, neither ever touches the simulation's memory directly" cheap enough not to think about. Shipping a live web view with zero JavaScript dependencies and zero build step, embedded as a Go string constant, is the kind of thing that is disproportionately easy in Go's specific combination of a capable stdlib and a static binary.
The honest cost: this whole milestone is essentially reimplementing a sliver of what an embedded dashboard framework or a proper time-series database plus Grafana gives you for free — retention, zooming, multiple metrics on one chart, alerting. For a single simulation you want to glance at, hand-rolled ANSI and a canvas are the right amount of engineering. For anything you would show to another team or keep historical data from, reach for the real tool; you would not want to grow this Frame/SSE setup into a general-purpose monitoring stack by hand.
Take the profile from Milestone 9 seriously and rewrite the hot path. Sensing becomes a lock-free read with no message at all, moves and drops become one-way messages, and writes are sharded across several owners. Then measure honestly.
Atomics for concurrent reads with single-owner writes, fixed-point arithmetic because Go has no atomic float, one-way messages, sharding by region, closing channels to drain them, and a genuine correctness bug found by a conservation test.
The profile said: channel operations, timers and scheduling dominate; simulation logic is invisible. So the optimisation strategy writes itself. Do fewer channel operations. Where do they come from?
| Operation | Frequency | Messages | Could it be fewer? |
|---|---|---|---|
| Sense | every tick, every ant | 2 (request + reply) | Yes: reads can be atomic |
| Move | ~95% of ticks | 2 | Yes: the new position is computable locally; only the pheromone deposit needs the world |
| Pickup | rare | 2 | No: only the owner can decide who gets the last unit |
| Drop | rare | 2 | Mostly: the ant knows it is at the nest |
So roughly 99% of the traffic is sensing and moving, and both can largely be eliminated. Three changes:
atomic.Int64, any goroutine may read it race-free at any time with no message. Writes stay single-owner, so there is no compare-and-swap loop and no contention between writers.Milestone 5 engine Milestone 11 engine
────────────────── ───────────────────
sense → message → owner → reply sense → atomic loads (no message)
move → message → owner → reply move → local compute
+ one-way deposit if carrying
pickup → message → owner → reply pickup → message → shard → reply
drop → message → owner → reply drop → one-way message → shard
1 owner, ~2 msgs/tick/ant N shards, ~0.05 msgs/tick/ant
The single-owner design had a property we are giving up: every world access was validated by one authority that saw everything in order. Now an ant computes its own position, so a buggy behaviour can put itself somewhere impossible; a drop is fire-and-forget, so nobody tells the ant it was rejected. We keep validation exactly where correctness demands it (a pickup still needs an authoritative answer about who got the unit) and drop it where the ant can be trusted.
Do this after a profile, never before. The Milestone 5 engine is the one I would ship if 3 million requests a second were enough, and it is the one to write first in any new system.
Faced with a profile dominated by channel and lock contention, a common first instinct is to reach for a general-purpose concurrent container — sync.Map, or a third-party lock-free hash map — and swap it in without changing the surrounding design. That often helps a little and rarely helps enough, because it treats the symptom (a data structure under contention) rather than the cause (a single serialisation point every operation must pass through).
The idiomatic Go move, and the one this milestone takes, is to change who owns what: atomic loads for the 99% of traffic that only reads, one-way messages for writes that need no answer, and sharding by the access pattern that actually exists (rows, because ants move locally) rather than by a generic hash. None of this needed a special data structure, only atomic.Int64, a channel, and a decision about ownership. The cost is that it is bespoke to this workload — a sync.Map would have taken five minutes to try, and this took an afternoon and a bug — and a workload without spatial locality would get nothing from the sharding half of it.
// pherScale turns float concentrations into fixed-point integers. Go has no
// atomic float64, so we store thousandths and convert on read.
const pherScale = 1000
// AtomicGrid is a grid that any goroutine may read without synchronisation
// and that only the owner of a row band may write. Reads use atomic loads,
// so they are race-free; writes stay single-owner, so no compare-and-swap
// loop is needed.
type AtomicGrid struct {
W, H int
cells []atomic.Int64
}
func (g *AtomicGrid) At(p Position) int64 {
if !g.InBounds(p) {
return 0
}
return g.cells[p.Y*g.W+p.X].Load()
}
func (g *AtomicGrid) AtFloat(p Position) float64 {
return float64(g.At(p)) / pherScale
}
// TakeOne removes a unit if there is one, reporting whether it succeeded.
func (g *AtomicGrid) TakeOne(p Position) bool {
if !g.InBounds(p) {
return false
}
c := &g.cells[p.Y*g.W+p.X]
if c.Load() <= 0 {
return false
}
c.Add(-1)
return true
}
Go has no atomic float64, and the reason is that the atomic instructions operate on integers. The standard workaround is fixed point: store thousandths as an integer and divide on read. You give up range and precision (a pheromone of 0.0004 is zero here) and gain a data type that atomic.Int64 can hold. The alternative, math.Float64bits with a CAS loop, gives you exact floats at the cost of a retry loop; fixed point is simpler and good enough for a concentration.
TakeOne is load-then-add rather than a compare-and-swap loop, and that is safe only because a single shard owns every write to this row. Two writers would race between the Load and the Add and could both take the last unit, driving the cell negative. The race detector would not catch it, because both operations are atomic; the food conservation test would. Atomics prevent torn reads and writes, not logical races. Write the ownership rule in a comment on the type, as this code does, because nothing else enforces it.
// shardFor maps a row to the goroutine that owns it.
func (e *FastEngine) shardFor(p world.Position) chan writeOp {
i := p.Y * len(e.shard) / e.cfg.Height
...
return e.shard[i]
}
// runShard owns one band of rows: it is the only writer of those cells.
func (e *FastEngine) runShard(ctx context.Context, i int) {
lo, hi := e.bandOf(i)
ticker := time.NewTicker(e.cfg.EvaporateEvery)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
e.pher.EvaporateBand(lo, hi, e.cfg.Evaporation)
case op, open := <-e.shard[i]:
if !open {
return // every ant has stopped and the queue is drained
}
switch op.kind {
case opDeposit:
e.pher.AddFloat(op.pos, e.cfg.Deposit, e.cfg.PherMax)
case opPickUp:
op.reply <- e.food.TakeOne(op.pos)
case opDrop:
e.M.Delivered.Inc()
}
}
}
}
Each shard now also evaporates its own band, which turns Milestone 6's single O(cells) pass on one goroutine into N parallel passes of O(cells/N). The horizontal-band split is chosen because ants move locally: an ant spends most of its life in one or two bands, so cross-shard traffic is low. A hash of the position would spread load perfectly and destroy that locality.
The sense path now needs no shard at all:
// senseFast reads the world directly. No message, no reply, no waiting.
func (e *FastEngine) senseFast(pos world.Position) Sense {
s := Sense{Pos: pos, Nest: e.nest, Food: int(e.food.At(pos))}
for i, d := range directions[:8] {
n := pos.Add(d)
s.Adjacent[i] = int(e.food.At(n))
s.Pher[i] = e.pher.AtFloat(n)
}
return s
}
Seventeen atomic loads and a struct return. Measured: 29 nanoseconds, against a median of 2 µs for the message-based version. Roughly seventy times faster for the single most frequent operation in the program.
// send is fire-and-forget with load shedding: if the owning shard is
// saturated, the deposit is dropped rather than stalling the ant.
func (e *FastEngine) send(ctx context.Context, op writeOp) {
select {
case e.shardFor(op.pos) <- op:
case <-ctx.Done():
default:
e.M.Dropped.Inc()
}
}
This is backpressure policy expressed in four lines, and the choice it encodes is: a lost pheromone deposit is better than a stalled ant. That is correct for pheromone, which is statistical and evaporates anyway, and it would be catastrophic for the pickup path, which is why pickups do not use this function. Shedding policy is per-message-type, not per-system, and deciding it message by message is the whole discipline.
The first version passed -race cleanly and failed this:
--- FAIL: TestFastEngineConservesFood (0.74s)
fast_test.go:22: food not conserved: 359 != 500
(ants 500 carrying 182 delivered 81 food left 96)
141 units of food had ceased to exist. No race, no panic, no error in any log. Instrumenting the ledger showed 336 successful pickups, 216 drops, 90 ants carrying: 30 units taken from the ground that no ant held and nobody delivered.
The cause was two shutdown bugs, both of them the kind that only appear at the boundary:
// BUG 1: when ctx is cancelled while we wait for the pickup answer,
// select may choose Done even though the reply is ready. The shard
// already took the food. We just threw it away.
select {
case ok := <-reply:
a.Carrying = ok
case <-ctx.Done():
return
}
Remember that select picks randomly among ready cases. At cancellation time, a couple of dozen ants were mid-pickup; for each, the world had already removed a unit of food, and half of them discarded the answer. The fix is to insist on collecting an answer you have already paid for:
case <-ctx.Done():
// The shard may already have taken a unit on our behalf.
// Collect the answer anyway, or that food vanishes from
// the ledger.
timer := time.NewTimer(e.cfg.RequestTimeout)
defer timer.Stop()
select {
case ok := <-reply:
a.Carrying = ok
case <-timer.C:
e.M.Lost.Inc()
}
return
Bug two was queued drops being discarded when the shards were cancelled with work still in their channels. The fix uses the one safe way to close a channel:
<-ctx.Done()
antWG.Wait() // every sender has stopped, so closing is now safe
for i := range e.shard {
close(e.shard[i])
}
shardWG.Wait()
Closing a channel that senders might still use panics, which is why Milestone 5 forbade it. Here it is not only safe but exactly right, because antWG.Wait() has already established that every sender has stopped. A closed channel keeps delivering its buffered values before reporting closed, so op, open := <-ch drains the queue and then exits. Close-after-senders-finish is the idiomatic way to say "process everything remaining, then stop".
Two lessons worth more than the speedup. First, a race detector cannot find a logic bug: every operation here was properly synchronised and the program was still wrong. Invariant tests find what sanitisers cannot. Second, shutdown is where work gets lost, in this project and in every queueing system, because that is the one moment when the pipeline is asked to stop while it is still full.
Food delivered in a fixed two-second window, same seed, same world, single-owner engine versus sharded:
500 ants: owner 827 delivered | sharded 1180 delivered | 1.4x
2000 ants: owner 811 delivered | sharded 1458 delivered | 1.8x
(an earlier run on a quieter machine: 1.4x and 2.9x)
Be careful about what this shows. My container has one core, so none of this gain comes from parallelism; it is all reduced overhead, which is exactly what the profile predicted. On a multi-core machine you should see considerably more, because sharded writes and parallel evaporation can then genuinely run at the same time. Run it yourself and compare with GOMAXPROCS set to 1, 2, 4 and 8; unlike the Milestone 4 mutex version, this one should improve rather than degrade.
Notice also that the gain grows with the number of ants (1.4× at 500, 1.8–2.9× at 2,000). That is the signature of removing a serialisation point: with few clients the owner was never the bottleneck, so removing it changes little.
go build -gcflags='-m' ./internal/sim/ 2>&1 | grep escapes and find one allocation in the ant path. Explain why it escapes to the heap and whether you can prevent it.1. With sensing free and most messages gone, the profile shifts to Decide (which is now genuinely a large share, since it is the only real computation left), the random number generator inside weightedStep, and EvaporateBand. That is a healthy profile: the program is finally spending its time on the simulation. The next optimisation would be algorithmic (lazy evaporation with timestamps instead of periodic sweeps), not mechanical.
2. reply := make(chan bool, 1) in the pickup path escapes, because the shard goroutine keeps a reference to it. It is allocated once per ant, not per operation, so it is not worth removing. Escape analysis output is dense; the useful discipline is to read it once per project and learn which patterns in your code allocate. The general rule in Go: a value escapes if its address outlives the function, which includes being stored in an interface, sent on a channel, or captured by a goroutine.
3. Too few shards and you are back to a serialisation point. Too many and you get: more goroutines than cores (so scheduling overhead with no parallelism gain), thinner bands (so more cross-band traffic as ants wander across boundaries), and N evaporation tickers firing independently. The optimum is usually close to GOMAXPROCS, which is why that is a sensible default, and the flat region around it is usually wide.
4. With bands, cross-shard traffic is low precisely because ants move one cell at a time; you should measure a small percentage. With a hash, essentially every operation crosses. The interesting part is that throughput may not differ much on a single machine, because a channel send costs the same either way. It matters enormously once shards live on different machines and a cross-shard operation becomes a network hop. Locality is cheap insurance that only pays out when you distribute, which is precisely Milestone 12.
The conservation test passes because every pickup is serialised through its shard's channel. Break that deliberately: change the pickup path so an ant calls e.food.TakeOne(pos) directly instead of sending opPickUp to the shard, and set a single food source with enough ants that several routinely stand on the same cell. Run go test -race: clean, because every access is a correctly-used atomic. Now run the conservation test by itself with -count=50: it fails intermittently, with more food taken from the ground than was ever delivered or lost. TakeOne's load-then-add is not atomic as a sequence, and this is what "atomics prevent torn reads and writes, not logical races" looks like when you trigger it instead of taking it on faith.
Decide, which was never the problem.-race proves correctness. It found nothing here; the conservation test found everything.WaitGroup proves they are done.float64 for pheromone?select's random choice caused food to disappear at shutdown.This is the milestone where Go's memory model and its atomics package earn their keep directly: atomic.Int64 gives every ant a race-free read of world state with no message, no lock, and a cost the profiler put at 29 nanoseconds. Go's happens-before guarantees around channel operations are also what let the shutdown fix (antWG.Wait() before closing the shard channels) be provably correct rather than merely untested. Few mainstream languages make "safe concurrent reads with no synchronisation primitive visible at the call site" this cheap and this explicit at once.
The honest cost sits right next to the win: Go has no atomic float64, which is why pheromone is stored as scaled integers, a workaround with real precision loss that a language with atomic floats (or one where you would reach for a different concurrency model entirely, like Erlang's isolated processes and message copying) would not force on you. And the conservation-test bug is a reminder that none of this machinery — not the race detector, not go vet, not the type system — checks the actual invariant you care about. Go gave us fast, correctly-synchronised primitives; it did not give us a correct program, and mistaking the first for the second is exactly how this bug shipped.
Run the world in one process and ants in another, talking over TCP. Kill the world, watch the ants wait and reconnect, and see what the network costs.
Interfaces as the seam between local and remote, encoding/gob, one goroutine per connection, connection ownership and why the client is not thread-safe, socket deadlines as the network form of a timeout, reconnection with jittered backoff, and a shutdown bug that only a network can give you.
This is the milestone that Milestone 5 was secretly building toward. Because the local protocol was already values in and values out, the change is mostly mechanical: give it a name as an interface, implement that interface twice.
// WorldService is everything an ant needs from the world. Both the local
// engine and the network client implement it, which is what lets an ant run
// unchanged against a world in another process.
type WorldService interface {
Sense(pos world.Position) (Sense, error)
Act(id int, pos world.Position, carrying bool, act Action) (ActResult, error)
}
Two methods. The local engine satisfies it with atomic reads and channel sends; the network client satisfies it with gob over a socket. Notice that both methods return an error, including the local one where nothing can fail. Designing the interface for the failing implementation is what makes the seam work: had Sense returned only a Sense, the remote version could not have reported a broken connection without panicking or lying.
A typical way to split this into two processes is HTTP plus JSON: a request per Sense/Act call, a REST-shaped handler, and libraries in every language that can speak to it immediately. It is the safe default for interoperability, and it is not free — a new TCP connection or HTTP/1.1 round trip per call (or connection pooling to avoid that), JSON encoding overhead, and a text format for what is fundamentally a few numbers.
The idiomatic Go version here is a raw net.Conn per ant carrying encoding/gob, which needs no schema, no code generation and no library, and sends compact binary values once the connection's type descriptions are established. The cost is exactly what the client's doc comment admits: this protocol cannot be spoken by anything that is not Go, there is no multiplexing, and production systems reach for gRPC (protobuf over HTTP/2) to get language-agnostic interoperability, streaming, and flow control that this hand-rolled version does not have.
process A: ants process B: the world
─────────────── ────────────────────
ant goroutine ─► netsim.Client ──TCP──► netsim.Server ─► FastEngine
▲ (one per ant) (one goroutine (shards)
│ per connection)
same loop as milestone 5
// Req and Resp are the wire format. They contain only values, which is
// exactly why the local protocol from milestone 5 could become a network
// protocol without redesign.
type Req struct {
Kind Kind
ID int
Pos world.Position
Carrying bool
Act sim.Action
}
type Resp struct {
Sense sim.Sense
Pos world.Position
Nest world.Position
OK bool
Err string // errors do not survive gob; send the text
}
encoding/gob is Go's native binary serialisation: self-describing, type-safe between Go programs, and requiring no schema file or code generation. It is a good fit here and a poor fit for anything that must interoperate with another language, where you would reach for protobuf or JSON. gob sends type descriptions once per connection and then compact values, so a long-lived connection is efficient.
Errors do not travel. error is an interface, and gob can only encode concrete types it has been told about, so an error field would fail to encode or arrive as something unusable. Sending the message as a string and reconstructing an error on the far side is the simple, standard answer. The cost is that errors.Is no longer works across the wire, which is a real loss: a remote ErrNoFood arrives as an unrelated error with the same text. Production protocols solve this with numeric error codes, which is worth knowing as the reason gRPC has a status enum.
// Server exposes a sim.WorldService over TCP. Each connection gets one
// goroutine, which is the standard Go server shape and costs almost nothing
// because goroutines are cheap.
func (s *Server) handle(conn net.Conn) {
defer conn.Close()
dec := gob.NewDecoder(conn)
enc := gob.NewEncoder(conn)
for {
var req Req
if err := dec.Decode(&req); err != nil {
if !errors.Is(err, io.EOF) && !errors.Is(err, net.ErrClosed) {
log.Printf("netsim: decode: %v", err)
}
return
}
resp := s.dispatch(req)
if err := enc.Encode(resp); err != nil {
return // the client went away; drop the connection
}
}
}
Blocking reads in a dedicated goroutine per connection is the whole Go server model, and it is the clearest demonstration of why cheap goroutines matter: this code reads like a single-threaded program from 1995 and scales to tens of thousands of connections. In a language with expensive threads you would be writing an event loop, a state machine per connection, and callbacks. Node, Python and Rust all have async machinery precisely to get back to code that looks like this.
One gob.Decoder per connection, created once and reused: a gob stream is stateful, because type descriptions are sent once, so creating a new decoder per message would be both wrong and slow.
My first version of Close was the obvious one:
func (s *Server) Close() {
s.ln.Close() // stop accepting
s.wg.Wait() // wait for the connection goroutines
}
The test suite hung until I killed it after 300 seconds. Closing the listener stops new connections and does nothing whatsoever to existing ones: each connection goroutine was parked in dec.Decode(&req), waiting for a request from a client that was still connected and simply not talking. wg.Wait() waited forever.
// Close stops accepting, hangs up on every client, and waits for the
// connection goroutines. Closing the listener alone is not enough: a
// goroutine blocked in Decode on a live connection never notices, and
// Close would wait forever.
func (s *Server) Close() {
if s.ln != nil {
s.ln.Close()
}
s.mu.Lock()
s.closing = true
for conn := range s.conns {
conn.Close() // unblocks the goroutine parked in Decode
}
s.mu.Unlock()
s.wg.Wait()
}
The fix requires a registry of live connections, and that registry is the one place in this entire project where a mutex is the right tool:
// mu guards conns. A connection registry is exactly the kind of small
// shared map a mutex is right for: no ownership to transfer, just a set
// that several goroutines add to and one goroutine walks at shutdown.
mu sync.Mutex
conns map[net.Conn]struct{}
closing bool
There is no ownership to hand over and no flow of work, just a set that accept-time adds to, connection-exit removes from, and shutdown iterates. A channel-based version would be strictly worse. This is what the Milestone 5 comparison meant by "use mutexes for shared state with simple invariants": the mutex is not a failure of the message-passing design, it is the right tool for a different job. The closing flag closes the race where a connection is accepted after shutdown began.
The general lesson is bigger than the fix. Cancelling a goroutine blocked on network I/O requires closing the connection or setting a deadline; a context alone does nothing, because the blocking call is in the kernel. This is the same fact as Milestone 7's "you cannot kill a goroutine", wearing different clothes.
// Client is a sim.WorldService backed by a TCP connection. It is NOT safe
// for concurrent use: the protocol is one request and one response in
// order, so each client belongs to exactly one goroutine. That restriction
// is the price of not putting request IDs on the wire.
Say this in the doc comment, because the compiler cannot. A stream protocol with no request IDs matches replies to requests by order, so two goroutines sharing a client would interleave their writes and each read the other's answer. The alternatives are a mutex around the whole round trip (correct, and it serialises every ant on one connection) or request IDs plus a demultiplexing reader goroutine (what gRPC does, and what Exercise 12 asks for). One connection per ant is the simplest thing that works, and it is what we do.
func (c *Client) roundTrip(req Req) (Resp, error) {
if err := c.connect(); err != nil {
return Resp{}, err
}
// A deadline on the socket is the network equivalent of the timeout we
// put on the local channel in milestone 5.
c.conn.SetDeadline(time.Now().Add(2 * time.Second))
if err := c.enc.Encode(req); err != nil {
c.drop()
return Resp{}, fmt.Errorf("%w: send: %v", ErrDisconnected, err)
}
var resp Resp
if err := c.dec.Decode(&resp); err != nil {
c.drop()
return Resp{}, fmt.Errorf("%w: receive: %v", ErrDisconnected, err)
}
if resp.Err != "" {
return resp, errors.New(resp.Err)
}
return resp, nil
}
Any I/O error drops the connection rather than trying to continue on it, because a gob stream that has lost sync cannot be repaired: the decoder's type state and the byte stream no longer agree. Reconnecting is cheap; resynchronising is impossible. Wrapping with %w keeps errors.Is(err, ErrDisconnected) working, which is how the ant distinguishes "the world is gone, wait and retry" from "the world said no".
// Backoff reports how long to wait before the next reconnection attempt,
// with full jitter so a thousand ants do not reconnect in lockstep.
func (c *Client) Backoff() time.Duration {
d := 10 * time.Millisecond << min(c.attempts, 6)
return time.Duration(c.rng.Int64N(int64(d) + 1))
}
Same reasoning as the supervisor backoff in Milestone 8, and now it matters more: when a server restarts, every client discovers it at the same instant, and unjittered retries would arrive as a synchronised wave that knocks it over again.
for ctx.Err() == nil {
sensed, err := c.Sense(a.Pos)
if err != nil {
if errors.Is(err, ErrDisconnected) {
select {
case <-time.After(c.Backoff()):
case <-ctx.Done():
}
continue
}
return
}
act := b.Decide(a, sensed, rng)
res, err := c.Act(a.ID, a.Pos, a.Carrying, act)
...
}
Sense, decide, act. Identical in structure to Milestone 5, with ErrTimeout handling replaced by ErrDisconnected handling. The behaviour code (TrailFollower) is byte-for-byte the same file that ran locally, because it only ever saw a Sense.
$ go test -race ./internal/netsim/
=== RUN TestColonyOverTheWire
net_test.go:94: delivered over TCP: 39, food left 458
--- PASS: TestColonyOverTheWire (0.71s)
=== RUN TestClientSurvivesServerRestart
--- PASS: TestClientSurvivesServerRestart (0.00s)
PASS
Forty remote ants foraging through forty TCP connections, delivering food to a world in a different object graph, race-clean. The second test kills the server mid-session, confirms the client reports ErrDisconnected, restarts the world at the same address, and confirms the client reconnects on its next attempt with no special handling.
sense latency: local 29ns, over loopback TCP 10.667µs (368x)
That is the number to carry out of this course. The same operation is 368 times slower over a loopback socket than in memory, and loopback is the best case: no switch, no cable, no congestion, no packet loss. Across a data centre add 200 µs or more; across a continent, 50 ms and up, which is five million times the local cost.
This is why "just make it a microservice" is a performance decision and not only an architectural one. Our sharded local engine does about 30 million senses per second per core; over loopback, about 100,000. Distribution buys you fault isolation, independent deployment and horizontal scale, and it costs you three orders of magnitude on every interaction you move across the boundary. The design question is always the same: which interactions cross? Here, the honest answer is that sensing should never cross a network at all, and a real distributed colony would replicate a read-only copy of the local grid to each ant process and only send mutations over the wire. That is a cache with a coherence protocol, and now you are writing a distributed system in earnest.
Go's case is strongest in this milestone and in Milestone 9. The standard library alone gave us a TCP server, binary serialisation, an HTTP server, a metrics endpoint, a profiler you can attach to a live process, and a test framework, with zero third-party dependencies in the entire project. The deployable artifact is one static binary. The concurrency model made "one goroutine per connection" the obvious and correct design rather than a scalability problem.
Honest counterweights. gob is Go-only, so this protocol cannot be consumed by anything else without rewriting the serialisation. Our hand-rolled protocol has no request IDs, no multiplexing, no flow control, no TLS and no authentication, all of which gRPC gives you for a schema file and a code generator, and a production system should use it rather than this. And the 368× penalty is not a Go number, it is a physics number: no language makes a socket as fast as a memory read.
The thing Go deserves credit for is that the distance between "single-process simulation" and "distributed system" turned out to be one interface and about 250 lines, because the design forced by cheap goroutines and channels in Milestone 5 was already the right shape.
Multiplex the connection. One TCP connection per ant does not scale to 50,000 ants: you run out of file descriptors and the server drowns in connection goroutines. Rewrite the client so that many ants share one connection.
Requirements. Add a request ID to Req and Resp. One writer goroutine and one reader goroutine per connection; the reader routes each response to the waiting caller. Callers still see the same blocking Sense/Act API. A caller whose context is cancelled must not leak its slot. The server may answer out of order.
Constraints. No goroutine per in-flight request beyond the caller's own. Bounded memory even if the server stops answering. -race clean with 1,000 ants on one connection.
Hints. What data structure maps IDs to waiting callers, and what protects it? What is the reply channel's buffer size, and why does that answer matter more here than it did in Milestone 5? What must happen to the map entry if the caller gives up?
type MuxClient struct {
conn net.Conn
enc *gob.Encoder
mu sync.Mutex // guards nextID and pending
nextID uint64
pending map[uint64]chan Resp
writeMu sync.Mutex // one writer at a time on the socket
}
func (c *MuxClient) roundTrip(ctx context.Context, req Req) (Resp, error) {
c.mu.Lock()
c.nextID++
id := c.nextID
ch := make(chan Resp, 1) // buffered: the reader must never block
c.pending[id] = ch
c.mu.Unlock()
defer func() { // always reclaim the slot
c.mu.Lock()
delete(c.pending, id)
c.mu.Unlock()
}()
req.ID = id
c.writeMu.Lock()
err := c.enc.Encode(req)
c.writeMu.Unlock()
if err != nil {
return Resp{}, fmt.Errorf("%w: send: %v", ErrDisconnected, err)
}
select {
case resp := <-ch:
return resp, nil
case <-ctx.Done():
return Resp{}, ctx.Err()
}
}
// readLoop is the only goroutine that reads the socket.
func (c *MuxClient) readLoop() {
for {
var resp Resp
if err := c.dec.Decode(&resp); err != nil {
c.failAll(err)
return
}
c.mu.Lock()
ch, ok := c.pending[resp.ID]
c.mu.Unlock()
if ok {
ch <- resp // never blocks: buffer 1, one send only
}
// unknown ID: the caller gave up. Dropping the response is correct.
}
}
Answers to the hints, which are the actual content of this exercise.
defer that deletes the entry is what keeps memory bounded: without it, every cancelled request leaks a map entry and a channel forever, and a server that stops responding would take the client down with it.writeMu serialises socket writes (a gob encoder is not safe for concurrent use, and interleaved frames would corrupt the stream); mu guards the map. Using one mutex for both would make every request wait for the socket write of every other. When lock contention appears, splitting locks by what they protect is the first thing to try.failAll on read error must close or signal every pending channel, or a thousand ants block forever on a dead connection.You have now implemented, in about eighty lines, the core of what HTTP/2, gRPC and every other multiplexed RPC protocol does. That is the right way to understand those systems: not as magic, but as this, plus flow control, plus TLS, plus a schema.
Reuse one plain Client (not the multiplexing MuxClient from the exercise) across two ant goroutines with no synchronisation at all, and run go test -race ./internal/netsim/. It fails immediately and reliably, pointing at concurrent writes through the same gob.Encoder — this is a case where the race detector is exactly the right tool, unlike Milestone 11's shard bug. Fix it the crude way, by wrapping the whole roundTrip call in a mutex, and the race disappears; but now watch throughput with two ants sharing one connection instead of two. Serialising away a race is easy. Getting the throughput back is the rest of Exercise 12.
gob.Encoder across goroutines without a mutex. Interleaved frames, corrupt stream, baffling decode errors.error in a gob struct. It will not encode. Send a string or a code.WorldService.Sense return an error even though the local implementation cannot fail?Close hang, and why is a mutex the right fix rather than a channel?errors.Is(err, world.ErrNoFood) across the network, and what do real protocols do about it?Decode. What can?antfarm/ 2,961 lines of Go, no dependencies
├── go.mod
├── cmd/antfarm/main.go flags, chaos parsing, signals, servers, wiring
└── internal/
├── sim/
│ ├── ant.go Ant, Action, ActionKind, Generation
│ ├── behaviour.go Behaviour, Forager, TrailFollower
│ ├── sense.go Sense, senseAt
│ ├── engine.go single-owner engine, supervisor, Frame, Run
│ ├── fast.go sharded engine, WorldService, atomic sensing
│ ├── sim.go Config, Chaos, Stats, the sequential Sim
│ ├── concurrent.go milestone 4's mutex versions, kept for contrast
│ └── *_test.go conservation, leaks, chaos, throughput, races
├── world/
│ ├── grid.go Position, Grid
│ ├── pheromone.go FloatGrid
│ ├── atomicgrid.go AtomicGrid, fixed-point pheromone
│ └── world.go World, food, nest, errors
├── metrics/
│ ├── metrics.go Counter, Gauge, Histogram, Prometheus output
│ └── http.go /metrics, /healthz, /debug/pprof
├── view/
│ ├── terminal.go ANSI renderer
│ ├── web.go SSE stream
│ └── page.go the embedded browser client
└── netsim/
├── proto.go Req, Resp
├── server.go TCP + gob server, connection registry
├── client.go reconnecting client, jittered backoff
└── net_test.go colony over TCP, survives a server restart
$ gofmt -l . && go vet ./... && go test -race ./...
ok github.com/yourname/antfarm/internal/netsim 0.705s
ok github.com/yourname/antfarm/internal/sim 12.426s
ok github.com/yourname/antfarm/internal/world 0.001s
$ git commit -am "milestone 12: the world moves to another process"
Continue