@cryptotaxi247 / kubo / commits / 8aed79cd9

fixed data races

Juan Batiz-Benet committed Oct 17, 2014 at 01:47 UTC 8aed79cd9794487294a655b1456c8652b6a31f1b
2 files changed +51 -37
crypto/spipe/pipe.go
+16 -20
@@ -34,34 +34,29 @@ type params struct {
34
35 // NewSecurePipe constructs a pipe with channels of a given buffer size.
36 func NewSecurePipe(ctx context.Context, bufsize int, local *peer.Peer,
37 - peers peer.Peerstore) (*SecurePipe, error) {
37 + peers peer.Peerstore, insecure Duplex) (*SecurePipe, error) {
38 +
39 + ctx, cancel := context.WithCancel(ctx)
40
41 sp := &SecurePipe{
42 Duplex: Duplex{
43 In: make(chan []byte, bufsize),
44 Out: make(chan []byte, bufsize),
45 },
44 - local: local,
45 - peers: peers,
46 - }
47 - return sp, nil
48 -}
46 + local: local,
47 + peers: peers,
48 + insecure: insecure,
49
50 -// Wrap creates a secure connection on top of an insecure duplex channel.
51 -func (s *SecurePipe) Wrap(ctx context.Context, insecure Duplex) error {
52 - if s.ctx != nil {
53 - return errors.New("Pipe in use")
50 + ctx: ctx,
51 + cancel: cancel,
52 }
53
56 - s.insecure = insecure
57 - s.ctx, s.cancel = context.WithCancel(ctx)
58 -
59 - if err := s.handshake(); err != nil {
60 - s.cancel()
61 - return err
54 + if err := sp.handshake(); err != nil {
55 + sp.Close()
56 + return nil, err
57 }
58
64 - return nil
59 + return sp, nil
60 }
61
62 // LocalPeer retrieves the local peer.
@@ -76,11 +71,12 @@ func (s *SecurePipe) RemotePeer() *peer.Peer {
71
72 // Close closes the secure pipe
73 func (s *SecurePipe) Close() error {
79 - if s.cancel == nil {
80 - return errors.New("pipe already closed")
74 + select {
75 + case <-s.ctx.Done():
76 + return errors.New("already closed")
77 + default:
78 }
79
80 s.cancel()
84 - s.cancel = nil
81 return nil
82 }
net/conn/conn.go
+35 -17
@@ -45,6 +45,7 @@ type singleConn struct {
45 // context + cancel
46 ctx context.Context
47 cancel context.CancelFunc
48 + closed chan struct{}
49
50 secure *spipe.SecurePipe
51 insecure *msgioPipe
@@ -66,6 +67,7 @@ func newSingleConn(ctx context.Context, local, remote *peer.Peer,
67 maconn: maconn,
68 ctx: ctx,
69 cancel: cancel,
70 + closed: make(chan struct{}),
71 insecure: newMsgioPipe(10),
72 msgpipe: msg.NewPipe(10),
73 }
@@ -92,20 +94,16 @@ func (c *singleConn) secureHandshake(peers peer.Peerstore) error {
94 return errors.New("Conn is already secured or being secured.")
95 }
96
95 - var err error
96 - c.secure, err = spipe.NewSecurePipe(c.ctx, 10, c.local, peers)
97 - if err != nil {
98 - return err
99 - }
100 -
97 // setup a Duplex pipe for spipe
98 insecure := spipe.Duplex{
99 In: c.insecure.incoming.MsgChan,
100 Out: c.insecure.outgoing.MsgChan,
101 }
102
107 - // Wrap actually performs the secure handshake, which takes multiple RTT
108 - if err := c.secure.Wrap(c.ctx, insecure); err != nil {
103 + // spipe performs the secure handshake, which takes multiple RTT
104 + var err error
105 + c.secure, err = spipe.NewSecurePipe(c.ctx, 10, c.local, peers, insecure)
106 + if err != nil {
107 return err
108 }
109
@@ -166,29 +164,33 @@ func (c *singleConn) waitToClose(ctx context.Context) {
164
165 // close underlying connection
166 c.maconn.Close()
169 - c.maconn = nil
167
168 // closing channels
169 c.insecure.outgoing.Close()
170 c.secure.Close()
171 close(c.msgpipe.Incoming)
172 + close(c.closed)
173 }
174
177 -// IsOpen returns whether this Conn is open or closed.
178 -func (c *singleConn) isOpen() bool {
179 - return c.maconn != nil
175 +// isClosed returns whether this Conn is open or closed.
176 +func (c *singleConn) isClosed() bool {
177 + select {
178 + case <-c.closed:
179 + return true
180 + default:
181 + return false
182 + }
183 }
184
185 // Close closes the connection, and associated channels.
186 func (c *singleConn) Close() error {
187 log.Debug("%s closing Conn with %s", c.local, c.remote)
185 - if !c.isOpen() {
186 - return fmt.Errorf("Already closed") // already closed
188 + if c.isClosed() {
189 + return fmt.Errorf("connection already closed")
190 }
191
192 // cancel context.
193 c.cancel()
191 - c.cancel = nil
194 return nil
195 }
196
@@ -278,6 +280,7 @@ type listener struct {
280 // ctx + cancel func
281 ctx context.Context
282 cancel context.CancelFunc
283 + closed chan struct{}
284 }
285
286 // waitToClose is needed to hand
@@ -286,8 +289,17 @@ func (l *listener) waitToClose() {
289 case <-l.ctx.Done():
290 }
291
289 - l.cancel = nil
292 l.Listener.Close()
293 + close(l.closed)
294 +}
295 +
296 +func (l *listener) isClosed() bool {
297 + select {
298 + case <-l.closed:
299 + return true
300 + default:
301 + return false
302 + }
303 }
304
305 func (l *listener) listen() {
@@ -312,7 +324,7 @@ func (l *listener) listen() {
324 if err != nil {
325
326 // if cancel is nil we're closed.
315 - if l.cancel == nil {
327 + if l.isClosed() {
328 return // done.
329 }
330
@@ -351,7 +363,12 @@ func (l *listener) Peerstore() peer.Peerstore {
363 // Close closes the listener.
364 // Any blocked Accept operations will be unblocked and return errors
365 func (l *listener) Close() error {
366 + if l.isClosed() {
367 + return errors.New("listener already closed")
368 + }
369 +
370 l.cancel()
371 + <-l.closed
372 return nil
373 }
374
@@ -371,6 +388,7 @@ func Listen(ctx context.Context, addr ma.Multiaddr, local *peer.Peer, peers peer
388 l := &listener{
389 ctx: ctx,
390 cancel: cancel,
391 + closed: make(chan struct{}),
392 Listener: ml,
393 maddr: addr,
394 peers: peers,