From 2e2d8c483d75402ebfc50b286494d72d2bf3bdd8 Mon Sep 17 00:00:00 2001 From: comicrime Date: Mon, 7 Sep 2026 00:05:00 +0300 Subject: [PATCH 1/2] + safegroup --- safegroup/go.mod | 3 ++ safegroup/safegroup.go | 109 +++++++++++++++++++++++++++++++++++++++++ 2 files changed, 112 insertions(+) create mode 100644 safegroup/go.mod create mode 100644 safegroup/safegroup.go diff --git a/safegroup/go.mod b/safegroup/go.mod new file mode 100644 index 0000000..4f40baa --- /dev/null +++ b/safegroup/go.mod @@ -0,0 +1,3 @@ +module github.com/draincloud/callpack/safegroup + +go 1.26.3 diff --git a/safegroup/safegroup.go b/safegroup/safegroup.go new file mode 100644 index 0000000..850ad91 --- /dev/null +++ b/safegroup/safegroup.go @@ -0,0 +1,109 @@ +package safegroup + +import ( + "context" + "errors" + "fmt" + "sync" +) + +var ErrPanic = errors.New("panic in a goroutine") + +type SafeGroup struct { + cancel func(error) + + wg sync.WaitGroup + + sem chan token + + errOnce sync.Once + err error +} + +type token struct{} + +func (g *SafeGroup) done() { + r := recover() + if r != nil { + g.errOnce.Do(func() { + g.err = fmt.Errorf("safegroup: %w: %s", ErrPanic, r) + if g.cancel != nil { + g.cancel(g.err) + } + }) + } + + if g.sem != nil { + <-g.sem + } + g.wg.Done() +} + +func WithContext(ctx context.Context) (*SafeGroup, context.Context) { + ctx, cancel := context.WithCancelCause(ctx) + return &SafeGroup{cancel: cancel}, ctx +} + +func (g *SafeGroup) Wait() error { + g.wg.Wait() + if g.cancel != nil { + g.cancel(g.err) + } + return g.err +} + +func (g *SafeGroup) Go(f func() error) { + if g.sem != nil { + g.sem <- token{} + } + + g.wg.Add(1) + go func() { + defer g.done() + + if err := f(); err != nil { + g.errOnce.Do(func() { + g.err = err + if g.cancel != nil { + g.cancel(g.err) + } + }) + } + }() +} + +func (g *SafeGroup) TryGo(f func() error) bool { + if g.sem != nil { + select { + case g.sem <- token{}: + default: + return false + } + } + + g.wg.Add(1) + go func() { + defer g.done() + + if err := f(); err != nil { + g.errOnce.Do(func() { + g.err = err + if g.cancel != nil { + g.cancel(g.err) + } + }) + } + }() + return true +} + +func (g *SafeGroup) SetLimit(n int) { + if n < 0 { + g.sem = nil + return + } + if active := len(g.sem); active != 0 { + panic(fmt.Errorf("safegroup: modify limit while %v goroutines in the group are still active", active)) + } + g.sem = make(chan token, n) +} From e8a5fc0f510fc2846e18084b260a395cd6c7c022 Mon Sep 17 00:00:00 2001 From: comicrime Date: Mon, 7 Sep 2026 00:11:30 +0300 Subject: [PATCH 2/2] +1 --- safegroup/go.sum | 0 1 file changed, 0 insertions(+), 0 deletions(-) create mode 100644 safegroup/go.sum diff --git a/safegroup/go.sum b/safegroup/go.sum new file mode 100644 index 0000000..e69de29