@cryptotaxi247 / kubo / commits / 8065b61c3

Added ContextCloser abstraction

Juan Batiz-Benet committed Oct 18, 2014 at 02:47 UTC 8065b61c305e7170c99e69e616e4267cf1e6b8a3
2 files changed +107 -64
net/conn/closer.go new
+83
@@ -0,0 +1,83 @@
1 +package conn
2 +
3 +import (
4 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5 +)
6 +
7 +// Wait is a readable channel to block on until it receives a signal.
8 +type Wait <-chan Signal
9 +
10 +// Signal is an empty channel
11 +type Signal struct{}
12 +
13 +// CloseFunc is a function used to close a ContextCloser
14 +type CloseFunc func() error
15 +
16 +// ContextCloser is an interface for services able to be opened and closed.
17 +type ContextCloser interface {
18 + Context() context.Context
19 +
20 + // Close is a method to call when you with to stop this ContextCloser
21 + Close() error
22 +
23 + // Done is a method to wait upon, like context.Context.Done
24 + Done() Wait
25 +}
26 +
27 +// contextCloser is an OpenCloser with a cancellable context
28 +type contextCloser struct {
29 + ctx context.Context
30 + cancel context.CancelFunc
31 +
32 + // called to close
33 + closeFunc CloseFunc
34 +
35 + // closed is released once the close function is done.
36 + closed chan Signal
37 +}
38 +
39 +// NewContextCloser constructs and returns a ContextCloser. It will call
40 +// cf CloseFunc before its Done() Wait signals fire.
41 +func NewContextCloser(ctx context.Context, cf CloseFunc) ContextCloser {
42 + ctx, cancel := context.WithCancel(ctx)
43 + c := &contextCloser{
44 + ctx: ctx,
45 + cancel: cancel,
46 + closeFunc: cf,
47 + closed: make(chan Signal),
48 + }
49 +
50 + go c.closeOnContextDone()
51 + return c
52 +}
53 +
54 +func (c *contextCloser) Context() context.Context {
55 + return c.ctx
56 +}
57 +
58 +func (c *contextCloser) Done() Wait {
59 + return c.closed
60 +}
61 +
62 +func (c *contextCloser) Close() error {
63 + select {
64 + case <-c.Done():
65 + panic("closed twice")
66 + default:
67 + }
68 +
69 + c.cancel() // release anyone waiting on the context
70 + err := c.closeFunc() // actually run the close logic
71 + close(c.closed) // relase everyone waiting on Done
72 + return err
73 +}
74 +
75 +func (c *contextCloser) closeOnContextDone() {
76 + <-c.ctx.Done()
77 + select {
78 + case <-c.Done():
79 + return // already closed
80 + default:
81 + }
82 + c.Close()
83 +}
net/conn/conn.go
+24 -64
@@ -41,35 +41,30 @@ type singleConn struct {
41 remote *peer.Peer
42 maconn manet.Conn
43
44 - // context + cancel
45 - ctx context.Context
46 - cancel context.CancelFunc
47 -
44 secure *spipe.SecurePipe
45 insecure *msgioPipe
46 +
47 + ContextCloser
48 }
49
50 // newConn constructs a new connection
51 func newSingleConn(ctx context.Context, local, remote *peer.Peer,
52 peers peer.Peerstore, maconn manet.Conn) (Conn, error) {
53
56 - ctx, cancel := context.WithCancel(ctx)
57 -
54 conn := &singleConn{
55 local: local,
56 remote: remote,
57 maconn: maconn,
62 - ctx: ctx,
63 - cancel: cancel,
58 insecure: newMsgioPipe(10),
59 }
60
61 + conn.ContextCloser = NewContextCloser(ctx, conn.close)
62 +
63 log.Info("newSingleConn: %v to %v", local, remote)
64
65 // setup the various io goroutines
66 go conn.insecure.outgoing.WriteTo(maconn)
67 go conn.insecure.incoming.ReadFrom(maconn, MaxMessageSize)
72 - go conn.waitToClose()
68
69 // perform secure handshake before returning this connection.
70 if err := conn.secureHandshake(peers); err != nil {
@@ -93,7 +88,7 @@ func (c *singleConn) secureHandshake(peers peer.Peerstore) error {
88 }
89
90 // spipe performs the secure handshake, which takes multiple RTT
96 - sp, err := spipe.NewSecurePipe(c.ctx, 10, c.local, peers, insecure)
91 + sp, err := spipe.NewSecurePipe(c.Context(), 10, c.local, peers, insecure)
92 if err != nil {
93 return err
94 }
@@ -114,42 +109,20 @@ func (c *singleConn) secureHandshake(peers peer.Peerstore) error {
109 return nil
110 }
111
117 -// waitToClose waits on the given context's Done before closing Conn.
118 -func (c *singleConn) waitToClose() {
119 - select {
120 - case <-c.ctx.Done():
121 - }
112 +// close is the internal close function, called by ContextCloser.Close
113 +func (c *singleConn) close() error {
114 + log.Debug("%s closing Conn with %s", c.local, c.remote)
115
116 // close underlying connection
124 - c.maconn.Close()
117 + err := c.maconn.Close()
118
119 // closing channels
120 c.insecure.outgoing.Close()
121 if c.secure != nil { // may never have gotten here.
122 c.secure.Close()
123 }
131 -}
124
133 -// isClosed returns whether this Conn is open or closed.
134 -func (c *singleConn) isClosed() bool {
135 - select {
136 - case <-c.ctx.Done():
137 - return true
138 - default:
139 - return false
140 - }
141 -}
142 -
143 -// Close closes the connection, and associated channels.
144 -func (c *singleConn) Close() error {
145 - log.Debug("%s closing Conn with %s", c.local, c.remote)
146 - if c.isClosed() {
147 - return fmt.Errorf("connection already closed")
148 - }
149 -
150 - // cancel context.
151 - c.cancel()
152 - return nil
125 + return err
126 }
127
128 // LocalPeer is the Peer on this side
@@ -235,23 +208,24 @@ type listener struct {
208 // Peerstore is the set of peers we know about locally
209 peers peer.Peerstore
210
238 - // ctx + cancel func
239 - ctx context.Context
240 - cancel context.CancelFunc
211 + // embedded ContextCloser
212 + ContextCloser
213 }
214
243 -// waitToClose is needed to hand
244 -func (l *listener) waitToClose() {
245 - select {
246 - case <-l.ctx.Done():
247 - }
215 +// disambiguate
216 +func (l *listener) Close() error {
217 + return l.ContextCloser.Close()
218 +}
219
249 - l.Listener.Close()
220 +// close called by ContextCloser.Close
221 +func (l *listener) close() error {
222 + log.Info("listener closing: %s %s", l.local, l.maddr)
223 + return l.Listener.Close()
224 }
225
226 func (l *listener) isClosed() bool {
227 select {
254 - case <-l.ctx.Done():
228 + case <-l.Done():
229 return true
230 default:
231 return false
@@ -266,7 +240,7 @@ func (l *listener) listen() {
240 // handle is a goroutine work function that handles the handshake.
241 // it's here only so that accepting new connections can happen quickly.
242 handle := func(maconn manet.Conn) {
269 - c, err := newSingleConn(l.ctx, l.local, nil, l.peers, maconn)
243 + c, err := newSingleConn(l.Context(), l.local, nil, l.peers, maconn)
244 if err != nil {
245 log.Error("Error accepting connection: %v", err)
246 } else {
@@ -316,22 +290,9 @@ func (l *listener) Peerstore() peer.Peerstore {
290 return l.peers
291 }
292
319 -// Close closes the listener.
320 -// Any blocked Accept operations will be unblocked and return errors
321 -func (l *listener) Close() error {
322 - if l.isClosed() {
323 - return errors.New("listener already closed")
324 - }
325 -
326 - l.cancel()
327 - return nil
328 -}
329 -
293 // Listen listens on the particular multiaddr, with given peer and peerstore.
294 func Listen(ctx context.Context, addr ma.Multiaddr, local *peer.Peer, peers peer.Peerstore) (Listener, error) {
295
333 - ctx, cancel := context.WithCancel(ctx)
334 -
296 ml, err := manet.Listen(addr)
297 if err != nil {
298 return nil, err
@@ -341,8 +302,6 @@ func Listen(ctx context.Context, addr ma.Multiaddr, local *peer.Peer, peers peer
302 chansize := 10
303
304 l := &listener{
344 - ctx: ctx,
345 - cancel: cancel,
305 Listener: ml,
306 maddr: addr,
307 peers: peers,
@@ -351,8 +310,9 @@ func Listen(ctx context.Context, addr ma.Multiaddr, local *peer.Peer, peers peer
310 chansize: chansize,
311 }
312
313 + l.ContextCloser = NewContextCloser(ctx, l.close)
314 +
315 go l.listen()
355 - go l.waitToClose()
316
317 return l, nil
318 }