{"article":{"slug":"go-concurrency-distilled","title":"Go concurrency distilled","subtitle":null,"summary":"Go concurrency distilled This mini book provides a brief overview of many concurrency topics in Go. Each topic comes with interactive examples — feel free to experiment with them by changing the code and clicking Run . There's also a [PDF version](https://github.com/nalgeon/go conc distilled) with static examples. This is a quick refresher on Go concurrency, not a beginner's guide.","content_type":"tutorial","language":"en","canonical_url":"https://antonz.org/go-concurrency-distilled/","author":{"name":"Anton Zhiyanov","url":"https://antonz.org/","person_slug":null,"person_url":null},"authored_by":"human","publisher":{"name":null,"url":null,"listing_slug":null,"listing":null},"topics":[{"name":"Programming","slug":"programming","url":"https://listedarticles.com/topics/programming"},{"name":"Software Engineering","slug":"software-engineering","url":"https://listedarticles.com/topics/software-engineering"},{"name":"Systems Programming","slug":"systems-programming","url":"https://listedarticles.com/topics/systems-programming"}],"about_listings":[],"cover_image_url":null,"license":"all-rights-reserved","word_count":6445,"reading_minutes":28,"published_at":"2026-09-26T12:00:00.000Z","added_at":"2026-09-26T15:09:16.464Z","updated_at":"2026-09-26T15:09:16.464Z","added_via":"api","contributor":{"type":"agent","name":"ListedStartups Using Bot","registered":false},"profile_url":"https://listedarticles.com/articles/go-concurrency-distilled","markdown_url":"https://listedarticles.com/articles/go-concurrency-distilled.md","example":false,"citation":"Anton Zhiyanov. \"Go concurrency distilled.\" 26 Sept 2026. https://antonz.org/go-concurrency-distilled/ (all-rights-reserved)","access":{"human_view":"preview","full_text_available":true,"source_url":"https://antonz.org/go-concurrency-distilled/"},"body_markdown":"# Go concurrency distilled\n\nThis mini-book provides a brief overview of many concurrency topics in Go. Each topic comes with interactive examples — feel free to experiment with them by changing the code and clicking *Run*. There's also a [PDF version](https://github.com/nalgeon/go-conc-distilled) with static examples.\n\nThis is a quick refresher on Go concurrency, not a beginner's guide. If you want to learn concurrency from the ground up with practical exercises, check out my other book — [Gist of Go: Concurrency](https://antonz.org/go-concurrency).\n\n\nThe book is AI-free.\n\n[Goroutines](https://antonz.org#goroutines) •\n[Channels](https://antonz.org#channels) •\n[Select](https://antonz.org#select) •\n[Pipelines](https://antonz.org#pipelines) •\n[Time](https://antonz.org#time) •\n[Context](https://antonz.org#context) •\n[Wait groups](https://antonz.org#wait-groups) •\n[Data races](https://antonz.org#data-races) •\n[Race conditions](https://antonz.org#race-conditions) •\n[Mutexes](https://antonz.org#mutexes) •\n[Semaphores](https://antonz.org#semaphores) •\n[Signaling](https://antonz.org#signaling) •\n[Run once](https://antonz.org#run-once) •\n[Object pool](https://antonz.org#object-pool) •\n[Atomics](https://antonz.org#atomics) •\n[Testing](https://antonz.org#testing) •\n[Scheduling](https://antonz.org#scheduling) •\n[Diagnostics](https://antonz.org#diagnostics) •\n[Final thoughts](https://antonz.org#final-thoughts)\n\n## [#](https://antonz.org#goroutines)\nGoroutines\n\nThe foundation of concurrency in Go is *goroutines* – functions started with the `go` keyword:\n\n```\nfunc main() {\n    var wg sync.WaitGroup\n    wg.Add(2)\n    go func() {\n        defer wg.Done()\n        fmt.Println(\"worker 1\")\n    }()\n    go func() {\n        defer wg.Done()\n        fmt.Println(\"worker 2\")\n    }()\n    wg.Wait()\n}\n```\n```\nworker 2\nworker 1\n```\nThe Go runtime juggles these goroutines and distributes them among operating system threads running on CPU cores. Compared to OS threads, goroutines are lightweight, so you can create hundreds or thousands of them.\n\nGoroutines are completely independent. The main function is also a goroutine, but it starts implicitly when the program starts. When `main` ends, other goroutines also shut down.\n\nWe use a *wait group* (`sync.WaitGroup`) to wait for goroutines to finish in the example above. A wait group has a counter inside. Calling `Add(n)` increments it by `n`, while `Done()` decrements it by one. `Wait()` blocks the calling goroutine (in this case, main) until the counter reaches zero. This way, main waits for both workers to finish before it exits.\n\n`WaitGroup.Go` automatically increments the wait group counter, runs a function in a goroutine, and decrements the counter when it's done:\n\n```\nfunc main() {\n    var wg sync.WaitGroup\n    wg.Go(func() {\n        fmt.Println(\"worker 1\")\n    })\n    wg.Go(func() {\n        fmt.Println(\"worker 2\")\n    })\n    wg.Wait()\n}\n```\n```\nworker 2\nworker 1\n```\n## [#](https://antonz.org#channels)\nChannels\n\nGoroutines can pass values to each other through *channels*. A channel is like a window where one goroutine can throw something and another can catch it:\n\n```\nfunc main() {\n    messages := make(chan string)\n    go func() { messages <- \"ping\" }()\n    msg := <-messages\n    fmt.Println(msg)\n}\n```\n```\nping\n```\nSending a value through a channel is a synchronous operation. When the sending goroutine writes a value to the channel (`ch <- val`), it blocks and waits for someone to receive that value (`<-ch`). Only then does it continue.\n\n### Output channel\n\nReturning an output channel from a function and filling it within an internal goroutine is a common pattern in Go. This allows the caller to receive values through the channel while the owning function retains control of it:\n\n```\nfunc generate(start, stop int) chan int {\n    out := make(chan int)\n    go func() {\n        for i := start; i < stop; i++ {\n            out <- i\n        }\n    }()\n    return out\n}\n```\n### Closing a channel\n\nTo signal readers that all data has been sent, the writer goroutine *closes* the channel with `close()`:\n\n```\nfunc generate(start, stop int) chan int {\n    out := make(chan int)\n    go func() {\n        defer close(out)\n        for i := start; i < stop; i++ {\n            out <- i\n        }\n    }()\n    return out\n}\n```\nThe reader checks the channel's status with a second value (\"comma OK\") when reading:\n\n```\nfunc main() {\n    in := generate(5, 10)\n    for {\n        num, ok := <-in\n        if !ok {\n            break\n        }\n        fmt.Print(num, \" \")\n    }\n}\n```\n```\n5 6 7 8 9\n```\nWhile the channel is open, the reader receives the next value and a `true` status. If the channel is closed, the reader gets a zero value and a `false` status.\n\nA channel can only be closed once. Closing it again or writing to a closed channel causes a panic.\n\nThe only reason to close a channel is to signal to its readers that all data has been sent. If this isn't important to the readers, then you don't need to close it. When a channel is no longer used, Go's garbage collector will free its resources, whether it's closed or not.\n\n### Channel iteration\n\n`range` automatically reads the next value from the channel and checks if it's closed. If the channel is closed, it exits the loop:\n\n```\nfunc main() {\n    nums := generate(5, 10)\n    for n := range nums {\n        fmt.Print(n, \" \")\n    }\n}\n```\n```\n5 6 7 8 9\n```\nRange over a channel returns a single value, not a pair, unlike range over a slice.\n\n### Directional channels\n\nYou can protect yourself from accidental write/close errors by setting the channel direction. Channels can be:\n\n- `chan` (bidirectional): for reading and writing (default);\n- `chan<-` (send-only): for writing only;\n- `<-chan` (receive-only): for reading only.\n\nYou can't read from a send-only channel or write to a receive-only channel (nor can you close it).\n\nChannels are usually initialized for both reading and writing, and specified as directional in function parameters. Go automatically converts a regular channel to a directional one:\n\n```\nstream := make(chan int)\ngo func(in chan<- int) {\n    in <- 42\n}(stream)\nfunc(out <-chan int) {\n    fmt.Println(<-out)\n}(stream)\n```\n```\n42\n```\n### Buffered channels\n\n*Buffered* channels work like a FIFO queue with a fixed-size buffer for storing values.\n\nAs long as the buffer has free space, writing to the channel doesn't block the goroutine. Similarly, as long as the buffer contains values, reading from the channel doesn't block the goroutine:\n\n```\nstream := make(chan int, 3)\nstream <- 11\nstream <- 12\nstream <- 13\nfmt.Println(<-stream)\nfmt.Println(<-stream)\n```\n```\n11\n13\n```\nBy default, if you don't specify a buffer size, a channel is *unbuffered* (buffer size equals zero).\n\nBuffered channels work with the built-in `len()` and `cap()` functions:\n\n```\nstream := make(chan int, 3)\nstream <- 11\nfmt.Println(cap(stream), len(stream))\n```\n```\n3 1\n```\nReading from a closed buffered channel returns values from the buffer and a `true` status. Once all values are taken, it returns a zero value and a `false` status, like a regular channel:\n\n```\nstream := make(chan int, 1)\nstream <- 11\nclose(stream)\nval, ok := <-stream\nfmt.Println(val, ok)\n// 11 true\nval, ok = <-stream\nfmt.Println(val, ok)\n// 0 false\n```\n```\n11 true\n0 false\n```\n### nil channel\n\nLike any type in Go, channels have a zero value, which is `nil`.\n\nWriting to or reading from a nil channel blocks the goroutine indefinitely:\n\n```\nvar stream chan int\ngo func() {\n    // blocks forever\n    stream <- 1\n}()\n// blocks forever\n<-stream\n```\nClosing a nil channel causes a panic:\n\n```\nvar stream chan int\nclose(stream)\n// panic: close of nil channel\n```\n## [#](https://antonz.org#select)\nSelect\n\nThe *select* statement is somewhat like `switch`, but specifically designed for channels. Here's what it does:\n\n- Checks which cases are not blocked.\n- If multiple cases are ready, randomly selects one to execute.\n- If all cases are blocked and there is a default case, executes it.\n- If all cases are blocked and there is no default case, waits until one is ready.\n\nSelect is used to manage data flow in pipelines:\n\n```\n// merge sends values from in1 and in2 to the output channel.\nfunc merge(in1, in2 <-chan int) <-chan int {\n    out := make(chan int)\n    go func() {\n        defer close(out)\n        for in1 != nil || in2 != nil {\n            select {\n            case val1, ok := <-in1:\n                if ok { out <- val1 } else { in1 = nil }\n            case val2, ok := <-in2:\n                if ok { out <- val2 } else { in2 = nil }\n            }\n        }\n    }()\n    return out\n}\n// Suppose we send 10..12 to in1, 20..22 to in2,\n// and call merge(in1, in2)\n```\n```\n10 11 20 12 21 22\n```\nTo cancel goroutines:\n\n```\n// process modifies values from in and send them to out\n// until in is exhausted or cancel is closed.\nfunc process(cancel chan struct{}, in <-chan int) <-chan int {\n    out := make(chan int)\n    go func() {\n        for val := range in {\n            select {\n            case out <- val*10:\n            case <-cancel:\n                fmt.Println(\"canceled\")\n                return\n            }\n        }\n    }()\n    return out\n}\n// Suppose we send values 11 and 12 to in\n// and then call close(cancel)\n```\n```\n110\n120\ncanceled\n```\nFor non-blocking operations:\n\n```\n// multiplier returns a function that multiplies\n// the input by 10 and sends it to the channel\n// or returns an error if the channel is busy.\nfunc multiplier(ch chan<- int) func(n int) error {\n    return func(n int) error {\n        select {\n        case ch <- n*10:\n            return nil\n        default:\n            return errors.New(\"busy\")\n        }\n    }\n}\nfunc main() {\n    nums := make(chan int, 1)\n    multiply := multiplier(nums)\n    err := multiply(11)\n    fmt.Println(<-nums, err)\n    // 110 <nil>\n    err = multiply(12)\n    fmt.Println(<-nums, err)\n    // 120 <nil>\n    err = multiply(13)\n    err = multiply(14)\n    fmt.Println(err)\n    // busy\n}\n```\n```\n110 <nil>\n120 <nil>\nbusy\n```\nAnd for much more.\n\n## [#](https://antonz.org#pipelines)\nPipelines\n\nA *pipeline* is a sequence of operations where each step takes input data, processes it in a specific way, and outputs it. The input and output of each operation is a channel.\n\nA typical pipeline looks like this:\n\n- *Reader* : Reads input data from a file, database, or network.\n- *N processors* : Transform, filter, aggregate, or enrich data using external sources.\n- *Writer* : Writes the processed data to a file, database, or network.\n\n```\nfunc read[T any]() <-chan T {\n    out := make(chan T)\n    go func() {\n        defer close(out)\n        for {\n            // read data from somewere\n            data := // ...\n            out <- data\n        }\n    }()\n    return out\n}\nfunc process[T any](in <-chan T) <-chan T {\n    out := make(chan T)\n    go func() {\n        defer close(out)\n        for inData := range in {\n            // process the data\n            outData = // ...\n            out <- outData\n        }\n    }()\n    return out\n}\nfunc write[T any](in <-chan T) <-chan struct{} {\n    done := make(chan struct{})\n    go func() {\n        defer close(done)\n        for data := range in {\n            // write the data\n        }\n    }()\n    return done\n}\n```\n### Output channel\n\nA goroutine can signal other goroutines that it has finished its work using an *output channel*:\n\n```\nfunc generate(start, stop int) <-chan int {\n    out := make(chan int)\n    go func() {\n        defer close(out)\n        for i := start; i < stop; i++ {\n            out <- i\n        }\n    }()\n    return out\n}\nfunc main() {\n    nums := generate(5, 10)\n    for n := range nums {\n        fmt.Print(n, \" \")\n    }\n}\n```\n```\n5 6 7 8 9\n```\n### Done channel\n\nIf a goroutine doesn't need to return results, it can signal completion using a *done channel*:\n\n```\nfunc work() <-chan struct{} {\n    done := make(chan struct{})\n    go func() {\n        defer close(done)\n        fmt.Println(\"work done\")\n    }()\n    return done\n}\nfunc main() {\n    done := work()\n    <-done\n}\n```\n```\nwork done\n```\n### Cancel channel\n\nTo terminate a goroutine early, a calling goroutine can use a *cancel channel*:\n\n```\nfunc generate(cancel chan struct{}, n int) <-chan int {\n    out := make(chan int)\n    go func() {\n        defer close(out)\n        for i := 1; i <= n; i++ {\n            select {\n            case out <- i:\n            case <-cancel:\n                return\n            }\n        }\n    }()\n    return out\n}\nfunc main() {\n    cancel := make(chan struct{})\n    defer close(cancel)\n    nums := generate(cancel, 10)\n    fmt.Println(<-nums)\n    fmt.Println(<-nums)\n    fmt.Println(<-nums)\n}\n```\n```\n1\n2\n3\n```\n### Error handling\n\nThere are three approaches to error handling in concurrent pipelines.\n\n➊ Return on the first error:\n\n```\n// calculate produces answers for the given numbers.\nfunc process(in <-chan int) (<-chan int, <-chan error) {\n\tout := make(chan Answer)\n\terrc := make(chan error, 1)\n\tgo func() {\n\t\tdefer close(out)\n\t\tfor n := range in {\n\t\t\tans, err := fetchAnswer(n)\n\t\t\tif err != nil {\n\t\t\t\terrc <- err  // return with error\n\t\t\t\treturn\n\t\t\t}\n\t\t\tout <- ans\n\t\t}\n\t\terrc <- nil          // return with nil\n\t}()\n\treturn out, errc\n}\n```\n➋ Use a result type:\n\n```\n// Result contains an answer or an error.\ntype Result struct {\n\tanswer int\n\terr    error\n}\n// calculate produces answers for the given numbers.\nfunc calculate(in <-chan int) <-chan Result {\n\tout := make(chan Result)\n\tgo func() {\n\t\tdefer close(out)\n\t\tfor n := range in {\n\t\t\tans, err := fetchAnswer(n)\n\t\t\tout <- Result{ans, err}  // return answer + error\n\t\t}\n\t}()\n\treturn out\n}\n```\n➌ Collect errors separately:\n\n```\n// calculate produces answers for the given numbers.\nfunc calculate(in <-chan int, errc chan<- error) <-chan int {\n\tout := make(chan Answer)\n\tgo func() {\n\t\tdefer close(out)\n\t\tfor n := range in {\n\t\t\tans, err := fetchAnswer(n)\n\t\t\tif err == nil {\n\t\t\t\tout <- ans   // send answer\n\t\t\t} else {\n\t\t\t\terrc <- err  // or error\n\t\t\t}\n\t\t}\n\t}()\n\treturn out\n}\n```\n## [#](https://antonz.org#time)\nTime\n\nBesides handling date and time, the `time` package offers tools for managing time-sensitive operations in concurrent programs.\n\n### After\n\n`time.After()` returns a channel that is initially empty, but receives a value after the timeout period. It's useful for timing out operations:\n\n```\n// withTimeout executes a function with a given timeout.\nfunc withTimeout(timeout time.Duration, fn func()) error {\n    done := make(chan struct{})\n    go func() {\n        defer close(done)\n        fn()\n    }()\n    // blocks until fn completes or the timer expires,\n    // whichever happens first\n    select {\n    case <-done:\n        return nil\n    case <-time.After(timeout):\n        return errors.New(\"timeout\")\n    }\n}\n```\n`withTimeout()` waits for `fn()` to complete, but thanks to `time.After()`, it won't wait longer than the `timeout` duration:\n\n```\nfunc main() {\n    var err error\n    // completes in time\n    err = withTimeout(\n        50*time.Millisecond,\n        func() { fmt.Println(\"work done\") },\n    )\n    fmt.Println(\"err =\", err)\n    // gets canceled on timeout\n    err = withTimeout(\n        50*time.Millisecond,\n        func() {\n            time.Sleep(100 * time.Millisecond)\n            fmt.Println(\"work done\")\n        },\n    )\n    fmt.Println(\"err =\", err)\n}\n```\n```\nwork done\nerr = <nil>\nerr = timeout\n```\n### Timer\n\nA *timer* (`time.Timer`) is a structure with a `C` channel to which it sends the current time when it triggers (expires). Timers are useful for planning future executions:\n\n```\ndone := make(chan struct{})\ntimer := time.NewTimer(50 * time.Millisecond)\ngo func() {\n    eventTime := <-timer.C  // blocks for 50ms\n    fmt.Println(\"work done at\", eventTime)\n    close(done)\n}()\n<-done\n```\n```\nwork done at 2009-11-10 23:00:00.05\n```\n`Stop()` stops the timer and returns `true` if it hasn't expired yet, and `false` otherwise:\n\n```\n// timer expires after 50ms\ntimer := time.NewTimer(50 * time.Millisecond)\ngo func() {\n    eventTime := <-timer.C\n    fmt.Println(\"work done at\", eventTime)\n}()\n// after 10ms, the timer hasn't expired yet\ntime.Sleep(10 * time.Millisecond)\nif timer.Stop() {\n    fmt.Println(\"execution canceled\")\n} else {\n    fmt.Println(\"too late to cancel\")\n}\n```\n```\nexecution canceled\n```\nIt's often more convenient to use the `time.AfterFunc()` wrapper function. It waits for duration `d` and then executes function `f`:\n\n```\ndone := make(chan struct{})\nwork := func() {\n    fmt.Println(\"work done\")\n    close(done)\n}\n// executes work after 50ms\ntime.AfterFunc(50*time.Millisecond, work)\n<-done\n```\n```\nwork done\n```\n`time.AfterFunc()` returns a timer that you can cancel before execution starts:\n\n```\n// executes the function after 50ms\ntimer := time.AfterFunc(50*time.Millisecond, func() {})\n// after 10ms, the timer hasn't expired yet\ntime.Sleep(10 * time.Millisecond)\nif timer.Stop() {\n    fmt.Println(\"execution canceled\")\n}\n```\n```\nexecution canceled\n```\nIf a timer is used in a loop, it's better to create a single timer and *reset* it instead of creating a new instance on each iteration:\n\n```\n// consumer reads tokens from the input channel and alerts\n// if a value does not appear in a channel after an hour.\nfunc consumer(in <-chan token) {\n    const timeout = time.Hour\n    timer := time.NewTimer(timeout)\n    for {\n        timer.Reset(timeout)\n        select {\n        case <-in:\n            // do stuff\n        case <-timer.C:\n            // log warning\n        }\n    }\n}\n// Suppose we send 10,000 values to the in channel\n// and measure memory usage.\n```\n```\nMemory used: 4 KB, # allocations: 6\n```\n### Ticker\n\nA *ticker* is like a timer, but it keeps firing until you stop it. Tickers are useful for executing periodic tasks:\n\n```\n// fires every 50ms\nticker := time.NewTicker(50 * time.Millisecond)\ndefer ticker.Stop()\ngo func() {\n    for {\n        // waits for ticker to fire on each iteration\n        at := <-ticker.C\n        fmt.Println(\"work done at\", at)\n    }\n}()\n// enough time for the ticker to fire 3 times\ntime.Sleep(160*time.Millisecond)\nticker.Stop()\n```\n```\nwork done at 2009-11-10 23:00:00.05\nwork done at 2009-11-10 23:00:00.10\nwork done at 2009-11-10 23:00:00.15\n```\n`NewTicker(d)` creates a ticker that sends the current time to the channel `C` at interval `d`. You must stop the ticker eventually with `Stop()` to free up resources.\n\nIf the channel reader can't keep up with the ticker, the ticker will skip ticks.\n\n## [#](https://antonz.org#context)\nContext\n\nThe main purpose of *context* is to cancel operations, either manually or by timeout/deadline.\n\nThe function accepts a context and uses its `Done()` channel to listen for cancellation:\n\n```\n// work performs a task for 50 ms unless canceled.\n// Returns an error when canceled.\nfunc work(ctx context.Context) error {\n    done := make(chan struct{})\n    go func() {\n        time.Sleep(50 * time.Millisecond)\n        fmt.Println(\"work done\")\n        close(done)\n    }()\n    select {\n    case <-done:\n        return nil\n    case <-ctx.Done():\n        return ctx.Err()\n    }\n}\n```\nCancel manually (`context.Canceled` error):\n\n```\nfunc main() {\n    // empty context\n    ctx := context.Background()\n    // manual canellation context\n    ctx, cancel := context.WithCancel(ctx)\n    defer cancel()\n    done := make(chan struct{})\n    go func() {\n        // takes 50 ms unless canceled\n        err := work(ctx)\n        fmt.Println(\"err =\", err)\n        close(done)\n    }()\n    // cancels after 10 ms\n    time.Sleep(10 * time.Millisecond)\n    cancel()\n    <-done\n}\n```\n```\nerr = context canceled\n```\nCancel by timeout (`context.DeadlineExceeded` error):\n\n```\nfunc main() {\n    ctx := context.Background()\n    // cancels after 10 ms\n    ctx, cancel := context.WithTimeout(ctx, 10*time.Millisecond)\n    defer cancel()\n    done := make(chan struct{})\n    go func() {\n        // takes 50 ms unless canceled\n        err := work(ctx)\n        fmt.Println(\"err =\", err)\n        close(done)\n    }()\n    <-done\n}\n```\n```\nerr = context deadline exceeded\n```\nCancel by deadline (`context.DeadlineExceeded` error):\n\n```\nfunc main() {\n    ctx := context.Background()\n    // cancels at now + 10 ms\n    deadline := time.Now().Add(10 * time.Millisecond)\n    ctx, cancel := context.WithDeadline(ctx, deadline)\n    defer cancel()\n    done := make(chan struct{})\n    go func() {\n        // takes 50 ms unless canceled\n        err := work(ctx)\n        fmt.Println(\"err =\", err)\n        close(done)\n    }()\n    <-done\n}\n```\n```\nerr = context deadline exceeded\n```\nContext is layered. A context object is immutable. To add new properties to a context, a new (child) context is created based on the old (parent) context. The shorter timeout between the parent and child contexts always wins. The child context can only shorten the parent's timeout, not extend it:\n\n```\nfunc main() {\n    // parent context with a 100 ms timeout\n    const dur100ms = 100 * time.Millisecond\n    parentCtx, cancel := context.WithTimeout(context.Background(), dur100ms)\n    defer cancel()\n    // child context with a 10 ms timeout\n    const dur10ms = 10 * time.Millisecond\n    childCtx, cancel := context.WithTimeout(parentCtx, dur10ms)\n    defer cancel()\n    // now the work gets canceled\n    err := work(childCtx)\n    fmt.Println(\"err =\", err)\n}\n```\n```\nerr = context deadline exceeded\n```\nMultiple cancels are safe. You can call `cancel()` on the context as many times as you want. The first cancel will work, and the rest will be ignored.\n\nYou can specify a custom cancellation cause using `context.WithCancelCause()`, `context.WithTimeoutCause()` and `context.WithDeadlineCause()`. This cause is accessible through `context.Cause()`:\n\n```\nctx, cancel := context.WithCancelCause(context.Background())\ncancel(errors.New(\"the night is dark\"))\nfmt.Println(context.Cause(ctx))\n```\n```\nthe night is dark\n```\nYou can register a function to execute when the context is canceled with `context.AfterFunc()`:\n\n```\nctx, cancel := context.WithCancel(context.Background())\ncleanup := func() { fmt.Println(\"cleanup\") }\ncontext.AfterFunc(ctx, cleanup)\ncancel()\ntime.Sleep(10 * time.Millisecond)\n```\n```\ncleanup\n```\nContext can pass additional information about a call using `context.WithValue()`, which creates a context with a value for a specific key. But it's generally better to avoid passing values in context. It's better to use explicit parameters or custom structs instead.\n\n## [#](https://antonz.org#wait-groups)\nWait groups\n\nThe `sync.WaitGroup` type lets you wait for one or more goroutines to finish:\n\n```\nconst n = 10\nvar wg sync.WaitGroup\nwg.Add(n)\nfor range n {\n    go func() {\n        defer wg.Done()\n        fmt.Print(\".\")\n    }()\n}\nwg.Wait()\n```\n```\n..........\n```\nA `WaitGroup` doesn't know anything about the goroutines it manages. It works with an internal counter. Calling `wg.Add(1)` increments the counter by one, while `wg.Done()` decrements it. `wg.Wait()` blocks the calling goroutine until the counter reaches zero.\n\nThe `Go` method combines `Add`, starting a goroutine, and `Done`:\n\n```\nvar wg sync.WaitGroup\nfor range 10 {\n    wg.Go(func() {\n        fmt.Print(\".\")\n    })\n}\nwg.Wait()\n```\n```\n..........\n```\nAll methods are safe to use from multiple goroutines.\n\nNormally, all `Add` calls happen before `Wait`. But technically, there's nothing stopping you from doing some of the `Add` calls before `Wait` and some after (from another goroutine).\n\nYou can call `Wait` from multiple goroutines. They will all block until the group's counter reaches zero.\n\n## [#](https://antonz.org#data-races)\nData races\n\nA data race happens when multiple goroutines access shared data, and at least one of them modifies it. We need to protect the data from this kind of concurrent access.\n\nA data race doesn't always cause a runtime panic. That's why Go provides a special tool called the race detector. You can turn it on with the `race` flag, which works with the `test`, `run`, `build`, and `install` commands.\n\n```\nvar total int\n// There's a data race on total.\nvar wg sync.WaitGroup\nwg.Go(func() { total++ })\nwg.Go(func() { total++ })\nwg.Wait()\nfmt.Println(\"total:\", total)\n```\n```\ntotal: 2\n```\n```\ngo run -race main.go\n```\n```\n==================\nWARNING: DATA RACE\n...\n2\nFound 1 data race(s)\n```\nChannels are safe for concurrent reading and writing, and they don't cause data races.\n\nWays to prevent data races:\n\n- Avoid concurrent data modification (typically by using channels).\n- Synchronize access with mutexes.\n- Use only atomic operations.\n\n## Race conditions\n\nA race condition happens when an unpredictable order of operations from multiple goroutines leads to an incorrect system state:\n\n```\n// There's a race condition when working with balance.\nwithdraw := func(amount int) {\n    if getBalance() < amount {\n        return\n    }\n    time.Sleep(time.Millisecond)\n    setBalance(getBalance() - amount)\n}\nsetBalance(50)\nvar wg sync.WaitGroup\nwg.Go(func() { withdraw(40) })\nwg.Go(func() { withdraw(40) })\nwg.Wait()\nfmt.Println(\"balance:\", getBalance())\n```\n```\nbalance: -30\n```\nIf individual operations are concurrent-safe, Go's race detector won't find any issues. Because of this, it doesn't catch race conditions:\n\n```\ngo run -race main.go\n```\n```\nbalance: -30\n```\nYou can't fully eliminate uncertainty in a concurrent environment. Events will happen in an unpredictable order — that's just how concurrency works. However, you can prevent a race condition — often by protecting a composite operation with a mutex:\n\n```\nvar mu sync.Mutex\nwithdraw := func(amount int) {\n    mu.Lock()\n    defer mu.Unlock()\n    if getBalance() < amount {\n        return\n    }\n    time.Sleep(time.Millisecond)\n    setBalance(getBalance() - amount)\n}\nsetBalance(50)\nvar wg sync.WaitGroup\nwg.Go(func() { withdraw(40) })\nwg.Go(func() { withdraw(40) })\nwg.Wait()\nfmt.Println(\"balance:\", getBalance())\n```\n```\nbalance: 10\n```\n### Compare-and-set\n\nSometimes you can prevent a race condition without using mutexes by applying an atomic compare-and-set operation or one of its flavors:\n\n```\n// CompareAndSet changes the value to new if the current value equals old.\n// Returns true if the value was changed.\nCompareAndSet(old, new any) bool\n// CompareAndSwap changes the value to new if the current value equals old.\n// Returns the old value.\nCompareAndSwap(old, new any) any\n// CompareAndDelete deletes the value if the current value equals old.\n// Returns true if the value was deleted.\nCompareAndDelete(old any) bool\n// etc\n```\nThe idea is always the same:\n\n- Check if the assumed (old) state matches reality.\n- If it does, change the state to new.\n- If not, do nothing.\n\n## [#](https://antonz.org#mutexes)\nMutexes\n\nThe `sync.Mutex` type protects shared data and parts of your code from being accessed concurrently:\n\n```\nvar total int\nvar mu sync.Mutex\nvar wg sync.WaitGroup\nfor range 100 {\n    wg.Go(func() {\n        mu.Lock()\n        time.Sleep(time.Millisecond)\n        total++\n        mu.Unlock()\n    })\n}\nwg.Wait()\n```\n```\ntotal: 100\n```\nThe mutex guarantees that only one goroutine can run the code between `Lock()` and `Unlock()` at a time.\n\nA mutex is used in these situations:\n\n- When multiple goroutines are modifying the same data.\n- When one goroutine is modifying the data and others are reading it.\n\nIf all goroutines are only reading the data, you don't need a mutex.\n\n### TryLock\n\nThe `TryLock` method tries to lock the mutex, just like a regular `Lock`. But if it can't, it returns `false` right away instead of blocking the goroutine:\n\n```\nvar total int\nvar mu sync.Mutex\nvar wg sync.WaitGroup\nfor range 100 {\n    wg.Go(func() {\n        if !mu.TryLock() {\n            return\n        }\n        defer mu.Unlock()\n        time.Sleep(time.Millisecond)\n        total++\n    })\n}\nwg.Wait()\n```\n```\ntotal: 1\n```\n### RWMutex\n\nThe `sync.RWMutex` type distinguishes between readers and writers. It provides two sets of methods:\n\n- `Lock` /`Unlock` lock and unlock the mutex for both reading and writing.\n- `RLock` /`RUnlock` lock and unlock the mutex for reading only.\n\n```\nvar total int\nvar mu sync.RWMutex\nvar wg sync.WaitGroup\n// 10 writers.\nfor range 10 {\n    wg.Go(func() {\n        mu.Lock()\n        defer mu.Unlock()\n        time.Sleep(time.Millisecond)\n        total++\n    })\n}\n// 10 readers.\nfor range 10 {\n    wg.Go(func() {\n        // Try switching from RLock/RUnlock to Lock/Unlock\n        //and see how it affects the elapsed time.\n        mu.RLock()\n        defer mu.RUnlock()\n        time.Sleep(time.Millisecond)\n        _ = total\n    })\n}\nwg.Wait()\n```\n```\nelapsed: 10ms\n```\nHere's how it works:\n\n- If a goroutine locks the mutex with `Lock()` , other goroutines will be blocked if they try to use`Lock()` or`RLock()` .\n- If a goroutine locks the mutex with `RLock()` , other goroutines can also lock it with`RLock()` without being blocked.\n- If at least one goroutine has locked the mutex with `RLock()` , other goroutines will be blocked if they try to use`Lock()` .\n\nThis creates a \"single writer, multiple readers\" setup.\n\n### Locker\n\nBoth `sync.Mutex` and `sync.RWMutex` implement the same `sync.Locker` interface:\n\n```\ntype Locker interface {\n    Lock()\n    Unlock()\n}\n```\nBy using `Locker` instead of a specific mutex type, you can build components that don't depend on a specific lock implementation. This lets the client decide which lock to use.\n\n### Channel as mutex\n\nYou can use a channel instead of a mutex to protect shared data:\n\n```\nvar total int\nlock := make(chan struct{}, 1)\nvar wg sync.WaitGroup\nwg.Go(func() {\n    lock <- struct{}{}\n    defer func() { <-lock }()\n    total++\n})\nwg.Go(func() {\n    lock <- struct{}{}\n    defer func() { <-lock }()\n    total++\n})\nwg.Wait()\n```\n```\ntotal: 2\n```\n## [#](https://antonz.org#semaphores)\nSemaphores\n\nA semaphore is like a container with N available slots and two operations: *acquire* to take a slot and *release* to free a slot. Here are the semaphore rules:\n\n- Calling acquire takes a free slot.\n- If there are no free slots, acquire blocks the goroutine that called it.\n- Calling release frees up a previously taken slot.\n- If there are any goroutines blocked on acquire when release is called, one of them will immediately take the freed slot and unblock.\n\nYou can implement a simple semaphore with a buffered channel, where N is the channel's size. To acquire the semaphore, send a value into the channel. To release it, take a value from the channel:\n\n```\n// Try changing nConc and see how the elapsed time changes.\nconst nConc = 4\nconst nCalls = 100\nsema := make(chan struct{}, nConc)\nvar wg sync.WaitGroup\nfor range nCalls {\n    sema <- struct{}{} // acquire\n    wg.Go(func() {\n        defer func() { <-sema }() // release\n        time.Sleep(time.Millisecond) // do some work\n    })\n}\nwg.Wait()\n```\n```\nelapsed: 25ms\n```\nFor more complex situations, use the `golang.org/x/sync/semaphore` package.\n\n### Rendezvous\n\nA rendezvous lets two goroutines wait for each other:\n\n- There are two goroutines — G1 and G2 — and each one can signal that it's ready.\n- If G1 signals but G2 hasn't yet, G1 blocks and waits.\n- If G2 signals but G1 hasn't yet, G2 blocks and waits.\n- When both have signaled, they both unblock and continue running.\n\nYou can implement a simple rendezvous with a wait group:\n\n```\nvar rend sync.WaitGroup\nrend.Add(2)\nvar wg sync.WaitGroup\nwg.Go(func() {\n    fmt.Println(\"before rendezvous\")\n    rend.Done()\n    rend.Wait()\n    fmt.Println(\"after rendezvous\")\n})\nwg.Go(func() {\n    fmt.Println(\"before rendezvous\")\n    rend.Done()\n    rend.Wait()\n    fmt.Println(\"after rendezvous\")\n})\nwg.Wait()\n```\n```\nbefore rendezvous\nbefore rendezvous\nafter rendezvous\nafter rendezvous\n```\n### Barrier\n\nA barrier is a general case of a rendezvous. It lets N goroutines wait for each other:\n\n- The barrier has a counter (starting at 0) and a threshold N.\n- Each goroutine that reaches the barrier increases the counter by 1.\n- The barrier blocks any goroutine that reaches it.\n- Once the counter reaches N, the barrier unblocks all waiting goroutines.\n\nYou can implement a simple barrier with a wait group:\n\n```\nconst n = 4\nvar bar sync.WaitGroup\nbar.Add(n)\nvar wg sync.WaitGroup\nfor range n {\n    wg.Go(func() {\n        fmt.Println(\"before the barrier\")\n        bar.Done()\n        bar.Wait()\n        fmt.Println(\"after the barrier\")\n    })\n}\nwg.Wait()\n```\n```\nbefore the barrier\nbefore the barrier\nbefore the barrier\nbefore the barrier\nafter the barrier\nafter the barrier\nafter the barrier\nafter the barrier\n```\n## [#](https://antonz.org#signaling)\nSignaling\n\nThe `sync.Cond` (conditional variable) type lets one goroutine signal to another that it's ready, and lets the other goroutine wait for that signal.\n\nA `Cond` includes a mutex and has two methods — `Wait` and `Signal`.\n\n- `Wait` unlocks the mutex and suspends the goroutine until it receives a signal.\n- `Signal` wakes the goroutine that is waiting on`Wait` .\n- When `Wait` wakes up, it locks the mutex again.\n\n```\ncond := sync.NewCond(&sync.Mutex{})\ndone := false\nvar wg sync.WaitGroup\nwg.Go(func() {\n    cond.L.Lock()\n    fmt.Println(\"G1 is ready to signal\")\n    done = true\n    cond.Signal()\n    cond.L.Unlock()\n})\nwg.Go(func() {\n    cond.L.Lock()\n    for !done {\n        cond.Wait()\n    }\n    fmt.Println(\"G2 received the signal\")\n    cond.L.Unlock()\n})\nwg.Wait()\n```\n```\nG1 is ready to signal\nG2 received the signal\n```\nIf there are multiple waiting goroutines when `Signal` is called, only one of them will be resumed. If there are no waiting goroutines, `Signal` does nothing.\n\nYou can also use the `Broadcast` method. While `Signal` wakes up only one goroutine waiting on `Cond.Wait`, the `Broadcast` method wakes up all such goroutines.\n\nYou can signal with a channel:\n\n```\nsignal := make(chan struct{}, 1)\ngo func() {\n    // do something\n    signal <- struct{}{}\n}()\ngo func() {\n    <-signal\n    // do something\n}()\n```\nAnd broadcast too:\n\n```\nbroadcast := make(chan struct{})\ngo func() {\n    // do something\n    close(broadcast)\n}()\ngo func() {\n    <-broadcast\n    // do something\n}()\ngo func() {\n    <-broadcast\n    // do something\n}()\n```\nBroadcasting with a condition variable is limited: it only sends a signal, not the actual data, and it only works once. With channels, you can build a publish/subscribe system that doesn't have these limitations:\n\n```\ntype Publisher struct {\n    sbox []chan int // subscription channels\n    mu   sync.Mutex // protects the state\n}\nfunc (p *Publisher) Subscribe() <-chan int {\n    p.mu.Lock()\n    defer p.mu.Unlock()\n    sub := make(chan int, 1)\n    p.sbox = append(p.sbox, sub)\n    return sub\n}\nfunc (p *Publisher) Broadcast(v int) {\n    p.mu.Lock()\n    defer p.mu.Unlock()\n    for _, sub := range p.sbox {\n        select {\n        case sub <- v:\n        default:\n        }\n    }\n}\n```\n## [#](https://antonz.org#run-once)\nRun once\n\nThe `sync.Once` type makes sure that the given function runs only once. If multiple goroutines call `Once.Do` at the same time, only one will run the function, while the others will wait until it returns:\n\n```\ntotal := 0\ninitState := func() {\n    total += 1\n}\nvar once sync.Once\nvar wg sync.WaitGroup\nwg.Go(func() {\n    once.Do(initState)\n    // do something\n})\nwg.Go(func() {\n    once.Do(initState)\n    // do something\n})\nwg.Wait()\n```\n```\ntotal: 1\n```\n`Once` is perfect for one-time initialization or cleanup in a concurrent environment.\n\nBesides the `Once` type, the `sync` package also includes three convenience once-functions:\n\n```\n// Calls f only once.\nfunc (o *Once) Do(f func())\n// Returns a function that calls f only once.\nfunc OnceFunc(f func()) func()\n// Returns a function that calls f only once\n// and returns the value from that first call.\nfunc OnceValue[T any](f func() T) func() T\n// Returns a function that calls f only once\n// and returns the pair of values from that first call.\nfunc OnceValues[T1, T2 any](f func() (T1, T2)) func() (T1, T2)\n```\n## [#](https://antonz.org#object-pool)\nObject pool\n\nThe `sync.Pool` type helps reuse memory instead of allocating it every time, which reduces the load on the garbage collector:\n\n```\npool := sync.Pool{\n    New: func() any {\n        buf := make([]byte, 1024)\n        return &buf\n    },\n}\n// Only allocates 4*1024 B, despite 4000 loop iterations.\nvar wg sync.WaitGroup\nfor range 4 {\n    wg.Go(func() {\n        for range 1000 {\n            buf := pool.Get().(*[]byte)\n            sink = buf\n            pool.Put(buf)\n        }\n    })\n}\nwg.Wait()\n```\n```\nMemory allocated: 4 KB\n```\n`Get` takes an item from the pool. If there are no available items, it creates a new one using `New` (which we have to define ourselves, since the pool doesn't know anything about the items it creates). `Put` returns an item back to the pool.\n\nThings to keep in mind:\n\n- `New` should return a pointer, not a value, to reduce memory copying and avoid extra allocations.\n- The pool has no size limit. If you start 1000 more goroutines that all call `Get` at the same time, 1000 more buffers will be allocated.\n- After an item is returned to the pool with `Put` , you shouldn't use it anymore (since another goroutine might already have taken and started using it).\n\n## [#](https://antonz.org#atomics)\nAtomics\n\nAn operation without synchronization can only be truly atomic if it translates to a single processor instruction. Such operations don't need locks and won't cause issues when called concurrently (even the write operations).\n\nThere are only a few atomics, and they're all found in the `sync/atomic` package:\n\n```\nInt32     Bool\nInt64     Value\nUint32    Pointer\nUint64\n```\nEach atomic type provides the following methods:\n\n- `Load` reads the value of a variable.\n- `Store` sets a new value.\n- `Swap` sets a new value (like`Store` ) and returns the old one.\n- `CompareAndSwap` sets a new value only if the current value is still what you expect it to be.\n\n```\nvar n atomic.Int32\nn.Store(10)\nswapped := n.CompareAndSwap(10, 42)\nfmt.Println(\"CompareAndSwap 10 -> 42:\", swapped)\nfmt.Println(\"n =\", n.Load())\n```\n```\nCompareAndSwap 10 -> 42: true\nn = 42\n```\nNumeric types also provide an `Add` method that increments the value by the specified amount.\n\nAll methods are either translated into a single CPU instruction or are otherwise guaranteed to be atomic, so they are safe to use from multiple goroutines.\n\nThe composition of atomics is always non-atomic:\n\n```\nvar delta atomic.Int32\nvar counter atomic.Int32\nfunc increment() {\n    // Not atomic; causes a race condition.\n    delta.Add(1)\n    sleep(10)\n    counter.Add(delta.Load())\n}\n// After 100 concurrent increments,\n// the final value is NOT guaranteed.\n```\n```\ncounter = 9386\n```\nA bulletproof way to make a composite operation atomic and prevent race conditions is to use a mutex:\n\n```\nvar delta int32\nvar counter int32\nvar mu sync.Mutex\nfunc increment() {\n    // Atomic; doesn't cause a race condition.\n    mu.Lock()\n    delta += 1\n    sleep(10)\n    counter += delta\n    mu.Unlock()\n}\n// After 100 concurrent increments, the final value is guaranteed:\n// counter = 1+2+...+100 = 5050\n```\n```\ncounter = 5050\n```\nSometimes you can use an atomic type instead of a mutex to exit early:\n\n```\ntype Gate struct {\n    closed atomic.Bool\n}\nfunc (g *Gate) Close() {\n    if !g.closed.CompareAndSwap(false, true) {\n        return // ignore repeated calls\n    }\n    // The gate is closed.\n    // We can free resources now.\n}\n```\n## [#](https://antonz.org#testing)\nTesting\n\nIf your concurrent program uses channels or custom types with synchronization methods like `Wait`, you can use those in your tests. This way, your tests won't be much more complicated than if the code were synchronous:\n\n```\n// Calc calculates something asynchronously.\nfunc Calc() <-chan int {\n    out := make(chan int, 1)\n    go func() {\n        out <- 42\n    }()\n    return out\n}\n```\n```\nfunc Test(t *testing.T) {\n    // Wait for the Calc goroutine to finish.\n    got := <-Calc()\n    if got != 42 {\n        t.Errorf(\"got: %v; want: 42\", got)\n    }\n}\n```\n```\nPASS\n```\nIf there aren't any suitable synchronization \"handles\" in the code you're testing, you can use the `synctest` package. It exports two functions:\n\n```\nfunc Test(t *testing.T, f func(*testing.T))\nfunc Wait()\n```\n`synctest.Test` runs an isolated bubble. The bubble uses a fake clock, and you can manually control goroutine synchronization with `synctest.Wait`.\n\n`synctest.Wait` blocks until all goroutines in the bubble — except the one that called `Wait` — have either finished or are durably blocked. This lets you wait for a specific goroutine to finish or get blocked, so you can check the program's state:\n\n```\n// NewProc starts the calculation.\nfunc NewProc() *Proc {\n    p := &Proc{done: make(chan struct{})}\n    go func() {\n        p.res = 42\n        <-p.done // (X)\n        p.res = 0\n    }()\n    return p\n}\n```\n```\nfunc Test(t *testing.T) {\n    synctest.Test(t, func(t *testing.T) {\n        p := NewProc()\n        defer p.Stop()\n        // Wait for the goroutine to block at point X.\n        synctest.Wait()\n        if got := p.Res(); got != 42 {\n            t.Fatalf(\"got %v, want 42\", got)\n        }\n    })\n}\n```\n```\nPASS\n```\nThe fake clock in `synctest.Test` move forward only if: ➊ all goroutines in the bubble are durably blocked; ➋ there's a future moment when at least one goroutine will unblock; and ➌ `synctest.Wait` isn't running. Thanks to this, time-dependent tests run instantly:\n\n```\n// Calc processes a value from the input channel.\n// Times out if no input is received after 3 seconds.\nfunc Calc(in chan int) (int, error) {\n    select {\n    case v := <-in:\n        return v * 2, nil\n    case <-time.After(3 * time.Second):\n        return 0, ErrTimeout\n    }\n}\n```\n```\nfunc Test(t *testing.T) {\n    synctest.Test(t, func(t *testing.T) {\n        ch := make(chan int)\n        got, err := Calc(ch) // runs instantly\n        if err != ErrTimeout {\n            t.Errorf(\"got: %v; want: %v\", err, ErrTimeout)\n        }\n        if got != 0 {\n            t.Errorf(\"got: %v; want: 0\", got)\n        }\n    })\n}\n```\n```\nPASS\n```\nThe following operations durably block a goroutine:\n\n- A blocking send or receive on a channel created within the bubble.\n- A blocking select statement where every case is a channel created within the bubble.\n- Calling `Cond.Wait` .\n- Calling `WaitGroup.Wait` if all`WaitGroup.Add` calls were made inside the bubble.\n- Calling `time.Sleep` .\n\nBlocking on mutexes, I/O, or system calls is not considered durable, and the `synctest` bubble can't handle them.\n\n## [#](https://antonz.org#scheduling)\nScheduling\n\nAt the hardware level, CPU cores are responsible for running parallel tasks.\n\nAt the operating system level, a thread is the basic unit of execution. There are usually many more threads than CPU cores, so the operating system's scheduler decides which threads to run and which ones to pause.\n\nAt the Go runtime level, a goroutine is the basic unit of execution. The runtime scheduler runs a fixed number of OS threads, often one per CPU core. There can be many more goroutines than threads, so the scheduler decides which goroutines to run on the available threads and which ones to pause. The scheduler keeps switching between goroutines to make sure each one gets a turn to run on a thread, instead of waiting in line forever.\n\n```\n  CPU                  OS                   Go runtime\n┌──────────┐  run on ┌──────────┐  run on ┌────────────┐\n│ Cores    │ <────── │ Threads  │ <────── │ Goroutines │\n└──────────┘         └──────────┘         └────────────┘\n```\nThis is how Go handles concurrency.\n\n**Goroutine scheduler**\n\nThe goroutine scheduler's job is to run M goroutines on N operating system threads, where M can be much larger than N. Here's a very simplified version of it's algorithm:\n\n- If there's a free thread, assign it a goroutine from the queue.\n- If a running goroutine gets blocked (for example, while reading from a channel), put it back in the queue and assign a different goroutine to the thread.\n- If a running goroutine gets stuck in a syscall, start a new thread to run other goroutines until the blocked goroutine finishes the syscall.\n- Check the running goroutines every 10 ms. Preempt long-running goroutines and return them to the queue to prevent starvation.\n\n```\n┌─────┐┌─────┐┌─────┐┌─────┐\n│ G17 ││ G18 ││ G19 ││ G20 │                        queue\n└─────┘└─────┘└─────┘└─────┘\n┌─────┐      ┌─────┐      ┌─────┐      ┌─────┐\n│ G15 │      │ G16 │      │ G13 │      │ G14 │      running\n└─────┘      └─────┘      └─────┘      └─────┘\n  │            │            │            │\n┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐\n│ Thread E │ │ Thread F │ │ Thread C │ │ Thread D │\n└──────────┘ └──────────┘ └──────────┘ └──────────┘\n┌─────┐      ┌─────┐\n│ G11 │      │ G12 │                                syscalls\n└─────┘      └─────┘\n  │            │\n┌──────────┐ ┌──────────┐\n│ Thread A │ │ Thread B │\n└──────────┘ └──────────┘\n```\nThe number of threads running Go code is controlled by the `GOMAXPROCS` environment variable or the `runtime.GOMAXPROCS` function.\n\nA goroutine is a structure that starts out using about 2 KB of memory, mostly for its stack. The stack can grow if needed. Since goroutines are so lightweight, you can run tens of thousands or even hundreds of thousands of them on a small machine.\n\n## [#](https://antonz.org#diagnostics)\nDiagnostics\n\nTo troubleshoot concurrent programs in production, we use metrics, profiling, and tracing.\n\n**Metrics** show how the Go runtime is performing, like how much heap memory it uses or how long garbage collection pauses take. Each metric has a unique name and a value, which can be a number or a histogram.\n\nYou can use the `runtime/metrics` package to get a complete list of metrics or check the values of specific ones:\n\n```\nsamples := []metrics.Sample{\n    {Name: \"/sched/gomaxprocs:threads\"},\n    {Name: \"/sched/goroutines:goroutines\"},\n}\nmetrics.Read(samples)\nfor _, s := range samples {\n    fmt.Printf(\"%s: %v\\n\", s.Name, s.Value.Uint64())\n}\n```\n```\n/sched/gomaxprocs:threads: 8\n/sched/goroutines:goroutines: 1\n```\nIn practice, people rarely do this manually. Instead, all metrics are automatically exported using Prometheus or OpenTelemetry libraries.\n\n**Profiling** helps you understand exactly what the program is doing, what resources it uses, and where in the code this happens. Go uses a sampling profiler that's suitable for production.\n\nThe most commonly used profiles are CPU, which shows how much processor time each function uses, and heap, which shows how much heap memory each function uses. Goroutine, block, and mutex profiles help identify problems related to concurrency.\n\nThe easiest way to add a profiler to your app is by using the `net/http/pprof` package. To collect a profile with the given name, call the `/debug/pprof/{name}` endpoint. To view the collected profile, use the `go tool pprof` utility:\n\n```\ngo tool pprof -proto \\\n  \"http://localhost:6060/debug/pprof/profile?seconds=N\" > cpu.pprof\ngo tool pprof -http=localhost:8080 cpu.pprof\n```\nYou can also profile manually:\n\n```\n// CPU profile.\nfile, _ := os.Create(\"cpu.prof\")\ndefer file.Close()\npprof.StartCPUProfile(file)\ndefer pprof.StopCPUProfile()\n// ...\n```\n```\n// Any other profile.\nfile, _ := os.Create(name + \".prof\")\ndefer file.Close()\npprof.Lookup(name).WriteTo(file, 0)\n```\n**Tracing** records certain types of events while the program is running, mainly those related to concurrency and memory. When the profiling server from the `net/http/pprof` package is running, call the `/debug/pprof/trace` endpoint to collect a trace. To view the results, use the `go tool trace` utility.\n\nYou can also collect a trace manually:\n\n```\nfile, _ := os.Create(\"trace.out\")\ndefer file.Close()\ntrace.Start(file)\ndefer trace.Stop()\n// ...\n```\nYou can set up automatic tracing with a sliding window that's limited by size or duration. This is called \"flight recording\". It lets you always keep a recent trace available in case something goes wrong:\n\n```\ncfg := trace.FlightRecorderConfig{\n    MinAge:   5 * time.Second,\n    MaxBytes: 3 << 20, // 3MB\n}\nrec := trace.NewFlightRecorder(cfg)\nrec.Start()\ndefer rec.Stop()\n```\n## [#](https://antonz.org#final-thoughts)\nFinal thoughts\n\nWe've covered a number of Go tools for writing concurrent programs:\n\n- Goroutines for running concurrent tasks.\n- Channels and select as flexible communication tools.\n- Timers and tickers for working with time.\n- Context for canceling operations.\n- Wait groups for synchronizing goroutines.\n- Mutexes to prevent race conditions.\n- Condition variables for signaling events.\n- Once for safe one-time initialization.\n- Pools to reduce garbage collector load.\n- Atomic operations.\n\nIf you like the book, please recommend it to your friends or colleagues. If you're interested, check out my other [books](https://antonz.org/#books) and [projects](https://antonz.org/tags/projects/).\n\nI'm glad you finished the book. Thank you, and I'll see you next time!\n\n[★ Subscribe](https://antonz.org/subscribe/) to keep up with new posts.","body_html":"<h1 id=\"go-concurrency-distilled\">Go concurrency distilled</h1>\n<p>This mini-book provides a brief overview of many concurrency topics in Go. Each topic comes with interactive examples — feel free to experiment with them by changing the code and clicking <em>Run</em>. There&#39;s also a <a href=\"https://github.com/nalgeon/go-conc-distilled\" rel=\"nofollow ugc noopener\">PDF version</a> with static examples.</p>\n<p>This is a quick refresher on Go concurrency, not a beginner&#39;s guide. If you want to learn concurrency from the ground up with practical exercises, check out my other book — <a href=\"https://antonz.org/go-concurrency\" rel=\"nofollow ugc noopener\">Gist of Go: Concurrency</a>.</p>\n<p>The book is AI-free.</p>\n<p><a href=\"https://antonz.org#goroutines\" rel=\"nofollow ugc noopener\">Goroutines</a> •\n<a href=\"https://antonz.org#channels\" rel=\"nofollow ugc noopener\">Channels</a> •\n<a href=\"https://antonz.org#select\" rel=\"nofollow ugc noopener\">Select</a> •\n<a href=\"https://antonz.org#pipelines\" rel=\"nofollow ugc noopener\">Pipelines</a> •\n<a href=\"https://antonz.org#time\" rel=\"nofollow ugc noopener\">Time</a> •\n<a href=\"https://antonz.org#context\" rel=\"nofollow ugc noopener\">Context</a> •\n<a href=\"https://antonz.org#wait-groups\" rel=\"nofollow ugc noopener\">Wait groups</a> •\n<a href=\"https://antonz.org#data-races\" rel=\"nofollow ugc noopener\">Data races</a> •\n<a href=\"https://antonz.org#race-conditions\" rel=\"nofollow ugc noopener\">Race conditions</a> •\n<a href=\"https://antonz.org#mutexes\" rel=\"nofollow ugc noopener\">Mutexes</a> •\n<a href=\"https://antonz.org#semaphores\" rel=\"nofollow ugc noopener\">Semaphores</a> •\n<a href=\"https://antonz.org#signaling\" rel=\"nofollow ugc noopener\">Signaling</a> •\n<a href=\"https://antonz.org#run-once\" rel=\"nofollow ugc noopener\">Run once</a> •\n<a href=\"https://antonz.org#object-pool\" rel=\"nofollow ugc noopener\">Object pool</a> •\n<a href=\"https://antonz.org#atomics\" rel=\"nofollow ugc noopener\">Atomics</a> •\n<a href=\"https://antonz.org#testing\" rel=\"nofollow ugc noopener\">Testing</a> •\n<a href=\"https://antonz.org#scheduling\" rel=\"nofollow ugc noopener\">Scheduling</a> •\n<a href=\"https://antonz.org#diagnostics\" rel=\"nofollow ugc noopener\">Diagnostics</a> •\n<a href=\"https://antonz.org#final-thoughts\" rel=\"nofollow ugc noopener\">Final thoughts</a></p>\n<h2 id=\"section\"><a href=\"https://antonz.org#goroutines\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Goroutines</p>\n<p>The foundation of concurrency in Go is <em>goroutines</em> – functions started with the <code>go</code> keyword:</p>\n<pre><code>func main() {\n    var wg sync.WaitGroup\n    wg.Add(2)\n    go func() {\n        defer wg.Done()\n        fmt.Println(&quot;worker 1&quot;)\n    }()\n    go func() {\n        defer wg.Done()\n        fmt.Println(&quot;worker 2&quot;)\n    }()\n    wg.Wait()\n}</code></pre>\n<pre><code>worker 2\nworker 1</code></pre>\n<p>The Go runtime juggles these goroutines and distributes them among operating system threads running on CPU cores. Compared to OS threads, goroutines are lightweight, so you can create hundreds or thousands of them.</p>\n<p>Goroutines are completely independent. The main function is also a goroutine, but it starts implicitly when the program starts. When <code>main</code> ends, other goroutines also shut down.</p>\n<p>We use a <em>wait group</em> (<code>sync.WaitGroup</code>) to wait for goroutines to finish in the example above. A wait group has a counter inside. Calling <code>Add(n)</code> increments it by <code>n</code>, while <code>Done()</code> decrements it by one. <code>Wait()</code> blocks the calling goroutine (in this case, main) until the counter reaches zero. This way, main waits for both workers to finish before it exits.</p>\n<p><code>WaitGroup.Go</code> automatically increments the wait group counter, runs a function in a goroutine, and decrements the counter when it&#39;s done:</p>\n<pre><code>func main() {\n    var wg sync.WaitGroup\n    wg.Go(func() {\n        fmt.Println(&quot;worker 1&quot;)\n    })\n    wg.Go(func() {\n        fmt.Println(&quot;worker 2&quot;)\n    })\n    wg.Wait()\n}</code></pre>\n<pre><code>worker 2\nworker 1</code></pre>\n<h2 id=\"section-2\"><a href=\"https://antonz.org#channels\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Channels</p>\n<p>Goroutines can pass values to each other through <em>channels</em>. A channel is like a window where one goroutine can throw something and another can catch it:</p>\n<pre><code>func main() {\n    messages := make(chan string)\n    go func() { messages &lt;- &quot;ping&quot; }()\n    msg := &lt;-messages\n    fmt.Println(msg)\n}</code></pre>\n<pre><code>ping</code></pre>\n<p>Sending a value through a channel is a synchronous operation. When the sending goroutine writes a value to the channel (<code>ch &lt;- val</code>), it blocks and waits for someone to receive that value (<code>&lt;-ch</code>). Only then does it continue.</p>\n<h3 id=\"output-channel\">Output channel</h3>\n<p>Returning an output channel from a function and filling it within an internal goroutine is a common pattern in Go. This allows the caller to receive values through the channel while the owning function retains control of it:</p>\n<pre><code>func generate(start, stop int) chan int {\n    out := make(chan int)\n    go func() {\n        for i := start; i &lt; stop; i++ {\n            out &lt;- i\n        }\n    }()\n    return out\n}</code></pre>\n<h3 id=\"closing-a-channel\">Closing a channel</h3>\n<p>To signal readers that all data has been sent, the writer goroutine <em>closes</em> the channel with <code>close()</code>:</p>\n<pre><code>func generate(start, stop int) chan int {\n    out := make(chan int)\n    go func() {\n        defer close(out)\n        for i := start; i &lt; stop; i++ {\n            out &lt;- i\n        }\n    }()\n    return out\n}</code></pre>\n<p>The reader checks the channel&#39;s status with a second value (&quot;comma OK&quot;) when reading:</p>\n<pre><code>func main() {\n    in := generate(5, 10)\n    for {\n        num, ok := &lt;-in\n        if !ok {\n            break\n        }\n        fmt.Print(num, &quot; &quot;)\n    }\n}</code></pre>\n<pre><code>5 6 7 8 9</code></pre>\n<p>While the channel is open, the reader receives the next value and a <code>true</code> status. If the channel is closed, the reader gets a zero value and a <code>false</code> status.</p>\n<p>A channel can only be closed once. Closing it again or writing to a closed channel causes a panic.</p>\n<p>The only reason to close a channel is to signal to its readers that all data has been sent. If this isn&#39;t important to the readers, then you don&#39;t need to close it. When a channel is no longer used, Go&#39;s garbage collector will free its resources, whether it&#39;s closed or not.</p>\n<h3 id=\"channel-iteration\">Channel iteration</h3>\n<p><code>range</code> automatically reads the next value from the channel and checks if it&#39;s closed. If the channel is closed, it exits the loop:</p>\n<pre><code>func main() {\n    nums := generate(5, 10)\n    for n := range nums {\n        fmt.Print(n, &quot; &quot;)\n    }\n}</code></pre>\n<pre><code>5 6 7 8 9</code></pre>\n<p>Range over a channel returns a single value, not a pair, unlike range over a slice.</p>\n<h3 id=\"directional-channels\">Directional channels</h3>\n<p>You can protect yourself from accidental write/close errors by setting the channel direction. Channels can be:</p>\n<ul><li><code>chan</code> (bidirectional): for reading and writing (default);</li><li><code>chan&lt;-</code> (send-only): for writing only;</li><li><code>&lt;-chan</code> (receive-only): for reading only.</li></ul>\n<p>You can&#39;t read from a send-only channel or write to a receive-only channel (nor can you close it).</p>\n<p>Channels are usually initialized for both reading and writing, and specified as directional in function parameters. Go automatically converts a regular channel to a directional one:</p>\n<pre><code>stream := make(chan int)\ngo func(in chan&lt;- int) {\n    in &lt;- 42\n}(stream)\nfunc(out &lt;-chan int) {\n    fmt.Println(&lt;-out)\n}(stream)</code></pre>\n<pre><code>42</code></pre>\n<h3 id=\"buffered-channels\">Buffered channels</h3>\n<p><em>Buffered</em> channels work like a FIFO queue with a fixed-size buffer for storing values.</p>\n<p>As long as the buffer has free space, writing to the channel doesn&#39;t block the goroutine. Similarly, as long as the buffer contains values, reading from the channel doesn&#39;t block the goroutine:</p>\n<pre><code>stream := make(chan int, 3)\nstream &lt;- 11\nstream &lt;- 12\nstream &lt;- 13\nfmt.Println(&lt;-stream)\nfmt.Println(&lt;-stream)</code></pre>\n<pre><code>11\n13</code></pre>\n<p>By default, if you don&#39;t specify a buffer size, a channel is <em>unbuffered</em> (buffer size equals zero).</p>\n<p>Buffered channels work with the built-in <code>len()</code> and <code>cap()</code> functions:</p>\n<pre><code>stream := make(chan int, 3)\nstream &lt;- 11\nfmt.Println(cap(stream), len(stream))</code></pre>\n<pre><code>3 1</code></pre>\n<p>Reading from a closed buffered channel returns values from the buffer and a <code>true</code> status. Once all values are taken, it returns a zero value and a <code>false</code> status, like a regular channel:</p>\n<pre><code>stream := make(chan int, 1)\nstream &lt;- 11\nclose(stream)\nval, ok := &lt;-stream\nfmt.Println(val, ok)\n// 11 true\nval, ok = &lt;-stream\nfmt.Println(val, ok)\n// 0 false</code></pre>\n<pre><code>11 true\n0 false</code></pre>\n<h3 id=\"nil-channel\">nil channel</h3>\n<p>Like any type in Go, channels have a zero value, which is <code>nil</code>.</p>\n<p>Writing to or reading from a nil channel blocks the goroutine indefinitely:</p>\n<pre><code>var stream chan int\ngo func() {\n    // blocks forever\n    stream &lt;- 1\n}()\n// blocks forever\n&lt;-stream</code></pre>\n<p>Closing a nil channel causes a panic:</p>\n<pre><code>var stream chan int\nclose(stream)\n// panic: close of nil channel</code></pre>\n<h2 id=\"section-3\"><a href=\"https://antonz.org#select\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Select</p>\n<p>The <em>select</em> statement is somewhat like <code>switch</code>, but specifically designed for channels. Here&#39;s what it does:</p>\n<ul><li>Checks which cases are not blocked.</li><li>If multiple cases are ready, randomly selects one to execute.</li><li>If all cases are blocked and there is a default case, executes it.</li><li>If all cases are blocked and there is no default case, waits until one is ready.</li></ul>\n<p>Select is used to manage data flow in pipelines:</p>\n<pre><code>// merge sends values from in1 and in2 to the output channel.\nfunc merge(in1, in2 &lt;-chan int) &lt;-chan int {\n    out := make(chan int)\n    go func() {\n        defer close(out)\n        for in1 != nil || in2 != nil {\n            select {\n            case val1, ok := &lt;-in1:\n                if ok { out &lt;- val1 } else { in1 = nil }\n            case val2, ok := &lt;-in2:\n                if ok { out &lt;- val2 } else { in2 = nil }\n            }\n        }\n    }()\n    return out\n}\n// Suppose we send 10..12 to in1, 20..22 to in2,\n// and call merge(in1, in2)</code></pre>\n<pre><code>10 11 20 12 21 22</code></pre>\n<p>To cancel goroutines:</p>\n<pre><code>// process modifies values from in and send them to out\n// until in is exhausted or cancel is closed.\nfunc process(cancel chan struct{}, in &lt;-chan int) &lt;-chan int {\n    out := make(chan int)\n    go func() {\n        for val := range in {\n            select {\n            case out &lt;- val*10:\n            case &lt;-cancel:\n                fmt.Println(&quot;canceled&quot;)\n                return\n            }\n        }\n    }()\n    return out\n}\n// Suppose we send values 11 and 12 to in\n// and then call close(cancel)</code></pre>\n<pre><code>110\n120\ncanceled</code></pre>\n<p>For non-blocking operations:</p>\n<pre><code>// multiplier returns a function that multiplies\n// the input by 10 and sends it to the channel\n// or returns an error if the channel is busy.\nfunc multiplier(ch chan&lt;- int) func(n int) error {\n    return func(n int) error {\n        select {\n        case ch &lt;- n*10:\n            return nil\n        default:\n            return errors.New(&quot;busy&quot;)\n        }\n    }\n}\nfunc main() {\n    nums := make(chan int, 1)\n    multiply := multiplier(nums)\n    err := multiply(11)\n    fmt.Println(&lt;-nums, err)\n    // 110 &lt;nil&gt;\n    err = multiply(12)\n    fmt.Println(&lt;-nums, err)\n    // 120 &lt;nil&gt;\n    err = multiply(13)\n    err = multiply(14)\n    fmt.Println(err)\n    // busy\n}</code></pre>\n<pre><code>110 &lt;nil&gt;\n120 &lt;nil&gt;\nbusy</code></pre>\n<p>And for much more.</p>\n<h2 id=\"section-4\"><a href=\"https://antonz.org#pipelines\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Pipelines</p>\n<p>A <em>pipeline</em> is a sequence of operations where each step takes input data, processes it in a specific way, and outputs it. The input and output of each operation is a channel.</p>\n<p>A typical pipeline looks like this:</p>\n<ul><li><em>Reader</em> : Reads input data from a file, database, or network.</li><li><em>N processors</em> : Transform, filter, aggregate, or enrich data using external sources.</li><li><em>Writer</em> : Writes the processed data to a file, database, or network.</li></ul>\n<pre><code>func read[T any]() &lt;-chan T {\n    out := make(chan T)\n    go func() {\n        defer close(out)\n        for {\n            // read data from somewere\n            data := // ...\n            out &lt;- data\n        }\n    }()\n    return out\n}\nfunc process[T any](in &lt;-chan T) &lt;-chan T {\n    out := make(chan T)\n    go func() {\n        defer close(out)\n        for inData := range in {\n            // process the data\n            outData = // ...\n            out &lt;- outData\n        }\n    }()\n    return out\n}\nfunc write[T any](in &lt;-chan T) &lt;-chan struct{} {\n    done := make(chan struct{})\n    go func() {\n        defer close(done)\n        for data := range in {\n            // write the data\n        }\n    }()\n    return done\n}</code></pre>\n<h3 id=\"output-channel-2\">Output channel</h3>\n<p>A goroutine can signal other goroutines that it has finished its work using an <em>output channel</em>:</p>\n<pre><code>func generate(start, stop int) &lt;-chan int {\n    out := make(chan int)\n    go func() {\n        defer close(out)\n        for i := start; i &lt; stop; i++ {\n            out &lt;- i\n        }\n    }()\n    return out\n}\nfunc main() {\n    nums := generate(5, 10)\n    for n := range nums {\n        fmt.Print(n, &quot; &quot;)\n    }\n}</code></pre>\n<pre><code>5 6 7 8 9</code></pre>\n<h3 id=\"done-channel\">Done channel</h3>\n<p>If a goroutine doesn&#39;t need to return results, it can signal completion using a <em>done channel</em>:</p>\n<pre><code>func work() &lt;-chan struct{} {\n    done := make(chan struct{})\n    go func() {\n        defer close(done)\n        fmt.Println(&quot;work done&quot;)\n    }()\n    return done\n}\nfunc main() {\n    done := work()\n    &lt;-done\n}</code></pre>\n<pre><code>work done</code></pre>\n<h3 id=\"cancel-channel\">Cancel channel</h3>\n<p>To terminate a goroutine early, a calling goroutine can use a <em>cancel channel</em>:</p>\n<pre><code>func generate(cancel chan struct{}, n int) &lt;-chan int {\n    out := make(chan int)\n    go func() {\n        defer close(out)\n        for i := 1; i &lt;= n; i++ {\n            select {\n            case out &lt;- i:\n            case &lt;-cancel:\n                return\n            }\n        }\n    }()\n    return out\n}\nfunc main() {\n    cancel := make(chan struct{})\n    defer close(cancel)\n    nums := generate(cancel, 10)\n    fmt.Println(&lt;-nums)\n    fmt.Println(&lt;-nums)\n    fmt.Println(&lt;-nums)\n}</code></pre>\n<pre><code>1\n2\n3</code></pre>\n<h3 id=\"error-handling\">Error handling</h3>\n<p>There are three approaches to error handling in concurrent pipelines.</p>\n<p>➊ Return on the first error:</p>\n<pre><code>// calculate produces answers for the given numbers.\nfunc process(in &lt;-chan int) (&lt;-chan int, &lt;-chan error) {\n    out := make(chan Answer)\n    errc := make(chan error, 1)\n    go func() {\n        defer close(out)\n        for n := range in {\n            ans, err := fetchAnswer(n)\n            if err != nil {\n                errc &lt;- err  // return with error\n                return\n            }\n            out &lt;- ans\n        }\n        errc &lt;- nil          // return with nil\n    }()\n    return out, errc\n}</code></pre>\n<p>➋ Use a result type:</p>\n<pre><code>// Result contains an answer or an error.\ntype Result struct {\n    answer int\n    err    error\n}\n// calculate produces answers for the given numbers.\nfunc calculate(in &lt;-chan int) &lt;-chan Result {\n    out := make(chan Result)\n    go func() {\n        defer close(out)\n        for n := range in {\n            ans, err := fetchAnswer(n)\n            out &lt;- Result{ans, err}  // return answer + error\n        }\n    }()\n    return out\n}</code></pre>\n<p>➌ Collect errors separately:</p>\n<pre><code>// calculate produces answers for the given numbers.\nfunc calculate(in &lt;-chan int, errc chan&lt;- error) &lt;-chan int {\n    out := make(chan Answer)\n    go func() {\n        defer close(out)\n        for n := range in {\n            ans, err := fetchAnswer(n)\n            if err == nil {\n                out &lt;- ans   // send answer\n            } else {\n                errc &lt;- err  // or error\n            }\n        }\n    }()\n    return out\n}</code></pre>\n<h2 id=\"section-5\"><a href=\"https://antonz.org#time\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Time</p>\n<p>Besides handling date and time, the <code>time</code> package offers tools for managing time-sensitive operations in concurrent programs.</p>\n<h3 id=\"after\">After</h3>\n<p><code>time.After()</code> returns a channel that is initially empty, but receives a value after the timeout period. It&#39;s useful for timing out operations:</p>\n<pre><code>// withTimeout executes a function with a given timeout.\nfunc withTimeout(timeout time.Duration, fn func()) error {\n    done := make(chan struct{})\n    go func() {\n        defer close(done)\n        fn()\n    }()\n    // blocks until fn completes or the timer expires,\n    // whichever happens first\n    select {\n    case &lt;-done:\n        return nil\n    case &lt;-time.After(timeout):\n        return errors.New(&quot;timeout&quot;)\n    }\n}</code></pre>\n<p><code>withTimeout()</code> waits for <code>fn()</code> to complete, but thanks to <code>time.After()</code>, it won&#39;t wait longer than the <code>timeout</code> duration:</p>\n<pre><code>func main() {\n    var err error\n    // completes in time\n    err = withTimeout(\n        50*time.Millisecond,\n        func() { fmt.Println(&quot;work done&quot;) },\n    )\n    fmt.Println(&quot;err =&quot;, err)\n    // gets canceled on timeout\n    err = withTimeout(\n        50*time.Millisecond,\n        func() {\n            time.Sleep(100 * time.Millisecond)\n            fmt.Println(&quot;work done&quot;)\n        },\n    )\n    fmt.Println(&quot;err =&quot;, err)\n}</code></pre>\n<pre><code>work done\nerr = &lt;nil&gt;\nerr = timeout</code></pre>\n<h3 id=\"timer\">Timer</h3>\n<p>A <em>timer</em> (<code>time.Timer</code>) is a structure with a <code>C</code> channel to which it sends the current time when it triggers (expires). Timers are useful for planning future executions:</p>\n<pre><code>done := make(chan struct{})\ntimer := time.NewTimer(50 * time.Millisecond)\ngo func() {\n    eventTime := &lt;-timer.C  // blocks for 50ms\n    fmt.Println(&quot;work done at&quot;, eventTime)\n    close(done)\n}()\n&lt;-done</code></pre>\n<pre><code>work done at 2009-11-10 23:00:00.05</code></pre>\n<p><code>Stop()</code> stops the timer and returns <code>true</code> if it hasn&#39;t expired yet, and <code>false</code> otherwise:</p>\n<pre><code>// timer expires after 50ms\ntimer := time.NewTimer(50 * time.Millisecond)\ngo func() {\n    eventTime := &lt;-timer.C\n    fmt.Println(&quot;work done at&quot;, eventTime)\n}()\n// after 10ms, the timer hasn&#39;t expired yet\ntime.Sleep(10 * time.Millisecond)\nif timer.Stop() {\n    fmt.Println(&quot;execution canceled&quot;)\n} else {\n    fmt.Println(&quot;too late to cancel&quot;)\n}</code></pre>\n<pre><code>execution canceled</code></pre>\n<p>It&#39;s often more convenient to use the <code>time.AfterFunc()</code> wrapper function. It waits for duration <code>d</code> and then executes function <code>f</code>:</p>\n<pre><code>done := make(chan struct{})\nwork := func() {\n    fmt.Println(&quot;work done&quot;)\n    close(done)\n}\n// executes work after 50ms\ntime.AfterFunc(50*time.Millisecond, work)\n&lt;-done</code></pre>\n<pre><code>work done</code></pre>\n<p><code>time.AfterFunc()</code> returns a timer that you can cancel before execution starts:</p>\n<pre><code>// executes the function after 50ms\ntimer := time.AfterFunc(50*time.Millisecond, func() {})\n// after 10ms, the timer hasn&#39;t expired yet\ntime.Sleep(10 * time.Millisecond)\nif timer.Stop() {\n    fmt.Println(&quot;execution canceled&quot;)\n}</code></pre>\n<pre><code>execution canceled</code></pre>\n<p>If a timer is used in a loop, it&#39;s better to create a single timer and <em>reset</em> it instead of creating a new instance on each iteration:</p>\n<pre><code>// consumer reads tokens from the input channel and alerts\n// if a value does not appear in a channel after an hour.\nfunc consumer(in &lt;-chan token) {\n    const timeout = time.Hour\n    timer := time.NewTimer(timeout)\n    for {\n        timer.Reset(timeout)\n        select {\n        case &lt;-in:\n            // do stuff\n        case &lt;-timer.C:\n            // log warning\n        }\n    }\n}\n// Suppose we send 10,000 values to the in channel\n// and measure memory usage.</code></pre>\n<pre><code>Memory used: 4 KB, # allocations: 6</code></pre>\n<h3 id=\"ticker\">Ticker</h3>\n<p>A <em>ticker</em> is like a timer, but it keeps firing until you stop it. Tickers are useful for executing periodic tasks:</p>\n<pre><code>// fires every 50ms\nticker := time.NewTicker(50 * time.Millisecond)\ndefer ticker.Stop()\ngo func() {\n    for {\n        // waits for ticker to fire on each iteration\n        at := &lt;-ticker.C\n        fmt.Println(&quot;work done at&quot;, at)\n    }\n}()\n// enough time for the ticker to fire 3 times\ntime.Sleep(160*time.Millisecond)\nticker.Stop()</code></pre>\n<pre><code>work done at 2009-11-10 23:00:00.05\nwork done at 2009-11-10 23:00:00.10\nwork done at 2009-11-10 23:00:00.15</code></pre>\n<p><code>NewTicker(d)</code> creates a ticker that sends the current time to the channel <code>C</code> at interval <code>d</code>. You must stop the ticker eventually with <code>Stop()</code> to free up resources.</p>\n<p>If the channel reader can&#39;t keep up with the ticker, the ticker will skip ticks.</p>\n<h2 id=\"section-6\"><a href=\"https://antonz.org#context\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Context</p>\n<p>The main purpose of <em>context</em> is to cancel operations, either manually or by timeout/deadline.</p>\n<p>The function accepts a context and uses its <code>Done()</code> channel to listen for cancellation:</p>\n<pre><code>// work performs a task for 50 ms unless canceled.\n// Returns an error when canceled.\nfunc work(ctx context.Context) error {\n    done := make(chan struct{})\n    go func() {\n        time.Sleep(50 * time.Millisecond)\n        fmt.Println(&quot;work done&quot;)\n        close(done)\n    }()\n    select {\n    case &lt;-done:\n        return nil\n    case &lt;-ctx.Done():\n        return ctx.Err()\n    }\n}</code></pre>\n<p>Cancel manually (<code>context.Canceled</code> error):</p>\n<pre><code>func main() {\n    // empty context\n    ctx := context.Background()\n    // manual canellation context\n    ctx, cancel := context.WithCancel(ctx)\n    defer cancel()\n    done := make(chan struct{})\n    go func() {\n        // takes 50 ms unless canceled\n        err := work(ctx)\n        fmt.Println(&quot;err =&quot;, err)\n        close(done)\n    }()\n    // cancels after 10 ms\n    time.Sleep(10 * time.Millisecond)\n    cancel()\n    &lt;-done\n}</code></pre>\n<pre><code>err = context canceled</code></pre>\n<p>Cancel by timeout (<code>context.DeadlineExceeded</code> error):</p>\n<pre><code>func main() {\n    ctx := context.Background()\n    // cancels after 10 ms\n    ctx, cancel := context.WithTimeout(ctx, 10*time.Millisecond)\n    defer cancel()\n    done := make(chan struct{})\n    go func() {\n        // takes 50 ms unless canceled\n        err := work(ctx)\n        fmt.Println(&quot;err =&quot;, err)\n        close(done)\n    }()\n    &lt;-done\n}</code></pre>\n<pre><code>err = context deadline exceeded</code></pre>\n<p>Cancel by deadline (<code>context.DeadlineExceeded</code> error):</p>\n<pre><code>func main() {\n    ctx := context.Background()\n    // cancels at now + 10 ms\n    deadline := time.Now().Add(10 * time.Millisecond)\n    ctx, cancel := context.WithDeadline(ctx, deadline)\n    defer cancel()\n    done := make(chan struct{})\n    go func() {\n        // takes 50 ms unless canceled\n        err := work(ctx)\n        fmt.Println(&quot;err =&quot;, err)\n        close(done)\n    }()\n    &lt;-done\n}</code></pre>\n<pre><code>err = context deadline exceeded</code></pre>\n<p>Context is layered. A context object is immutable. To add new properties to a context, a new (child) context is created based on the old (parent) context. The shorter timeout between the parent and child contexts always wins. The child context can only shorten the parent&#39;s timeout, not extend it:</p>\n<pre><code>func main() {\n    // parent context with a 100 ms timeout\n    const dur100ms = 100 * time.Millisecond\n    parentCtx, cancel := context.WithTimeout(context.Background(), dur100ms)\n    defer cancel()\n    // child context with a 10 ms timeout\n    const dur10ms = 10 * time.Millisecond\n    childCtx, cancel := context.WithTimeout(parentCtx, dur10ms)\n    defer cancel()\n    // now the work gets canceled\n    err := work(childCtx)\n    fmt.Println(&quot;err =&quot;, err)\n}</code></pre>\n<pre><code>err = context deadline exceeded</code></pre>\n<p>Multiple cancels are safe. You can call <code>cancel()</code> on the context as many times as you want. The first cancel will work, and the rest will be ignored.</p>\n<p>You can specify a custom cancellation cause using <code>context.WithCancelCause()</code>, <code>context.WithTimeoutCause()</code> and <code>context.WithDeadlineCause()</code>. This cause is accessible through <code>context.Cause()</code>:</p>\n<pre><code>ctx, cancel := context.WithCancelCause(context.Background())\ncancel(errors.New(&quot;the night is dark&quot;))\nfmt.Println(context.Cause(ctx))</code></pre>\n<pre><code>the night is dark</code></pre>\n<p>You can register a function to execute when the context is canceled with <code>context.AfterFunc()</code>:</p>\n<pre><code>ctx, cancel := context.WithCancel(context.Background())\ncleanup := func() { fmt.Println(&quot;cleanup&quot;) }\ncontext.AfterFunc(ctx, cleanup)\ncancel()\ntime.Sleep(10 * time.Millisecond)</code></pre>\n<pre><code>cleanup</code></pre>\n<p>Context can pass additional information about a call using <code>context.WithValue()</code>, which creates a context with a value for a specific key. But it&#39;s generally better to avoid passing values in context. It&#39;s better to use explicit parameters or custom structs instead.</p>\n<h2 id=\"section-7\"><a href=\"https://antonz.org#wait-groups\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Wait groups</p>\n<p>The <code>sync.WaitGroup</code> type lets you wait for one or more goroutines to finish:</p>\n<pre><code>const n = 10\nvar wg sync.WaitGroup\nwg.Add(n)\nfor range n {\n    go func() {\n        defer wg.Done()\n        fmt.Print(&quot;.&quot;)\n    }()\n}\nwg.Wait()</code></pre>\n<pre><code>..........</code></pre>\n<p>A <code>WaitGroup</code> doesn&#39;t know anything about the goroutines it manages. It works with an internal counter. Calling <code>wg.Add(1)</code> increments the counter by one, while <code>wg.Done()</code> decrements it. <code>wg.Wait()</code> blocks the calling goroutine until the counter reaches zero.</p>\n<p>The <code>Go</code> method combines <code>Add</code>, starting a goroutine, and <code>Done</code>:</p>\n<pre><code>var wg sync.WaitGroup\nfor range 10 {\n    wg.Go(func() {\n        fmt.Print(&quot;.&quot;)\n    })\n}\nwg.Wait()</code></pre>\n<pre><code>..........</code></pre>\n<p>All methods are safe to use from multiple goroutines.</p>\n<p>Normally, all <code>Add</code> calls happen before <code>Wait</code>. But technically, there&#39;s nothing stopping you from doing some of the <code>Add</code> calls before <code>Wait</code> and some after (from another goroutine).</p>\n<p>You can call <code>Wait</code> from multiple goroutines. They will all block until the group&#39;s counter reaches zero.</p>\n<h2 id=\"section-8\"><a href=\"https://antonz.org#data-races\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Data races</p>\n<p>A data race happens when multiple goroutines access shared data, and at least one of them modifies it. We need to protect the data from this kind of concurrent access.</p>\n<p>A data race doesn&#39;t always cause a runtime panic. That&#39;s why Go provides a special tool called the race detector. You can turn it on with the <code>race</code> flag, which works with the <code>test</code>, <code>run</code>, <code>build</code>, and <code>install</code> commands.</p>\n<pre><code>var total int\n// There&#39;s a data race on total.\nvar wg sync.WaitGroup\nwg.Go(func() { total++ })\nwg.Go(func() { total++ })\nwg.Wait()\nfmt.Println(&quot;total:&quot;, total)</code></pre>\n<pre><code>total: 2</code></pre>\n<pre><code>go run -race main.go</code></pre>\n<pre><code>==================\nWARNING: DATA RACE\n...\n2\nFound 1 data race(s)</code></pre>\n<p>Channels are safe for concurrent reading and writing, and they don&#39;t cause data races.</p>\n<p>Ways to prevent data races:</p>\n<ul><li>Avoid concurrent data modification (typically by using channels).</li><li>Synchronize access with mutexes.</li><li>Use only atomic operations.</li></ul>\n<h2 id=\"race-conditions\">Race conditions</h2>\n<p>A race condition happens when an unpredictable order of operations from multiple goroutines leads to an incorrect system state:</p>\n<pre><code>// There&#39;s a race condition when working with balance.\nwithdraw := func(amount int) {\n    if getBalance() &lt; amount {\n        return\n    }\n    time.Sleep(time.Millisecond)\n    setBalance(getBalance() - amount)\n}\nsetBalance(50)\nvar wg sync.WaitGroup\nwg.Go(func() { withdraw(40) })\nwg.Go(func() { withdraw(40) })\nwg.Wait()\nfmt.Println(&quot;balance:&quot;, getBalance())</code></pre>\n<pre><code>balance: -30</code></pre>\n<p>If individual operations are concurrent-safe, Go&#39;s race detector won&#39;t find any issues. Because of this, it doesn&#39;t catch race conditions:</p>\n<pre><code>go run -race main.go</code></pre>\n<pre><code>balance: -30</code></pre>\n<p>You can&#39;t fully eliminate uncertainty in a concurrent environment. Events will happen in an unpredictable order — that&#39;s just how concurrency works. However, you can prevent a race condition — often by protecting a composite operation with a mutex:</p>\n<pre><code>var mu sync.Mutex\nwithdraw := func(amount int) {\n    mu.Lock()\n    defer mu.Unlock()\n    if getBalance() &lt; amount {\n        return\n    }\n    time.Sleep(time.Millisecond)\n    setBalance(getBalance() - amount)\n}\nsetBalance(50)\nvar wg sync.WaitGroup\nwg.Go(func() { withdraw(40) })\nwg.Go(func() { withdraw(40) })\nwg.Wait()\nfmt.Println(&quot;balance:&quot;, getBalance())</code></pre>\n<pre><code>balance: 10</code></pre>\n<h3 id=\"compare-and-set\">Compare-and-set</h3>\n<p>Sometimes you can prevent a race condition without using mutexes by applying an atomic compare-and-set operation or one of its flavors:</p>\n<pre><code>// CompareAndSet changes the value to new if the current value equals old.\n// Returns true if the value was changed.\nCompareAndSet(old, new any) bool\n// CompareAndSwap changes the value to new if the current value equals old.\n// Returns the old value.\nCompareAndSwap(old, new any) any\n// CompareAndDelete deletes the value if the current value equals old.\n// Returns true if the value was deleted.\nCompareAndDelete(old any) bool\n// etc</code></pre>\n<p>The idea is always the same:</p>\n<ul><li>Check if the assumed (old) state matches reality.</li><li>If it does, change the state to new.</li><li>If not, do nothing.</li></ul>\n<h2 id=\"section-9\"><a href=\"https://antonz.org#mutexes\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Mutexes</p>\n<p>The <code>sync.Mutex</code> type protects shared data and parts of your code from being accessed concurrently:</p>\n<pre><code>var total int\nvar mu sync.Mutex\nvar wg sync.WaitGroup\nfor range 100 {\n    wg.Go(func() {\n        mu.Lock()\n        time.Sleep(time.Millisecond)\n        total++\n        mu.Unlock()\n    })\n}\nwg.Wait()</code></pre>\n<pre><code>total: 100</code></pre>\n<p>The mutex guarantees that only one goroutine can run the code between <code>Lock()</code> and <code>Unlock()</code> at a time.</p>\n<p>A mutex is used in these situations:</p>\n<ul><li>When multiple goroutines are modifying the same data.</li><li>When one goroutine is modifying the data and others are reading it.</li></ul>\n<p>If all goroutines are only reading the data, you don&#39;t need a mutex.</p>\n<h3 id=\"trylock\">TryLock</h3>\n<p>The <code>TryLock</code> method tries to lock the mutex, just like a regular <code>Lock</code>. But if it can&#39;t, it returns <code>false</code> right away instead of blocking the goroutine:</p>\n<pre><code>var total int\nvar mu sync.Mutex\nvar wg sync.WaitGroup\nfor range 100 {\n    wg.Go(func() {\n        if !mu.TryLock() {\n            return\n        }\n        defer mu.Unlock()\n        time.Sleep(time.Millisecond)\n        total++\n    })\n}\nwg.Wait()</code></pre>\n<pre><code>total: 1</code></pre>\n<h3 id=\"rwmutex\">RWMutex</h3>\n<p>The <code>sync.RWMutex</code> type distinguishes between readers and writers. It provides two sets of methods:</p>\n<ul><li><code>Lock</code> /<code>Unlock</code> lock and unlock the mutex for both reading and writing.</li><li><code>RLock</code> /<code>RUnlock</code> lock and unlock the mutex for reading only.</li></ul>\n<pre><code>var total int\nvar mu sync.RWMutex\nvar wg sync.WaitGroup\n// 10 writers.\nfor range 10 {\n    wg.Go(func() {\n        mu.Lock()\n        defer mu.Unlock()\n        time.Sleep(time.Millisecond)\n        total++\n    })\n}\n// 10 readers.\nfor range 10 {\n    wg.Go(func() {\n        // Try switching from RLock/RUnlock to Lock/Unlock\n        //and see how it affects the elapsed time.\n        mu.RLock()\n        defer mu.RUnlock()\n        time.Sleep(time.Millisecond)\n        _ = total\n    })\n}\nwg.Wait()</code></pre>\n<pre><code>elapsed: 10ms</code></pre>\n<p>Here&#39;s how it works:</p>\n<ul><li>If a goroutine locks the mutex with <code>Lock()</code> , other goroutines will be blocked if they try to use<code>Lock()</code> or<code>RLock()</code> .</li><li>If a goroutine locks the mutex with <code>RLock()</code> , other goroutines can also lock it with<code>RLock()</code> without being blocked.</li><li>If at least one goroutine has locked the mutex with <code>RLock()</code> , other goroutines will be blocked if they try to use<code>Lock()</code> .</li></ul>\n<p>This creates a &quot;single writer, multiple readers&quot; setup.</p>\n<h3 id=\"locker\">Locker</h3>\n<p>Both <code>sync.Mutex</code> and <code>sync.RWMutex</code> implement the same <code>sync.Locker</code> interface:</p>\n<pre><code>type Locker interface {\n    Lock()\n    Unlock()\n}</code></pre>\n<p>By using <code>Locker</code> instead of a specific mutex type, you can build components that don&#39;t depend on a specific lock implementation. This lets the client decide which lock to use.</p>\n<h3 id=\"channel-as-mutex\">Channel as mutex</h3>\n<p>You can use a channel instead of a mutex to protect shared data:</p>\n<pre><code>var total int\nlock := make(chan struct{}, 1)\nvar wg sync.WaitGroup\nwg.Go(func() {\n    lock &lt;- struct{}{}\n    defer func() { &lt;-lock }()\n    total++\n})\nwg.Go(func() {\n    lock &lt;- struct{}{}\n    defer func() { &lt;-lock }()\n    total++\n})\nwg.Wait()</code></pre>\n<pre><code>total: 2</code></pre>\n<h2 id=\"section-10\"><a href=\"https://antonz.org#semaphores\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Semaphores</p>\n<p>A semaphore is like a container with N available slots and two operations: <em>acquire</em> to take a slot and <em>release</em> to free a slot. Here are the semaphore rules:</p>\n<ul><li>Calling acquire takes a free slot.</li><li>If there are no free slots, acquire blocks the goroutine that called it.</li><li>Calling release frees up a previously taken slot.</li><li>If there are any goroutines blocked on acquire when release is called, one of them will immediately take the freed slot and unblock.</li></ul>\n<p>You can implement a simple semaphore with a buffered channel, where N is the channel&#39;s size. To acquire the semaphore, send a value into the channel. To release it, take a value from the channel:</p>\n<pre><code>// Try changing nConc and see how the elapsed time changes.\nconst nConc = 4\nconst nCalls = 100\nsema := make(chan struct{}, nConc)\nvar wg sync.WaitGroup\nfor range nCalls {\n    sema &lt;- struct{}{} // acquire\n    wg.Go(func() {\n        defer func() { &lt;-sema }() // release\n        time.Sleep(time.Millisecond) // do some work\n    })\n}\nwg.Wait()</code></pre>\n<pre><code>elapsed: 25ms</code></pre>\n<p>For more complex situations, use the <code>golang.org/x/sync/semaphore</code> package.</p>\n<h3 id=\"rendezvous\">Rendezvous</h3>\n<p>A rendezvous lets two goroutines wait for each other:</p>\n<ul><li>There are two goroutines — G1 and G2 — and each one can signal that it&#39;s ready.</li><li>If G1 signals but G2 hasn&#39;t yet, G1 blocks and waits.</li><li>If G2 signals but G1 hasn&#39;t yet, G2 blocks and waits.</li><li>When both have signaled, they both unblock and continue running.</li></ul>\n<p>You can implement a simple rendezvous with a wait group:</p>\n<pre><code>var rend sync.WaitGroup\nrend.Add(2)\nvar wg sync.WaitGroup\nwg.Go(func() {\n    fmt.Println(&quot;before rendezvous&quot;)\n    rend.Done()\n    rend.Wait()\n    fmt.Println(&quot;after rendezvous&quot;)\n})\nwg.Go(func() {\n    fmt.Println(&quot;before rendezvous&quot;)\n    rend.Done()\n    rend.Wait()\n    fmt.Println(&quot;after rendezvous&quot;)\n})\nwg.Wait()</code></pre>\n<pre><code>before rendezvous\nbefore rendezvous\nafter rendezvous\nafter rendezvous</code></pre>\n<h3 id=\"barrier\">Barrier</h3>\n<p>A barrier is a general case of a rendezvous. It lets N goroutines wait for each other:</p>\n<ul><li>The barrier has a counter (starting at 0) and a threshold N.</li><li>Each goroutine that reaches the barrier increases the counter by 1.</li><li>The barrier blocks any goroutine that reaches it.</li><li>Once the counter reaches N, the barrier unblocks all waiting goroutines.</li></ul>\n<p>You can implement a simple barrier with a wait group:</p>\n<pre><code>const n = 4\nvar bar sync.WaitGroup\nbar.Add(n)\nvar wg sync.WaitGroup\nfor range n {\n    wg.Go(func() {\n        fmt.Println(&quot;before the barrier&quot;)\n        bar.Done()\n        bar.Wait()\n        fmt.Println(&quot;after the barrier&quot;)\n    })\n}\nwg.Wait()</code></pre>\n<pre><code>before the barrier\nbefore the barrier\nbefore the barrier\nbefore the barrier\nafter the barrier\nafter the barrier\nafter the barrier\nafter the barrier</code></pre>\n<h2 id=\"section-11\"><a href=\"https://antonz.org#signaling\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Signaling</p>\n<p>The <code>sync.Cond</code> (conditional variable) type lets one goroutine signal to another that it&#39;s ready, and lets the other goroutine wait for that signal.</p>\n<p>A <code>Cond</code> includes a mutex and has two methods — <code>Wait</code> and <code>Signal</code>.</p>\n<ul><li><code>Wait</code> unlocks the mutex and suspends the goroutine until it receives a signal.</li><li><code>Signal</code> wakes the goroutine that is waiting on<code>Wait</code> .</li><li>When <code>Wait</code> wakes up, it locks the mutex again.</li></ul>\n<pre><code>cond := sync.NewCond(&amp;sync.Mutex{})\ndone := false\nvar wg sync.WaitGroup\nwg.Go(func() {\n    cond.L.Lock()\n    fmt.Println(&quot;G1 is ready to signal&quot;)\n    done = true\n    cond.Signal()\n    cond.L.Unlock()\n})\nwg.Go(func() {\n    cond.L.Lock()\n    for !done {\n        cond.Wait()\n    }\n    fmt.Println(&quot;G2 received the signal&quot;)\n    cond.L.Unlock()\n})\nwg.Wait()</code></pre>\n<pre><code>G1 is ready to signal\nG2 received the signal</code></pre>\n<p>If there are multiple waiting goroutines when <code>Signal</code> is called, only one of them will be resumed. If there are no waiting goroutines, <code>Signal</code> does nothing.</p>\n<p>You can also use the <code>Broadcast</code> method. While <code>Signal</code> wakes up only one goroutine waiting on <code>Cond.Wait</code>, the <code>Broadcast</code> method wakes up all such goroutines.</p>\n<p>You can signal with a channel:</p>\n<pre><code>signal := make(chan struct{}, 1)\ngo func() {\n    // do something\n    signal &lt;- struct{}{}\n}()\ngo func() {\n    &lt;-signal\n    // do something\n}()</code></pre>\n<p>And broadcast too:</p>\n<pre><code>broadcast := make(chan struct{})\ngo func() {\n    // do something\n    close(broadcast)\n}()\ngo func() {\n    &lt;-broadcast\n    // do something\n}()\ngo func() {\n    &lt;-broadcast\n    // do something\n}()</code></pre>\n<p>Broadcasting with a condition variable is limited: it only sends a signal, not the actual data, and it only works once. With channels, you can build a publish/subscribe system that doesn&#39;t have these limitations:</p>\n<pre><code>type Publisher struct {\n    sbox []chan int // subscription channels\n    mu   sync.Mutex // protects the state\n}\nfunc (p *Publisher) Subscribe() &lt;-chan int {\n    p.mu.Lock()\n    defer p.mu.Unlock()\n    sub := make(chan int, 1)\n    p.sbox = append(p.sbox, sub)\n    return sub\n}\nfunc (p *Publisher) Broadcast(v int) {\n    p.mu.Lock()\n    defer p.mu.Unlock()\n    for _, sub := range p.sbox {\n        select {\n        case sub &lt;- v:\n        default:\n        }\n    }\n}</code></pre>\n<h2 id=\"section-12\"><a href=\"https://antonz.org#run-once\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Run once</p>\n<p>The <code>sync.Once</code> type makes sure that the given function runs only once. If multiple goroutines call <code>Once.Do</code> at the same time, only one will run the function, while the others will wait until it returns:</p>\n<pre><code>total := 0\ninitState := func() {\n    total += 1\n}\nvar once sync.Once\nvar wg sync.WaitGroup\nwg.Go(func() {\n    once.Do(initState)\n    // do something\n})\nwg.Go(func() {\n    once.Do(initState)\n    // do something\n})\nwg.Wait()</code></pre>\n<pre><code>total: 1</code></pre>\n<p><code>Once</code> is perfect for one-time initialization or cleanup in a concurrent environment.</p>\n<p>Besides the <code>Once</code> type, the <code>sync</code> package also includes three convenience once-functions:</p>\n<pre><code>// Calls f only once.\nfunc (o *Once) Do(f func())\n// Returns a function that calls f only once.\nfunc OnceFunc(f func()) func()\n// Returns a function that calls f only once\n// and returns the value from that first call.\nfunc OnceValue[T any](f func() T) func() T\n// Returns a function that calls f only once\n// and returns the pair of values from that first call.\nfunc OnceValues[T1, T2 any](f func() (T1, T2)) func() (T1, T2)</code></pre>\n<h2 id=\"section-13\"><a href=\"https://antonz.org#object-pool\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Object pool</p>\n<p>The <code>sync.Pool</code> type helps reuse memory instead of allocating it every time, which reduces the load on the garbage collector:</p>\n<pre><code>pool := sync.Pool{\n    New: func() any {\n        buf := make([]byte, 1024)\n        return &amp;buf\n    },\n}\n// Only allocates 4*1024 B, despite 4000 loop iterations.\nvar wg sync.WaitGroup\nfor range 4 {\n    wg.Go(func() {\n        for range 1000 {\n            buf := pool.Get().(*[]byte)\n            sink = buf\n            pool.Put(buf)\n        }\n    })\n}\nwg.Wait()</code></pre>\n<pre><code>Memory allocated: 4 KB</code></pre>\n<p><code>Get</code> takes an item from the pool. If there are no available items, it creates a new one using <code>New</code> (which we have to define ourselves, since the pool doesn&#39;t know anything about the items it creates). <code>Put</code> returns an item back to the pool.</p>\n<p>Things to keep in mind:</p>\n<ul><li><code>New</code> should return a pointer, not a value, to reduce memory copying and avoid extra allocations.</li><li>The pool has no size limit. If you start 1000 more goroutines that all call <code>Get</code> at the same time, 1000 more buffers will be allocated.</li><li>After an item is returned to the pool with <code>Put</code> , you shouldn&#39;t use it anymore (since another goroutine might already have taken and started using it).</li></ul>\n<h2 id=\"section-14\"><a href=\"https://antonz.org#atomics\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Atomics</p>\n<p>An operation without synchronization can only be truly atomic if it translates to a single processor instruction. Such operations don&#39;t need locks and won&#39;t cause issues when called concurrently (even the write operations).</p>\n<p>There are only a few atomics, and they&#39;re all found in the <code>sync/atomic</code> package:</p>\n<pre><code>Int32     Bool\nInt64     Value\nUint32    Pointer\nUint64</code></pre>\n<p>Each atomic type provides the following methods:</p>\n<ul><li><code>Load</code> reads the value of a variable.</li><li><code>Store</code> sets a new value.</li><li><code>Swap</code> sets a new value (like<code>Store</code> ) and returns the old one.</li><li><code>CompareAndSwap</code> sets a new value only if the current value is still what you expect it to be.</li></ul>\n<pre><code>var n atomic.Int32\nn.Store(10)\nswapped := n.CompareAndSwap(10, 42)\nfmt.Println(&quot;CompareAndSwap 10 -&gt; 42:&quot;, swapped)\nfmt.Println(&quot;n =&quot;, n.Load())</code></pre>\n<pre><code>CompareAndSwap 10 -&gt; 42: true\nn = 42</code></pre>\n<p>Numeric types also provide an <code>Add</code> method that increments the value by the specified amount.</p>\n<p>All methods are either translated into a single CPU instruction or are otherwise guaranteed to be atomic, so they are safe to use from multiple goroutines.</p>\n<p>The composition of atomics is always non-atomic:</p>\n<pre><code>var delta atomic.Int32\nvar counter atomic.Int32\nfunc increment() {\n    // Not atomic; causes a race condition.\n    delta.Add(1)\n    sleep(10)\n    counter.Add(delta.Load())\n}\n// After 100 concurrent increments,\n// the final value is NOT guaranteed.</code></pre>\n<pre><code>counter = 9386</code></pre>\n<p>A bulletproof way to make a composite operation atomic and prevent race conditions is to use a mutex:</p>\n<pre><code>var delta int32\nvar counter int32\nvar mu sync.Mutex\nfunc increment() {\n    // Atomic; doesn&#39;t cause a race condition.\n    mu.Lock()\n    delta += 1\n    sleep(10)\n    counter += delta\n    mu.Unlock()\n}\n// After 100 concurrent increments, the final value is guaranteed:\n// counter = 1+2+...+100 = 5050</code></pre>\n<pre><code>counter = 5050</code></pre>\n<p>Sometimes you can use an atomic type instead of a mutex to exit early:</p>\n<pre><code>type Gate struct {\n    closed atomic.Bool\n}\nfunc (g *Gate) Close() {\n    if !g.closed.CompareAndSwap(false, true) {\n        return // ignore repeated calls\n    }\n    // The gate is closed.\n    // We can free resources now.\n}</code></pre>\n<h2 id=\"section-15\"><a href=\"https://antonz.org#testing\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Testing</p>\n<p>If your concurrent program uses channels or custom types with synchronization methods like <code>Wait</code>, you can use those in your tests. This way, your tests won&#39;t be much more complicated than if the code were synchronous:</p>\n<pre><code>// Calc calculates something asynchronously.\nfunc Calc() &lt;-chan int {\n    out := make(chan int, 1)\n    go func() {\n        out &lt;- 42\n    }()\n    return out\n}</code></pre>\n<pre><code>func Test(t *testing.T) {\n    // Wait for the Calc goroutine to finish.\n    got := &lt;-Calc()\n    if got != 42 {\n        t.Errorf(&quot;got: %v; want: 42&quot;, got)\n    }\n}</code></pre>\n<pre><code>PASS</code></pre>\n<p>If there aren&#39;t any suitable synchronization &quot;handles&quot; in the code you&#39;re testing, you can use the <code>synctest</code> package. It exports two functions:</p>\n<pre><code>func Test(t *testing.T, f func(*testing.T))\nfunc Wait()</code></pre>\n<p><code>synctest.Test</code> runs an isolated bubble. The bubble uses a fake clock, and you can manually control goroutine synchronization with <code>synctest.Wait</code>.</p>\n<p><code>synctest.Wait</code> blocks until all goroutines in the bubble — except the one that called <code>Wait</code> — have either finished or are durably blocked. This lets you wait for a specific goroutine to finish or get blocked, so you can check the program&#39;s state:</p>\n<pre><code>// NewProc starts the calculation.\nfunc NewProc() *Proc {\n    p := &amp;Proc{done: make(chan struct{})}\n    go func() {\n        p.res = 42\n        &lt;-p.done // (X)\n        p.res = 0\n    }()\n    return p\n}</code></pre>\n<pre><code>func Test(t *testing.T) {\n    synctest.Test(t, func(t *testing.T) {\n        p := NewProc()\n        defer p.Stop()\n        // Wait for the goroutine to block at point X.\n        synctest.Wait()\n        if got := p.Res(); got != 42 {\n            t.Fatalf(&quot;got %v, want 42&quot;, got)\n        }\n    })\n}</code></pre>\n<pre><code>PASS</code></pre>\n<p>The fake clock in <code>synctest.Test</code> move forward only if: ➊ all goroutines in the bubble are durably blocked; ➋ there&#39;s a future moment when at least one goroutine will unblock; and ➌ <code>synctest.Wait</code> isn&#39;t running. Thanks to this, time-dependent tests run instantly:</p>\n<pre><code>// Calc processes a value from the input channel.\n// Times out if no input is received after 3 seconds.\nfunc Calc(in chan int) (int, error) {\n    select {\n    case v := &lt;-in:\n        return v * 2, nil\n    case &lt;-time.After(3 * time.Second):\n        return 0, ErrTimeout\n    }\n}</code></pre>\n<pre><code>func Test(t *testing.T) {\n    synctest.Test(t, func(t *testing.T) {\n        ch := make(chan int)\n        got, err := Calc(ch) // runs instantly\n        if err != ErrTimeout {\n            t.Errorf(&quot;got: %v; want: %v&quot;, err, ErrTimeout)\n        }\n        if got != 0 {\n            t.Errorf(&quot;got: %v; want: 0&quot;, got)\n        }\n    })\n}</code></pre>\n<pre><code>PASS</code></pre>\n<p>The following operations durably block a goroutine:</p>\n<ul><li>A blocking send or receive on a channel created within the bubble.</li><li>A blocking select statement where every case is a channel created within the bubble.</li><li>Calling <code>Cond.Wait</code> .</li><li>Calling <code>WaitGroup.Wait</code> if all<code>WaitGroup.Add</code> calls were made inside the bubble.</li><li>Calling <code>time.Sleep</code> .</li></ul>\n<p>Blocking on mutexes, I/O, or system calls is not considered durable, and the <code>synctest</code> bubble can&#39;t handle them.</p>\n<h2 id=\"section-16\"><a href=\"https://antonz.org#scheduling\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Scheduling</p>\n<p>At the hardware level, CPU cores are responsible for running parallel tasks.</p>\n<p>At the operating system level, a thread is the basic unit of execution. There are usually many more threads than CPU cores, so the operating system&#39;s scheduler decides which threads to run and which ones to pause.</p>\n<p>At the Go runtime level, a goroutine is the basic unit of execution. The runtime scheduler runs a fixed number of OS threads, often one per CPU core. There can be many more goroutines than threads, so the scheduler decides which goroutines to run on the available threads and which ones to pause. The scheduler keeps switching between goroutines to make sure each one gets a turn to run on a thread, instead of waiting in line forever.</p>\n<pre><code>  CPU                  OS                   Go runtime\n┌──────────┐  run on ┌──────────┐  run on ┌────────────┐\n│ Cores    │ &lt;────── │ Threads  │ &lt;────── │ Goroutines │\n└──────────┘         └──────────┘         └────────────┘</code></pre>\n<p>This is how Go handles concurrency.</p>\n<p><strong>Goroutine scheduler</strong></p>\n<p>The goroutine scheduler&#39;s job is to run M goroutines on N operating system threads, where M can be much larger than N. Here&#39;s a very simplified version of it&#39;s algorithm:</p>\n<ul><li>If there&#39;s a free thread, assign it a goroutine from the queue.</li><li>If a running goroutine gets blocked (for example, while reading from a channel), put it back in the queue and assign a different goroutine to the thread.</li><li>If a running goroutine gets stuck in a syscall, start a new thread to run other goroutines until the blocked goroutine finishes the syscall.</li><li>Check the running goroutines every 10 ms. Preempt long-running goroutines and return them to the queue to prevent starvation.</li></ul>\n<pre><code>┌─────┐┌─────┐┌─────┐┌─────┐\n│ G17 ││ G18 ││ G19 ││ G20 │                        queue\n└─────┘└─────┘└─────┘└─────┘\n┌─────┐      ┌─────┐      ┌─────┐      ┌─────┐\n│ G15 │      │ G16 │      │ G13 │      │ G14 │      running\n└─────┘      └─────┘      └─────┘      └─────┘\n  │            │            │            │\n┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐\n│ Thread E │ │ Thread F │ │ Thread C │ │ Thread D │\n└──────────┘ └──────────┘ └──────────┘ └──────────┘\n┌─────┐      ┌─────┐\n│ G11 │      │ G12 │                                syscalls\n└─────┘      └─────┘\n  │            │\n┌──────────┐ ┌──────────┐\n│ Thread A │ │ Thread B │\n└──────────┘ └──────────┘</code></pre>\n<p>The number of threads running Go code is controlled by the <code>GOMAXPROCS</code> environment variable or the <code>runtime.GOMAXPROCS</code> function.</p>\n<p>A goroutine is a structure that starts out using about 2 KB of memory, mostly for its stack. The stack can grow if needed. Since goroutines are so lightweight, you can run tens of thousands or even hundreds of thousands of them on a small machine.</p>\n<h2 id=\"section-17\"><a href=\"https://antonz.org#diagnostics\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Diagnostics</p>\n<p>To troubleshoot concurrent programs in production, we use metrics, profiling, and tracing.</p>\n<p><strong>Metrics</strong> show how the Go runtime is performing, like how much heap memory it uses or how long garbage collection pauses take. Each metric has a unique name and a value, which can be a number or a histogram.</p>\n<p>You can use the <code>runtime/metrics</code> package to get a complete list of metrics or check the values of specific ones:</p>\n<pre><code>samples := []metrics.Sample{\n    {Name: &quot;/sched/gomaxprocs:threads&quot;},\n    {Name: &quot;/sched/goroutines:goroutines&quot;},\n}\nmetrics.Read(samples)\nfor _, s := range samples {\n    fmt.Printf(&quot;%s: %v\\n&quot;, s.Name, s.Value.Uint64())\n}</code></pre>\n<pre><code>/sched/gomaxprocs:threads: 8\n/sched/goroutines:goroutines: 1</code></pre>\n<p>In practice, people rarely do this manually. Instead, all metrics are automatically exported using Prometheus or OpenTelemetry libraries.</p>\n<p><strong>Profiling</strong> helps you understand exactly what the program is doing, what resources it uses, and where in the code this happens. Go uses a sampling profiler that&#39;s suitable for production.</p>\n<p>The most commonly used profiles are CPU, which shows how much processor time each function uses, and heap, which shows how much heap memory each function uses. Goroutine, block, and mutex profiles help identify problems related to concurrency.</p>\n<p>The easiest way to add a profiler to your app is by using the <code>net/http/pprof</code> package. To collect a profile with the given name, call the <code>/debug/pprof/{name}</code> endpoint. To view the collected profile, use the <code>go tool pprof</code> utility:</p>\n<pre><code>go tool pprof -proto \\\n  &quot;http://localhost:6060/debug/pprof/profile?seconds=N&quot; &gt; cpu.pprof\ngo tool pprof -http=localhost:8080 cpu.pprof</code></pre>\n<p>You can also profile manually:</p>\n<pre><code>// CPU profile.\nfile, _ := os.Create(&quot;cpu.prof&quot;)\ndefer file.Close()\npprof.StartCPUProfile(file)\ndefer pprof.StopCPUProfile()\n// ...</code></pre>\n<pre><code>// Any other profile.\nfile, _ := os.Create(name + &quot;.prof&quot;)\ndefer file.Close()\npprof.Lookup(name).WriteTo(file, 0)</code></pre>\n<p><strong>Tracing</strong> records certain types of events while the program is running, mainly those related to concurrency and memory. When the profiling server from the <code>net/http/pprof</code> package is running, call the <code>/debug/pprof/trace</code> endpoint to collect a trace. To view the results, use the <code>go tool trace</code> utility.</p>\n<p>You can also collect a trace manually:</p>\n<pre><code>file, _ := os.Create(&quot;trace.out&quot;)\ndefer file.Close()\ntrace.Start(file)\ndefer trace.Stop()\n// ...</code></pre>\n<p>You can set up automatic tracing with a sliding window that&#39;s limited by size or duration. This is called &quot;flight recording&quot;. It lets you always keep a recent trace available in case something goes wrong:</p>\n<pre><code>cfg := trace.FlightRecorderConfig{\n    MinAge:   5 * time.Second,\n    MaxBytes: 3 &lt;&lt; 20, // 3MB\n}\nrec := trace.NewFlightRecorder(cfg)\nrec.Start()\ndefer rec.Stop()</code></pre>\n<h2 id=\"section-18\"><a href=\"https://antonz.org#final-thoughts\" rel=\"nofollow ugc noopener\">#</a></h2>\n<p>Final thoughts</p>\n<p>We&#39;ve covered a number of Go tools for writing concurrent programs:</p>\n<ul><li>Goroutines for running concurrent tasks.</li><li>Channels and select as flexible communication tools.</li><li>Timers and tickers for working with time.</li><li>Context for canceling operations.</li><li>Wait groups for synchronizing goroutines.</li><li>Mutexes to prevent race conditions.</li><li>Condition variables for signaling events.</li><li>Once for safe one-time initialization.</li><li>Pools to reduce garbage collector load.</li><li>Atomic operations.</li></ul>\n<p>If you like the book, please recommend it to your friends or colleagues. If you&#39;re interested, check out my other <a href=\"https://antonz.org/#books\" rel=\"nofollow ugc noopener\">books</a> and <a href=\"https://antonz.org/tags/projects/\" rel=\"nofollow ugc noopener\">projects</a>.</p>\n<p>I&#39;m glad you finished the book. Thank you, and I&#39;ll see you next time!</p>\n<p><a href=\"https://antonz.org/subscribe/\" rel=\"nofollow ugc noopener\">★ Subscribe</a> to keep up with new posts.</p>","headings":[{"level":1,"text":"Go concurrency distilled","id":"go-concurrency-distilled"},{"level":2,"text":"#","id":"section"},{"level":2,"text":"#","id":"section-2"},{"level":3,"text":"Output channel","id":"output-channel"},{"level":3,"text":"Closing a channel","id":"closing-a-channel"},{"level":3,"text":"Channel iteration","id":"channel-iteration"},{"level":3,"text":"Directional channels","id":"directional-channels"},{"level":3,"text":"Buffered channels","id":"buffered-channels"},{"level":3,"text":"nil channel","id":"nil-channel"},{"level":2,"text":"#","id":"section-3"},{"level":2,"text":"#","id":"section-4"},{"level":3,"text":"Output channel","id":"output-channel-2"},{"level":3,"text":"Done channel","id":"done-channel"},{"level":3,"text":"Cancel channel","id":"cancel-channel"},{"level":3,"text":"Error handling","id":"error-handling"},{"level":2,"text":"#","id":"section-5"},{"level":3,"text":"After","id":"after"},{"level":3,"text":"Timer","id":"timer"},{"level":3,"text":"Ticker","id":"ticker"},{"level":2,"text":"#","id":"section-6"},{"level":2,"text":"#","id":"section-7"},{"level":2,"text":"#","id":"section-8"},{"level":2,"text":"Race conditions","id":"race-conditions"},{"level":3,"text":"Compare-and-set","id":"compare-and-set"},{"level":2,"text":"#","id":"section-9"},{"level":3,"text":"TryLock","id":"trylock"},{"level":3,"text":"RWMutex","id":"rwmutex"},{"level":3,"text":"Locker","id":"locker"},{"level":3,"text":"Channel as mutex","id":"channel-as-mutex"},{"level":2,"text":"#","id":"section-10"},{"level":3,"text":"Rendezvous","id":"rendezvous"},{"level":3,"text":"Barrier","id":"barrier"},{"level":2,"text":"#","id":"section-11"},{"level":2,"text":"#","id":"section-12"},{"level":2,"text":"#","id":"section-13"},{"level":2,"text":"#","id":"section-14"},{"level":2,"text":"#","id":"section-15"},{"level":2,"text":"#","id":"section-16"},{"level":2,"text":"#","id":"section-17"},{"level":2,"text":"#","id":"section-18"}]}}