@cryptotaxi247 / kubo / commits / d26fd5818

ctx closer races #270

Juan Batiz-Benet committed Nov 5, 2014 at 07:45 UTC d26fd581827f44a7e3e7075e84fe96779beb02aa
6 files changed +11 -12
net/conn/conn.go
+2 -2
@@ -89,13 +89,13 @@ func newSingleConn(ctx context.Context, local, remote peer.Peer,
89 log.Info("newSingleConn: %v to %v", local, remote)
90
91 // setup the various io goroutines
92 + conn.Children().Add(1)
93 go func() {
93 - conn.Children().Add(1)
94 conn.msgio.outgoing.WriteTo(maconn)
95 conn.Children().Done()
96 }()
97 + conn.Children().Add(1)
98 go func() {
98 - conn.Children().Add(1)
99 conn.msgio.incoming.ReadFrom(maconn, MaxMessageSize)
100 conn.Children().Done()
101 }()
net/conn/listen.go
+1 -1
@@ -47,7 +47,6 @@ func (l *listener) close() error {
47 }
48
49 func (l *listener) listen() {
50 - l.Children().Add(1)
50 defer l.Children().Done()
51
52 // handle at most chansize concurrent handshakes
@@ -143,6 +142,7 @@ func Listen(ctx context.Context, addr ma.Multiaddr, local peer.Peer, peers peer.
142 ctx2, _ := context.WithCancel(ctx)
143 l.ContextCloser = ctxc.NewContextCloser(ctx2, l.close)
144
145 + l.Children().Add(1)
146 go l.listen()
147
148 return l, nil
net/conn/multiconn.go
+4 -4
@@ -57,6 +57,8 @@ func NewMultiConn(ctx context.Context, local, remote peer.Peer, conns []Conn) (*
57 if conns != nil && len(conns) > 0 {
58 c.Add(conns...)
59 }
60 +
61 + c.Children().Add(1)
62 go c.fanOut()
63 return c, nil
64 }
@@ -81,6 +83,8 @@ func (c *MultiConn) Add(conns ...Conn) {
83 }
84
85 c.conns[c2.ID()] = c2
86 + c.Children().Add(1)
87 + c2.Children().Add(1) // yep, on the child too.
88 go c.fanInSingle(c2)
89 log.Infof("MultiConn: added %s", c2)
90 }
@@ -134,7 +138,6 @@ func CloseConns(conns ...Conn) {
138 // fanOut is the multiplexor out -- it sends outgoing messages over the
139 // underlying single connections.
140 func (c *MultiConn) fanOut() {
137 - c.Children().Add(1)
141 defer c.Children().Done()
142
143 i := 0
@@ -165,9 +168,6 @@ func (c *MultiConn) fanOut() {
168 // fanInSingle is a multiplexor in -- it receives incoming messages over the
169 // underlying single connections.
170 func (c *MultiConn) fanInSingle(child Conn) {
168 - c.Children().Add(1)
169 - child.Children().Add(1) // yep, on the child too.
170 -
171 // cleanup all data associated with this child Connection.
172 defer func() {
173 log.Infof("closing: %s", child)
net/swarm/conn.go
+2 -4
@@ -139,6 +139,8 @@ func (s *Swarm) connSetup(c conn.Conn) (conn.Conn, error) {
139 s.connsLock.Unlock()
140
141 // kick off reader goroutine
142 + s.Children().Add(1)
143 + mc.Children().Add(1) // child of Conn as well.
144 go s.fanInSingle(mc)
145 log.Debugf("added new multiconn: %s", mc)
146 } else {
@@ -154,7 +156,6 @@ func (s *Swarm) connSetup(c conn.Conn) (conn.Conn, error) {
156
157 // Handles the unwrapping + sending of messages to the right connection.
158 func (s *Swarm) fanOut() {
157 - s.Children().Add(1)
159 defer s.Children().Done()
160
161 i := 0
@@ -194,9 +195,6 @@ func (s *Swarm) fanOut() {
195 // Handles the receiving + wrapping of messages, per conn.
196 // Consider using reflect.Select with one goroutine instead of n.
197 func (s *Swarm) fanInSingle(c conn.Conn) {
197 - s.Children().Add(1)
198 - c.Children().Add(1) // child of Conn as well.
199 -
198 // cleanup all data associated with this child Connection.
199 defer func() {
200 // remove it from the map.
net/swarm/swarm.go
+1
@@ -83,6 +83,7 @@ func NewSwarm(ctx context.Context, listenAddrs []ma.Multiaddr, local peer.Peer,
83 // ContextCloser for proper child management.
84 s.ContextCloser = ctxc.NewContextCloser(ctx, s.close)
85
86 + s.Children().Add(1)
87 go s.fanOut()
88 return s, s.listen(listenAddrs)
89 }
util/ctxcloser/closer.go
+1 -1
@@ -120,6 +120,7 @@ func NewContextCloser(ctx context.Context, cf CloseFunc) ContextCloser {
120 closed: make(chan struct{}),
121 }
122
123 + c.Children().Add(1) // we're a child goroutine, to be waited upon.
124 go c.closeOnContextDone()
125 return c
126 }
@@ -176,7 +177,6 @@ func (c *contextCloser) closeLogic() {
177 // we need to go through the Close motions anyway. Hence all the sync
178 // stuff all over the place...
179 func (c *contextCloser) closeOnContextDone() {
179 - c.Children().Add(1) // we're a child goroutine, to be waited upon.
180 <-c.Context().Done() // wait until parent (context) is done.
181 c.internalClose()
182 c.Children().Done()