Async
Structured concurrency: Promise.all and errgroup
How Node.js's Promise.all/allSettled/any combined with AbortController compare to Go's errgroup, sync.WaitGroup, and channels: cancelling sibling tasks when one fails, limiting concurrency, collecting every result, and taking the fastest one.
- Minimum versions
- Node.js ≥ 17.2Go ≥ 1.25
- Verified on
- Node.js 24.12.0Go 1.27.1
Node.js code is an ES module: save it as .mjs or set "type": "module" in package.json.
"Structured concurrency" is a simple rule: every child task started inside a scope has to finish before that scope ends — no task "leaks" out, and a child's failure is reported back to the parent scope. Node.js ships the Promise.all/allSettled/any/race combinators, but they only wait — they don't cancel anything: when Promise.all rejects because one promise failed, the others keep running to completion. To make sibling tasks stop, you have to pass an AbortSignal in yourself (see the cancellation post). Go has no Promise type; the equivalent patterns are assembled from goroutines, sync.WaitGroup, channels, and context — with errgroup (from golang.org/x/sync, outside the standard library) being the closest thing to "Promise.all with cancellation".
The first error cancels the siblings (Promise.all + AbortController / errgroup)
import { setTimeout as sleep } from "node:timers/promises";
// work simulates a job that takes `ms` milliseconds, may fail, and stops
// early when the signal is aborted.
async function work({ name, ms, fail }, signal) {
await sleep(ms, undefined, { signal });
if (fail) throw new Error(`${name} failed`);
return `${name} ok`;
}
// all: like Promise.all, but the first error aborts the sibling tasks, and it
// only returns once EVERY task has actually stopped (like errgroup.Wait). Task i
// records its outcome in statuses[i], like the statuses slice main allocates in Go.
async function all(tasks, statuses) {
const controller = new AbortController();
const settled = await Promise.allSettled(
tasks.map(async (task, i) => {
try {
const value = await work(task, controller.signal);
statuses[i] = "ok";
return value;
} catch (err) {
statuses[i] = err.name === "AbortError" ? "cancelled" : "failed";
controller.abort(err); // later aborts are no-ops: reason keeps the first error
throw err;
}
}),
);
if (controller.signal.aborted) throw controller.signal.reason;
return settled.map((s) => s.value);
}
const tasks = [
{ name: "A", ms: 100 },
{ name: "B", ms: 50, fail: true }, // B fails at ~50ms, before A and C finish
{ name: "C", ms: 300 },
];
const statuses = new Array(tasks.length);
try {
await all(tasks, statuses); // waits for EVERY task to stop, then throws the first error
} catch (err) {
console.log(statuses); // → [ 'cancelled', 'failed', 'cancelled' ]
console.log("error:", err.message); // → error: B failed
}// needs a go.mod: go mod init example.com/app && go get golang.org/x/sync
package main
import (
"context"
"errors"
"fmt"
"time"
"golang.org/x/sync/errgroup"
)
type task struct {
name string
d time.Duration
fail bool
}
// work simulates a job that takes t.d, may fail, and stops early when ctx is cancelled.
func work(ctx context.Context, t task) (string, error) {
timer := time.NewTimer(t.d)
defer timer.Stop()
select {
case <-timer.C:
case <-ctx.Done():
return "", ctx.Err()
}
if t.fail {
return "", fmt.Errorf("%s failed", t.name)
}
return t.name + " ok", nil
}
func main() {
tasks := []task{
{name: "A", d: 100 * time.Millisecond},
{name: "B", d: 50 * time.Millisecond, fail: true}, // B fails at ~50ms, before A and C finish
{name: "C", d: 300 * time.Millisecond},
}
// ctx is cancelled as soon as a goroutine in g returns the first error
g, ctx := errgroup.WithContext(context.Background())
statuses := make([]string, len(tasks))
for i, t := range tasks {
g.Go(func() error {
_, err := work(ctx, t)
switch {
case errors.Is(err, context.Canceled):
statuses[i] = "cancelled"
case err != nil:
statuses[i] = "failed"
default:
statuses[i] = "ok"
}
return err
})
}
err := g.Wait() // wait for EVERY goroutine to stop, then return the first error
fmt.Println(statuses) // → [cancelled failed cancelled]
fmt.Println("error:", err) // → error: B failed
}Limiting how many tasks run at once (mapLimit / g.SetLimit)
import { setTimeout as sleep } from "node:timers/promises";
// mapLimit: runs fn on every item but never more than `limit` at a time —
// Node.js has nothing built in, so build it from `limit` "workers" pulling from one queue.
async function mapLimit(items, limit, fn) {
const results = new Array(items.length);
let next = 0;
let failed = false;
async function worker() {
while (!failed && next < items.length) {
const i = next++; // safe: JS is single-threaded, nothing runs between the read and the increment
try {
results[i] = await fn(items[i]);
} catch (err) {
failed = true; // other workers see this flag and stop taking new items
throw err;
}
}
}
await Promise.all(Array.from({ length: limit }, () => worker()));
return results;
}
let inFlight = 0;
let maxInFlight = 0;
const squares = await mapLimit([1, 2, 3, 4, 5, 6], 2, async (n) => {
inFlight++;
maxInFlight = Math.max(maxInFlight, inFlight);
await sleep(20); // simulate I/O
inFlight--;
return n * n;
});
console.log(squares); // → [ 1, 4, 9, 16, 25, 36 ]
console.log(maxInFlight); // → 2 (at most)// needs a go.mod: go mod init example.com/app && go get golang.org/x/sync
package main
import (
"fmt"
"sync"
"time"
"golang.org/x/sync/errgroup"
)
func main() {
items := []int{1, 2, 3, 4, 5, 6}
squares := make([]int, len(items))
var mu sync.Mutex
inFlight, maxInFlight := 0, 0
var g errgroup.Group // the zero value works when you don't need a ctx
g.SetLimit(2) // g.Go blocks until fewer than 2 goroutines are running
for i, n := range items {
g.Go(func() error {
mu.Lock()
inFlight++
maxInFlight = max(maxInFlight, inFlight)
mu.Unlock()
time.Sleep(20 * time.Millisecond) // simulate I/O
mu.Lock()
inFlight--
mu.Unlock()
squares[i] = n * n
return nil
})
}
if err := g.Wait(); err != nil {
fmt.Println("error:", err)
return
}
fmt.Println(squares) // → [1 4 9 16 25 36]
fmt.Println(maxInFlight) // → 2 (at most; a slow machine may show less)
}Collecting every result, errors included (Promise.allSettled / WaitGroup + errors.Join)
import { setTimeout as sleep } from "node:timers/promises";
async function check(name, ms, fail = false) {
await sleep(ms);
if (fail) throw new Error(`${name} unreachable`);
return `${name} healthy`;
}
// allSettled never rejects: it waits for every promise, even if some fail
const settled = await Promise.allSettled([
check("db", 50),
check("cache", 20, true),
check("queue", 30, true),
]);
for (const s of settled) {
console.log(s.status, s.status === "fulfilled" ? s.value : s.reason.message);
}
// → fulfilled db healthy
// → rejected cache unreachable
// → rejected queue unreachablepackage main
import (
"errors"
"fmt"
"sync"
"time"
)
func check(name string, d time.Duration, fail bool) (string, error) {
time.Sleep(d)
if fail {
return "", fmt.Errorf("%s unreachable", name)
}
return name + " healthy", nil
}
func main() {
type spec struct {
name string
d time.Duration
fail bool
}
specs := []spec{
{"db", 50 * time.Millisecond, false},
{"cache", 20 * time.Millisecond, true},
{"queue", 30 * time.Millisecond, true},
}
// each goroutine writes to its own slot: no lock needed, order preserved
values := make([]string, len(specs))
errs := make([]error, len(specs))
var wg sync.WaitGroup
for i, s := range specs {
wg.Go(func() {
values[i], errs[i] = check(s.name, s.d, s.fail)
})
}
wg.Wait() // wait for all of them, cancel nobody — like allSettled
for i := range specs {
if errs[i] != nil {
fmt.Println("rejected", errs[i])
} else {
fmt.Println("fulfilled", values[i])
}
}
// errors.Join merges every non-nil error into one (nil if there are none)
if err := errors.Join(errs...); err != nil {
fmt.Printf("%q\n", err.Error())
}
}
// → fulfilled db healthy
// → rejected cache unreachable
// → rejected queue unreachable
// → "cache unreachable\nqueue unreachable"Taking the first success (Promise.any / channel + cancel)
import { setTimeout as sleep } from "node:timers/promises";
async function fetchFrom({ name, ms, fail }, signal) {
await sleep(ms, undefined, { signal });
if (fail) throw new Error(`${name} down`);
return `data from ${name}`;
}
async function fastest(mirrors) {
const controller = new AbortController();
try {
// any: take the first SUCCESSFUL result, ignoring failures
return await Promise.any(mirrors.map((m) => fetchFrom(m, controller.signal)));
} finally {
controller.abort(); // cancel the requests still running — Promise.any doesn't do this itself
}
}
console.log(await fastest([
{ name: "eu", ms: 300 },
{ name: "us", ms: 100 },
{ name: "asia", ms: 50, fail: true }, // fails first, but any ignores it
])); // → data from us
try {
await fastest([{ name: "eu", ms: 30, fail: true }, { name: "us", ms: 10, fail: true }]);
} catch (err) {
console.log(err.name, err.errors.map((e) => e.message)); // → AggregateError [ 'eu down', 'us down' ]
}package main
import (
"context"
"errors"
"fmt"
"time"
)
type mirror struct {
name string
d time.Duration
fail bool
}
func fetchFrom(ctx context.Context, m mirror) (string, error) {
timer := time.NewTimer(m.d)
defer timer.Stop()
select {
case <-timer.C:
case <-ctx.Done():
return "", ctx.Err()
}
if m.fail {
return "", fmt.Errorf("%s down", m.name)
}
return "data from " + m.name, nil
}
// fastest returns the first SUCCESSFUL result, like Promise.any; if every
// attempt fails it returns all the errors joined together, like AggregateError.
func fastest(mirrors []mirror) (string, error) {
ctx, cancel := context.WithCancel(context.Background())
defer cancel() // cancel the losing goroutines once there's a winner
type result struct {
data string
err error
}
// buffered with room for every goroutine: losers send and exit instead of
// blocking forever with nobody receiving (a goroutine leak)
results := make(chan result, len(mirrors))
for _, m := range mirrors {
go func() {
data, err := fetchFrom(ctx, m)
results <- result{data, err}
}()
}
errs := make([]error, 0, len(mirrors))
for range mirrors {
r := <-results
if r.err == nil {
return r.data, nil
}
errs = append(errs, r.err)
}
return "", errors.Join(errs...)
}
func main() {
data, err := fastest([]mirror{
{"eu", 300 * time.Millisecond, false},
{"us", 100 * time.Millisecond, false},
{"asia", 50 * time.Millisecond, true}, // fails first, but is ignored
})
fmt.Println(data, err) // → data from us <nil>
_, err = fastest([]mirror{
{"eu", 30 * time.Millisecond, true},
{"us", 10 * time.Millisecond, true},
})
fmt.Printf("%q\n", err.Error()) // → "us down\neu down" (in the order the errors arrived)
}Key differences
| Node.js | Go | |
|---|---|---|
| First error rejects immediately (no waiting, no cancelling) | Promise.all |
no direct equivalent |
| Wait for all, then return the first error | Promise.allSettled + rethrow (the all function above) |
errgroup.Group.Wait() |
| First error cancels siblings | AbortController passed in by hand |
errgroup.WithContext |
| Limit concurrent tasks | hand-rolled (or a library like p-limit) |
g.SetLimit(n) |
| Collect every result and error | Promise.allSettled |
sync.WaitGroup + errors.Join |
| First success | Promise.any → AggregateError |
buffered channel + cancel() → errors.Join |
| First to settle | Promise.race |
first value received from the channel |