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_fan.go | |
| parent | 970b9f872254612b673c3a09579ed902cade3c99 (diff) | |
set up and add sync package
Diffstat (limited to 'sync/chan_fan.go')
| -rw-r--r-- | sync/chan_fan.go | 79 |
1 files changed, 79 insertions, 0 deletions
diff --git a/sync/chan_fan.go b/sync/chan_fan.go new file mode 100644 index 0000000..13957ed --- /dev/null +++ b/sync/chan_fan.go @@ -0,0 +1,79 @@ +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 +} |
