summaryrefslogtreecommitdiff
path: root/sync/chan_fan.go
diff options
context:
space:
mode:
Diffstat (limited to 'sync/chan_fan.go')
-rw-r--r--sync/chan_fan.go79
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
-}