@cryptotaxi247 / kubo / commits / 23081430a

Fixed panic on closer

Juan Batiz-Benet committed Oct 19, 2014 at 01:50 UTC 23081430a2cc402045f90f4c7065716bb962a62a
5 files changed +118 -50
net/conn/closer.go
+100 -32
@@ -1,29 +1,71 @@
1 package conn
2
3 import (
4 - "errors"
4 + "sync"
5
6 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7 )
8
9 -// Wait is a readable channel to block on until it receives a signal.
10 -type Wait <-chan Signal
11 -
12 -// Signal is an empty channel
13 -type Signal struct{}
14 -
9 // CloseFunc is a function used to close a ContextCloser
10 type CloseFunc func() error
11
12 // ContextCloser is an interface for services able to be opened and closed.
13 +// It has a parent Context, and Children. But ContextCloser is not a proper
14 +// "tree" like the Context tree. It is more like a Context-WaitGroup hybrid.
15 +// It models a main object with a few children objects -- and, unlike the
16 +// context -- concerns itself with the parent-child closing semantics:
17 +//
18 +// - Can define a CloseFunc (func() error) to be run at Close time.
19 +// - Children call Children().Add(1) to be waited upon
20 +// - Children can select on <-Closing() to know when they should shut down.
21 +// - Close() will wait until all children call Children().Done()
22 +// - <-Closed() signals when the service is completely closed.
23 +//
24 +// ContextCloser can be embedded into the main object itself. In that case,
25 +// the closeFunc (if a member function) has to be set after the struct
26 +// is intialized:
27 +//
28 +// type service struct {
29 +// ContextCloser
30 +// net.Conn
31 +// }
32 +//
33 +// func (s *service) close() error {
34 +// return s.Conn.Close()
35 +// }
36 +//
37 +// func newService(ctx context.Context, c net.Conn) *service {
38 +// s := &service{c}
39 +// s.ContextCloser = NewContextCloser(ctx, s.close)
40 +// return s
41 +// }
42 +//
43 type ContextCloser interface {
44 +
45 + // Context is the context of this ContextCloser. It is "sort of" a parent.
46 Context() context.Context
47
22 - // Close is a method to call when you with to stop this ContextCloser
48 + // Children is a sync.Waitgroup for all children goroutines that should
49 + // shut down completely before this service is said to be "closed".
50 + // Follows the semantics of WaitGroup:
51 + // Children().Add(1) // add one more dependent child
52 + // Children().Done() // child signals it is done
53 + Children() *sync.WaitGroup
54 +
55 + // Close is a method to call when you wish to stop this ContextCloser
56 Close() error
57
25 - // Done is a method to wait upon, like context.Context.Done
26 - Done() Wait
58 + // Closing is a signal to wait upon, like Context.Done().
59 + // It fires when the object should be closing (but hasn't yet fully closed).
60 + // The primary use case is for child goroutines who need to know when
61 + // they should shut down. (equivalent to Context().Done())
62 + Closing() <-chan struct{}
63 +
64 + // Closed is a method to wait upon, like Context.Done().
65 + // It fires when the entire object is fully closed.
66 + // The primary use case is for external listeners who need to know when
67 + // this object is completly done, and all its children closed.
68 + Closed() <-chan struct{}
69 }
70
71 // contextCloser is an OpenCloser with a cancellable context
@@ -31,11 +73,20 @@ type contextCloser struct {
73 ctx context.Context
74 cancel context.CancelFunc
75
34 - // called to close
76 + // called to run the close logic.
77 closeFunc CloseFunc
78
79 // closed is released once the close function is done.
38 - closed chan Signal
80 + closed chan struct{}
81 +
82 + // wait group for child goroutines
83 + children sync.WaitGroup
84 +
85 + // sync primitive to ensure the close logic is only called once.
86 + closeOnce sync.Once
87 +
88 + // error to return to clients of Close().
89 + closeErr error
90 }
91
92 // NewContextCloser constructs and returns a ContextCloser. It will call
@@ -46,7 +97,7 @@ func NewContextCloser(ctx context.Context, cf CloseFunc) ContextCloser {
97 ctx: ctx,
98 cancel: cancel,
99 closeFunc: cf,
49 - closed: make(chan Signal),
100 + closed: make(chan struct{}),
101 }
102
103 go c.closeOnContextDone()
@@ -57,30 +108,47 @@ func (c *contextCloser) Context() context.Context {
108 return c.ctx
109 }
110
60 -func (c *contextCloser) Done() Wait {
61 - return c.closed
111 +func (c *contextCloser) Children() *sync.WaitGroup {
112 + return &c.children
113 }
114
115 +// Close is the external close function. it's a wrapper around internalClose
116 +// that waits on Closed()
117 func (c *contextCloser) Close() error {
65 - select {
66 - case <-c.Done():
67 - // panic("closed twice")
68 - return errors.New("closed twice")
69 - default:
70 - }
118 + c.internalClose()
119 + <-c.Closed() // wait until we're totally done.
120 + return c.closeErr
121 +}
122
72 - err := c.closeFunc() // actually run the close logic
73 - close(c.closed) // relase everyone waiting on Done
74 - c.cancel() // release anyone waiting on the context
75 - return err
123 +func (c *contextCloser) Closing() <-chan struct{} {
124 + return c.Context().Done()
125 }
126
127 +func (c *contextCloser) Closed() <-chan struct{} {
128 + return c.closed
129 +}
130 +
131 +func (c *contextCloser) internalClose() {
132 + go c.closeOnce.Do(c.closeLogic)
133 +}
134 +
135 +// the _actual_ close process.
136 +func (c *contextCloser) closeLogic() {
137 + // this function should only be called once (hence the sync.Once).
138 + // and it will panic at the bottom (on close(c.closed)) otherwise.
139 +
140 + c.cancel() // signal that we're shutting down (Closing)
141 + c.closeErr = c.closeFunc() // actually run the close logic
142 + c.children.Wait() // wait till all children are done.
143 + close(c.closed) // signal that we're shut down (Closed)
144 +}
145 +
146 +// if parent context is shut down before we call Close explicitly,
147 +// we need to go through the Close motions anyway. Hence all the sync
148 +// stuff all over the place...
149 func (c *contextCloser) closeOnContextDone() {
79 - <-c.ctx.Done()
80 - select {
81 - case <-c.Done():
82 - return // already closed
83 - default:
84 - }
85 - c.Close()
150 + c.Children().Add(1) // we're a child goroutine, to be waited upon.
151 + <-c.Context().Done() // wait until parent (context) is done.
152 + c.internalClose()
153 + c.Children().Done()
154 }
net/conn/conn.go
+1 -1
@@ -218,7 +218,7 @@ func (l *listener) close() error {
218
219 func (l *listener) isClosed() bool {
220 select {
221 - case <-l.Done():
221 + case <-l.Closed():
222 return true
223 default:
224 return false
net/conn/conn_test.go
+8 -8
@@ -18,9 +18,9 @@ func TestClose(t *testing.T) {
18 c1, c2 := setupConn(t, ctx, "/ip4/127.0.0.1/tcp/1234", "/ip4/127.0.0.1/tcp/2345")
19
20 select {
21 - case <-c1.Done():
21 + case <-c1.Closed():
22 t.Fatal("done before close")
23 - case <-c2.Done():
23 + case <-c2.Closed():
24 t.Fatal("done before close")
25 default:
26 }
@@ -28,7 +28,7 @@ func TestClose(t *testing.T) {
28 c1.Close()
29
30 select {
31 - case <-c1.Done():
31 + case <-c1.Closed():
32 default:
33 t.Fatal("not done after cancel")
34 }
@@ -36,7 +36,7 @@ func TestClose(t *testing.T) {
36 c2.Close()
37
38 select {
39 - case <-c2.Done():
39 + case <-c2.Closed():
40 default:
41 t.Fatal("not done after cancel")
42 }
@@ -50,9 +50,9 @@ func TestCancel(t *testing.T) {
50 c1, c2 := setupConn(t, ctx, "/ip4/127.0.0.1/tcp/1234", "/ip4/127.0.0.1/tcp/2345")
51
52 select {
53 - case <-c1.Done():
53 + case <-c1.Closed():
54 t.Fatal("done before close")
55 - case <-c2.Done():
55 + case <-c2.Closed():
56 t.Fatal("done before close")
57 default:
58 }
@@ -64,13 +64,13 @@ func TestCancel(t *testing.T) {
64 // test that cancel called Close.
65
66 select {
67 - case <-c1.Done():
67 + case <-c1.Closed():
68 default:
69 t.Fatal("not done after cancel")
70 }
71
72 select {
73 - case <-c2.Done():
73 + case <-c2.Closed():
74 default:
75 t.Fatal("not done after cancel")
76 }
net/conn/handshake.go
+1 -1
@@ -31,7 +31,7 @@ func VersionHandshake(ctx context.Context, c Conn) error {
31 case <-ctx.Done():
32 return ctx.Err()
33
34 - case <-c.Done():
34 + case <-c.Closed():
35 return errors.New("remote closed connection during version exchange")
36
37 case data, ok := <-c.In():
net/conn/secure_conn_test.go
+8 -8
@@ -37,9 +37,9 @@ func TestSecureClose(t *testing.T) {
37 c2 = setupSecureConn(t, c2)
38
39 select {
40 - case <-c1.Done():
40 + case <-c1.Closed():
41 t.Fatal("done before close")
42 - case <-c2.Done():
42 + case <-c2.Closed():
43 t.Fatal("done before close")
44 default:
45 }
@@ -47,7 +47,7 @@ func TestSecureClose(t *testing.T) {
47 c1.Close()
48
49 select {
50 - case <-c1.Done():
50 + case <-c1.Closed():
51 default:
52 t.Fatal("not done after cancel")
53 }
@@ -55,7 +55,7 @@ func TestSecureClose(t *testing.T) {
55 c2.Close()
56
57 select {
58 - case <-c2.Done():
58 + case <-c2.Closed():
59 default:
60 t.Fatal("not done after cancel")
61 }
@@ -72,9 +72,9 @@ func TestSecureCancel(t *testing.T) {
72 c2 = setupSecureConn(t, c2)
73
74 select {
75 - case <-c1.Done():
75 + case <-c1.Closed():
76 t.Fatal("done before close")
77 - case <-c2.Done():
77 + case <-c2.Closed():
78 t.Fatal("done before close")
79 default:
80 }
@@ -86,13 +86,13 @@ func TestSecureCancel(t *testing.T) {
86 // test that cancel called Close.
87
88 select {
89 - case <-c1.Done():
89 + case <-c1.Closed():
90 default:
91 t.Fatal("not done after cancel")
92 }
93
94 select {
95 - case <-c2.Done():
95 + case <-c2.Closed():
96 default:
97 t.Fatal("not done after cancel")
98 }