@cryptotaxi247 / kubo / commits / afe85ce1c

add in basic bandwidth tracking to the muxer

Jeromy committed Oct 13, 2014 at 10:42 UTC afe85ce1c8cbef0a06e8e26297a955a3d47582ec
1 file changed +27
net/mux/mux.go
+27
@@ -36,6 +36,12 @@ type Muxer struct {
36 ctx context.Context
37 wg sync.WaitGroup
38
39 + bwiLock sync.Mutex
40 + bwIn uint64
41 +
42 + bwoLock sync.Mutex
43 + bwOut uint64
44 +
45 *msg.Pipe
46 }
47
@@ -76,6 +82,17 @@ func (m *Muxer) Start(ctx context.Context) error {
82 return nil
83 }
84
85 +func (m *Muxer) GetBandwidthTotals() (in uint64, out uint64) {
86 + m.bwiLock.Lock()
87 + in = m.bwIn
88 + m.bwiLock.Unlock()
89 +
90 + m.bwoLock.Lock()
91 + out = m.bwOut
92 + m.bwoLock.Unlock()
93 + return
94 +}
95 +
96 // Stop stops muxer activity.
97 func (m *Muxer) Stop() {
98 if m.cancel == nil {
@@ -125,6 +142,11 @@ func (m *Muxer) handleIncomingMessages() {
142 // handleIncomingMessage routes message to the appropriate protocol.
143 func (m *Muxer) handleIncomingMessage(m1 msg.NetMessage) {
144
145 + m.bwiLock.Lock()
146 + // TODO: compensate for overhead
147 + m.bwIn += uint64(len(m1.Data()))
148 + m.bwiLock.Unlock()
149 +
150 data, pid, err := unwrapData(m1.Data())
151 if err != nil {
152 log.Error("muxer de-serializing error: %v", err)
@@ -173,6 +195,11 @@ func (m *Muxer) handleOutgoingMessage(pid ProtocolID, m1 msg.NetMessage) {
195 return
196 }
197
198 + m.bwoLock.Lock()
199 + // TODO: compensate for overhead
200 + m.bwOut += uint64(len(data))
201 + m.bwoLock.Unlock()
202 +
203 m2 := msg.New(m1.Peer(), data)
204 select {
205 case m.GetPipe().Outgoing <- m2: