concurrency-semaphore-pool
The implementation uses a buffered channel as a semaphore to bound concurrency, and a pending-map reorder buffer to emit results in submission order despite out-of-order completion. Context cancellation propagates to all in-flight workers, and the first job error cancels the remaining work.
┌──────────────────┐
jobs[] ──►│ semaphore (chan)│◄── at most N tokens
└──────┬───────────┘
│ launch goroutines
▼
┌──────────────┐
│ jobFn(ctx) │ (each acquires + releases a token)
└──────┬───────┘
│ result{idx, val, err}
▼
┌──────────────┐
│ pending map │ reorder by index
└──────┬───────┘
▼
results[] in order
pool.gopackage semaphorepool
import (
"context"
"sync"
)
// SemaphorePool is a generic worker pool that runs at most maxWorkers
// concurrent goroutines.
type SemaphorePool[T, R any] struct {
maxWorkers int
jobFn func(context.Context, T) (R, error)
}
// New creates a pool that runs at most n goroutines concurrently.
func New[T, R any](n int, fn func(context.Context, T) (R, error)) *SemaphorePool[T, R] {
return &SemaphorePool[T, R]{maxWorkers: n, jobFn: fn}
}
type result[T any] struct {
idx int
val T
err error
}
// Process runs all jobs with bounded concurrency and returns results
// in the same order as the input slice.
func (p *SemaphorePool[T, R]) Process(ctx context.Context, jobs []T) ([]R, error) {
if len(jobs) == 0 {
return nil, nil
}
n := len(jobs)
results := make([]R, n)
sem := make(chan struct{}, p.maxWorkers)
resCh := make(chan result[R], n)
ctx, cancel := context.WithCancel(ctx)
defer cancel()
var wg sync.WaitGroup
var firstErr error
var setErrOnce sync.Once
// Submit all jobs, acquiring semaphore tokens
for i, job := range jobs {
if ctx.Err() != nil {
break
}
select {
case sem <- struct{}{}:
case <-ctx.Done():
break
}
wg.Add(1)
go func(idx int, job T) {
defer func() {
<-sem
wg.Done()
}()
if ctx.Err() != nil {
resCh <- result[R]{idx: idx, err: ctx.Err()}
return
}
val, err := p.jobFn(ctx, job)
if err != nil {
setErrOnce.Do(func() {
firstErr = err
cancel()
})
}
resCh <- result[R]{idx: idx, val: val, err: err}
}(i, job)
}
go func() {
wg.Wait()
close(resCh)
}()
// Reorder results using a pending map keyed by index
pending := make(map[int]result[R])
nextIdx := 0
for r := range resCh {
pending[r.idx] = r
for {
pr, ok := pending[nextIdx]
if !ok {
break
}
results[pr.idx] = pr.val
if pr.err != nil && firstErr == nil {
firstErr = pr.err
}
delete(pending, nextIdx)
nextIdx++
}
}
if firstErr == nil {
firstErr = ctx.Err()
}
return results, firstErr
}
| Concern | Solution |
|---|---|
| Concurrency limit | sem <- struct{}{} acquire / <-sem release — a buffered channel acts as a counting semaphore |
| Result ordering | Each worker tags its output with the input index; a pending map[int]result + nextIdx counter drains in sequential order |
| Context cancellation | Derive a child context via WithCancel; on external cancel or first job error, call cancel() so all workers see ctx.Done() |
| First-error semantics | sync.Once captures the first non-nil error and cancels the rest; partial results are still returned |
| Generics | SemaphorePool[T, R] works with any job input/output types |
### Test summary (13 tests, all pass) ``` === RUN TestEmptyJobs --- PASS === RUN TestSingleJob --- PASS === RUN TestOrderPreserved --- PASS === RUN TestOrderWithVaryingDelays --- PASS === RUN TestConcurrencyLimit --- PASS === RUN TestContextCancel --- PASS === RUN TestContextCancelDuringWork --- PASS === RUN TestFirstErrorCancelsOthers --- PASS === RUN TestDeadlineExceeded --- PASS === RUN TestManyJobs --- PASS === RUN TestSingleConcurrency --- PASS === RUN TestAllJobsError --- PASS === RUN TestGenericTypes --- PASS ``` ### Edge cases covered | Test | What it verifies | |---|---| | `TestEmptyJobs` | Nil and empty input slices return nil/empty without calling jobFn | | `TestOrderPreserved` | 10 jobs with varying sleep times still return results in submission order | | `TestOrderWithVaryingDelays` | Descending sleep times `[5,4,3,2,1]` ms — workers finish in reverse order, results stay correct | | `TestConcurrencyLimit` | Measures `maxSeen` concurrent goroutines; never exceeds the limit of 3 | | `TestSingleConcurrency` | `maxWorkers=1` forces sequential execution; maxSeen=1 | | `TestContextCancel` | Context cancelled before `Process` begins — all results are zero and error is `context.Canceled` | | `TestContextCancelDuringWork` | Short 50ms timeout with 200ms jobs — returns partial results + deadline error | | `TestFirstErrorCancelsOthers` | Job at index 2 fails — first error propagates, remaining workers cancelled | | `TestDeadlineExceeded` | Explicit `WithTimeout(50ms)` vs 200ms jobs — returns `DeadlineExceeded` | | `TestManyJobs` | 1000 jobs with concurrency 10 — all results correct and in order | | `TestAllJobsError` | Every job returns an error — first error returned, all result slots filled | | `TestGenericTypes` | `SemaphorePool[string, output]` — generic parameters work with struct types | ---
{"model": "gemini-2.5-pro", "problem_class": "concurrency-semaphore-pool", "result": "passed", "tests": 13}