Go

Implement a bounded, cancellable fan-out that returns the first error and does not leak goroutines.

Question 175HardGo 1.22 to 1.25

Requirements: limit concurrency, stop the remaining work at the first error, wait for every goroutine to finish, and never block a sender forever. Without errgroup:

func FanOut[T any](ctx context.Context, items []T, limit int,
	fn func(context.Context, T) error) error {
	ctx, cancel := context.WithCancelCause(ctx)
	defer cancel(nil)

	sem := make(chan struct{}, limit)
	var wg sync.WaitGroup
loop:
	for _, it := range items {
		select {
		case sem <- struct{}{}:
		case <-ctx.Done():
			break loop // stop launching new work
		}
		wg.Go(func() {
			defer func() { <-sem }()
			if err := fn(ctx, it); err != nil {
				cancel(err) // first cause wins; later calls are no-ops
			}
		})
	}
	wg.Wait()
	if err := context.Cause(ctx); err != nil && ctx.Err() != nil {
		return err
	}
	return nil
}

Key points:

  • The semaphore is acquired before go. This bounds the number of live goroutines, not just the number doing work.
  • WithCancelCause records the first error without a mutex.
  • fn must respect ctx, otherwise cancellation cannot stop it.
  • wg.Wait() guarantees there are no stray goroutines when the function returns.

In real code, errgroup.WithContext plus SetLimit does exactly this.

More on Goroutines & the Scheduler

All 35 Goroutines & the Scheduler questions