conn: raw []byte, not msg
This commit actually removes the previously introduced chan net.NetMessage, in favor of raw []byte. It plays nicer with crypto/spipe, and it makes more sense in the context of a "single connection", i.e. I already know the peer I'm talking to, from the connection. The NetMessage additional Peer is useful swarm and up.
Juan Batiz-Benet committed
Oct 18, 2014 at 01:56 UTC
7a7bf8d839053b07aca70a4c3f07254d6d93f50d
3 files changed
+23
-64
net/conn/conn.go
+6
-49
@@ -10,7 +10,6 @@ import (
10
manet "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr/net"
11
12
spipe "github.com/jbenet/go-ipfs/crypto/spipe"
13
- msg "github.com/jbenet/go-ipfs/net/message"
13
peer "github.com/jbenet/go-ipfs/peer"
14
u "github.com/jbenet/go-ipfs/util"
15
)
@@ -48,12 +47,8 @@ type singleConn struct {
47
48
secure *spipe.SecurePipe
49
insecure *msgioPipe
51
- msgpipe *msg.Pipe
50
}
51
54
-// Map maps Keys (Peer.IDs) to Connections.
55
-type Map map[u.Key]Conn
56
-
52
// newConn constructs a new connection
53
func newSingleConn(ctx context.Context, local, remote *peer.Peer,
54
peers peer.Peerstore, maconn manet.Conn) (Conn, error) {
@@ -67,7 +62,6 @@ func newSingleConn(ctx context.Context, local, remote *peer.Peer,
62
ctx: ctx,
63
cancel: cancel,
64
insecure: newMsgioPipe(10),
70
- msgpipe: msg.NewPipe(10),
65
}
66
67
log.Info("newSingleConn: %v to %v", local, remote)
@@ -117,45 +111,9 @@ func (c *singleConn) secureHandshake(peers peer.Peerstore) error {
111
panic("peers not being constructed correctly.")
112
}
113
120
- // silly we have to do it this way.
121
- go c.unwrapOutMsgs()
122
- go c.wrapInMsgs()
123
-
114
return nil
115
}
116
127
-// unwrapOutMsgs sends just the raw data of a message through secure
128
-func (c *singleConn) unwrapOutMsgs() {
129
- for {
130
- select {
131
- case <-c.ctx.Done():
132
- return
133
- case m, more := <-c.msgpipe.Outgoing:
134
- if !more {
135
- return
136
- }
137
-
138
- c.secure.Out <- m.Data()
139
- }
140
- }
141
-}
142
-
143
-// wrapInMsgs wraps a message
144
-func (c *singleConn) wrapInMsgs() {
145
- for {
146
- select {
147
- case <-c.ctx.Done():
148
- return
149
- case d, more := <-c.secure.In:
150
- if !more {
151
- return
152
- }
153
-
154
- c.msgpipe.Incoming <- msg.New(c.remote, d)
155
- }
156
- }
157
-}
158
-
117
// waitToClose waits on the given context's Done before closing Conn.
118
func (c *singleConn) waitToClose() {
119
select {
@@ -170,7 +128,6 @@ func (c *singleConn) waitToClose() {
128
if c.secure != nil { // may never have gotten here.
129
c.secure.Close()
130
}
173
- close(c.msgpipe.Incoming)
131
}
132
133
// isClosed returns whether this Conn is open or closed.
@@ -205,14 +162,14 @@ func (c *singleConn) RemotePeer() *peer.Peer {
162
return c.remote
163
}
164
208
-// MsgIn returns a readable message channel
209
-func (c *singleConn) MsgIn() <-chan msg.NetMessage {
210
- return c.msgpipe.Incoming
165
+// In returns a readable message channel
166
+func (c *singleConn) In() <-chan []byte {
167
+ return c.secure.In
168
}
169
213
-// MsgOut returns a writable message channel
214
-func (c *singleConn) MsgOut() chan<- msg.NetMessage {
215
- return c.msgpipe.Outgoing
170
+// Out returns a writable message channel
171
+func (c *singleConn) Out() chan<- []byte {
172
+ return c.secure.Out
173
}
174
175
// Dialer is an object that can open connections. We could have a "convenience"
net/conn/conn_test.go
+9
-10
@@ -4,7 +4,6 @@ import (
4
"testing"
5
6
ci "github.com/jbenet/go-ipfs/crypto"
7
- msg "github.com/jbenet/go-ipfs/net/message"
7
peer "github.com/jbenet/go-ipfs/peer"
8
9
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
@@ -50,8 +49,8 @@ func echo(ctx context.Context, c Conn) {
49
select {
50
case <-ctx.Done():
51
return
53
- case m := <-c.MsgIn():
54
- c.MsgOut() <- m
52
+ case m := <-c.In():
53
+ c.Out() <- m
54
}
55
}
56
}
@@ -93,19 +92,19 @@ func TestDialer(t *testing.T) {
92
}
93
94
// fmt.Println("sending")
96
- c.MsgOut() <- msg.New(p2, []byte("beep"))
97
- c.MsgOut() <- msg.New(p2, []byte("boop"))
95
+ c.Out() <- []byte("beep")
96
+ c.Out() <- []byte("boop")
97
99
- out := <-c.MsgIn()
98
+ out := <-c.In()
99
// fmt.Println("recving", string(out))
101
- data := string(out.Data())
100
+ data := string(out)
101
if data != "beep" {
102
t.Error("unexpected conn output", data)
103
}
104
106
- out = <-c.MsgIn()
107
- data = string(out.Data())
108
- if string(out.Data()) != "boop" {
105
+ out = <-c.In()
106
+ data = string(out)
107
+ if string(out) != "boop" {
108
t.Error("unexpected conn output", data)
109
}
110
net/conn/interface.go
+8
-5
@@ -1,12 +1,15 @@
1
package conn
2
3
import (
4
- msg "github.com/jbenet/go-ipfs/net/message"
4
peer "github.com/jbenet/go-ipfs/peer"
5
+ u "github.com/jbenet/go-ipfs/util"
6
7
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
8
)
9
10
+// Map maps Keys (Peer.IDs) to Connections.
11
+type Map map[u.Key]Conn
12
+
13
// Conn is a generic message-based Peer-to-Peer connection.
14
type Conn interface {
15
@@ -16,11 +19,11 @@ type Conn interface {
19
// RemotePeer is the Peer on the remote side
20
RemotePeer() *peer.Peer
21
19
- // MsgIn returns a readable message channel
20
- MsgIn() <-chan msg.NetMessage
22
+ // In returns a readable message channel
23
+ In() <-chan []byte
24
22
- // MsgOut returns a writable message channel
23
- MsgOut() chan<- msg.NetMessage
25
+ // Out returns a writable message channel
26
+ Out() chan<- []byte
27
28
// Close ends the connection
29
Close() error