diff options
| author | DJ O'Leary <dijitol@proton.me> | 2026-07-14 02:48:22 +0200 |
|---|---|---|
| committer | DJ O'Leary <dijitol@proton.me> | 2026-07-14 02:48:22 +0200 |
| commit | 03878d151215277d948b0a1ce31c65d9305d3aa4 (patch) | |
| tree | ac3136fa6f85f85c54a74e7acbd7c83c325bd594 /pkg/sync/chan_pipeline.go | |
| parent | 783ab544c03991419a7bf7ac033bdf5d34ff9cc3 (diff) | |
refactor(pkg): move packages that can be imported by others to pkg directory
Diffstat (limited to 'pkg/sync/chan_pipeline.go')
| -rw-r--r-- | pkg/sync/chan_pipeline.go | 48 |
1 files changed, 48 insertions, 0 deletions
diff --git a/pkg/sync/chan_pipeline.go b/pkg/sync/chan_pipeline.go new file mode 100644 index 0000000..6b5f37f --- /dev/null +++ b/pkg/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 +} |
