@cryptotaxi247 / kubo / commits / c2a228f65

use ContextCloser better (listener fix)

Juan Batiz-Benet committed Oct 19, 2014 at 02:26 UTC c2a228f650c2a8d62ca71661d0e53d96c9cbed10
5 files changed +24 -14
net/conn/conn.go
+16 -13
@@ -65,8 +65,16 @@ func newSingleConn(ctx context.Context, local, remote *peer.Peer,
65 log.Info("newSingleConn: %v to %v", local, remote)
66
67 // setup the various io goroutines
68 - go conn.msgio.outgoing.WriteTo(maconn)
69 - go conn.msgio.incoming.ReadFrom(maconn, MaxMessageSize)
68 + go func() {
69 + conn.Children().Add(1)
70 + conn.msgio.outgoing.WriteTo(maconn)
71 + conn.Children().Done()
72 + }()
73 + go func() {
74 + conn.Children().Add(1)
75 + conn.msgio.incoming.ReadFrom(maconn, MaxMessageSize)
76 + conn.Children().Done()
77 + }()
78
79 // version handshake
80 ctxT, _ := context.WithTimeout(ctx, HandshakeTimeout)
@@ -216,16 +224,9 @@ func (l *listener) close() error {
224 return l.Listener.Close()
225 }
226
219 -func (l *listener) isClosed() bool {
220 - select {
221 - case <-l.Closed():
222 - return true
223 - default:
224 - return false
225 - }
226 -}
227 -
227 func (l *listener) listen() {
228 + l.Children().Add(1)
229 + defer l.Children().Done()
230
231 // handle at most chansize concurrent handshakes
232 sem := make(chan struct{}, l.chansize)
@@ -254,9 +255,11 @@ func (l *listener) listen() {
255 maconn, err := l.Listener.Accept()
256 if err != nil {
257
257 - // if cancel is nil we're closed.
258 - if l.isClosed() {
258 + // if closing, we should exit.
259 + select {
260 + case <-l.Closing():
261 return // done.
262 + default:
263 }
264
265 log.Error("Failed to accept connection: %v", err)
net/conn/conn_test.go
+3
@@ -13,6 +13,7 @@ import (
13 )
14
15 func TestClose(t *testing.T) {
16 + // t.Skip("Skipping in favor of another test")
17
18 ctx, cancel := context.WithCancel(context.Background())
19 c1, c2 := setupConn(t, ctx, "/ip4/127.0.0.1/tcp/1234", "/ip4/127.0.0.1/tcp/2345")
@@ -45,6 +46,7 @@ func TestClose(t *testing.T) {
46 }
47
48 func TestCancel(t *testing.T) {
49 + // t.Skip("Skipping in favor of another test")
50
51 ctx, cancel := context.WithCancel(context.Background())
52 c1, c2 := setupConn(t, ctx, "/ip4/127.0.0.1/tcp/1234", "/ip4/127.0.0.1/tcp/2345")
@@ -78,6 +80,7 @@ func TestCancel(t *testing.T) {
80 }
81
82 func TestCloseLeak(t *testing.T) {
83 + // t.Skip("Skipping in favor of another test")
84
85 var wg sync.WaitGroup
86
net/conn/dial_test.go
+1
@@ -93,6 +93,7 @@ func setupConn(t *testing.T, ctx context.Context, a1, a2 string) (a, b Conn) {
93 }
94
95 func TestDialer(t *testing.T) {
96 + // t.Skip("Skipping in favor of another test")
97
98 p1, err := setupPeer("/ip4/127.0.0.1/tcp/1234")
99 if err != nil {
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.Closed():
34 + case <-c.Closing():
35 return errors.New("remote closed connection during version exchange")
36
37 case data, ok := <-c.In():
net/conn/secure_conn_test.go
+3
@@ -29,6 +29,7 @@ func setupSecureConn(t *testing.T, c Conn) Conn {
29 }
30
31 func TestSecureClose(t *testing.T) {
32 + // t.Skip("Skipping in favor of another test")
33
34 ctx, cancel := context.WithCancel(context.Background())
35 c1, c2 := setupConn(t, ctx, "/ip4/127.0.0.1/tcp/1234", "/ip4/127.0.0.1/tcp/2345")
@@ -64,6 +65,7 @@ func TestSecureClose(t *testing.T) {
65 }
66
67 func TestSecureCancel(t *testing.T) {
68 + // t.Skip("Skipping in favor of another test")
69
70 ctx, cancel := context.WithCancel(context.Background())
71 c1, c2 := setupConn(t, ctx, "/ip4/127.0.0.1/tcp/1234", "/ip4/127.0.0.1/tcp/2345")
@@ -100,6 +102,7 @@ func TestSecureCancel(t *testing.T) {
102 }
103
104 func TestSecureCloseLeak(t *testing.T) {
105 + // t.Skip("Skipping in favor of another test")
106
107 var wg sync.WaitGroup
108