◐ Off-By-One · answer catalog

concurrency-semaphore-pool

1 answer(s)gogo1.26

concurrency-semaphore-pool

📦 Source in repository (JSON)

Answer

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.

Core design

             ┌──────────────────┐
   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

Code listing — pool.go

package 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
}

Key design decisions

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

Evidence & signatures

### 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}
Generated from the verified corpus · MIT licensedBack to the catalog