programming · Level 3
Go: build a concurrent service
Use goroutines, channels and cancellation to build a small service you can explain and test.

Start with the essentials
The short answer
A concurrent service coordinates work that can overlap. This lab accepts a small batch of integers, squares them through two workers and returns results in input order. Two requests can be admitted at once; excess requests receive 503. A deadline or client cancellation asks workers to stop, and the handler waits for them before returning. The example is local teaching code, with no external provider or public deployment.
What you will learn
- Read Go declarations, slices, methods and explicit errors.
- Bound jobs and workers while preserving input order.
- Propagate cancellation and wait for workers to finish.
- Validate a small HTTP JSON request before doing work.
- Separate request admission from worker concurrency.
- Test local HTTP behaviour and report unverified operational limits.
Who it is for
Builders who can read a small program and want a first practical Go service, including people moving from Python.
Before you start
- Read functions, loops and a basic test assertion.
- Understand an HTTP request and response at a high level.
- Have a local Go toolchain available for the lab.
Read a sample · Chapter 01 of 06
Read the contract before the goroutines
Give a small service a precise job and a deliberately small input.
Create a new folder containing go.mod, batch.go, server.go, main.go and server_test.go. Every file is printed in full across this workbook. Where a file has several code blocks, concatenate them in order with one blank line between blocks. Keep package main only at the start of each Go file. Use a dedicated folder, not an existing application package. The module name is a local teaching identifier; it does not cause a connection to that domain.
module example.com/mickai/concurrent-service
go 1.26.0The module declares Go 1.26.0 and the release ran with the installed Go 1.26.2 toolchain. This is a reproducibility record, not a recommendation that this patch is current for a public service. Review the supported toolchain before real deployment. No external module is imported. The validation run set GOTOOLCHAIN=local and GOPROXY=off, so it neither downloaded a toolchain nor fetched a dependency. Your local environment may use different settings.
POST /batch accepts application/json containing numbers, a list of one to eight integers. Values range from zero to one million inclusive. The response contains results in the same order. For example, input [3,1,2] produces [9,1,4]. A value of one million squares to one trillion, which fits int64. The low-level worker helper accepts an empty list, but the HTTP contract rejects it; these are different interfaces.
The fictional operation waits about 100 milliseconds before squaring each value. That delay models waiting for work, not the cost of multiplication, a measured external service or an AI provider. With two workers, eight jobs need four waiting rounds under favourable scheduling. The 500-millisecond budget includes their execution and scheduling after body parsing. No exact duration or throughput is promised. Operating-system delays can still make a request miss its budget.
In Go, := declares and initialises a local variable, while var can declare its type separately. A slice such as []int64 is a view over ordered elements. A function returning (int64, error) makes success and failure explicit: callers examine the error before using the value. Nil is the absence value for an error or reference-like type. The blank identifier _ means that a parameter or result is deliberately unused; it does not hide a failure safely by itself.
The service is intentionally small enough to read from input to output. A worker operation receives a context and one integer. It must honour cancellation and return an error when unfinished. It must not panic or launch untracked background work. Those requirements are part of our own operation contract. Replacing the operation with an arbitrary library call requires checking whether that call can meet them.
Try it yourself · Activity 01
15 minWrite the request contract
Predict the result of three fictional batches before reading the implementation.
- Calculate [3,1,2] and [0,1000000].
- Decide whether an empty HTTP list and the value 1000001 should be accepted.
- State which timing claims the artificial delay cannot establish.
Can a new reader understand your service input, output and failure conditions without reading a worker implementation?
Worked answer
The results are [9,1,4] and [0,1000000000000]. The empty HTTP list and 1000001 both receive 400. Two workers can overlap waiting, but this example does not measure a provider, CPU parallel speed-up or production capacity. The 500-millisecond work budget may expire if the machine is delayed.
Read a sample · Chapter 02 of 06
Bound work and keep its original order
Use channels to distribute ownership instead of sharing an append operation.
Save batch.go from the following blocks. The function type operation makes the worker dependency explicit and allows tests to substitute a fast or blocked operation. RunBatch remains unexported because its name starts with a lowercase letter. The helper validates worker count and job count before starting goroutines. It does not validate numeric range; that is the HTTP parser responsibility in this lab.
package main
import (
"context"
"errors"
"sync"
"time"
)
type operation func(context.Context, int64) (int64, error)func runBatch(parent context.Context, values []int64,
workers int, work operation) ([]int64, error) {
if workers < 1 || workers > 8 || len(values) > 8 || work == nil {
return nil, errors.New("invalid batch configuration")
}
ctx, cancel := context.WithCancel(parent)
defer cancel()
jobs := make(chan int, len(values))
for i := range values {
jobs <- i
}
close(jobs)
results := make([]int64, len(values))
var wg sync.WaitGroup
var once sync.Once
var firstErr error
for w := 0; w < workers; w++ {
wg.Add(1)
go func() {
defer wg.Done()
for i := range jobs {
if ctx.Err() != nil {
return
}
value, err := work(ctx, values[i])
if err != nil {
once.Do(func() { firstErr = err; cancel() })
return
}
results[i] = value
}
}()
}
wg.Wait()
if firstErr != nil {
return nil, firstErr
}
if err := ctx.Err(); err != nil {
return nil, err
}
return results, nil
}func squareAfterWait(ctx context.Context, n int64) (int64, error) {
timer := time.NewTimer(100 * time.Millisecond)
defer timer.Stop()
select {
case <-ctx.Done():
return 0, ctx.Err()
case <-timer.C:
return n * n, nil
}
}The buffered jobs channel contains indexes rather than values. Its capacity is at most eight, because validation already bounded the input. The creating goroutine enqueues all indexes and closes the channel before starting workers. That works only because this queue has room for every job. Changing it to an unbuffered channel while keeping this sequence would block the sender before any receiver exists. A streaming producer would require a different lifecycle.
Each index is received by one worker. That worker writes only results[i], and the result slice was allocated once before they started. No worker appends, resizes the slice or writes another index. The caller reads it after the WaitGroup joins all workers. These ownership rules are why the design does not need a mutex around the output slice. A shared map, counter or append would have a different synchronisation requirement.
Input order and completion order are different. If job 1 finishes before job 0, its result still goes into slot 1. The response waits for the entire batch, so the client receives ordered results only on success. When any operation fails, the partial result slice is discarded and the error is returned. The first error recorded by sync.Once wins; if several errors occur together, its identity can depend on scheduling.
The producer owns closing jobs. Workers only receive from it. They never close it, send another job or depend on a result channel being drained. Closing a channel is a communication event, not a request that goroutines stop immediately. Here a worker checks the context before beginning each new job and the operation also waits on context cancellation. A task already in progress must cooperate.
The WaitGroup counter is increased before each goroutine starts. Done runs through defer when that worker returns. The caller then waits before examining the first error or output. Calling cancel alone would not establish that all operations have returned. Keeping the join explicit also prevents the handler from releasing its request slot while a cooperative worker still owns work.
Try it yourself · Activity 02
20 minTrace three indexes
Draw two worker lanes for input [3,1,2].
- Label each queue item with its input index.
- Let the second item finish first and place each result in its original slot.
- Explain what breaks if results[i] becomes an unsynchronised append to one shared slice.
Who owns closing each channel, and can every send finish when a request is cancelled?
Worked answer
The queue contains indexes 0,1,2. Results belong in slots 0=9, 1=1 and 2=4, regardless of completion order. A shared append changes the slice header and can race; even a correctly locked append would normally record completion order rather than input order. Preallocated, disjoint slots plus a join meet this batch contract.
Read a sample · Chapter 03 of 06
Validate input and admit a bounded number of requests
Worker limits and request limits control different amounts of work.
Start server.go with these three blocks. Service holds the shared request-slot channel, worker count, work budget and operation. NewService creates one instance with two request slots and two workers per admitted batch. All HTTP calls on that same instance share its slot channel. Creating a new service for each incoming request would destroy this admission boundary.
package main
import (
"context"
"encoding/json"
"errors"
"io"
"mime"
"net/http"
"strconv"
"strings"
"time"
)
type service struct {
slots chan struct{}
workers int
budget time.Duration
work operation
}func newService(work operation) *service {
return &service{make(chan struct{}, 2), 2,
500 * time.Millisecond, work}
}func numbers(w http.ResponseWriter, r *http.Request) ([]int64, int) {
media, _, err := mime.ParseMediaType(r.Header.Get("Content-Type"))
if err != nil || media != "application/json" {
return nil, http.StatusUnsupportedMediaType
}
r.Body = http.MaxBytesReader(w, r.Body, 1024)
defer r.Body.Close()
var input struct {
Numbers []json.RawMessage `json:"numbers"`
}
decoder := json.NewDecoder(r.Body)
decoder.DisallowUnknownFields()
err = decoder.Decode(&input)
if err == nil {
var extra any
err = decoder.Decode(&extra)
if err == io.EOF {
err = nil
} else if err == nil {
err = errors.New("extra JSON value")
}
}
if err != nil {
var tooLarge *http.MaxBytesError
if errors.As(err, &tooLarge) {
return nil, http.StatusRequestEntityTooLarge
}
return nil, http.StatusBadRequest
}
if len(input.Numbers) < 1 || len(input.Numbers) > 8 {
return nil, http.StatusBadRequest
}
values := make([]int64, len(input.Numbers))
for i, raw := range input.Numbers {
n, err := strconv.ParseInt(strings.TrimSpace(string(raw)), 10, 64)
if err != nil || n < 0 || n > 1000000 {
return nil, http.StatusBadRequest
}
values[i] = n
}
return values, 0
}The parser accepts application/json, including a valid charset parameter. It limits the request body to 1,024 bytes, rejects unknown struct fields and insists that a second decoder call reaches the end of the body. That second read rejects another JSON value and also counts trailing whitespace against the byte cap. A malformed short document receives 400. A read that actually encounters the byte limit receives 413; not every invalid oversized document necessarily reaches that read before another error is found.
RawMessage preserves each number token until ParseInt examines it. This avoids an intermediate float conversion. Fractions, exponent notation, quoted numbers, booleans and null are rejected even where another parser might coerce them. Negative values and values above one million are also rejected. JSON negative zero is accepted as zero. The limits are teaching choices for one tiny endpoint, not a complete validation policy for an unrelated service.
This version uses encoding/json with its established v1 behaviour. DisallowUnknownFields does not reject duplicate keys or enforce case-sensitive field names. For example, differently cased numbers can match the field, and a later duplicate list replaces an earlier list. Do not label this parser a canonical or hostile-input JSON validator. A signed payload, gateway or second parser can interpret ambiguous input differently and would need a deliberately stricter contract.
Admission happens before body decoding in the handler shown next. At most two batch requests on one service instance can hold slots, including requests still reading their body. Each then runs at most two operations, giving at most four operation calls in flight across those admitted batches. This does not cap TCP connections, HTTP handler goroutines, other endpoints, other service instances or the entire process memory. The standard HTTP server still accepts and dispatches requests.
The server timeouts in main.go place finite read and write deadlines around connections. Our 500-millisecond context is created after successful parsing, so it is not the complete request deadline. If you move parsing into a different server or call ServeHTTP directly, those server-level timeouts may be absent. Keep body size, body read duration, admission and operation duration as separate controls when reviewing the boundary.
Try it yourself · Activity 03
20 minSeparate four limits
Review one service instance receiving three batches at the same time.
- Calculate the maximum admitted batches and operation calls.
- Explain when a slow request body occupies a slot.
- Classify [1.5], [null], an unknown field and a second JSON object.
Does a number called max concurrency in your own code constrain tasks, admitted requests, connections or something else?
Worked answer
Two batches can hold slots and at most four operations can run across them. A slow body occupies one of those slots before it is fully parsed; server read timeouts matter here. The third request receives 503 if both slots are occupied. Each listed invalid body receives 400 when the relevant malformed content is reached within the body cap. Duplicate keys and case-insensitive names remain accepted v1 behaviours.
Read a sample · Chapter 04 of 06
Propagate failure and join cancellation
Return one complete result or an explicit failure.
Follow the work, keep the order
For input [3,1,2], the fixture holds value 3 until the other two
operations have recorded their results. Recorded ranks describe operation log events just before return, not exact function-return
instants or the behaviour of every run. Each job still writes its original
input index. After the worker join, the result is [9,1,4].
| Recorded rank | Input index | Input value | Square |
|---|---|---|---|
| 1 | 1 | 1 | 1 |
| 2 | 2 | 2 | 4 |
| 3 | 0 | 3 | 9 |
A second check sends two two-item batches, A and B, to one service instance.
All four operations wait on a test barrier. C reaches the same handler while
both request slots are occupied and receives 503 with
Retry-After: 1. Releasing the barrier lets A and B return 200,
leaving zero active operations and zero occupied slots.
| Observation | Occupied slots | Active operations | Response status |
|---|---|---|---|
| A + B held | 2 | 4 | Pending |
| C while full | 2 | 4 | 503 |
| A + B released | 0 | 0 | 200 / 200 |
Four is the maximum number of operation calls across the admitted batches on this instance. It is not a limit on TCP connections, all HTTP handler goroutines, other endpoints or a whole deployment. The standard server can still receive C; the application refuses its batch admission.
| Observation | Occupied slots | Active operations | Handler status |
|---|---|---|---|
| Before cancel | 1 | 2 | Pending |
| After join + return | 0 | 0 | 408 |
The third check cancels a request context after both operations enter. They return, the worker join completes and the handler releases its slot before the final observation. An in-memory response recorder observes 408, the lab's chosen policy; a disconnected network client may receive no response. Cancellation is cooperative. A worker that ignores its context can block the join and retain the slot.
All three checks use a five-second test budget to isolate the control flow from the lab's normal 500-millisecond work budget. Test operations and barriers replace the fictional timer wait; the published implementation is unchanged. These are local behaviour checks, not performance measurements. Race-detector tests could not run with cgo disabled, and OS-signal shutdown remains untested. Neither result is implied by this figure.
Append the final ServeHTTP block to server.go. A method on *service receives a pointer to its shared configuration. That method satisfies http.Handler through its name and signature. GET /health returns a small status response without taking a batch slot. Unknown paths return 404; /batch with the wrong method returns 405 and Allow: POST. Health is a local process response, not a downstream readiness check.
func (s *service) ServeHTTP(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/health" && r.Method == http.MethodGet {
w.Header().Set("Content-Type", "application/json")
io.WriteString(w, "{\"status\":\"ok\"}\n")
return
}
if r.URL.Path != "/batch" {
http.NotFound(w, r)
return
}
if r.Method != http.MethodPost {
w.Header().Set("Allow", http.MethodPost)
http.Error(w, "POST required", http.StatusMethodNotAllowed)
return
}
select {
case s.slots <- struct{}{}:
defer func() { <-s.slots }()
default:
w.Header().Set("Retry-After", "1")
http.Error(w, "busy", http.StatusServiceUnavailable)
return
}
values, status := numbers(w, r)
if status != 0 {
http.Error(w, "invalid request", status)
return
}
ctx, cancel := context.WithTimeout(r.Context(), s.budget)
defer cancel()
results, err := runBatch(ctx, values, s.workers, s.work)
if err != nil {
status = http.StatusInternalServerError
if errors.Is(err, context.DeadlineExceeded) ||
errors.Is(err, context.Canceled) {
status = http.StatusRequestTimeout
}
http.Error(w, "batch not completed", status)
return
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(struct {
Results []int64 `json:"results"`
}{results})
}The non-blocking select tries to put a token into slots. If neither slot is free, default returns 503 with Retry-After: 1. There is no application admission wait queue. A caller should not treat that header as an instruction to retry forever or to retry every failed operation. Our batch has no external effect, but a different task could need an idempotency key before a retry is safe.
After admission, defer removes the token on every normal return path: invalid input, failed work and successful response. Context.WithTimeout inherits request cancellation and adds the work budget. Calling its cancel function releases associated resources. It does not interrupt arbitrary Go code. The operation must observe ctx.Done or use an API that honours the context. RunBatch then joins the workers before the handler returns.
The lab maps an expired or cancelled batch to 408 and another operation error to 500. This is an explicit teaching policy, not a universal recommendation for every HTTP API. A client that already cancelled its connection may never receive that response. The response text deliberately does not expose the internal error string. A production observability design would need useful internal diagnostics, controlled disclosure and request correlation.
The operation uses a timer and select to make its wait cancellable. When the timer and cancellation are both ready, select can choose either branch. RunBatch checks context again after joining, so a cancelled batch is not returned as a normal success merely because a task finished near the boundary. Cancellation that arrives after that final check can still race with response writing; no local function makes a network write atomic with a remote client decision.
The small JSON encoder write has no recovery path here. If the client disconnects while writing the response, trying to write a second error body would not repair a partially transmitted response. More generally, this implementation does not recover panics, interrupt a non-cooperative operation or account for goroutines created secretly inside an operation. Those are reasons to keep the injected work contract explicit and to review new dependencies before substituting them.
Try it yourself · Activity 04
15 minFollow a cancelled batch
The first two jobs are waiting when the client cancels.
- Trace request context, child context and each operation.
- Explain why the request slot is released after the join.
- Describe what happens if an operation ignores its context forever.
Can every dependency in a worker stop under the same cancellation contract, or are some calls only timed out from the caller perspective?
Worked answer
Request cancellation reaches the child context. Cooperative operations return an error, each worker calls Done, and Wait completes. RunBatch returns the failure; the handler releases its request token through defer. A permanently blocked non-cooperative operation prevents the join and retains the slot. CancelFunc cannot kill that goroutine, so the current design cannot impose a hard deadline on arbitrary work.
Read a sample · Chapter 05 of 06
Start locally and test lifecycle rules
Build one executable, then test the worker rules without a fixed public port.
Save main.go. The executable listens on 127.0.0.1:8088, so its intended use is the local machine. ReadHeaderTimeout and ReadTimeout are two seconds, WriteTimeout is two seconds and IdleTimeout is fifteen seconds. MaxHeaderBytes constrains request-header parsing. These settings are not capacity measurements, TLS, authentication or protection against every form of resource exhaustion.
package main
import (
"context"
"errors"
"log"
"net/http"
"os"
"os/signal"
"time"
)func main() {
stop, cancel := signal.NotifyContext(context.Background(), os.Interrupt)
defer cancel()
server := &http.Server{
Addr: "127.0.0.1:8088", Handler: newService(squareAfterWait),
ReadHeaderTimeout: 2 * time.Second,
ReadTimeout: 2 * time.Second,
WriteTimeout: 2 * time.Second,
IdleTimeout: 15 * time.Second,
MaxHeaderBytes: 4096,
}
done := make(chan error, 1)
go func() { done <- server.ListenAndServe() }()
select {
case err := <-done:
if !errors.Is(err, http.ErrServerClosed) {
log.Print(err)
}
case <-stop.Done():
ctx, release := context.WithTimeout(context.Background(), time.Second)
defer release()
if err := server.Shutdown(ctx); err != nil {
log.Print(err)
server.Close()
}
<-done
}
}Run go fmt ./..., go test ./..., go vet ./... and go build . from the new folder after all files are saved. Run go run . only when ready to occupy the local port. Ctrl+C requests shutdown through signal.NotifyContext. Shutdown has a one-second budget; if it fails, Close is attempted. The release compiled this path but did not exercise the operating-system signal or graceful-drain sequence. The network tests instead use httptest.NewServer and its automatically chosen local port.
Start server_test.go with the following blocks, through TestSquareWait. All remaining blocks appear in the next chapter. The request helper invokes a handler with an in-memory ResponseRecorder, while await uses a one-second deadline to catch a missing test event. Neither is a production request implementation. Contexts, channels and atomic counters make the worker tests observable without relying on the fictional 100-millisecond operation.
package main
import (
"context"
"errors"
"fmt"
"io"
"net/http"
"net/http/httptest"
"reflect"
"strings"
"sync/atomic"
"testing"
"time"
)
const jsonMedia = "application/json"func fastSquare(_ context.Context, n int64) (int64, error) {
return n * n, nil
}func request(s http.Handler, method, path, body, media string,
) *httptest.ResponseRecorder {
r := httptest.NewRequest(method, path, strings.NewReader(body))
r.Header.Set("Content-Type", media)
w := httptest.NewRecorder()
s.ServeHTTP(w, r)
return w
}func await(t *testing.T, ch <-chan struct{}) {
t.Helper()
select {
case <-ch:
case <-time.After(time.Second):
t.Fatal("worker did not reach barrier")
}
}func TestOrder(t *testing.T) {
second := make(chan struct{})
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
work := func(ctx context.Context, n int64) (int64, error) {
if n == 1 {
select {
case <-second:
case <-ctx.Done():
return 0, ctx.Err()
}
} else if n == 2 {
close(second)
}
return n * n, nil
}
got, err := runBatch(ctx, []int64{1, 2, 3}, 2, work)
if err != nil || !reflect.DeepEqual(got, []int64{1, 4, 9}) {
t.Fatalf("ordered results: %v, %v", got, err)
}
}func TestWorkerBound(t *testing.T) {
var active, peak atomic.Int32
started := make(chan struct{}, 8)
gate := make(chan struct{})
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
work := func(ctx context.Context, n int64) (int64, error) {
now := active.Add(1)
defer active.Add(-1)
for old := peak.Load(); now > old; old = peak.Load() {
if peak.CompareAndSwap(old, now) {
break
}
}
started <- struct{}{}
select {
case <-gate:
return n * n, nil
case <-ctx.Done():
return 0, ctx.Err()
}
}
done := make(chan error, 1)
go func() {
_, err := runBatch(ctx, []int64{1, 2, 3, 4, 5, 6, 7, 8}, 2, work)
done <- err
}()
await(t, started)
await(t, started)
// Keep both workers occupied during this bounded observation.
select {
case <-started:
t.Error("more than two workers entered the operation")
case <-time.After(20 * time.Millisecond):
}
close(gate)
if err := <-done; err != nil || peak.Load() != 2 || active.Load() != 0 {
t.Fatalf("err=%v peak=%d active=%d", err, peak.Load(), active.Load())
}
}func TestCancellationJoinsWorkers(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
started := make(chan struct{}, 2)
var active atomic.Int32
work := func(ctx context.Context, _ int64) (int64, error) {
active.Add(1)
defer active.Add(-1)
started <- struct{}{}
<-ctx.Done()
return 0, ctx.Err()
}
done := make(chan error, 1)
go func() {
_, err := runBatch(ctx, []int64{1, 2, 3}, 2, work)
done <- err
}()
await(t, started)
await(t, started)
cancel()
select {
case err := <-done:
if !errors.Is(err, context.Canceled) || active.Load() != 0 {
t.Fatalf("err=%v active=%d", err, active.Load())
}
case <-time.After(time.Second):
t.Fatal("cancellation did not join workers")
}
}func TestFirstErrorCancelsPeer(t *testing.T) {
want := errors.New("fictional operation failed")
peer := make(chan struct{})
peerDone := make(chan struct{})
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
work := func(ctx context.Context, n int64) (int64, error) {
if n == 1 {
select {
case <-peer:
return 0, want
case <-ctx.Done():
return 0, ctx.Err()
}
}
close(peer)
<-ctx.Done()
close(peerDone)
return 0, ctx.Err()
}
got, err := runBatch(ctx, []int64{1, 2}, 2, work)
if got != nil || !errors.Is(err, want) {
t.Fatalf("got=%v err=%v", got, err)
}
select {
case <-peerDone:
default:
t.Fatal("returned before the peer stopped")
}
}func TestBatchConfiguration(t *testing.T) {
ctx := context.Background()
for _, workers := range []int{0, 9} {
_, err := runBatch(ctx, []int64{1}, workers, fastSquare)
if err == nil {
t.Fatalf("accepted %d workers", workers)
}
}
if _, err := runBatch(ctx, make([]int64, 9), 2, fastSquare); err == nil {
t.Fatal("accepted nine jobs")
}
if _, err := runBatch(ctx, []int64{1}, 2, nil); err == nil {
t.Fatal("accepted nil operation")
}
got, err := runBatch(ctx, nil, 2, fastSquare)
if len(got) != 0 || err != nil {
t.Fatalf("empty internal batch: %v %v", got, err)
}
}func TestSquareWait(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
cancel()
if _, err := squareAfterWait(ctx, 3); !errors.Is(err, context.Canceled) {
t.Fatalf("cancelled wait: %v", err)
}
ctx, cancel = context.WithTimeout(context.Background(), time.Second)
defer cancel()
if got, err := squareAfterWait(ctx, 3); got != 9 || err != nil {
t.Fatalf("square=%d err=%v", got, err)
}
}TestOrder blocks the first item until a later item reaches a channel signal. This forces a useful ordering situation without demanding a scheduler-specific completion duration. TestWorkerBound holds two operations at a barrier and observes a short interval in which a third must not enter; it also checks peak activity and zero activity after return. That interval is a bounded observation, not a proof about every possible execution. A global go test timeout catches a larger deadlock.
TestCancellationJoinsWorkers cancels only after both workers report that they started. It expects context.Canceled and zero active operations when RunBatch returns. TestFirstErrorCancelsPeer makes one worker fail after its peer starts, then confirms the peer exits before return. TestBatchConfiguration checks invalid worker counts, excess jobs, a nil operation and the valid empty internal batch. TestSquareWait also executes the real timer operation on success and with an already-cancelled context. Each assertion targets a promise made by this lab.
Try it yourself · Activity 05
15 minExplain the shutdown boundary
Compare stopping a worker, stopping admission and stopping the HTTP server.
- Identify which function joins worker goroutines.
- Identify which object limits batch admission.
- State exactly which shutdown path this release has and has not exercised.
Which observation would demonstrate that an in-flight request is allowed to finish when your real process receives its shutdown signal?
Worked answer
RunBatch joins its workers through WaitGroup. The service instance holds the request-slot channel. The executable contains a Ctrl+C-triggered Shutdown path with a one-second budget and a Close fallback. Compilation and local server/client tests pass, but the operating-system signal and graceful-drain sequence were not exercised. Releasing a request slot and shutting down the whole server are separate events.
Read a sample · Chapter 06 of 06
Test the HTTP boundary and read the evidence
Use failures to check what the happy path leaves hidden.
Append these remaining blocks to server_test.go. TestHTTPContract covers twenty request cases. TestAdmissionLimits keeps two requests occupied while checking the third receives 503 and the health route remains responsive. TestSlotsReleased repeats both successful and invalid requests to catch leaked admission tokens. Deadline and operation-error tests check both response behaviour and released resources.
func TestHTTPContract(t *testing.T) {
cases := []struct {
name, method, path, body, media string
status int
}{
{"good", "POST", "/batch", `{"numbers":[0,2,1000000]}`, jsonMedia, 200},
{"charset", "POST", "/batch", `{"numbers":[1]}`,
"application/json; charset=utf-8", 200},
{"method", "GET", "/batch", "", "", 405},
{"path", "POST", "/missing", "", "", 404},
{"health", "GET", "/health", "", "", 200},
{"media", "POST", "/batch", `{"numbers":[1]}`, "text/plain", 415},
{"empty", "POST", "/batch", `{"numbers":[]}`, jsonMedia, 400},
{"missing", "POST", "/batch", `{}`, jsonMedia, 400},
{"null", "POST", "/batch", `{"numbers":[null]}`, jsonMedia, 400},
{"fraction", "POST", "/batch", `{"numbers":[1.5]}`, jsonMedia, 400},
{"exponent", "POST", "/batch", `{"numbers":[1e2]}`, jsonMedia, 400},
{"string", "POST", "/batch", `{"numbers":["1"]}`, jsonMedia, 400},
{"boolean", "POST", "/batch", `{"numbers":[true]}`, jsonMedia, 400},
{"negative", "POST", "/batch", `{"numbers":[-1]}`, jsonMedia, 400},
{"range", "POST", "/batch", `{"numbers":[1000001]}`, jsonMedia, 400},
{"count", "POST", "/batch",
`{"numbers":[1,2,3,4,5,6,7,8,9]}`, jsonMedia, 400},
{"extra", "POST", "/batch",
`{"numbers":[1],"other":0}`, jsonMedia, 400},
{"second", "POST", "/batch", `{"numbers":[1]} {}`, jsonMedia, 400},
{"malformed", "POST", "/batch", `{"numbers":[`, jsonMedia, 400},
{"large", "POST", "/batch",
`{"numbers":[1]}` + strings.Repeat(" ", 1024), jsonMedia, 413},
}
s := newService(fastSquare)
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
w := request(s, tc.method, tc.path, tc.body, tc.media)
if w.Code != tc.status {
t.Fatalf("status=%d body=%s", w.Code, w.Body)
}
expected := "{\"results\":[0,4,1000000000000]}\n"
if tc.name == "good" && w.Body.String() != expected {
t.Fatalf("results=%s", w.Body)
}
if tc.name == "method" && w.Header().Get("Allow") != "POST" {
t.Fatal("missing Allow header")
}
})
}
}func TestAdmissionLimits(t *testing.T) {
started := make(chan struct{}, 2)
gate := make(chan struct{})
s := newService(func(ctx context.Context, n int64) (int64, error) {
started <- struct{}{}
select {
case <-gate:
return n * n, nil
case <-ctx.Done():
return 0, ctx.Err()
}
})
done := make(chan int, 2)
for i := 0; i < 2; i++ {
go func() {
w := request(s, "POST", "/batch", `{"numbers":[2]}`, jsonMedia)
done <- w.Code
}()
}
await(t, started)
await(t, started)
w := request(s, "POST", "/batch", `{"numbers":[2]}`, jsonMedia)
if w.Code != 503 || w.Header().Get("Retry-After") != "1" {
t.Errorf("overload=%d headers=%v", w.Code, w.Header())
}
if request(s, "GET", "/health", "", "").Code != 200 {
t.Error("health unavailable during admitted work")
}
close(gate)
for i := 0; i < 2; i++ {
if status := <-done; status != 200 {
t.Errorf("admitted status=%d", status)
}
}
}func TestSlotsReleased(t *testing.T) {
s := newService(fastSquare)
for i := 0; i < 6; i++ {
w := request(s, "POST", "/batch", `{"numbers":[2]}`, jsonMedia)
if w.Code != 200 {
t.Fatalf("request %d: status %d", i, w.Code)
}
request(s, "POST", "/batch", `{}`, jsonMedia)
}
if len(s.slots) != 0 {
t.Fatal("admission slot retained")
}
}func TestDeadline(t *testing.T) {
var active atomic.Int32
s := newService(func(ctx context.Context, _ int64) (int64, error) {
active.Add(1)
defer active.Add(-1)
<-ctx.Done()
return 0, ctx.Err()
})
s.budget = 20 * time.Millisecond
w := request(s, "POST", "/batch", `{"numbers":[1,2]}`, jsonMedia)
if w.Code != 408 || active.Load() != 0 || len(s.slots) != 0 {
t.Fatalf("status=%d active=%d slots=%d",
w.Code, active.Load(), len(s.slots))
}
}func TestOperationError(t *testing.T) {
s := newService(func(context.Context, int64) (int64, error) {
return 0, errors.New("private diagnostic detail")
})
w := request(s, "POST", "/batch", `{"numbers":[2]}`, jsonMedia)
leaked := strings.Contains(w.Body.String(), "private")
if w.Code != 500 || leaked || len(s.slots) != 0 {
t.Fatalf("status=%d body=%s slots=%d", w.Code, w.Body, len(s.slots))
}
}func TestNetworkCancellation(t *testing.T) {
started, stopped := make(chan struct{}), make(chan struct{})
s := newService(func(ctx context.Context, _ int64) (int64, error) {
close(started)
<-ctx.Done()
close(stopped)
return 0, ctx.Err()
})
s.budget = 5 * time.Second
server := httptest.NewServer(s)
defer server.Close()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
r, err := http.NewRequestWithContext(ctx, "POST",
server.URL+"/batch", strings.NewReader(`{"numbers":[1]}`))
if err != nil {
t.Fatal(err)
}
r.Header.Set("Content-Type", jsonMedia)
client := &http.Client{Timeout: 2 * time.Second}
done := make(chan error, 1)
go func() {
response, err := client.Do(r)
if response != nil {
response.Body.Close()
}
done <- err
}()
await(t, started)
cancel()
if err := <-done; !errors.Is(err, context.Canceled) {
t.Fatalf("client error=%v", err)
}
await(t, stopped)
}func Example_service() {
server := httptest.NewServer(newService(fastSquare))
defer server.Close()
client := &http.Client{Timeout: 2 * time.Second}
response, err := client.Post(server.URL+"/batch", jsonMedia,
strings.NewReader(`{"numbers":[3,1,2]}`))
if err != nil {
panic(err)
}
defer response.Body.Close()
body, err := io.ReadAll(response.Body)
if err != nil {
panic(err)
}
fmt.Println(response.StatusCode)
fmt.Print(string(body))
// Output:
// 200
// {"results":[9,1,4]}
}TestNetworkCancellation starts an actual local HTTP server. It waits until the operation is entered, cancels the client request and expects the client error to identify context cancellation. It also waits for the server operation to stop. This test sets the work budget to five seconds but requires the server stop signal within one second of client cancellation. The test therefore cannot pass by waiting for that five-second work deadline. It still does not establish a production cancellation-latency target. Example_service starts another actual server, sends [3,1,2] and checks the exact output shown in its Output comment.
The verified suite contains twelve top-level Test functions, twenty HTTP subtests and one executable Example. The complete suite passed once with verbose output and then twenty repeated runs. Go vet and a normal executable build also passed. Repetition can expose scheduling-dependent bugs but does not explore every interleaving. The release has no production-load, public-network, authentication or operating-system signal acceptance result.
Run go test -race ./... when your environment supports it. In this release it stopped before running tests with the message that -race requires cgo and CGO_ENABLED=1. The official Go race-detector documentation also describes platform and C-toolchain requirements. A disabled detector is not a clean race report. Enabling an environment variable alone is not proof that all required native tooling is installed. Even a passing detector reports only races encountered in executed paths.
Release review deliberately changed three copied implementations and ran focused tests. Rotating the result index caused TestOrder to fail. Removing the upper numeric bound caused the range case to fail. Removing the deferred slot release caused TestSlotsReleased to fail. Each failure came from a behavioural assertion, not a compiler error. The original files stayed unchanged. These experiments show that those three defects are detected by the relevant tests.
Before adding real work, decide how it is authorised, cancelled, retried and observed. Do not simply replace multiplication with an external send and assume a retry is harmless. Decide whether admission should be per tenant, per process or across a deployment, and measure the work and body-reading limits on the intended system. Preserve the small lab as an executable reference for ownership and response ordering while expanding the implementation.
Try it yourself · Activity 06
20 minCatch a leaked slot
Use a separate copy of this local five-file lab.
- Run the original tests successfully.
- Replace the deferred receive from s.slots with an empty deferred function.
- Run go test -run TestSlotsReleased and restore the original before rerunning all tests.
Does your test prove an observable contract, or would the same assertion pass if the important operation silently did nothing?
Worked answer
Without the receive, admitted requests leave tokens in the channel. Once both slots are occupied, a later sequential request receives 503 even though earlier handlers have returned. TestSlotsReleased detects this. Restore defer func() { <-s.slots }() and require the suite to pass. This demonstrates one admission lifecycle defect; it does not substitute for the unavailable race-detector run.
Keep learning
The complete workbook
Read a complete Go program across a module file and four source files. Implement a bounded queue, explicit channel ownership, joined cancellation and a small JSON contract. Check ordering, overload, deadlines, validation and a real HTTP disconnect. Six worked activities connect syntax and concurrency to observable behaviour.
- 01Read the contract before the goroutinesRead here · 1 exercise
Give a small service a precise job and a deliberately small input.
- 02Bound work and keep its original orderRead here · 1 exercise
Use channels to distribute ownership instead of sharing an append operation.
- 03Validate input and admit a bounded number of requestsRead here · 1 exercise
Worker limits and request limits control different amounts of work.
- 04Propagate failure and join cancellationRead here · 1 exercise
Return one complete result or an explicit failure.
- 05Start locally and test lifecycle rulesRead here · 1 exercise
Build one executable, then test the worker rules without a fixed public port.
- 06Test the HTTP boundary and read the evidenceRead here · 1 exercise
Use failures to check what the happy path leaves hidden.
Also inside: a 8-point checklist, a glossary of 8 terms and 10 questions and answers to test yourself. 6 hands-on exercises, each with a worked answer at the back where the workbook gives one.
No login, no card, no account. Before the download we ask you to follow Mickai (two quick links). Free to download and use for personal learning, study groups and inside your own team. Please do not resell the workbooks or republish them as your own. Link people to trust-agent.ai instead.
Test yourself
Questions and answers
Does this call an AI provider?
No. The operation is a fictional cancellable wait followed by integer squaring; no model or external account is involved.
Why is output ordered?
Each queue item is an input index and its worker writes to that preallocated result slot. The caller reads only after joining all workers.
Who closes the jobs channel?
The producing goroutine closes it after enqueuing every bounded job. Workers only receive.
How many operations can run together?
One service instance admits two batches with two workers each, so at most four operation calls can be in flight across those batches.
Are duplicate JSON keys rejected?
No. The encoding/json v1 parser retains duplicate-key and case-insensitive matching behaviour. This is not a canonical JSON validator.
Does cancel stop arbitrary code?
No. Cancellation is cooperative. A blocked operation that ignores its context can prevent the join and retain a request slot.
Does the budget include body parsing?
No. The work context is created after parsing. Body size and server read timeouts are separate controls.
Was Ctrl+C shutdown verified?
The executable builds, but the operating-system signal and graceful-drain sequence were not exercised in release testing.
Did the race detector pass?
No race-detector result exists. The command could not run tests because cgo was disabled in this environment.
What did the tests establish?
Twelve top-level tests, twenty HTTP subtests and one executable network example pass, including twenty repeated suite runs. Three intentional behavioural defects are caught. This remains bounded local verification.
When you have finished
Get your certificate of completion
Type your name and download a certificate for this workbook as a PDF, ready to print or to add to LinkedIn. It is made on your own device, so your name is never sent to us. It is a self-declared certificate, not an accredited qualification.
Learn the language
Key terms
- Goroutine
- A function executing concurrently with other goroutines under the Go runtime.
- Channel
- A typed communication mechanism used here to hand each queued index to one worker.
- Worker pool
- A fixed number of workers receiving multiple jobs from a shared queue.
- Context
- A value carrying cancellation and deadline signals across cooperating operations.
- Join
- Waiting until owned concurrent work has returned before continuing.
- Admission control
- A rule that decides whether a new request may begin holding application resources.
6 of the workbook's 8 terms. The complete glossary is in the workbook.
Follow the evidence
Sources and checks
Facts last checked: .
Examples in this workbook were run on: Go 1.26.2 on Windows 11, windows/amd64. Standard library only, GOTOOLCHAIN=local and GOPROXY=off. Twelve tests, twenty HTTP cases and one executable network example pass, with twenty repeated suite runs. Build and go vet pass. The race detector could not run because cgo is disabled. OS-signal shutdown was not exercised. (2026-09-27).
These workbooks use AI assistance. See how the workbooks are made.
- Effective Go: core language and concurrency conventionsThe Go Authors
- Go 1.26.2 context packageThe Go Authors
- Go 1.26.2 net/http packageThe Go Authors
- Go 1.26.2 httptest packageThe Go Authors
- Go race detector: use and requirementsThe Go Authors
- Go 1.26.2 encoding/json and parser considerationsThe Go Authors
Created by Mickarle Wagstaff-Irons - Micky Irons with the Mickai team. Published by Mickai LTD. Last updated 27 September 2026.
NextKeep going
Where to go next
Recommended for you
Node.js: build and test a small API
Build a local read-only HTTP API with Node.js: explicit routes, strict query validation, predictable errors and executable tests using built-in modules.
Recommended for you
Test-driven development with an AI coding assistant
Use failing examples to guide an AI coding assistant, test boundary cases, catch deliberate defects and hand over a small Python change with clear evidence.