From 03878d151215277d948b0a1ce31c65d9305d3aa4 Mon Sep 17 00:00:00 2001 From: DJ O'Leary Date: Tue, 14 Jul 2026 02:48:22 +0200 Subject: refactor(pkg): move packages that can be imported by others to pkg directory --- pkg/sync/chan_fan.go | 79 ++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 79 insertions(+) create mode 100644 pkg/sync/chan_fan.go (limited to 'pkg/sync/chan_fan.go') diff --git a/pkg/sync/chan_fan.go b/pkg/sync/chan_fan.go new file mode 100644 index 0000000..13957ed --- /dev/null +++ b/pkg/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 +} -- cgit v1.2.3