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/README.md | 10 ++++++ pkg/sync/chan.go | 48 ++++++++++++++++++++++++++++ pkg/sync/chan_fan.go | 79 +++++++++++++++++++++++++++++++++++++++++++++++ pkg/sync/chan_pipeline.go | 48 ++++++++++++++++++++++++++++ sync/README.md | 10 ------ sync/chan.go | 48 ---------------------------- sync/chan_fan.go | 79 ----------------------------------------------- sync/chan_pipeline.go | 48 ---------------------------- 8 files changed, 185 insertions(+), 185 deletions(-) create mode 100644 pkg/sync/README.md create mode 100644 pkg/sync/chan.go create mode 100644 pkg/sync/chan_fan.go create mode 100644 pkg/sync/chan_pipeline.go delete mode 100644 sync/README.md delete mode 100644 sync/chan.go delete mode 100644 sync/chan_fan.go delete mode 100644 sync/chan_pipeline.go 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 // , 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 +} diff --git a/sync/README.md b/sync/README.md deleted file mode 100644 index 93a4526..0000000 --- a/sync/README.md +++ /dev/null @@ -1,10 +0,0 @@ -# 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/sync/chan.go b/sync/chan.go deleted file mode 100644 index 91501f9..0000000 --- a/sync/chan.go +++ /dev/null @@ -1,48 +0,0 @@ -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 // , 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/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 -} diff --git a/sync/chan_pipeline.go b/sync/chan_pipeline.go deleted file mode 100644 index 6b5f37f..0000000 --- a/sync/chan_pipeline.go +++ /dev/null @@ -1,48 +0,0 @@ -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 -} -- cgit v1.2.3