diff --git a/app/app.go b/app/app.go new file mode 100644 index 0000000..4e1c7da --- /dev/null +++ b/app/app.go @@ -0,0 +1,53 @@ +package app + +import ( + "context" + "fmt" + "log/slog" + + "github.com/draincloud/logger" + "golang.org/x/sync/errgroup" +) + +type Runnable interface { + Run(ctx context.Context) error +} + +type App struct { + name string + runnables []Runnable +} + +func NewApp( + name string, + runnables ...Runnable, +) *App { + return &App{ + name: name, + runnables: runnables, + } +} + +func (a *App) Run(ctx context.Context) error { + ctx = logger.WithAttrs(ctx, slog.String("app", a.name)) + logger.Warn(ctx, "[App][Run] sstarting app") + + eg, egCtx := errgroup.WithContext(ctx) + + runCtx, cancel := context.WithCancel(egCtx) + defer cancel() + + for _, r := range a.runnables { + eg.Go(func() error { + defer cancel() + + return r.Run(runCtx) + }) + } + + if err := eg.Wait(); err != nil { + return fmt.Errorf("[app][Run] %s: %w", a.name, err) + } + + return nil +} diff --git a/app/go.mod b/app/go.mod new file mode 100644 index 0000000..b5920e3 --- /dev/null +++ b/app/go.mod @@ -0,0 +1,15 @@ +module github.com/draincloud/callpack/app + +go 1.26.3 + +require ( + github.com/draincloud/logger v0.0.5 + golang.org/x/sync v0.22.0 +) + +require ( + github.com/fatih/color v1.18.0 // indirect + github.com/mattn/go-colorable v0.1.13 // indirect + github.com/mattn/go-isatty v0.0.20 // indirect + golang.org/x/sys v0.25.0 // indirect +) diff --git a/app/go.sum b/app/go.sum new file mode 100644 index 0000000..3fcd09b --- /dev/null +++ b/app/go.sum @@ -0,0 +1,15 @@ +github.com/draincloud/logger v0.0.5 h1:4saeda/sm5E6S0lnfwCIdgmqZYiIGuTIUN0HpIlm3js= +github.com/draincloud/logger v0.0.5/go.mod h1:Z/GP5qHAC+MrhTMs5FckdQahiP2MjugaeQxjaMMaqWY= +github.com/fatih/color v1.18.0 h1:S8gINlzdQ840/4pfAwic/ZE0djQEH3wM94VfqLTZcOM= +github.com/fatih/color v1.18.0/go.mod h1:4FelSpRwEGDpQ12mAdzqdOukCy4u8WUtOY6lkT/6HfU= +github.com/mattn/go-colorable v0.1.13 h1:fFA4WZxdEF4tXPZVKMLwD8oUnCTTo08duU7wxecdEvA= +github.com/mattn/go-colorable v0.1.13/go.mod h1:7S9/ev0klgBDR4GtXTXX8a3vIGJpMovkB8vQcUbaXHg= +github.com/mattn/go-isatty v0.0.16/go.mod h1:kYGgaQfpe5nmfYZH+SKPsOc2e4SrIfOl2e/yFXSvRLM= +github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= +github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek= +golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.25.0 h1:r+8e+loiHxRqhXVl6ML1nO3l1+oFoWbnlu2Ehimmi34= +golang.org/x/sys v0.25.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= diff --git a/app/run_test.go b/app/run_test.go new file mode 100644 index 0000000..3db236c --- /dev/null +++ b/app/run_test.go @@ -0,0 +1,57 @@ +package app_test + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/draincloud/callpack/app" +) + +type runnableFunc func(ctx context.Context) error + +func (f runnableFunc) Run(ctx context.Context) error { return f(ctx) } + +func TestRunReturnsWhenRunnableExitsCleanly(t *testing.T) { + t.Parallel() + + blocked := runnableFunc(func(ctx context.Context) error { + <-ctx.Done() + + return nil + }) + oneShot := runnableFunc(func(context.Context) error { return nil }) + + if err := runWithin(t, time.Second, app.NewApp("test", blocked, oneShot)); err != nil { + t.Fatalf("Run: %v", err) + } +} + +func TestRunWrapsRunnableError(t *testing.T) { + t.Parallel() + + errBoom := errors.New("boom") + failing := runnableFunc(func(context.Context) error { return errBoom }) + + err := runWithin(t, time.Second, app.NewApp("test", failing)) + if !errors.Is(err, errBoom) { + t.Fatalf("Run: got %v, want %v", err, errBoom) + } +} + +func runWithin(t *testing.T, d time.Duration, a *app.App) error { + t.Helper() + + done := make(chan error, 1) + go func() { done <- a.Run(t.Context()) }() + + select { + case err := <-done: + return err + case <-time.After(d): + t.Fatal("Run did not return") + + return nil + } +} diff --git a/closer/closer.go b/closer/closer.go new file mode 100644 index 0000000..10b58f9 --- /dev/null +++ b/closer/closer.go @@ -0,0 +1,75 @@ +package closer + +import ( + "context" + "errors" + "sync" + + "github.com/draincloud/logger" +) + +type CloseFunc func(ctx context.Context) error + +// ErrClosed is returned by Add once Close has started; the cleanup is not registered. +var ErrClosed = errors.New("[closer] already closed") + +var globalCloser = &Closer{} + +type Closer struct { + mu sync.Mutex + closed bool + done chan struct{} + err error + closeFns []CloseFunc +} + +func (c *Closer) Add(fn CloseFunc) error { + c.mu.Lock() + defer c.mu.Unlock() + + if c.closed { + return ErrClosed + } + c.closeFns = append(c.closeFns, fn) + + return nil +} + +// Close runs the registered cleanups in reverse registration order, so dependents +// are torn down before the resources they use. Concurrent callers block until the +// first Close finishes and receive its result. +func (c *Closer) Close(ctx context.Context) error { + c.mu.Lock() + if c.closed { + done := c.done + c.mu.Unlock() + <-done + + return c.err + } + c.closed = true + c.done = make(chan struct{}) + done, fns := c.done, c.closeFns + c.mu.Unlock() + + var commonErr error + for i := len(fns) - 1; i >= 0; i-- { + if err := fns[i](ctx); err != nil { + logger.Error(ctx, "[closer][Close] error at close func call", logger.Err(err)) + commonErr = errors.Join(commonErr, err) + } + } + + c.err = commonErr + close(done) + + return commonErr +} + +func Add(fn CloseFunc) error { + return globalCloser.Add(fn) +} + +func Close(ctx context.Context) error { + return globalCloser.Close(ctx) +} diff --git a/closer/closer_test.go b/closer/closer_test.go new file mode 100644 index 0000000..be003f3 --- /dev/null +++ b/closer/closer_test.go @@ -0,0 +1,126 @@ +package closer_test + +import ( + "context" + "errors" + "slices" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/draincloud/callpack/closer" +) + +func TestCloserAddConcurrent(t *testing.T) { + t.Parallel() + + const n = 50 + c := &closer.Closer{} + var calls atomic.Int64 + var wg sync.WaitGroup + wg.Add(n) + for range n { + go func() { + defer wg.Done() + if err := c.Add(func(context.Context) error { + calls.Add(1) + + return nil + }); err != nil { + t.Error(err) + } + }() + } + wg.Wait() + + if err := c.Close(t.Context()); err != nil { + t.Fatalf("Close: %v", err) + } + if got := calls.Load(); got != n { + t.Fatalf("ran %d cleanups, want %d", got, n) + } +} + +func TestCloserAddAfterClose(t *testing.T) { + t.Parallel() + + c := &closer.Closer{} + if err := c.Close(t.Context()); err != nil { + t.Fatalf("Close: %v", err) + } + + if err := c.Add(func(context.Context) error { return nil }); !errors.Is(err, closer.ErrClosed) { + t.Fatalf("Add after Close: got %v, want ErrClosed", err) + } +} + +func TestCloserClosesInReverseOrder(t *testing.T) { + t.Parallel() + + c := &closer.Closer{} + var order []int + for i := range 3 { + if err := c.Add(func(context.Context) error { + order = append(order, i) + + return nil + }); err != nil { + t.Fatal(err) + } + } + + if err := c.Close(t.Context()); err != nil { + t.Fatalf("Close: %v", err) + } + if want := []int{2, 1, 0}; !slices.Equal(order, want) { + t.Fatalf("close order %v, want %v", order, want) + } +} + +func TestCloserConcurrentCloseWaits(t *testing.T) { + t.Parallel() + + errCleanup := errors.New("cleanup failed") + released := make(chan struct{}) + var done atomic.Bool + + c := &closer.Closer{} + if err := c.Add(func(context.Context) error { + <-released + done.Store(true) + + return errCleanup + }); err != nil { + t.Fatal(err) + } + + first := make(chan error, 1) + go func() { first <- c.Close(t.Context()) }() + + // Let the first Close reach the blocking cleanup before the second starts. + time.Sleep(10 * time.Millisecond) + second := make(chan error, 1) + go func() { second <- c.Close(t.Context()) }() + + select { + case err := <-second: + t.Fatalf("second Close returned %v while cleanups were still running", err) + case <-time.After(50 * time.Millisecond): + } + + close(released) + for _, ch := range []chan error{first, second} { + select { + case err := <-ch: + if !errors.Is(err, errCleanup) { + t.Fatalf("Close: got %v, want %v", err, errCleanup) + } + case <-time.After(time.Second): + t.Fatal("Close did not return") + } + } + if !done.Load() { + t.Fatal("cleanup did not finish") + } +} diff --git a/closer/go.mod b/closer/go.mod new file mode 100644 index 0000000..b23a30b --- /dev/null +++ b/closer/go.mod @@ -0,0 +1,12 @@ +module github.com/draincloud/callpack/closer + +go 1.26.3 + +require github.com/draincloud/logger v0.0.5 + +require ( + github.com/fatih/color v1.18.0 // indirect + github.com/mattn/go-colorable v0.1.13 // indirect + github.com/mattn/go-isatty v0.0.20 // indirect + golang.org/x/sys v0.25.0 // indirect +) diff --git a/closer/go.sum b/closer/go.sum new file mode 100644 index 0000000..e6720bd --- /dev/null +++ b/closer/go.sum @@ -0,0 +1,13 @@ +github.com/draincloud/logger v0.0.5 h1:4saeda/sm5E6S0lnfwCIdgmqZYiIGuTIUN0HpIlm3js= +github.com/draincloud/logger v0.0.5/go.mod h1:Z/GP5qHAC+MrhTMs5FckdQahiP2MjugaeQxjaMMaqWY= +github.com/fatih/color v1.18.0 h1:S8gINlzdQ840/4pfAwic/ZE0djQEH3wM94VfqLTZcOM= +github.com/fatih/color v1.18.0/go.mod h1:4FelSpRwEGDpQ12mAdzqdOukCy4u8WUtOY6lkT/6HfU= +github.com/mattn/go-colorable v0.1.13 h1:fFA4WZxdEF4tXPZVKMLwD8oUnCTTo08duU7wxecdEvA= +github.com/mattn/go-colorable v0.1.13/go.mod h1:7S9/ev0klgBDR4GtXTXX8a3vIGJpMovkB8vQcUbaXHg= +github.com/mattn/go-isatty v0.0.16/go.mod h1:kYGgaQfpe5nmfYZH+SKPsOc2e4SrIfOl2e/yFXSvRLM= +github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= +github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.25.0 h1:r+8e+loiHxRqhXVl6ML1nO3l1+oFoWbnlu2Ehimmi34= +golang.org/x/sys v0.25.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= diff --git a/go.work b/go.work index 0ce97a0..d189fae 100644 --- a/go.work +++ b/go.work @@ -1,7 +1,9 @@ -go 1.25.0 +go 1.26.3 use ( + ./app ./caller + ./closer ./integration ./registry )