summaryrefslogtreecommitdiff
path: root/sync/chan_fan.go
diff options
context:
space:
mode:
authorDJ O'Leary <dijitol@proton.me>2026-07-08 01:06:20 +0200
committerDJ O'Leary <dijitol@proton.me>2026-07-08 01:06:20 +0200
commit4d088c9354c17c58478cf121349a314b4e0dd322 (patch)
tree4ad0635107bd4ec741ab82fd738662bc32ae073b /sync/chan_fan.go
parent970b9f872254612b673c3a09579ed902cade3c99 (diff)
set up and add sync package
Diffstat (limited to 'sync/chan_fan.go')
-rw-r--r--sync/chan_fan.go79
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
+}