handshake3 (addrs)
Juan Batiz-Benet committed
Oct 22, 2014 at 05:25 UTC
701035d5b01e1d8b069f9b7e2ea9e35e85b64b05
5 files changed
+129
-6
net/conn/conn.go
+2
-2
@@ -79,9 +79,9 @@ func newSingleConn(ctx context.Context, local, remote peer.Peer,
79
80
// version handshake
81
ctxT, _ := context.WithTimeout(ctx, HandshakeTimeout)
82
- if err := VersionHandshake(ctxT, conn); err != nil {
82
+ if err := Handshake1(ctxT, conn); err != nil {
83
conn.Close()
84
- return nil, fmt.Errorf("Version handshake: %s", err)
84
+ return nil, fmt.Errorf("Handshake1 failed: %s", err)
85
}
86
87
return conn, nil
net/conn/handshake.go
+46
-2
@@ -11,9 +11,9 @@ import (
11
proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
12
)
13
14
-// VersionHandshake exchanges local and remote versions and compares them
14
+// Handshake1 exchanges local and remote versions and compares them
15
// closes remote and returns an error in case of major difference
16
-func VersionHandshake(ctx context.Context, c Conn) error {
16
+func Handshake1(ctx context.Context, c Conn) error {
17
rpeer := c.RemotePeer()
18
lpeer := c.LocalPeer()
19
@@ -57,3 +57,47 @@ func VersionHandshake(ctx context.Context, c Conn) error {
57
log.Debug("%s version handshake compatible %s", lpeer, rpeer)
58
return nil
59
}
60
+
61
+// Handshake3 exchanges local and remote service information
62
+func Handshake3(ctx context.Context, c Conn) error {
63
+ rpeer := c.RemotePeer()
64
+ lpeer := c.LocalPeer()
65
+
66
+ var remoteH, localH *hspb.Handshake3
67
+ localH = handshake.Handshake3Msg(lpeer)
68
+ localB, err := proto.Marshal(localH)
69
+ if err != nil {
70
+ return err
71
+ }
72
+
73
+ c.Out() <- localB
74
+ log.Debug("Handshake1: sent to %s", rpeer)
75
+
76
+ select {
77
+ case <-ctx.Done():
78
+ return ctx.Err()
79
+
80
+ case <-c.Closing():
81
+ return errors.New("Handshake3: error remote connection closed")
82
+
83
+ case remoteB, ok := <-c.In():
84
+ if !ok {
85
+ return fmt.Errorf("Handshake3 error receiving from conn: %v", rpeer)
86
+ }
87
+
88
+ remoteH = new(hspb.Handshake3)
89
+ err = proto.Unmarshal(remoteB, remoteH)
90
+ if err != nil {
91
+ return fmt.Errorf("Handshake3 could not decode remote msg: %q", err)
92
+ }
93
+
94
+ log.Debug("Handshake3 received from %s", rpeer)
95
+ }
96
+
97
+ if err := handshake.Handshake3UpdatePeer(rpeer, remoteH); err != nil {
98
+ log.Error("Handshake3 failed to update %s", rpeer)
99
+ return err
100
+ }
101
+
102
+ return nil
103
+}
net/handshake/handshake3.go
new
+54
@@ -0,0 +1,54 @@
1
+package handshake
2
+
3
+import (
4
+ "fmt"
5
+
6
+ pb "github.com/jbenet/go-ipfs/net/handshake/pb"
7
+ peer "github.com/jbenet/go-ipfs/peer"
8
+ u "github.com/jbenet/go-ipfs/util"
9
+
10
+ ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
11
+)
12
+
13
+var log = u.Logger("handshake")
14
+
15
+// Handshake3Msg constructs a Handshake3 msg.
16
+func Handshake3Msg(localPeer peer.Peer) *pb.Handshake3 {
17
+ var msg pb.Handshake3
18
+ // don't need publicKey after secure channel.
19
+ // msg.PublicKey = localPeer.PubKey().Bytes()
20
+
21
+ // addresses
22
+ addrs := localPeer.Addresses()
23
+ msg.ListenAddrs = make([][]byte, len(addrs))
24
+ for i, a := range addrs {
25
+ msg.ListenAddrs[i] = a.Bytes()
26
+ }
27
+
28
+ // services
29
+ // srv := localPeer.Services()
30
+ // msg.Services = make([]mux.ProtocolID, len(srv))
31
+ // for i, pid := range srv {
32
+ // msg.Services[i] = pid
33
+ // }
34
+
35
+ return &msg
36
+}
37
+
38
+// Handshake3UpdatePeer updates a remote peer with the information in the
39
+// handshake3 msg we received from them.
40
+func Handshake3UpdatePeer(remotePeer peer.Peer, msg *pb.Handshake3) error {
41
+
42
+ // addresses
43
+ for _, a := range msg.GetListenAddrs() {
44
+ addr, err := ma.NewMultiaddrBytes(a)
45
+ if err != nil {
46
+ err = fmt.Errorf("remote peer address not a multiaddr: %s", err)
47
+ log.Error("Handshake3: error %s", err)
48
+ return err
49
+ }
50
+ remotePeer.AddAddress(addr)
51
+ }
52
+
53
+ return nil
54
+}
net/swarm/conn.go
+8
@@ -7,6 +7,7 @@ import (
7
conn "github.com/jbenet/go-ipfs/net/conn"
8
msg "github.com/jbenet/go-ipfs/net/message"
9
10
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
11
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
12
)
13
@@ -99,6 +100,13 @@ func (s *Swarm) connSetup(c conn.Conn) (conn.Conn, error) {
100
// addresses should be figured out through the DHT.
101
// c.Remote.AddAddress(c.Conn.RemoteMultiaddr())
102
103
+ // handshake3
104
+ ctxT, _ := context.WithTimeout(c.Context(), conn.HandshakeTimeout)
105
+ if err := conn.Handshake3(ctxT, c); err != nil {
106
+ c.Close()
107
+ return nil, fmt.Errorf("Handshake3 failed: %s", err)
108
+ }
109
+
110
// add to conns
111
s.connsLock.Lock()
112
peer/peer.go
+19
-2
@@ -1,6 +1,7 @@
1
package peer
2
3
import (
4
+ "bytes"
5
"errors"
6
"fmt"
7
"sync"
@@ -9,10 +10,9 @@ import (
10
b58 "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-base58"
11
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
12
mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
13
+
14
ic "github.com/jbenet/go-ipfs/crypto"
15
u "github.com/jbenet/go-ipfs/util"
14
-
15
- "bytes"
16
)
17
18
var log = u.Logger("peer")
@@ -71,6 +71,10 @@ type Peer interface {
71
// NetAddress returns the first Multiaddr found for a given network.
72
NetAddress(n string) ma.Multiaddr
73
74
+ // Services returns the peer's services
75
+ // Services() []mux.ProtocolID
76
+ // SetServices([]mux.ProtocolID)
77
+
78
// Priv/PubKey returns the peer's Private Key
79
PrivKey() ic.PrivKey
80
PubKey() ic.PubKey
@@ -92,6 +96,7 @@ type Peer interface {
96
type peer struct {
97
id ID
98
addresses []ma.Multiaddr
99
+ // services []mux.ProtocolID
100
101
privKey ic.PrivKey
102
pubKey ic.PubKey
@@ -163,6 +168,18 @@ func (p *peer) NetAddress(n string) ma.Multiaddr {
168
return nil
169
}
170
171
+// func (p *peer) Services() []mux.ProtocolID {
172
+// p.RLock()
173
+// defer p.RUnlock()
174
+// return p.services
175
+// }
176
+//
177
+// func (p *peer) SetServices(s []mux.ProtocolID) {
178
+// p.Lock()
179
+// defer p.Unlock()
180
+// p.services = s
181
+// }
182
+
183
// GetLatency retrieves the current latency measurement.
184
func (p *peer) GetLatency() (out time.Duration) {
185
p.RLock()