summaryrefslogtreecommitdiff
path: root/pkg/sync
diff options
context:
space:
mode:
authorDJ O'Leary <dijitol@proton.me>2026-07-14 02:48:22 +0200
committerDJ O'Leary <dijitol@proton.me>2026-07-14 02:48:22 +0200
commit03878d151215277d948b0a1ce31c65d9305d3aa4 (patch)
treeac3136fa6f85f85c54a74e7acbd7c83c325bd594 /pkg/sync
parent783ab544c03991419a7bf7ac033bdf5d34ff9cc3 (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.md10
-rw-r--r--pkg/sync/chan.go48
-rw-r--r--pkg/sync/chan_fan.go79
-rw-r--r--pkg/sync/chan_pipeline.go48
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
+}