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 }