Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 53 additions & 0 deletions app/app.go
Original file line number Diff line number Diff line change
@@ -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
}
15 changes: 15 additions & 0 deletions app/go.mod
Original file line number Diff line number Diff line change
@@ -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
)
15 changes: 15 additions & 0 deletions app/go.sum
Original file line number Diff line number Diff line change
@@ -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=
57 changes: 57 additions & 0 deletions app/run_test.go
Original file line number Diff line number Diff line change
@@ -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
}
}
75 changes: 75 additions & 0 deletions closer/closer.go
Original file line number Diff line number Diff line change
@@ -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)
}
126 changes: 126 additions & 0 deletions closer/closer_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
12 changes: 12 additions & 0 deletions closer/go.mod
Original file line number Diff line number Diff line change
@@ -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
)
13 changes: 13 additions & 0 deletions closer/go.sum
Original file line number Diff line number Diff line change
@@ -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=
4 changes: 3 additions & 1 deletion go.work
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
go 1.25.0
go 1.26.3

use (
./app
./caller
./closer
./integration
./registry
)
Loading