@cryptotaxi247 / kubo / commits / ec481b5ad

refactor(net/mux) move proto to internal pb package

Brian Tiger Chow committed Oct 22, 2014 at 04:23 UTC ec481b5ad434937167a9cbeddc4a680526fc243c
5 files changed +26 -18
net/mux/internal/pb/Makefile renamed
net/mux/internal/pb/mux.pb.go renamed
net/mux/internal/pb/mux.proto renamed
net/mux/mux.go
+15 -8
@@ -5,6 +5,7 @@ import (
5 "sync"
6
7 msg "github.com/jbenet/go-ipfs/net/message"
8 + pb "github.com/jbenet/go-ipfs/net/mux/internal/pb"
9 u "github.com/jbenet/go-ipfs/util"
10
11 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
@@ -13,6 +14,12 @@ import (
14
15 var log = u.Logger("muxer")
16
17 +var (
18 + ProtocolID_Routing = pb.ProtocolID_Routing
19 + ProtocolID_Exchange = pb.ProtocolID_Exchange
20 + ProtocolID_Diagnostic = pb.ProtocolID_Diagnostic
21 +)
22 +
23 // Protocol objects produce + consume raw data. They are added to the Muxer
24 // with a ProtocolID, which is added to outgoing payloads. Muxer properly
25 // encapsulates and decapsulates when interfacing with its Protocols. The
@@ -22,7 +29,7 @@ type Protocol interface {
29 }
30
31 // ProtocolMap maps ProtocolIDs to Protocols.
25 -type ProtocolMap map[ProtocolID]Protocol
32 +type ProtocolMap map[pb.ProtocolID]Protocol
33
34 // Muxer is a simple multiplexor that reads + writes to Incoming and Outgoing
35 // channels. It multiplexes various protocols, wrapping and unwrapping data
@@ -107,7 +114,7 @@ func (m *Muxer) Stop() {
114 }
115
116 // AddProtocol adds a Protocol with given ProtocolID to the Muxer.
110 -func (m *Muxer) AddProtocol(p Protocol, pid ProtocolID) error {
117 +func (m *Muxer) AddProtocol(p Protocol, pid pb.ProtocolID) error {
118 if _, found := m.Protocols[pid]; found {
119 return errors.New("Another protocol already using this ProtocolID")
120 }
@@ -170,7 +177,7 @@ func (m *Muxer) handleIncomingMessage(m1 msg.NetMessage) {
177
178 // handleOutgoingMessages consumes the messages on the proto.Outgoing channel,
179 // wraps them and sends them out.
173 -func (m *Muxer) handleOutgoingMessages(pid ProtocolID, proto Protocol) {
180 +func (m *Muxer) handleOutgoingMessages(pid pb.ProtocolID, proto Protocol) {
181 defer m.wg.Done()
182
183 for {
@@ -188,7 +195,7 @@ func (m *Muxer) handleOutgoingMessages(pid ProtocolID, proto Protocol) {
195 }
196
197 // handleOutgoingMessage wraps out a message and sends it out the
191 -func (m *Muxer) handleOutgoingMessage(pid ProtocolID, m1 msg.NetMessage) {
198 +func (m *Muxer) handleOutgoingMessage(pid pb.ProtocolID, m1 msg.NetMessage) {
199 data, err := wrapData(m1.Data(), pid)
200 if err != nil {
201 log.Error("muxer serializing error: %v", err)
@@ -208,9 +215,9 @@ func (m *Muxer) handleOutgoingMessage(pid ProtocolID, m1 msg.NetMessage) {
215 }
216 }
217
211 -func wrapData(data []byte, pid ProtocolID) ([]byte, error) {
218 +func wrapData(data []byte, pid pb.ProtocolID) ([]byte, error) {
219 // Marshal
213 - pbm := new(PBProtocolMessage)
220 + pbm := new(pb.PBProtocolMessage)
221 pbm.ProtocolID = &pid
222 pbm.Data = data
223 b, err := proto.Marshal(pbm)
@@ -221,9 +228,9 @@ func wrapData(data []byte, pid ProtocolID) ([]byte, error) {
228 return b, nil
229 }
230
224 -func unwrapData(data []byte) ([]byte, ProtocolID, error) {
231 +func unwrapData(data []byte) ([]byte, pb.ProtocolID, error) {
232 // Unmarshal
226 - pbm := new(PBProtocolMessage)
233 + pbm := new(pb.PBProtocolMessage)
234 err := proto.Unmarshal(data, pbm)
235 if err != nil {
236 return nil, 0, err
net/mux/mux_test.go
+11 -10
@@ -8,6 +8,7 @@ import (
8
9 mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
10 msg "github.com/jbenet/go-ipfs/net/message"
11 + pb "github.com/jbenet/go-ipfs/net/mux/internal/pb"
12 peer "github.com/jbenet/go-ipfs/peer"
13
14 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
@@ -37,7 +38,7 @@ func testMsg(t *testing.T, m msg.NetMessage, data []byte) {
38 }
39 }
40
40 -func testWrappedMsg(t *testing.T, m msg.NetMessage, pid ProtocolID, data []byte) {
41 +func testWrappedMsg(t *testing.T, m msg.NetMessage, pid pb.ProtocolID, data []byte) {
42 data2, pid2, err := unwrapData(m.Data())
43 if err != nil {
44 t.Error(err)
@@ -57,8 +58,8 @@ func TestSimpleMuxer(t *testing.T) {
58 // setup
59 p1 := &TestProtocol{Pipe: msg.NewPipe(10)}
60 p2 := &TestProtocol{Pipe: msg.NewPipe(10)}
60 - pid1 := ProtocolID_Test
61 - pid2 := ProtocolID_Routing
61 + pid1 := pb.ProtocolID_Test
62 + pid2 := pb.ProtocolID_Routing
63 mux1 := NewMuxer(ProtocolMap{
64 pid1: p1,
65 pid2: p2,
@@ -108,8 +109,8 @@ func TestSimultMuxer(t *testing.T) {
109 // setup
110 p1 := &TestProtocol{Pipe: msg.NewPipe(10)}
111 p2 := &TestProtocol{Pipe: msg.NewPipe(10)}
111 - pid1 := ProtocolID_Test
112 - pid2 := ProtocolID_Identify
112 + pid1 := pb.ProtocolID_Test
113 + pid2 := pb.ProtocolID_Identify
114 mux1 := NewMuxer(ProtocolMap{
115 pid1: p1,
116 pid2: p2,
@@ -127,7 +128,7 @@ func TestSimultMuxer(t *testing.T) {
128 counts := [2][2][2]int{}
129
130 // run producers at every end sending incrementing messages
130 - produceOut := func(pid ProtocolID, size int) {
131 + produceOut := func(pid pb.ProtocolID, size int) {
132 limiter := time.Tick(speed)
133 for i := 0; i < size; i++ {
134 <-limiter
@@ -139,7 +140,7 @@ func TestSimultMuxer(t *testing.T) {
140 }
141 }
142
142 - produceIn := func(pid ProtocolID, size int) {
143 + produceIn := func(pid pb.ProtocolID, size int) {
144 limiter := time.Tick(speed)
145 for i := 0; i < size; i++ {
146 <-limiter
@@ -175,7 +176,7 @@ func TestSimultMuxer(t *testing.T) {
176 }
177 }
178
178 - consumeIn := func(pid ProtocolID) {
179 + consumeIn := func(pid pb.ProtocolID) {
180 for {
181 select {
182 case m := <-mux1.Protocols[pid].GetPipe().Incoming:
@@ -217,8 +218,8 @@ func TestStopping(t *testing.T) {
218 // setup
219 p1 := &TestProtocol{Pipe: msg.NewPipe(10)}
220 p2 := &TestProtocol{Pipe: msg.NewPipe(10)}
220 - pid1 := ProtocolID_Test
221 - pid2 := ProtocolID_Identify
221 + pid1 := pb.ProtocolID_Test
222 + pid2 := pb.ProtocolID_Identify
223 mux1 := NewMuxer(ProtocolMap{
224 pid1: p1,
225 pid2: p2,