diff options
| author | DJ O'Leary <dijitol@proton.me> | 2026-07-14 02:48:22 +0200 |
|---|---|---|
| committer | DJ O'Leary <dijitol@proton.me> | 2026-07-14 02:48:22 +0200 |
| commit | 03878d151215277d948b0a1ce31c65d9305d3aa4 (patch) | |
| tree | ac3136fa6f85f85c54a74e7acbd7c83c325bd594 /pkg/sync | |
| parent | 783ab544c03991419a7bf7ac033bdf5d34ff9cc3 (diff) | |
refactor(pkg): move packages that can be imported by others to pkg directory
Diffstat (limited to 'pkg/sync')
| -rw-r--r-- | pkg/sync/README.md | 10 | ||||
| -rw-r--r-- | pkg/sync/chan.go | 48 | ||||
| -rw-r--r-- | pkg/sync/chan_fan.go | 79 | ||||
| -rw-r--r-- | pkg/sync/chan_pipeline.go | 48 |
4 files changed, 185 insertions, 0 deletions
diff --git a/pkg/sync/README.md b/pkg/sync/README.md new file mode 100644 index 0000000..93a4526 --- /dev/null +++ b/pkg/sync/README.md @@ -0,0 +1,10 @@ +# Sync + +## Resources + +- https://gobyexample.com/channels +- https://go.dev/tour/concurrency/2 +- https://www.oreilly.com/library/view/concurrency-in-go/9781491941294/ +- https://youtu.be/gz4DKBKe58A?si=Ks5VhvUMNaa2NQ60 +- https://go.dev/doc/effective_go#channels +- https://dave.cheney.net/2014/03/19/channel-axioms diff --git a/pkg/sync/chan.go b/pkg/sync/chan.go new file mode 100644 index 0000000..91501f9 --- /dev/null +++ b/pkg/sync/chan.go @@ -0,0 +1,48 @@ +package sync + +import "context" + +func OrDone[T any](ctx context.Context, in <-chan T) <-chan T { + if in == nil { + // we only expect this to happen during development + panic("in-channel is nil") + } + + out := make(chan T) + go func() { + defer close(out) + for { + v, ok := CancelOrReceive(ctx, in) + if !ok { + return + } + + if !CancelOrSend(ctx, out, v) { + return + } + } + }() + return out +} + +// CancelOrReceive is trying to solve the same problem as OrDone but does so +// without an extra goroutine. +func CancelOrReceive[T any](ctx context.Context, in <-chan T) (v T, ok bool) { + select { + case <-ctx.Done(): + return // <zero value of type T>, false + case v, ok = <-in: + return // v, ok + } +} + +// CancelOrSend is a helper function to simplify sending values to channels +// while keeping the context in mind. +func CancelOrSend[T any](ctx context.Context, out chan<- T, val T) bool { + select { + case <-ctx.Done(): + return false + case out <- val: + return true + } +} 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 +} diff --git a/pkg/sync/chan_pipeline.go b/pkg/sync/chan_pipeline.go new file mode 100644 index 0000000..6b5f37f --- /dev/null +++ b/pkg/sync/chan_pipeline.go @@ -0,0 +1,48 @@ +package sync + +import ( + "context" + "sync" +) + +func Pipe[IN, OUT any]( + ctx context.Context, + in <-chan IN, + transform func(context.Context, IN) OUT, + workerCount int, + buffer int, +) <-chan OUT { + if in == nil { + // we only expect this to happen during development + panic("in-channel is nil") + } + if workerCount < 1 { + workerCount = 1 + } + if buffer < 0 { + buffer = 0 + } + + wg := sync.WaitGroup{} + + out := make(chan OUT, buffer) + + for range workerCount { + wg.Go(func() { + for v := range OrDone(ctx, in) { + w := transform(ctx, v) + + if !CancelOrSend(ctx, out, w) { + return + } + } + }) + } + + go func() { + wg.Wait() + close(out) + }() + + return out +} |
