diff options
| author | DJ O'Leary <dijitol@proton.me> | 2026-07-08 01:06:20 +0200 |
|---|---|---|
| committer | DJ O'Leary <dijitol@proton.me> | 2026-07-08 01:06:20 +0200 |
| commit | 4d088c9354c17c58478cf121349a314b4e0dd322 (patch) | |
| tree | 4ad0635107bd4ec741ab82fd738662bc32ae073b /sync/chan_pipeline.go | |
| parent | 970b9f872254612b673c3a09579ed902cade3c99 (diff) | |
set up and add sync package
Diffstat (limited to 'sync/chan_pipeline.go')
| -rw-r--r-- | sync/chan_pipeline.go | 48 |
1 files changed, 48 insertions, 0 deletions
diff --git a/sync/chan_pipeline.go b/sync/chan_pipeline.go new file mode 100644 index 0000000..6b5f37f --- /dev/null +++ b/sync/chan_pipeline.go @@ -0,0 +1,48 @@ +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 +} |
