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 /sync/chan_fan.go | |
| parent | 783ab544c03991419a7bf7ac033bdf5d34ff9cc3 (diff) | |
refactor(pkg): move packages that can be imported by others to pkg directory
Diffstat (limited to 'sync/chan_fan.go')
| -rw-r--r-- | sync/chan_fan.go | 79 |
1 files changed, 0 insertions, 79 deletions
diff --git a/sync/chan_fan.go b/sync/chan_fan.go deleted file mode 100644 index 13957ed..0000000 --- a/sync/chan_fan.go +++ /dev/null @@ -1,79 +0,0 @@ -package sync - -import ( - "context" - "sync" -) - -func FanOut[T any](ctx context.Context, in <-chan T, channelCount int, buffer int) []<-chan T { - if in == nil { - // we only expect this to happen during development - panic("in-channel is nil") - } - if channelCount < 1 { - channelCount = 1 - } - if buffer < 0 { - buffer = 0 - } - - chs := make([]chan T, channelCount) - for i := range channelCount { - chs[i] = make(chan T, buffer) - } - - go func() { - defer func() { - for _, out := range chs { - close(out) - } - }() - - for v := range OrDone(ctx, in) { - for _, out := range chs { - // N.b. this will block if *any* of the out channels has a full buffer - if !CancelOrSend(ctx, out, v) { - return - } - } - } - }() - - outs := make([]<-chan T, channelCount) - for i, out := range chs { - outs[i] = out - } - - return outs -} - -func FanIn[T any](ctx context.Context, buffer int, ins ...<-chan T) <-chan T { - if buffer < 0 { - buffer = 0 - } - - wg := sync.WaitGroup{} - - out := make(chan T, buffer) - - for _, in := range ins { - wg.Go(func() { - if in == nil { - return - } - - for v := range OrDone(ctx, in) { - if !CancelOrSend(ctx, out, v) { - return - } - } - }) - } - - go func() { - wg.Wait() - close(out) - }() - - return out -} |
