Skip to content

Commit 8e46022

Browse files
authored
Merge pull request #48 from pior/naming
Simplify runnable naming and add FuncNamed
2 parents dcb5cce + 4e179c2 commit 8e46022

14 files changed

Lines changed: 130 additions & 110 deletions

README.md

Lines changed: 14 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -103,7 +103,7 @@ Components with dependencies will be stopped before their dependencies.
103103

104104
Example with three components:
105105
```go
106-
g := runnable.Manager(nil)
106+
g := runnable.NewManager()
107107
g.Add(jobQueue)
108108
g.Add(httpServer, jobQueue) // jobs is a dependency
109109
g.Add(monitor)
@@ -115,20 +115,20 @@ runnable.Run(g.Build())
115115
<summary markdown="span">Logs of a demo app</summary>
116116

117117
```shell
118-
$ go run ./cmd/example
119-
[RUNNABLE] 2020/10/22 22:42:26 INFO manager: main.JobQueue started
120-
[RUNNABLE] 2020/10/22 22:42:26 INFO manager: runnable.httpServer started
121-
[RUNNABLE] 2020/10/22 22:42:26 INFO manager: main.Monitor started
118+
$ go run ./examples/example
119+
level=INFO msg=started runnable=manager/JobQueue
120+
level=INFO msg=started runnable=manager/httpserver
121+
level=INFO msg=started runnable=manager/Monitor
122122
...
123-
^C[RUNNABLE] 2020/10/22 22:42:34 INFO signal: received signal interrupt
124-
[RUNNABLE] 2020/10/22 22:42:34 INFO manager: starting shutdown (context cancelled)
125-
[RUNNABLE] 2020/10/22 22:42:34 INFO manager: runnable.httpServer cancelled
126-
[RUNNABLE] 2020/10/22 22:42:34 INFO manager: main.Monitor cancelled
127-
[RUNNABLE] 2020/10/22 22:42:34 INFO manager: main.Monitor stopped
128-
[RUNNABLE] 2020/10/22 22:42:34 INFO manager: runnable.httpServer stopped
129-
[RUNNABLE] 2020/10/22 22:42:34 INFO manager: main.JobQueue cancelled
130-
[RUNNABLE] 2020/10/22 22:42:34 INFO manager: main.JobQueue stopped
131-
[RUNNABLE] 2020/10/22 22:42:34 INFO manager: shutdown complete
123+
^Clevel=INFO msg="received signal" runnable=signal/manager signal=interrupt
124+
level=INFO msg="starting shutdown" runnable=manager reason="context cancelled"
125+
level=INFO msg=cancelled runnable=manager/httpserver
126+
level=INFO msg=cancelled runnable=manager/Monitor
127+
level=INFO msg=stopped runnable=manager/Monitor
128+
level=INFO msg=stopped runnable=manager/httpserver
129+
level=INFO msg=cancelled runnable=manager/JobQueue
130+
level=INFO msg=stopped runnable=manager/JobQueue
131+
level=INFO msg="shutdown complete" runnable=manager
132132
```
133133

134134
</details>

closer.go

Lines changed: 10 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -6,40 +6,41 @@ import (
66

77
// Closer returns a runnable intended to call a Close method on shutdown.
88
func Closer(c interface{ Close() }) Runnable {
9-
return &closer{baseWrapper{"closer", c}, c, func(ctx context.Context) error {
9+
return &closer{"closer/" + runnableName(c), func(ctx context.Context) error {
1010
c.Close()
1111
return nil
1212
}}
1313
}
1414

15-
// Closer returns a runnable intended to call a Close method on shutdown.
15+
// CloserErr returns a runnable intended to call a Close method on shutdown.
1616
func CloserErr(c interface{ Close() error }) Runnable {
17-
return &closer{baseWrapper{"closer", c}, c, func(ctx context.Context) error {
17+
return &closer{"closer/" + runnableName(c), func(ctx context.Context) error {
1818
return c.Close()
1919
}}
2020
}
2121

22-
// Closer returns a runnable intended to call a Close method on shutdown.
22+
// CloserCtx returns a runnable intended to call a Close method on shutdown.
2323
func CloserCtx(c interface{ Close(context.Context) }) Runnable {
24-
return &closer{baseWrapper{"closer", c}, c, func(ctx context.Context) error {
24+
return &closer{"closer/" + runnableName(c), func(ctx context.Context) error {
2525
c.Close(ctx)
2626
return nil
2727
}}
2828
}
2929

30-
// Closer returns a runnable intended to call a Close method on shutdown.
30+
// CloserCtxErr returns a runnable intended to call a Close method on shutdown.
3131
func CloserCtxErr(c interface{ Close(context.Context) error }) Runnable {
32-
return &closer{baseWrapper{"closer", c}, c, func(ctx context.Context) error {
32+
return &closer{"closer/" + runnableName(c), func(ctx context.Context) error {
3333
return c.Close(ctx)
3434
}}
3535
}
3636

3737
type closer struct {
38-
baseWrapper
39-
c any
38+
name string
4039
closeFn func(context.Context) error
4140
}
4241

42+
func (c *closer) runnableName() string { return c.name }
43+
4344
func (c *closer) Run(ctx context.Context) error {
4445
<-ctx.Done()
4546
err := c.closeFn(ctx)

every.go

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -8,18 +8,20 @@ import (
88
// Every returns a runnable that will periodically run the runnable passed in argument.
99
func Every(runnable Runnable, period time.Duration) Runnable {
1010
return &every{
11-
baseWrapper{"every-" + period.String(), runnable},
12-
runnable,
13-
period,
11+
name: "every-" + period.String() + "/" + runnableName(runnable),
12+
runnable: runnable,
13+
period: period,
1414
}
1515
}
1616

1717
type every struct {
18-
baseWrapper
18+
name string
1919
runnable Runnable
2020
period time.Duration
2121
}
2222

23+
func (e *every) runnableName() string { return e.name }
24+
2325
func (e *every) Run(ctx context.Context) (err error) {
2426
ticker := time.NewTicker(e.period)
2527

example_test.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -65,7 +65,7 @@ func Example() {
6565
}
6666
g.Add(runnable.HTTPServer(server), jobs)
6767

68-
task := runnable.Func(func(ctx context.Context) error {
68+
task := runnable.FuncNamed("enqueue", func(ctx context.Context) error {
6969
_, _ = http.Post("http://127.0.0.1:8080/?id=1", "test/plain", nil)
7070
_, _ = http.Post("http://127.0.0.1:8080/?id=2", "test/plain", nil)
7171
_, _ = http.Post("http://127.0.0.1:8080/?id=3", "test/plain", nil)
@@ -81,12 +81,12 @@ func Example() {
8181

8282
// level=INFO msg=started runnable=manager/Jobs
8383
// level=INFO msg=started runnable=manager/httpserver
84-
// level=INFO msg=started runnable=manager/RunnableFunc
84+
// level=INFO msg=started runnable=manager/enqueue
8585
// level=INFO msg=started runnable=manager/every-1h0m0s/CleanupTask
8686
// level=INFO msg=listening runnable=httpserver addr=127.0.0.1:8080
8787
// Starting job 1
8888
// Completed job 1
8989
// ...
90-
// level=INFO msg="starting shutdown" runnable=manager reason="RunnableFunc died"
90+
// level=INFO msg="starting shutdown" runnable=manager reason="enqueue died"
9191
// level=INFO msg="shutdown complete" runnable=manager
9292
}

func.go

Lines changed: 18 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,16 +1,30 @@
11
package runnable
22

3-
import "context"
3+
import (
4+
"context"
5+
"reflect"
6+
"runtime"
7+
)
48

9+
// RunnableFunc is a function that implements the Runnable contract.
510
type RunnableFunc func(context.Context) error
611

712
type funcRunnable struct {
8-
baseWrapper
9-
fn RunnableFunc
13+
name string
14+
fn RunnableFunc
1015
}
1116

17+
func (f *funcRunnable) runnableName() string { return f.name }
18+
19+
// Func returns a Runnable from a function. The name is derived from the function using reflection.
1220
func Func(fn RunnableFunc) Runnable {
13-
return &funcRunnable{baseWrapper{"", fn}, fn}
21+
name := runtime.FuncForPC(reflect.ValueOf(fn).Pointer()).Name()
22+
return &funcRunnable{name, fn}
23+
}
24+
25+
// FuncNamed returns a named Runnable from a function.
26+
func FuncNamed(name string, fn RunnableFunc) Runnable {
27+
return &funcRunnable{name, fn}
1428
}
1529

1630
func (f *funcRunnable) Run(ctx context.Context) error {

manager.go

Lines changed: 18 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,12 @@ type AppManager interface {
1616
// ManagerOption configures the behavior of a Manager.
1717
type ManagerOption func(*manager)
1818

19+
func ManagerName(name string) ManagerOption {
20+
return func(m *manager) {
21+
m.name = name
22+
}
23+
}
24+
1925
func ManagerShutdownTimeout(dur time.Duration) ManagerOption {
2026
return func(m *manager) {
2127
m.shutdownTimeout = dur
@@ -26,6 +32,7 @@ func ManagerShutdownTimeout(dur time.Duration) ManagerOption {
2632
// Runnables can declare a dependency on another runnable. Dependencies are started first and stopped last.
2733
func NewManager(opts ...ManagerOption) AppManager {
2834
m := &manager{
35+
name: "manager",
2936
shutdownTimeout: 10 * time.Second,
3037
}
3138

@@ -40,10 +47,13 @@ func NewManager(opts ...ManagerOption) AppManager {
4047
}
4148

4249
type manager struct {
50+
name string
4351
containers []*managerContainer
4452
shutdownTimeout time.Duration
4553
}
4654

55+
func (m *manager) runnableName() string { return m.name }
56+
4757
func (m *manager) Add(runnable Runnable, dependencies ...Runnable) {
4858
container := m.insertRunnable(runnable)
4959
for _, dep := range dependencies {
@@ -80,15 +90,15 @@ func (m *manager) Run(ctx context.Context) error {
8090
// run the runnables in Go routines.
8191
for _, c := range m.containers {
8292
c.launch(completedChan, dying)
83-
logger.Info("started", "runnable", "manager/"+c.name())
93+
logger.Info("started", "runnable", m.runnableName()+"/"+c.name())
8494
}
8595

8696
// block until group is cancelled, or a runnable dies.
8797
select {
8898
case <-ctx.Done():
89-
logger.Info("starting shutdown", "runnable", "manager", "reason", "context cancelled")
99+
logger.Info("starting shutdown", "runnable", m.runnableName(), "reason", "context cancelled")
90100
case c := <-dying:
91-
logger.Info("starting shutdown", "runnable", "manager", "reason", c.name()+" died")
101+
logger.Info("starting shutdown", "runnable", m.runnableName(), "reason", c.name()+" died")
92102
}
93103

94104
// starting shutdown
@@ -117,7 +127,7 @@ func (m *manager) Run(ctx context.Context) error {
117127
}
118128

119129
if !cancelled.contains(c) {
120-
logger.Info("cancelled", "runnable", "manager/"+c.name())
130+
logger.Info("cancelled", "runnable", m.runnableName()+"/"+c.name())
121131
c.shutdown()
122132
cancelled.insert(c)
123133
}
@@ -129,9 +139,9 @@ func (m *manager) Run(ctx context.Context) error {
129139
completed.insert(c)
130140

131141
if c.err == nil || errors.Is(c.err, context.Canceled) {
132-
logger.Info("stopped", "runnable", "manager/"+c.name())
142+
logger.Info("stopped", "runnable", m.runnableName()+"/"+c.name())
133143
} else {
134-
logger.Info("stopped with error", "runnable", "manager/"+c.name(), "error", c.err)
144+
logger.Info("stopped with error", "runnable", m.runnableName()+"/"+c.name(), "error", c.err)
135145
}
136146

137147
if len(completed) == len(m.containers) {
@@ -146,15 +156,15 @@ func (m *manager) Run(ctx context.Context) error {
146156
errs := []string{}
147157
for _, c := range m.containers {
148158
if !completed.contains(c) {
149-
logger.Info("still running", "runnable", "manager/"+c.name())
159+
logger.Info("still running", "runnable", m.runnableName()+"/"+c.name())
150160
errs = append(errs, fmt.Sprintf("%s is still running", c.name()))
151161
}
152162
if c.err != nil && !errors.Is(c.err, context.Canceled) {
153163
errs = append(errs, fmt.Sprintf("%s crashed with %+v", c.name(), c.err))
154164
}
155165
}
156166

157-
logger.Info("shutdown complete", "runnable", "manager")
167+
logger.Info("shutdown complete", "runnable", m.runnableName())
158168

159169
if len(errs) != 0 {
160170
return fmt.Errorf("manager: %s", strings.Join(errs, ", "))

manager_container.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@ func (c *managerContainer) insertUser(container *managerContainer) {
3232
}
3333

3434
func (c *managerContainer) name() string {
35-
return findName(c.runnable)
35+
return runnableName(c.runnable)
3636
}
3737

3838
func (c *managerContainer) launch(completed chan *managerContainer, dying chan *managerContainer) {

recover.go

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -22,15 +22,17 @@ func (e *PanicError) Unwrap() error {
2222

2323
// Recover returns a runnable that recovers when a runnable panics and return an error to represent this panic.
2424
func Recover(runnable Runnable) Runnable {
25-
return &RecoverRunner{baseWrapper{"recover", runnable}, runnable}
25+
return &recoverRunner{"recover/" + runnableName(runnable), runnable}
2626
}
2727

28-
type RecoverRunner struct {
29-
baseWrapper
28+
type recoverRunner struct {
29+
name string
3030
runnable Runnable
3131
}
3232

33-
func (r *RecoverRunner) Run(ctx context.Context) (err error) {
33+
func (r *recoverRunner) runnableName() string { return r.name }
34+
35+
func (r *recoverRunner) Run(ctx context.Context) (err error) {
3436
defer func() {
3537
if value := recover(); value != nil {
3638
err = &PanicError{value}

restart.go

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -51,21 +51,23 @@ func Restart(runnable Runnable, opts ...RestartOption) Runnable {
5151
if cfg.crashBackoffDelayFn == nil {
5252
cfg.crashBackoffDelayFn = crashBackoffDelay
5353
}
54-
return &restart{baseWrapper{"restart", runnable}, runnable, cfg}
54+
return &restart{"restart/" + runnableName(runnable), runnable, cfg}
5555
}
5656

5757
type restart struct {
58-
baseWrapper
58+
name string
5959
runnable Runnable
6060
cfg restartConfig
6161
}
6262

63+
func (r *restart) runnableName() string { return r.name }
64+
6365
func (r *restart) Run(ctx context.Context) error {
6466
restartCount := 0
6567
crashCount := 0
6668

6769
for {
68-
logger.Info("starting", "runnable", findName(r), "restart", restartCount, "crash", crashCount)
70+
logger.Info("starting", "runnable", r.name, "restart", restartCount, "crash", crashCount)
6971
err := r.runnable.Run(ctx)
7072
isCrash := err != nil
7173

@@ -74,12 +76,12 @@ func (r *restart) Run(ctx context.Context) error {
7476
}
7577

7678
if r.cfg.restartLimit > 0 && restartCount >= r.cfg.restartLimit {
77-
logger.Info("not restarting", "runnable", findName(r), "reason", "restart limit", "limit", r.cfg.restartLimit)
79+
logger.Info("not restarting", "runnable", r.name, "reason", "restart limit", "limit", r.cfg.restartLimit)
7880
return err
7981
}
8082

8183
if r.cfg.crashLimit > 0 && crashCount >= r.cfg.crashLimit {
82-
logger.Info("not restarting", "runnable", findName(r), "reason", "crash limit", "limit", r.cfg.crashLimit)
84+
logger.Info("not restarting", "runnable", r.name, "reason", "crash limit", "limit", r.cfg.crashLimit)
8385
return err
8486
}
8587

0 commit comments

Comments
 (0)