package sync import ( "context" "sync" ) func Pipe[IN, OUT any]( ctx context.Context, in <-chan IN, transform func(context.Context, IN) OUT, workerCount int, buffer int, ) <-chan OUT { if in == nil { // we only expect this to happen during development panic("in-channel is nil") } if workerCount < 1 { workerCount = 1 } if buffer < 0 { buffer = 0 } wg := sync.WaitGroup{} out := make(chan OUT, buffer) for range workerCount { wg.Go(func() { for v := range OrDone(ctx, in) { w := transform(ctx, v) if !CancelOrSend(ctx, out, w) { return } } }) } go func() { wg.Wait() close(out) }() return out }