added msg counters to logs
Juan Batiz-Benet committed
Oct 19, 2014 at 05:49 UTC
29ab6dec6026b073af42edb5f3b3f244dd306e1c
3 files changed
+23
-5
net/conn/multiconn.go
+11
@@ -132,6 +132,7 @@ func (c *MultiConn) fanOut() {
132
c.Children().Add(1)
133
defer c.Children().Done()
134
135
+ i := 0
136
for {
137
select {
138
case <-c.Closing():
@@ -140,6 +141,7 @@ func (c *MultiConn) fanOut() {
141
// send data out through our "best connection"
142
case m, more := <-c.duplex.Out:
143
if !more {
144
+ log.Info("%s out channel closed", c)
145
return
146
}
147
sc := c.BestConn()
@@ -147,6 +149,9 @@ func (c *MultiConn) fanOut() {
149
// maybe this should be a logged error, not a panic.
150
panic("sending out multiconn without any live connection")
151
}
152
+
153
+ i++
154
+ log.Info("%s sending (%d)", sc, i)
155
sc.Out() <- m
156
}
157
}
@@ -160,6 +165,8 @@ func (c *MultiConn) fanInSingle(child Conn) {
165
166
// cleanup all data associated with this child Connection.
167
defer func() {
168
+ log.Info("closing: %s", child)
169
+
170
// in case it still is in the map, remove it.
171
c.Lock()
172
delete(c.conns, child.ID())
@@ -174,6 +181,7 @@ func (c *MultiConn) fanInSingle(child Conn) {
181
}
182
}()
183
184
+ i := 0
185
for {
186
select {
187
case <-c.Closing(): // multiconn closing
@@ -184,8 +192,11 @@ func (c *MultiConn) fanInSingle(child Conn) {
192
193
case m, more := <-child.In(): // receiving data
194
if !more {
195
+ log.Info("%s in channel closed", child)
196
return // closed
197
}
198
+ i++
199
+ log.Info("%s received (%d)", child, i)
200
c.duplex.In <- m
201
}
202
}
net/swarm/conn.go
+8
-3
@@ -136,6 +136,7 @@ func (s *Swarm) fanOut() {
136
s.Children().Add(1)
137
defer s.Children().Done()
138
139
+ i := 0
140
for {
141
select {
142
case <-s.Closing():
@@ -143,6 +144,7 @@ func (s *Swarm) fanOut() {
144
145
case msg, ok := <-s.Outgoing:
146
if !ok {
147
+ log.Info("%s outgoing channel closed", s)
148
return
149
}
150
@@ -157,8 +159,8 @@ func (s *Swarm) fanOut() {
159
continue
160
}
161
160
- // log.Debug("[peer: %s] Sent message [to = %s]", s.local, msg.Peer())
161
-
162
+ i++
163
+ log.Debug("%s sent message to %s (%d)", s.local, msg.Peer(), i)
164
// queue it in the connection's buffer
165
c.Out() <- msg.Data()
166
}
@@ -182,6 +184,7 @@ func (s *Swarm) fanInSingle(c conn.Conn) {
184
c.Children().Done() // child of Conn as well.
185
}()
186
187
+ i := 0
188
for {
189
select {
190
case <-s.Closing(): // Swarm closing
@@ -192,9 +195,11 @@ func (s *Swarm) fanInSingle(c conn.Conn) {
195
196
case data, ok := <-c.In():
197
if !ok {
198
+ log.Info("%s in channel closed", c)
199
return // channel closed.
200
}
197
- // log.Debug("[peer: %s] Received message [from = %s]", s.local, c.Peer)
201
+ i++
202
+ log.Debug("%s received message from %s (%d)", s.local, c.RemotePeer(), i)
203
s.Incoming <- msg.New(c.RemotePeer(), data)
204
}
205
}
net/swarm/swarm_test.go
+4
-2
@@ -16,6 +16,7 @@ import (
16
)
17
18
func pong(ctx context.Context, swarm *Swarm) {
19
+ i := 0
20
for {
21
select {
22
case <-ctx.Done():
@@ -23,7 +24,8 @@ func pong(ctx context.Context, swarm *Swarm) {
24
case m1 := <-swarm.Incoming:
25
if bytes.Equal(m1.Data(), []byte("ping")) {
26
m2 := msg.New(m1.Peer(), []byte("pong"))
26
- log.Debug("%s pong %s", swarm.local, m1.Peer())
27
+ i++
28
+ log.Debug("%s pong %s (%d)", swarm.local, m1.Peer(), i)
29
swarm.Outgoing <- m2
30
}
31
}
@@ -132,7 +134,7 @@ func SubtestSwarm(t *testing.T, addrs []string, MsgNum int) {
134
135
for k := 0; k < MsgNum; k++ {
136
for _, p := range *peers {
135
- log.Debug("%s ping %s", s1.local, p)
137
+ log.Debug("%s ping %s (%d)", s1.local, p, k)
138
s1.Outgoing <- msg.New(p, []byte("ping"))
139
}
140
}