@cryptotaxi247 / kubo / commits / 12fd3555a

added dep github.com/jbenet/goprocess

Juan Batiz-Benet committed Jan 6, 2015 at 08:53 UTC 12fd3555a121f145002b7a805e2bab44df56d940
11 files changed +1374
Godeps/Godeps.json
+4
@@ -142,6 +142,10 @@
142 "ImportPath": "github.com/jbenet/go-random",
143 "Rev": "2e83344e7dc7898f94501665af34edd4aa95a013"
144 },
145 + {
146 + "ImportPath": "github.com/jbenet/goprocess",
147 + "Rev": "162148a58668ca38b0b8f0459ccc6ca88e32f1f4"
148 + },
149 {
150 "ImportPath": "github.com/kr/binarydist",
151 "Rev": "9955b0ab8708602d411341e55fffd7e0700f86bd"
Godeps/_workspace/src/github.com/jbenet/goprocess/README.md new
+130
@@ -0,0 +1,130 @@
1 +# goprocess - lifecycles in go
2 +
3 +(Based on https://github.com/jbenet/go-ctxgroup)
4 +
5 +- Godoc: https://godoc.org/github.com/jbenet/goprocess
6 +
7 +`goprocess` introduces a way to manage process lifecycles in go. It is
8 +much like [go.net/context](https://godoc.org/code.google.com/p/go.net/context)
9 +(it actually uses a Context), but it is more like a Context-WaitGroup hybrid.
10 +`goprocess` is about being able to start and stop units of work, which may
11 +receive `Close` signals from many clients. Think of it like a UNIX process
12 +tree, but inside go.
13 +
14 +`goprocess` seeks to minimally affect your objects, so you can use it
15 +with both embedding or composition. At the heart of `goprocess` is the
16 +`Process` interface:
17 +
18 +```Go
19 +// Process is the basic unit of work in goprocess. It defines a computation
20 +// with a lifecycle:
21 +// - running (before calling Close),
22 +// - closing (after calling Close at least once),
23 +// - closed (after Close returns, and all teardown has _completed_).
24 +//
25 +// More specifically, it fits this:
26 +//
27 +// p := WithTeardown(tf) // new process is created, it is now running.
28 +// p.AddChild(q) // can register children **before** Closing.
29 +// go p.Close() // blocks until done running teardown func.
30 +// <-p.Closing() // would now return true.
31 +// <-p.childrenDone() // wait on all children to be done
32 +// p.teardown() // runs the user's teardown function tf.
33 +// p.Close() // now returns, with error teardown returned.
34 +// <-p.Closed() // would now return true.
35 +//
36 +// Processes can be arranged in a process "tree", where children are
37 +// automatically Closed if their parents are closed. (Note, it is actually
38 +// a Process DAG, children may have multiple parents). A process may also
39 +// optionally wait for another to fully Close before beginning to Close.
40 +// This makes it easy to ensure order of operations and proper sequential
41 +// teardown of resurces. For example:
42 +//
43 +// p1 := goprocess.WithTeardown(func() error {
44 +// fmt.Println("closing 1")
45 +// })
46 +// p2 := goprocess.WithTeardown(func() error {
47 +// fmt.Println("closing 2")
48 +// })
49 +// p3 := goprocess.WithTeardown(func() error {
50 +// fmt.Println("closing 3")
51 +// })
52 +//
53 +// p1.AddChild(p2)
54 +// p2.AddChild(p3)
55 +//
56 +//
57 +// go p1.Close()
58 +// go p2.Close()
59 +// go p3.Close()
60 +//
61 +// // Output:
62 +// // closing 3
63 +// // closing 2
64 +// // closing 1
65 +//
66 +// Process is modelled after the UNIX processes group idea, and heavily
67 +// informed by sync.WaitGroup and go.net/context.Context.
68 +//
69 +// In the function documentation of this interface, `p` always refers to
70 +// the self Process.
71 +type Process interface {
72 +
73 + // WaitFor makes p wait for q before exiting. Thus, p will _always_ close
74 + // _after_ q. Note well: a waiting cycle is deadlock.
75 + //
76 + // If q is already Closed, WaitFor calls p.Close()
77 + // If p is already Closing or Closed, WaitFor panics. This is the same thing
78 + // as calling Add(1) _after_ calling Done() on a wait group. Calling WaitFor
79 + // on an already-closed process is a programming error likely due to bad
80 + // synchronization
81 + WaitFor(q Process)
82 +
83 + // AddChildNoWait registers child as a "child" of Process. As in UNIX,
84 + // when parent is Closed, child is Closed -- child may Close beforehand.
85 + // This is the equivalent of calling:
86 + //
87 + // go func(parent, child Process) {
88 + // <-parent.Closing()
89 + // child.Close()
90 + // }(p, q)
91 + //
92 + // Note: the naming of functions is `AddChildNoWait` and `AddChild` (instead
93 + // of `AddChild` and `AddChildWaitFor`) because:
94 + // - it is the more common operation,
95 + // - explicitness is helpful in the less common case (no waiting), and
96 + // - usual "child" semantics imply parent Processes should wait for children.
97 + AddChildNoWait(q Process)
98 +
99 + // AddChild is the equivalent of calling:
100 + // parent.AddChildNoWait(q)
101 + // parent.WaitFor(q)
102 + AddChild(q Process)
103 +
104 + // Go creates a new process, adds it as a child, and spawns the ProcessFunc f
105 + // in its own goroutine. It is equivalent to:
106 + //
107 + // GoChild(p, f)
108 + //
109 + // It is useful to construct simple asynchronous workers, children of p.
110 + Go(f ProcessFunc) Process
111 +
112 + // Close ends the process. Close blocks until the process has completely
113 + // shut down, and any teardown has run _exactly once_. The returned error
114 + // is available indefinitely: calling Close twice returns the same error.
115 + // If the process has already been closed, Close returns immediately.
116 + Close() error
117 +
118 + // Closing is a signal to wait upon. The returned channel is closed
119 + // _after_ Close has been called at least once, but teardown may or may
120 + // not be done yet. The primary use case of Closing is for children who
121 + // need to know when a parent is shutting down, and therefore also shut
122 + // down.
123 + Closing() <-chan struct{}
124 +
125 + // Closed is a signal to wait upon. The returned channel is closed
126 + // _after_ Close has completed; teardown has finished. The primary use case
127 + // of Closed is waiting for a Process to Close without _causing_ the Close.
128 + Closed() <-chan struct{}
129 +}
130 +```
Godeps/_workspace/src/github.com/jbenet/goprocess/context/context.go new
+82
@@ -0,0 +1,82 @@
1 +package goprocessctx
2 +
3 +import (
4 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5 + goprocess "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/goprocess"
6 +)
7 +
8 +// WithContext constructs and returns a Process that respects
9 +// given context. It is the equivalent of:
10 +//
11 +// func ProcessWithContext(ctx context.Context) goprocess.Process {
12 +// p := goprocess.WithParent(goprocess.Background())
13 +// go func() {
14 +// <-ctx.Done()
15 +// p.Close()
16 +// }()
17 +// return p
18 +// }
19 +//
20 +func WithContext(ctx context.Context) goprocess.Process {
21 + if ctx == nil {
22 + panic("nil Context")
23 + }
24 +
25 + p := goprocess.WithParent(goprocess.Background())
26 + go func() {
27 + <-ctx.Done()
28 + p.Close()
29 + }()
30 + return p
31 +}
32 +
33 +// WaitForContext makes p WaitFor ctx. When Closing, p waits for
34 +// ctx.Done(), before being Closed(). It is simply:
35 +//
36 +// p.WaitFor(goprocess.WithContext(ctx))
37 +//
38 +func WaitForContext(ctx context.Context, p goprocess.Process) {
39 + p.WaitFor(WithContext(ctx))
40 +}
41 +
42 +// WithProcessClosing returns a context.Context derived from ctx that
43 +// is cancelled as p is Closing (after: <-p.Closing()). It is simply:
44 +//
45 +// func WithProcessClosing(ctx context.Context, p goprocess.Process) context.Context {
46 +// ctx, cancel := context.WithCancel(ctx)
47 +// go func() {
48 +// <-p.Closing()
49 +// cancel()
50 +// }()
51 +// return ctx
52 +// }
53 +//
54 +func WithProcessClosing(ctx context.Context, p goprocess.Process) context.Context {
55 + ctx, cancel := context.WithCancel(ctx)
56 + go func() {
57 + <-p.Closing()
58 + cancel()
59 + }()
60 + return ctx
61 +}
62 +
63 +// WithProcessClosed returns a context.Context that is cancelled
64 +// after Process p is Closed. It is the equivalent of:
65 +//
66 +// func WithProcessClosed(ctx context.Context, p goprocess.Process) context.Context {
67 +// ctx, cancel := context.WithCancel(ctx)
68 +// go func() {
69 +// <-p.Closed()
70 +// cancel()
71 +// }()
72 +// return ctx
73 +// }
74 +//
75 +func WithProcessClosed(ctx context.Context, p goprocess.Process) context.Context {
76 + ctx, cancel := context.WithCancel(ctx)
77 + go func() {
78 + <-p.Closed()
79 + cancel()
80 + }()
81 + return ctx
82 +}
Godeps/_workspace/src/github.com/jbenet/goprocess/example_test.go new
+37
@@ -0,0 +1,37 @@
1 +package goprocess_test
2 +
3 +import (
4 + "fmt"
5 + "time"
6 +
7 + "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/goprocess"
8 +)
9 +
10 +func ExampleGo() {
11 + p := goprocess.Go(func(p goprocess.Process) {
12 + ticker := time.Tick(200 * time.Millisecond)
13 + for {
14 + select {
15 + case <-ticker:
16 + fmt.Println("tick")
17 + case <-p.Closing():
18 + fmt.Println("closing")
19 + return
20 + }
21 + }
22 + })
23 +
24 + <-time.After(1100 * time.Millisecond)
25 + p.Close()
26 + fmt.Println("closed")
27 + <-time.After(100 * time.Millisecond)
28 +
29 + // Output:
30 + // tick
31 + // tick
32 + // tick
33 + // tick
34 + // tick
35 + // closing
36 + // closed
37 +}
Godeps/_workspace/src/github.com/jbenet/goprocess/goprocess.go new
+262
@@ -0,0 +1,262 @@
1 +// Package goprocess introduces a Process abstraction that allows simple
2 +// organization, and orchestration of work. It is much like a WaitGroup,
3 +// and much like a context.Context, but also ensures safe **exactly-once**,
4 +// and well-ordered teardown semantics.
5 +package goprocess
6 +
7 +import (
8 + "os"
9 + "os/signal"
10 +)
11 +
12 +// Process is the basic unit of work in goprocess. It defines a computation
13 +// with a lifecycle:
14 +// - running (before calling Close),
15 +// - closing (after calling Close at least once),
16 +// - closed (after Close returns, and all teardown has _completed_).
17 +//
18 +// More specifically, it fits this:
19 +//
20 +// p := WithTeardown(tf) // new process is created, it is now running.
21 +// p.AddChild(q) // can register children **before** Closing.
22 +// go p.Close() // blocks until done running teardown func.
23 +// <-p.Closing() // would now return true.
24 +// <-p.childrenDone() // wait on all children to be done
25 +// p.teardown() // runs the user's teardown function tf.
26 +// p.Close() // now returns, with error teardown returned.
27 +// <-p.Closed() // would now return true.
28 +//
29 +// Processes can be arranged in a process "tree", where children are
30 +// automatically Closed if their parents are closed. (Note, it is actually
31 +// a Process DAG, children may have multiple parents). A process may also
32 +// optionally wait for another to fully Close before beginning to Close.
33 +// This makes it easy to ensure order of operations and proper sequential
34 +// teardown of resurces. For example:
35 +//
36 +// p1 := goprocess.WithTeardown(func() error {
37 +// fmt.Println("closing 1")
38 +// })
39 +// p2 := goprocess.WithTeardown(func() error {
40 +// fmt.Println("closing 2")
41 +// })
42 +// p3 := goprocess.WithTeardown(func() error {
43 +// fmt.Println("closing 3")
44 +// })
45 +//
46 +// p1.AddChild(p2)
47 +// p2.AddChild(p3)
48 +//
49 +//
50 +// go p1.Close()
51 +// go p2.Close()
52 +// go p3.Close()
53 +//
54 +// // Output:
55 +// // closing 3
56 +// // closing 2
57 +// // closing 1
58 +//
59 +// Process is modelled after the UNIX processes group idea, and heavily
60 +// informed by sync.WaitGroup and go.net/context.Context.
61 +//
62 +// In the function documentation of this interface, `p` always refers to
63 +// the self Process.
64 +type Process interface {
65 +
66 + // WaitFor makes p wait for q before exiting. Thus, p will _always_ close
67 + // _after_ q. Note well: a waiting cycle is deadlock.
68 + //
69 + // If q is already Closed, WaitFor calls p.Close()
70 + // If p is already Closing or Closed, WaitFor panics. This is the same thing
71 + // as calling Add(1) _after_ calling Done() on a wait group. Calling WaitFor
72 + // on an already-closed process is a programming error likely due to bad
73 + // synchronization
74 + WaitFor(q Process)
75 +
76 + // AddChildNoWait registers child as a "child" of Process. As in UNIX,
77 + // when parent is Closed, child is Closed -- child may Close beforehand.
78 + // This is the equivalent of calling:
79 + //
80 + // go func(parent, child Process) {
81 + // <-parent.Closing()
82 + // child.Close()
83 + // }(p, q)
84 + //
85 + // Note: the naming of functions is `AddChildNoWait` and `AddChild` (instead
86 + // of `AddChild` and `AddChildWaitFor`) because:
87 + // - it is the more common operation,
88 + // - explicitness is helpful in the less common case (no waiting), and
89 + // - usual "child" semantics imply parent Processes should wait for children.
90 + AddChildNoWait(q Process)
91 +
92 + // AddChild is the equivalent of calling:
93 + // parent.AddChildNoWait(q)
94 + // parent.WaitFor(q)
95 + AddChild(q Process)
96 +
97 + // Go is much like `go`, as it runs a function in a newly spawned goroutine.
98 + // The neat part of Process.Go is that the Process object you call it on will:
99 + // * construct a child Process, and call AddChild(child) on it
100 + // * spawn a goroutine, and call the given function
101 + // * Close the child when the function exits.
102 + // This way, you can rest assured each goroutine you spawn has its very own
103 + // Process context, and that it will be closed when the function exits.
104 + // It is the function's responsibility to respect the Closing of its Process,
105 + // namely it should exit (return) when <-Closing() is ready. It is basically:
106 + //
107 + // func (p Process) Go(f ProcessFunc) Process {
108 + // child := WithParent(p)
109 + // go func () {
110 + // f(child)
111 + // child.Close()
112 + // }()
113 + // }
114 + //
115 + // It is useful to construct simple asynchronous workers, children of p.
116 + Go(f ProcessFunc)
117 +
118 + // Close ends the process. Close blocks until the process has completely
119 + // shut down, and any teardown has run _exactly once_. The returned error
120 + // is available indefinitely: calling Close twice returns the same error.
121 + // If the process has already been closed, Close returns immediately.
122 + Close() error
123 +
124 + // Closing is a signal to wait upon. The returned channel is closed
125 + // _after_ Close has been called at least once, but teardown may or may
126 + // not be done yet. The primary use case of Closing is for children who
127 + // need to know when a parent is shutting down, and therefore also shut
128 + // down.
129 + Closing() <-chan struct{}
130 +
131 + // Closed is a signal to wait upon. The returned channel is closed
132 + // _after_ Close has completed; teardown has finished. The primary use case
133 + // of Closed is waiting for a Process to Close without _causing_ the Close.
134 + Closed() <-chan struct{}
135 +}
136 +
137 +// TeardownFunc is a function used to cleanup state at the end of the
138 +// lifecycle of a Process.
139 +type TeardownFunc func() error
140 +
141 +var nilTeardownFunc = func() error { return nil }
142 +
143 +// ProcessFunc is a function that takes a process. Its main use case is goprocess.Go,
144 +// which spawns a ProcessFunc in its own goroutine, and returns a corresponding
145 +// Process object.
146 +type ProcessFunc func(proc Process)
147 +
148 +var nilProcessFunc = func(Process) {}
149 +
150 +// Go is much like `go`: it runs a function in a newly spawned goroutine. The neat
151 +// part of Go is that it provides Process object to communicate between the
152 +// function and the outside world. Thus, callers can easily WaitFor, or Close the
153 +// function. It is the function's responsibility to respect the Closing of its Process,
154 +// namely it should exit (return) when <-Closing() is ready. It is simply:
155 +//
156 +// func Go(f ProcessFunc) Process {
157 +// p := WithParent(Background())
158 +// p.Go(f)
159 +// return p
160 +// }
161 +//
162 +// Note that a naive implementation of Go like the following would not work:
163 +//
164 +// func Go(f ProcessFunc) Process {
165 +// return Background().Go(f)
166 +// }
167 +//
168 +// This is because having the process you
169 +func Go(f ProcessFunc) Process {
170 + return GoChild(Background(), f)
171 +}
172 +
173 +// GoChild is like Go, but it registers the returned Process as a child of parent,
174 +// **before** spawning the goroutine, which ensures proper synchronization with parent.
175 +// It is somewhat like
176 +//
177 +// func GoChild(parent Process, f ProcessFunc) Process {
178 +// p := WithParent(parent)
179 +// p.Go(f)
180 +// return p
181 +// }
182 +//
183 +// And it is similar to the classic WaitGroup use case:
184 +//
185 +// func WaitGroupGo(wg sync.WaitGroup, child func()) {
186 +// wg.Add(1)
187 +// go func() {
188 +// child()
189 +// wg.Done()
190 +// }()
191 +// }
192 +//
193 +func GoChild(parent Process, f ProcessFunc) Process {
194 + p := WithParent(parent)
195 + p.Go(f)
196 + return p
197 +}
198 +
199 +// Spawn is an alias of `Go`. In many contexts, Spawn is a
200 +// well-known Process launching word, which fits our use case.
201 +var Spawn = Go
202 +
203 +// SpawnChild is an alias of `GoChild`. In many contexts, Spawn is a
204 +// well-known Process launching word, which fits our use case.
205 +var SpawnChild = GoChild
206 +
207 +// WithTeardown constructs and returns a Process with a TeardownFunc.
208 +// TeardownFunc tf will be called **exactly-once** when Process is
209 +// Closing, after all Children have fully closed, and before p is Closed.
210 +// In fact, Process p will not be Closed until tf runs and exits.
211 +// See lifecycle in Process doc.
212 +func WithTeardown(tf TeardownFunc) Process {
213 + if tf == nil {
214 + panic("nil tf TeardownFunc")
215 + }
216 + return newProcess(tf)
217 +}
218 +
219 +// WithParent constructs and returns a Process with a given parent.
220 +func WithParent(parent Process) Process {
221 + if parent == nil {
222 + panic("nil parent Process")
223 + }
224 + q := newProcess(nil)
225 + parent.AddChild(q)
226 + return q
227 +}
228 +
229 +// WithSignals returns a Process that will Close() when any given signal fires.
230 +// This is useful to bind Process trees to syscall.SIGTERM, SIGKILL, etc.
231 +func WithSignals(sig ...os.Signal) Process {
232 + p := WithParent(Background())
233 + c := make(chan os.Signal)
234 + signal.Notify(c, sig...)
235 + go func() {
236 + <-c
237 + signal.Stop(c)
238 + p.Close()
239 + }()
240 + return p
241 +}
242 +
243 +// Background returns the "background" Process: a statically allocated
244 +// process that can _never_ close. It also never enters Closing() state.
245 +// Calling Background().Close() will hang indefinitely.
246 +func Background() Process {
247 + return background
248 +}
249 +
250 +// background is the background process
251 +var background = &unclosable{Process: newProcess(nil)}
252 +
253 +// unclosable is a process that _cannot_ be closed. calling Close simply hangs.
254 +type unclosable struct {
255 + Process
256 +}
257 +
258 +func (p *unclosable) Close() error {
259 + var hang chan struct{}
260 + <-hang // hang forever
261 + return nil
262 +}
Godeps/_workspace/src/github.com/jbenet/goprocess/goprocess_test.go new
+433
@@ -0,0 +1,433 @@
1 +package goprocess
2 +
3 +import (
4 + "syscall"
5 + "testing"
6 + "time"
7 +)
8 +
9 +type tree struct {
10 + Process
11 + c []tree
12 +}
13 +
14 +func setupHierarchy(p Process) tree {
15 + t := func(n Process, ts ...tree) tree {
16 + return tree{n, ts}
17 + }
18 +
19 + a := WithParent(p)
20 + b1 := WithParent(a)
21 + b2 := WithParent(a)
22 + c1 := WithParent(b1)
23 + c2 := WithParent(b1)
24 + c3 := WithParent(b2)
25 + c4 := WithParent(b2)
26 +
27 + return t(a, t(b1, t(c1), t(c2)), t(b2, t(c3), t(c4)))
28 +}
29 +
30 +func TestClosingClosed(t *testing.T) {
31 +
32 + a := WithParent(Background())
33 + b := WithParent(a)
34 +
35 + Q := make(chan string, 3)
36 +
37 + go func() {
38 + <-a.Closing()
39 + Q <- "closing"
40 + b.Close()
41 + }()
42 +
43 + go func() {
44 + <-a.Closed()
45 + Q <- "closed"
46 + }()
47 +
48 + go func() {
49 + a.Close()
50 + Q <- "closed"
51 + }()
52 +
53 + if q := <-Q; q != "closing" {
54 + t.Error("order incorrect. closing not first")
55 + }
56 + if q := <-Q; q != "closed" {
57 + t.Error("order incorrect. closing not first")
58 + }
59 + if q := <-Q; q != "closed" {
60 + t.Error("order incorrect. closing not first")
61 + }
62 +}
63 +
64 +func TestChildFunc(t *testing.T) {
65 + a := WithParent(Background())
66 +
67 + wait1 := make(chan struct{})
68 + wait2 := make(chan struct{})
69 + wait3 := make(chan struct{})
70 + wait4 := make(chan struct{})
71 + go func() {
72 + a.Close()
73 + wait4 <- struct{}{}
74 + }()
75 +
76 + a.Go(func(process Process) {
77 + wait1 <- struct{}{}
78 + <-wait2
79 + wait3 <- struct{}{}
80 + })
81 +
82 + <-wait1
83 + select {
84 + case <-wait3:
85 + t.Error("should not be closed yet")
86 + case <-wait4:
87 + t.Error("should not be closed yet")
88 + case <-a.Closed():
89 + t.Error("should not be closed yet")
90 + default:
91 + }
92 +
93 + wait2 <- struct{}{}
94 +
95 + select {
96 + case <-wait3:
97 + case <-time.After(time.Second):
98 + t.Error("should be closed now")
99 + }
100 +
101 + select {
102 + case <-wait4:
103 + case <-time.After(time.Second):
104 + t.Error("should be closed now")
105 + }
106 +}
107 +
108 +func TestTeardownCalledOnce(t *testing.T) {
109 + a := setupHierarchy(Background())
110 +
111 + onlyOnce := func() func() error {
112 + count := 0
113 + return func() error {
114 + count++
115 + if count > 1 {
116 + t.Error("called", count, "times")
117 + }
118 + return nil
119 + }
120 + }
121 +
122 + setTeardown := func(t tree, tf TeardownFunc) {
123 + t.Process.(*process).teardown = tf
124 + }
125 +
126 + setTeardown(a, onlyOnce())
127 + setTeardown(a.c[0], onlyOnce())
128 + setTeardown(a.c[0].c[0], onlyOnce())
129 + setTeardown(a.c[0].c[1], onlyOnce())
130 + setTeardown(a.c[1], onlyOnce())
131 + setTeardown(a.c[1].c[0], onlyOnce())
132 + setTeardown(a.c[1].c[1], onlyOnce())
133 +
134 + a.c[0].c[0].Close()
135 + a.c[0].c[0].Close()
136 + a.c[0].c[0].Close()
137 + a.c[0].c[0].Close()
138 + a.c[0].Close()
139 + a.c[0].Close()
140 + a.c[0].Close()
141 + a.c[0].Close()
142 + a.Close()
143 + a.Close()
144 + a.Close()
145 + a.Close()
146 + a.c[1].Close()
147 + a.c[1].Close()
148 + a.c[1].Close()
149 + a.c[1].Close()
150 +}
151 +
152 +func TestOnClosed(t *testing.T) {
153 +
154 + Q := make(chan string, 10)
155 + p := WithParent(Background())
156 + a := setupHierarchy(p)
157 +
158 + go onClosedStr(Q, "0", a.c[0])
159 + go onClosedStr(Q, "10", a.c[1].c[0])
160 + go onClosedStr(Q, "", a)
161 + go onClosedStr(Q, "00", a.c[0].c[0])
162 + go onClosedStr(Q, "1", a.c[1])
163 + go onClosedStr(Q, "01", a.c[0].c[1])
164 + go onClosedStr(Q, "11", a.c[1].c[1])
165 +
166 + go p.Close()
167 +
168 + testStrs(t, Q, "00", "01", "10", "11")
169 + testStrs(t, Q, "00", "01", "10", "11")
170 + testStrs(t, Q, "00", "01", "10", "11")
171 + testStrs(t, Q, "00", "01", "10", "11")
172 + testStrs(t, Q, "0", "1")
173 + testStrs(t, Q, "0", "1")
174 + testStrs(t, Q, "")
175 +}
176 +
177 +func TestWaitFor(t *testing.T) {
178 +
179 + Q := make(chan string, 5)
180 + a := WithParent(Background())
181 + b := WithParent(Background())
182 + c := WithParent(Background())
183 + d := WithParent(Background())
184 + e := WithParent(Background())
185 +
186 + go onClosedStr(Q, "a", a)
187 + go onClosedStr(Q, "b", b)
188 + go onClosedStr(Q, "c", c)
189 + go onClosedStr(Q, "d", d)
190 + go onClosedStr(Q, "e", e)
191 +
192 + testNone(t, Q)
193 + a.WaitFor(b)
194 + a.WaitFor(c)
195 + b.WaitFor(d)
196 + e.WaitFor(d)
197 + testNone(t, Q)
198 +
199 + go a.Close() // should do nothing.
200 + testNone(t, Q)
201 +
202 + go e.Close()
203 + testNone(t, Q)
204 +
205 + d.Close()
206 + testStrs(t, Q, "d", "e")
207 + testStrs(t, Q, "d", "e")
208 +
209 + c.Close()
210 + testStrs(t, Q, "c")
211 +
212 + b.Close()
213 + testStrs(t, Q, "a", "b")
214 + testStrs(t, Q, "a", "b")
215 +}
216 +
217 +func TestAddChildNoWait(t *testing.T) {
218 +
219 + Q := make(chan string, 5)
220 + a := WithParent(Background())
221 + b := WithParent(Background())
222 + c := WithParent(Background())
223 + d := WithParent(Background())
224 + e := WithParent(Background())
225 +
226 + go onClosedStr(Q, "a", a)
227 + go onClosedStr(Q, "b", b)
228 + go onClosedStr(Q, "c", c)
229 + go onClosedStr(Q, "d", d)
230 + go onClosedStr(Q, "e", e)
231 +
232 + testNone(t, Q)
233 + a.AddChildNoWait(b)
234 + a.AddChildNoWait(c)
235 + b.AddChildNoWait(d)
236 + e.AddChildNoWait(d)
237 + testNone(t, Q)
238 +
239 + b.Close()
240 + testStrs(t, Q, "b", "d")
241 + testStrs(t, Q, "b", "d")
242 +
243 + a.Close()
244 + testStrs(t, Q, "a", "c")
245 + testStrs(t, Q, "a", "c")
246 +
247 + e.Close()
248 + testStrs(t, Q, "e")
249 +}
250 +
251 +func TestAddChild(t *testing.T) {
252 +
253 + a := WithParent(Background())
254 + b := WithParent(Background())
255 + c := WithParent(Background())
256 + d := WithParent(Background())
257 + e := WithParent(Background())
258 + Q := make(chan string, 5)
259 +
260 + go onClosedStr(Q, "a", a)
261 + go onClosedStr(Q, "b", b)
262 + go onClosedStr(Q, "c", c)
263 + go onClosedStr(Q, "d", d)
264 + go onClosedStr(Q, "e", e)
265 +
266 + testNone(t, Q)
267 + a.AddChild(b)
268 + a.AddChild(c)
269 + b.AddChild(d)
270 + e.AddChild(d)
271 + testNone(t, Q)
272 +
273 + go b.Close()
274 + testNone(t, Q)
275 + d.Close()
276 + testStrs(t, Q, "b", "d")
277 + testStrs(t, Q, "b", "d")
278 +
279 + go a.Close()
280 + testNone(t, Q)
281 + c.Close()
282 + testStrs(t, Q, "a", "c")
283 + testStrs(t, Q, "a", "c")
284 +
285 + e.Close()
286 + testStrs(t, Q, "e")
287 +}
288 +
289 +func TestGoChildrenClose(t *testing.T) {
290 +
291 + var a, b, c, d, e Process
292 +
293 + var ready = make(chan struct{})
294 + var bWait = make(chan struct{})
295 + var cWait = make(chan struct{})
296 + var dWait = make(chan struct{})
297 + var eWait = make(chan struct{})
298 +
299 + a = WithParent(Background())
300 + a.Go(func(p Process) {
301 + b = p
302 + b.Go(func(p Process) {
303 + c = p
304 + ready <- struct{}{}
305 + <-cWait
306 + })
307 + ready <- struct{}{}
308 + <-bWait
309 + })
310 + a.Go(func(p Process) {
311 + d = p
312 + d.Go(func(p Process) {
313 + e = p
314 + ready <- struct{}{}
315 + <-eWait
316 + })
317 + ready <- struct{}{}
318 + <-dWait
319 + })
320 +
321 + <-ready
322 + <-ready
323 + <-ready
324 + <-ready
325 +
326 + Q := make(chan string, 5)
327 +
328 + go onClosedStr(Q, "a", a)
329 + go onClosedStr(Q, "b", b)
330 + go onClosedStr(Q, "c", c)
331 + go onClosedStr(Q, "d", d)
332 + go onClosedStr(Q, "e", e)
333 +
334 + testNone(t, Q)
335 + go a.Close()
336 + testNone(t, Q)
337 +
338 + go b.Close()
339 + testNone(t, Q)
340 +
341 + c.Close()
342 + testStrs(t, Q, "b", "c")
343 + testStrs(t, Q, "b", "c")
344 +
345 + e.Close()
346 + testStrs(t, Q, "e")
347 +
348 + d.Close()
349 + <-a.Closed()
350 + testStrs(t, Q, "a", "d")
351 + testStrs(t, Q, "a", "d")
352 +}
353 +
354 +func TestBackground(t *testing.T) {
355 + // test it hangs indefinitely:
356 + b := Background()
357 +
358 + go b.Close()
359 + go func() {
360 + b.Close()
361 + }()
362 +
363 + select {
364 + case <-b.Closing():
365 + t.Error("b.Closing() closed :(")
366 + default:
367 + }
368 +}
369 +
370 +func TestWithSignals(t *testing.T) {
371 + p := WithSignals(syscall.SIGABRT)
372 + testNotClosed(t, p)
373 +
374 + syscall.Kill(syscall.Getpid(), syscall.SIGABRT)
375 + testClosed(t, p)
376 +}
377 +
378 +func testClosing(t *testing.T, p Process) {
379 + select {
380 + case <-p.Closing():
381 + case <-time.After(50 * time.Millisecond):
382 + t.Fatal("should be closing")
383 + }
384 +}
385 +
386 +func testNotClosing(t *testing.T, p Process) {
387 + select {
388 + case <-p.Closing():
389 + t.Fatal("should not be closing")
390 + case <-p.Closed():
391 + t.Fatal("should not be closed")
392 + default:
393 + }
394 +}
395 +
396 +func testClosed(t *testing.T, p Process) {
397 + select {
398 + case <-p.Closed():
399 + case <-time.After(50 * time.Millisecond):
400 + t.Fatal("should be closed")
401 + }
402 +}
403 +
404 +func testNotClosed(t *testing.T, p Process) {
405 + select {
406 + case <-p.Closed():
407 + t.Fatal("should not be closed")
408 + case <-time.After(50 * time.Millisecond):
409 + }
410 +}
411 +
412 +func testNone(t *testing.T, c <-chan string) {
413 + select {
414 + case <-c:
415 + t.Fatal("none should be closed")
416 + default:
417 + }
418 +}
419 +
420 +func testStrs(t *testing.T, Q <-chan string, ss ...string) {
421 + s1 := <-Q
422 + for _, s2 := range ss {
423 + if s1 == s2 {
424 + return
425 + }
426 + }
427 + t.Error("context not in group:", s1, ss)
428 +}
429 +
430 +func onClosedStr(Q chan<- string, s string, p Process) {
431 + <-p.Closed()
432 + Q <- s
433 +}
Godeps/_workspace/src/github.com/jbenet/goprocess/impl-goroutines.go new
+114
@@ -0,0 +1,114 @@
1 +// +build ignore
2 +
3 +// WARNING: this implementation is not correct.
4 +// here only for historical purposes.
5 +
6 +package goprocess
7 +
8 +import (
9 + "sync"
10 +)
11 +
12 +// process implements Process
13 +type process struct {
14 + children sync.WaitGroup // wait group for child goroutines
15 + teardown TeardownFunc // called to run the teardown logic.
16 + closing chan struct{} // closed once close starts.
17 + closed chan struct{} // closed once close is done.
18 + closeOnce sync.Once // ensure close is only called once.
19 + closeErr error // error to return to clients of Close()
20 +}
21 +
22 +// newProcess constructs and returns a Process.
23 +// It will call tf TeardownFunc exactly once:
24 +// **after** all children have fully Closed,
25 +// **after** entering <-Closing(), and
26 +// **before** <-Closed().
27 +func newProcess(tf TeardownFunc) *process {
28 + if tf == nil {
29 + tf = nilTeardownFunc
30 + }
31 +
32 + return &process{
33 + teardown: tf,
34 + closed: make(chan struct{}),
35 + closing: make(chan struct{}),
36 + }
37 +}
38 +
39 +func (p *process) WaitFor(q Process) {
40 + p.children.Add(1) // p waits on q to be done
41 + go func(p *process, q Process) {
42 + <-q.Closed() // wait until q is closed
43 + p.children.Done() // p done waiting on q
44 + }(p, q)
45 +}
46 +
47 +func (p *process) AddChildNoWait(child Process) {
48 + go func(p, child Process) {
49 + <-p.Closing() // wait until p is closing
50 + child.Close() // close child
51 + }(p, child)
52 +}
53 +
54 +func (p *process) AddChild(child Process) {
55 + select {
56 + case <-p.Closing():
57 + panic("attempt to add child to closing or closed process")
58 + default:
59 + }
60 +
61 + p.children.Add(1) // p waits on child to be done
62 + go func(p *process, child Process) {
63 + <-p.Closing() // wait until p is closing
64 + child.Close() // close child and wait
65 + p.children.Done() // p done waiting on child
66 + }(p, child)
67 +}
68 +
69 +func (p *process) Go(f ProcessFunc) Process {
70 + select {
71 + case <-p.Closing():
72 + panic("attempt to add child to closing or closed process")
73 + default:
74 + }
75 +
76 + // this is very similar to AddChild, but also runs the func
77 + // in the child. we replicate it here to save one goroutine.
78 + child := newProcessGoroutines(nil)
79 + child.children.Add(1) // child waits on func to be done
80 + p.AddChild(child)
81 + go func() {
82 + f(child)
83 + child.children.Done() // wait on child's children to be done.
84 + child.Close() // close to tear down.
85 + }()
86 + return child
87 +}
88 +
89 +// Close is the external close function.
90 +// it's a wrapper around internalClose that waits on Closed()
91 +func (p *process) Close() error {
92 + p.closeOnce.Do(p.doClose)
93 + <-p.Closed() // sync.Once should block, but this checks chan is closed too
94 + return p.closeErr
95 +}
96 +
97 +func (p *process) Closing() <-chan struct{} {
98 + return p.closing
99 +}
100 +
101 +func (p *process) Closed() <-chan struct{} {
102 + return p.closed
103 +}
104 +
105 +// the _actual_ close process.
106 +func (p *process) doClose() {
107 + // this function should only be called once (hence the sync.Once).
108 + // and it will panic (on closing channels) otherwise.
109 +
110 + close(p.closing) // signal that we're shutting down (Closing)
111 + p.children.Wait() // wait till all children are done (before teardown)
112 + p.closeErr = p.teardown() // actually run the close logic (ok safe to teardown)
113 + close(p.closed) // signal that we're shut down (Closed)
114 +}
Godeps/_workspace/src/github.com/jbenet/goprocess/impl-mutex.go new
+127
@@ -0,0 +1,127 @@
1 +package goprocess
2 +
3 +import (
4 + "sync"
5 +)
6 +
7 +// process implements Process
8 +type process struct {
9 + children []Process // process to close with us
10 + waitfors []Process // process to only wait for
11 + teardown TeardownFunc // called to run the teardown logic.
12 + closing chan struct{} // closed once close starts.
13 + closed chan struct{} // closed once close is done.
14 + closeErr error // error to return to clients of Close()
15 +
16 + sync.Mutex
17 +}
18 +
19 +// newProcess constructs and returns a Process.
20 +// It will call tf TeardownFunc exactly once:
21 +// **after** all children have fully Closed,
22 +// **after** entering <-Closing(), and
23 +// **before** <-Closed().
24 +func newProcess(tf TeardownFunc) *process {
25 + if tf == nil {
26 + tf = nilTeardownFunc
27 + }
28 +
29 + return &process{
30 + teardown: tf,
31 + closed: make(chan struct{}),
32 + closing: make(chan struct{}),
33 + }
34 +}
35 +
36 +func (p *process) WaitFor(q Process) {
37 + p.Lock()
38 +
39 + select {
40 + case <-p.Closed():
41 + panic("Process cannot wait after being closed")
42 + default:
43 + }
44 +
45 + p.waitfors = append(p.waitfors, q)
46 + p.Unlock()
47 +}
48 +
49 +func (p *process) AddChildNoWait(child Process) {
50 + p.Lock()
51 +
52 + select {
53 + case <-p.Closed():
54 + panic("Process cannot add children after being closed")
55 + default:
56 + }
57 +
58 + p.children = append(p.children, child)
59 + p.Unlock()
60 +}
61 +
62 +func (p *process) AddChild(child Process) {
63 + p.Lock()
64 +
65 + select {
66 + case <-p.Closed():
67 + panic("Process cannot add children after being closed")
68 + default:
69 + }
70 +
71 + p.waitfors = append(p.waitfors, child)
72 + p.children = append(p.children, child)
73 + p.Unlock()
74 +}
75 +
76 +func (p *process) Go(f ProcessFunc) {
77 + child := newProcess(nil)
78 + p.AddChild(child)
79 + go func() {
80 + f(child)
81 + child.Close() // close to tear down.
82 + }()
83 +}
84 +
85 +// Close is the external close function.
86 +// it's a wrapper around internalClose that waits on Closed()
87 +func (p *process) Close() error {
88 + p.Lock()
89 + defer p.Unlock()
90 +
91 + // if already closed, get out.
92 + select {
93 + case <-p.Closed():
94 + return p.closeErr
95 + default:
96 + }
97 +
98 + p.doClose()
99 + return p.closeErr
100 +}
101 +
102 +func (p *process) Closing() <-chan struct{} {
103 + return p.closing
104 +}
105 +
106 +func (p *process) Closed() <-chan struct{} {
107 + return p.closed
108 +}
109 +
110 +// the _actual_ close process.
111 +func (p *process) doClose() {
112 + // this function is only be called once (protected by p.Lock()).
113 + // and it will panic (on closing channels) otherwise.
114 +
115 + close(p.closing) // signal that we're shutting down (Closing)
116 +
117 + for _, c := range p.children {
118 + go c.Close() // force all children to shut down
119 + }
120 +
121 + for _, w := range p.waitfors {
122 + <-w.Closed() // wait till all waitfors are fully closed (before teardown)
123 + }
124 +
125 + p.closeErr = p.teardown() // actually run the close logic (ok safe to teardown)
126 + close(p.closed) // signal that we're shut down (Closed)
127 +}
Godeps/_workspace/src/github.com/jbenet/goprocess/ratelimit/README.md new
+4
@@ -0,0 +1,4 @@
1 +# goprocess/ratelimit - ratelimit children creation
2 +
3 +- goprocess: https://github.com/jbenet/goprocess
4 +- Godoc: https://godoc.org/github.com/jbenet/goprocess/ratelimit
Godeps/_workspace/src/github.com/jbenet/goprocess/ratelimit/ratelimit.go new
+83
@@ -0,0 +1,83 @@
1 +// Package ratelimit is part of github.com/jbenet/goprocess.
2 +// It provides a simple process that ratelimits child creation.
3 +// This is done internally with a channel/semaphore.
4 +// So the call `RateLimiter.LimitedGo` may block until another
5 +// child is Closed().
6 +package ratelimit
7 +
8 +import (
9 + process "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/goprocess"
10 +)
11 +
12 +// RateLimiter limits the spawning of children. It does so
13 +// with an internal semaphore. Note that Go will continue
14 +// to be the unlimited process.Process.Go, and ONLY the
15 +// added function `RateLimiter.LimitedGo` will honor the
16 +// limit. This is to improve readability and avoid confusion
17 +// for the reader, particularly if code changes over time.
18 +type RateLimiter struct {
19 + process.Process
20 +
21 + limiter chan struct{}
22 +}
23 +
24 +func NewRateLimiter(parent process.Process, limit int) *RateLimiter {
25 + proc := process.WithParent(parent)
26 + return &RateLimiter{Process: proc, limiter: LimitChan(limit)}
27 +}
28 +
29 +// LimitedGo creates a new process, adds it as a child, and spawns the
30 +// ProcessFunc f in its own goroutine, but may block according to the
31 +// internal rate limit. It is equivalent to:
32 +//
33 +// func(f process.ProcessFunc) {
34 +// <-limitch
35 +// p.Go(func (child process.Process) {
36 +// f(child)
37 +// f.Close() // make sure its children close too!
38 +// limitch<- struct{}{}
39 +// })
40 +/// }
41 +//
42 +// It is useful to construct simple asynchronous workers, children of p,
43 +// and rate limit their creation, to avoid spinning up too many, too fast.
44 +// This is great for providing backpressure to producers.
45 +func (rl *RateLimiter) LimitedGo(f process.ProcessFunc) {
46 +
47 + <-rl.limiter
48 + rl.Go(func(child process.Process) {
49 +
50 + // call the function as rl.Go would.
51 + f(child)
52 +
53 + // this close is here because the child may have spawned
54 + // children of its own, and our rate limiter should capture that.
55 + // we have two options:
56 + // * this approach (which is what process.Go itself does), or
57 + // * spawn another goroutine that waits on <-child.Closed()
58 + //
59 + // go func() {
60 + // <-child.Closed()
61 + // rl.limiter <- struct{}{}
62 + // }()
63 + //
64 + // This approach saves a goroutine. It is fine to call child.Close()
65 + // multiple times.
66 + child.Close()
67 +
68 + // after it's done.
69 + rl.limiter <- struct{}{}
70 + })
71 +}
72 +
73 +// LimitChan returns a rate-limiting channel. it is the usual, simple,
74 +// golang-idiomatic rate-limiting semaphore. This function merely
75 +// initializes it with certain buffer size, and sends that many values,
76 +// so it is ready to be used.
77 +func LimitChan(limit int) chan struct{} {
78 + limitch := make(chan struct{}, limit)
79 + for i := 0; i < limit; i++ {
80 + limitch <- struct{}{}
81 + }
82 + return limitch
83 +}
Godeps/_workspace/src/github.com/jbenet/goprocess/ratelimit/ratelimit_test.go new
+98
@@ -0,0 +1,98 @@
1 +package ratelimit
2 +
3 +import (
4 + "testing"
5 + "time"
6 +
7 + process "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/goprocess"
8 +)
9 +
10 +func TestRateLimitLimitedGoBlocks(t *testing.T) {
11 + numChildren := 6
12 +
13 + t.Logf("create a rate limiter with limit of %d", numChildren/2)
14 + rl := NewRateLimiter(process.Background(), numChildren/2)
15 +
16 + doneSpawning := make(chan struct{})
17 + childClosing := make(chan struct{})
18 +
19 + t.Log("spawn 6 children with LimitedGo.")
20 + go func() {
21 + for i := 0; i < numChildren; i++ {
22 + rl.LimitedGo(func(child process.Process) {
23 + // hang until we drain childClosing
24 + childClosing <- struct{}{}
25 + })
26 + t.Logf("spawned %d", i)
27 + }
28 + close(doneSpawning)
29 + }()
30 +
31 + t.Log("should have blocked.")
32 + select {
33 + case <-doneSpawning:
34 + t.Error("did not block")
35 + case <-time.After(time.Millisecond): // for scheduler
36 + t.Log("blocked")
37 + }
38 +
39 + t.Logf("drain %d children so they close", numChildren/2)
40 + for i := 0; i < numChildren/2; i++ {
41 + t.Logf("closing %d", i)
42 + <-childClosing // consume child cloing
43 + t.Logf("closed %d", i)
44 + }
45 +
46 + t.Log("should be done spawning.")
47 + select {
48 + case <-doneSpawning:
49 + case <-time.After(time.Millisecond): // for scheduler
50 + t.Error("still blocked...")
51 + }
52 +
53 + t.Logf("drain %d children so they close", numChildren/2)
54 + for i := 0; i < numChildren/2; i++ {
55 + <-childClosing
56 + t.Logf("closed %d", i)
57 + }
58 +
59 + rl.Close() // ensure everyone's closed.
60 +}
61 +
62 +func TestRateLimitGoDoesntBlock(t *testing.T) {
63 + numChildren := 6
64 +
65 + t.Logf("create a rate limiter with limit of %d", numChildren/2)
66 + rl := NewRateLimiter(process.Background(), numChildren/2)
67 +
68 + doneSpawning := make(chan struct{})
69 + childClosing := make(chan struct{})
70 +
71 + t.Log("spawn 6 children with usual Process.Go.")
72 + go func() {
73 + for i := 0; i < numChildren; i++ {
74 + rl.Go(func(child process.Process) {
75 + // hang until we drain childClosing
76 + childClosing <- struct{}{}
77 + })
78 + t.Logf("spawned %d", i)
79 + }
80 + close(doneSpawning)
81 + }()
82 +
83 + t.Log("should not have blocked.")
84 + select {
85 + case <-doneSpawning:
86 + t.Log("did not block")
87 + case <-time.After(time.Millisecond): // for scheduler
88 + t.Error("process.Go blocked. it should not.")
89 + }
90 +
91 + t.Log("drain children so they close")
92 + for i := 0; i < numChildren; i++ {
93 + <-childClosing
94 + t.Logf("closed %d", i)
95 + }
96 +
97 + rl.Close() // ensure everyone's closed.
98 +}